| 1 | // EduBerza market simulation bot.
|
|---|
| 2 | //
|
|---|
| 3 | // Walks a small price for every active market, inserts rows into
|
|---|
| 4 | // market_trades once per tick, and upserts the current 1m candle.
|
|---|
| 5 | //
|
|---|
| 6 | // Run from the repo root so the relative .env path resolves:
|
|---|
| 7 | //
|
|---|
| 8 | // go run ./bots/...
|
|---|
| 9 | package main
|
|---|
| 10 |
|
|---|
| 11 | import (
|
|---|
| 12 | "bufio"
|
|---|
| 13 | "database/sql"
|
|---|
| 14 | "flag"
|
|---|
| 15 | "fmt"
|
|---|
| 16 | "log"
|
|---|
| 17 | "math/rand"
|
|---|
| 18 | "os"
|
|---|
| 19 | "strings"
|
|---|
| 20 | "time"
|
|---|
| 21 |
|
|---|
| 22 | _ "github.com/lib/pq"
|
|---|
| 23 | )
|
|---|
| 24 |
|
|---|
| 25 | type market struct {
|
|---|
| 26 | id string
|
|---|
| 27 | symbol string
|
|---|
| 28 | price float64
|
|---|
| 29 | }
|
|---|
| 30 |
|
|---|
| 31 | func main() {
|
|---|
| 32 | interval := flag.Duration("interval", 3*time.Second, "seconds between price ticks")
|
|---|
| 33 | flag.Parse()
|
|---|
| 34 |
|
|---|
| 35 | loadEnv(".env")
|
|---|
| 36 | dsn := fmt.Sprintf(
|
|---|
| 37 | "host=%s port=%s user=%s password=%s dbname=%s sslmode=disable options='--search_path=project,public'",
|
|---|
| 38 | env("DBHOST", "localhost"),
|
|---|
| 39 | env("DBPORT", "5432"),
|
|---|
| 40 | env("DBUSER", "postgres"),
|
|---|
| 41 | env("DBPASSWORD", ""),
|
|---|
| 42 | env("DBNAME", "postgres"),
|
|---|
| 43 | )
|
|---|
| 44 | db, err := sql.Open("postgres", dsn)
|
|---|
| 45 | if err != nil {
|
|---|
| 46 | log.Fatalf("open db: %v", err)
|
|---|
| 47 | }
|
|---|
| 48 | defer db.Close()
|
|---|
| 49 | if err := db.Ping(); err != nil {
|
|---|
| 50 | log.Fatalf("ping: %v", err)
|
|---|
| 51 | }
|
|---|
| 52 |
|
|---|
| 53 | markets, err := loadMarkets(db)
|
|---|
| 54 | if err != nil {
|
|---|
| 55 | log.Fatalf("load markets: %v", err)
|
|---|
| 56 | }
|
|---|
| 57 | if len(markets) == 0 {
|
|---|
| 58 | log.Fatal("no active markets found - run `go run ./server -init` first")
|
|---|
| 59 | }
|
|---|
| 60 |
|
|---|
| 61 | log.Printf("bot started. simulating %d markets every %s", len(markets), *interval)
|
|---|
| 62 | rnd := rand.New(rand.NewSource(time.Now().UnixNano()))
|
|---|
| 63 |
|
|---|
| 64 | for {
|
|---|
| 65 | for i := range markets {
|
|---|
| 66 | m := &markets[i]
|
|---|
| 67 | // random walk: ±0.3% per tick
|
|---|
| 68 | drift := (rnd.Float64() - 0.5) * 0.006
|
|---|
| 69 | m.price = m.price * (1 + drift)
|
|---|
| 70 | if m.price <= 0 {
|
|---|
| 71 | m.price = 0.000001
|
|---|
| 72 | }
|
|---|
| 73 | qty := rnd.Float64()*0.5 + 0.01
|
|---|
| 74 |
|
|---|
| 75 | side := "buy"
|
|---|
| 76 | if rnd.Float64() < 0.5 {
|
|---|
| 77 | side = "sell"
|
|---|
| 78 | }
|
|---|
| 79 |
|
|---|
| 80 | if err := insertTick(db, m.id, m.price, qty, side); err != nil {
|
|---|
| 81 | log.Printf("insert tick %s: %v", m.symbol, err)
|
|---|
| 82 | continue
|
|---|
| 83 | }
|
|---|
| 84 | log.Printf(" %-8s %.6f qty=%.4f side=%s", m.symbol, m.price, qty, side)
|
|---|
| 85 | }
|
|---|
| 86 | // P7 background job: the prices just moved, so fill any resting
|
|---|
| 87 | // limit order the new market price has reached.
|
|---|
| 88 | var filled int
|
|---|
| 89 | if err := db.QueryRow(`SELECT fill_marketable_orders()`).Scan(&filled); err != nil {
|
|---|
| 90 | log.Printf("fill_marketable_orders: %v", err)
|
|---|
| 91 | } else if filled > 0 {
|
|---|
| 92 | log.Printf(" filled %d resting limit order(s) at the new market price", filled)
|
|---|
| 93 | }
|
|---|
| 94 | time.Sleep(*interval)
|
|---|
| 95 | }
|
|---|
| 96 | }
|
|---|
| 97 |
|
|---|
| 98 | func loadMarkets(db *sql.DB) ([]market, error) {
|
|---|
| 99 | rows, err := db.Query(`
|
|---|
| 100 | SELECT m.id, c.symbol, COALESCE(lp.price, 100)
|
|---|
| 101 | FROM markets m
|
|---|
| 102 | JOIN crypto c ON c.id = m.crypto_id
|
|---|
| 103 | LEFT JOIN v_latest_prices lp ON lp.market_id = m.id
|
|---|
| 104 | WHERE m.is_active = true
|
|---|
| 105 | ORDER BY c.symbol`)
|
|---|
| 106 | if err != nil {
|
|---|
| 107 | return nil, err
|
|---|
| 108 | }
|
|---|
| 109 | defer rows.Close()
|
|---|
| 110 | var out []market
|
|---|
| 111 | for rows.Next() {
|
|---|
| 112 | var m market
|
|---|
| 113 | if err := rows.Scan(&m.id, &m.symbol, &m.price); err != nil {
|
|---|
| 114 | return nil, err
|
|---|
| 115 | }
|
|---|
| 116 | out = append(out, m)
|
|---|
| 117 | }
|
|---|
| 118 | return out, nil
|
|---|
| 119 | }
|
|---|
| 120 |
|
|---|
| 121 | func insertTick(db *sql.DB, marketID string, price, qty float64, side string) error {
|
|---|
| 122 | tx, err := db.Begin()
|
|---|
| 123 | if err != nil {
|
|---|
| 124 | return err
|
|---|
| 125 | }
|
|---|
| 126 | defer tx.Rollback()
|
|---|
| 127 |
|
|---|
| 128 | if _, err := tx.Exec(
|
|---|
| 129 | `INSERT INTO market_trades (market_id, executed_at, price, quantity, side, source)
|
|---|
| 130 | VALUES ($1, now(), $2, $3, $4, 'simulation')`,
|
|---|
| 131 | marketID, price, qty, side,
|
|---|
| 132 | ); err != nil {
|
|---|
| 133 | return err
|
|---|
| 134 | }
|
|---|
| 135 |
|
|---|
| 136 | // upsert the current 1m candle
|
|---|
| 137 | if _, err := tx.Exec(`
|
|---|
| 138 | INSERT INTO market_candles (market_id, timeframe, open, high, low, close, volume, candle_time)
|
|---|
| 139 | VALUES ($1, '1m', $2, $2, $2, $2, $3, date_trunc('minute', now()))
|
|---|
| 140 | ON CONFLICT (market_id, timeframe, candle_time) DO UPDATE
|
|---|
| 141 | SET high = GREATEST(market_candles.high, EXCLUDED.close),
|
|---|
| 142 | low = LEAST( market_candles.low, EXCLUDED.close),
|
|---|
| 143 | close = EXCLUDED.close,
|
|---|
| 144 | volume = market_candles.volume + EXCLUDED.volume`,
|
|---|
| 145 | marketID, price, qty,
|
|---|
| 146 | ); err != nil {
|
|---|
| 147 | return err
|
|---|
| 148 | }
|
|---|
| 149 | return tx.Commit()
|
|---|
| 150 | }
|
|---|
| 151 |
|
|---|
| 152 | func loadEnv(path string) {
|
|---|
| 153 | f, err := os.Open(path)
|
|---|
| 154 | if err != nil {
|
|---|
| 155 | return
|
|---|
| 156 | }
|
|---|
| 157 | defer f.Close()
|
|---|
| 158 | s := bufio.NewScanner(f)
|
|---|
| 159 | for s.Scan() {
|
|---|
| 160 | line := strings.TrimSpace(s.Text())
|
|---|
| 161 | if line == "" || strings.HasPrefix(line, "#") {
|
|---|
| 162 | continue
|
|---|
| 163 | }
|
|---|
| 164 | parts := strings.SplitN(line, "=", 2)
|
|---|
| 165 | if len(parts) == 2 {
|
|---|
| 166 | os.Setenv(strings.TrimSpace(parts[0]), strings.TrimSpace(parts[1]))
|
|---|
| 167 | }
|
|---|
| 168 | }
|
|---|
| 169 | }
|
|---|
| 170 |
|
|---|
| 171 | func env(k, def string) string {
|
|---|
| 172 | if v := os.Getenv(k); v != "" {
|
|---|
| 173 | return v
|
|---|
| 174 | }
|
|---|
| 175 | return def
|
|---|
| 176 | }
|
|---|