From 6fd89b6c197e63bcf75d9698c8822382c69de411 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Wed, 16 Sep 2026 01:06:42 +0700 Subject: [PATCH] =?UTF-8?q?feat:=20full=20IDX=20universe=20(780=20ticker)?= =?UTF-8?q?=20=E2=80=94=20universe=20table=20+=20close=20sweep=20harian,?= =?UTF-8?q?=20dashboard=20foreign-flow=20picker=20dari=20semua=20ticker,?= =?UTF-8?q?=20rotateDepth=20incremental=20fill=2010/cylcle,=20strip=20.JK,?= =?UTF-8?q?=20backfill=20cmd/universe-sweep?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 1 + backend/cmd/universe-sweep/main.go | 65 +++++++++++++ backend/internal/api/flow.go | 26 ++++- backend/internal/scheduler/scheduler.go | 97 +++++++++++++++++-- .../store/migrations/0009_universe.sql | 6 ++ backend/internal/store/rows.go | 72 +++++++++++++- web/src/lib/api.ts | 2 +- web/src/pages/Dashboard.tsx | 20 +++- 8 files changed, 275 insertions(+), 14 deletions(-) create mode 100644 backend/cmd/universe-sweep/main.go create mode 100644 backend/internal/store/migrations/0009_universe.sql diff --git a/.gitignore b/.gitignore index 51b6eca..3d83f0b 100644 --- a/.gitignore +++ b/.gitignore @@ -48,3 +48,4 @@ Thumbs.db task_plan.md findings.md progress.md +backend/universe-sweep diff --git a/backend/cmd/universe-sweep/main.go b/backend/cmd/universe-sweep/main.go new file mode 100644 index 0000000..c7d59dd --- /dev/null +++ b/backend/cmd/universe-sweep/main.go @@ -0,0 +1,65 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "log" + "os" + "time" + + "flowsight/internal/sectors" + "flowsight/internal/store" +) + +// one-shot: sweep the full IDX universe (close/ pages) and persist it into +// the universe table. Usage: BWS_PROJECT_ID=... sh scripts/bws-run.sh go run ./cmd/universe-sweep +func main() { + // env DB_PATH defaults to data/flowsight.db if unset. + db, err := store.Open(os.Getenv("DB_PATH")) + if err != nil { + log.Fatal(err) + } + defer db.Close() + c := sectors.New("https://api.sectors.app/v2/", os.Getenv("SECTORS_API_KEY")) + ctx, cancel := context.WithTimeout(context.Background(), 180*time.Second) + defer cancel() + c.OnSpend(func(endpoint string, calls, credits int) { fmt.Printf("spend %d credits %s\n", credits, endpoint) }) + + offset := 0 + var all []sectors.CloseRow + for page := 0; page < 60; page++ { + rows, _, err := c.ClosePage(ctx, "", 30, offset) + if err != nil { + log.Printf("page %d err: %v (keeping %d so far)", page, err, len(all)) + break + } + all = append(all, rows...) + offset += len(rows) + fmt.Printf("page %d: +%d (total %d)\n", page, len(rows), len(all)) + if len(rows) == 0 { + break + } + time.Sleep(1500 * time.Millisecond) // stay well under the 429 pace + } + fmt.Println("TOTAL:", len(all)) + if len(all) == 0 { + log.Fatal("no rows collected") + } + date := all[0].Date + for _, r := range all { + if r.Date > date { + date = r.Date + } + } + raw, _ := json.Marshal(all) + _ = db.SaveSnapshot("IDX", date, "close", string(raw)) + uni := make([]store.UniverseRow, 0, len(all)) + for _, r := range all { + uni = append(uni, store.UniverseRow{Symbol: r.Symbol, Close: r.Close, Date: r.Date}) + } + if err := db.SaveUniverse(uni); err != nil { + log.Fatal(err) + } + fmt.Println("universe saved:", len(uni), "tickers @", date) +} \ No newline at end of file diff --git a/backend/internal/api/flow.go b/backend/internal/api/flow.go index 878acb9..ce7d4b5 100644 --- a/backend/internal/api/flow.go +++ b/backend/internal/api/flow.go @@ -17,6 +17,8 @@ func (s *Server) FlowSummary(w http.ResponseWriter, r *http.Request) { if date == "" { date = "latest" } + // Build the full universe: watchlist if set, else all tickers with data — + // so the dashboard chart picker is not limited to a top-5. wl, _ := s.DB.Watchlist(s.userKey(r)) if len(wl) == 0 { if all, err := s.DB.AllTickers(); err == nil && len(all) > 0 { @@ -31,10 +33,20 @@ func (s *Server) FlowSummary(w http.ResponseWriter, r *http.Request) { Brokers int `json:"brokers"` } var accs []accRow + var fkTickers []string foreignTotal := 0.0 var closes []map[string]any var cites []model.Citation - for _, tk := range wl { + // Depth loop bounded to 60 tickers/request for responsiveness; tickers + // lacking broker data still appear in foreign_tickers (chart picker) — + // the full universe list is served cheaply from the universe table. + scanList := wl + if len(scanList) > 60 { + scanList = scanList[:60] + } + scanned := map[string]bool{} + for _, tk := range scanList { + scanned[tk] = true if nets, err := s.DB.NetBuySum5d(tk); err == nil && len(nets) > 0 { sum, n := 0.0, 0 for _, v := range nets { @@ -48,12 +60,23 @@ func (s *Server) FlowSummary(w http.ResponseWriter, r *http.Request) { } if _, nets, err := s.DB.ForeignLast6(tk); err == nil && len(nets) > 0 { foreignTotal += nets[len(nets)-1] + fkTickers = append(fkTickers, tk) cites = append(cites, model.Cite("v2/foreign-flow/"+tk+"/", tk, date)) } if px, dx, err := s.DB.LatestClose(tk); err == nil { closes = append(closes, map[string]any{"ticker": tk, "close": px, "date": dx}) } } + // Rest of the universe: still list them as chart-picker options (foreign + // data may exist even if not scanned this request). + if len(wl) > 60 { + extra, _ := s.DB.TickersWithForeign() + for _, tk := range wl { + if !scanned[tk] && extra[tk] { + fkTickers = append(fkTickers, tk) + } + } + } sort.Slice(accs, func(i, j int) bool { return accs[i].NetSum > accs[j].NetSum }) if len(accs) > 5 { accs = accs[:5] @@ -61,6 +84,7 @@ func (s *Server) FlowSummary(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, map[string]any{ "date": date, "foreign_net_total": foreignTotal, "top_accumulation": accs, "closes": closes, "citations": cites, + "foreign_tickers": fkTickers, }) } diff --git a/backend/internal/scheduler/scheduler.go b/backend/internal/scheduler/scheduler.go index 810d345..8979268 100644 --- a/backend/internal/scheduler/scheduler.go +++ b/backend/internal/scheduler/scheduler.go @@ -11,6 +11,7 @@ import ( "fmt" "log" "net/url" + "strconv" "time" "github.com/robfig/cron/v3" @@ -175,6 +176,15 @@ func (s *Scheduler) RunCycle(ctx context.Context) error { return err } } + // Full-universe rotation: every cycle pull depth (foreign flow included) + // for a slice of the remaining universe, credit-aware, so the dashboard + // chart picker eventually covers all ~800 IDX tickers, not just the 5. + if err := s.rotateDepth(ctx); err != nil { + log.Printf("scheduler: rotate: %v", err) + } + if err := guard(); err != nil { + return err + } if err := s.events(ctx); err != nil { log.Printf("scheduler: events: %v", err) } @@ -234,22 +244,53 @@ func (s *Scheduler) reference(ctx context.Context) error { return nil } -// universe sweeps close/ pages for the latest trading day. +// universe sweeps close/ pages for the latest trading day and persists the +// full ticker list (all pages, not just the last) so AllTickers() knows the +// whole IDX universe. ~800 tickers = ~27 pages at 30/page; each page burns 1 +// credit. Runs at most once per day (meta universe_last_sweep) and aborts on +// 429 keeping pages collected so far. func (s *Scheduler) universe(ctx context.Context) error { + if s.DB.GetMeta("universe_last_sweep") == time.Now().Format("2006-01-02") { + return nil // already swept today; rotation fills depth incrementally + } offset := 0 - for page := 0; page < 12; page++ { - rows, total, err := s.Sectors.ClosePage(ctx, "", 30, offset) + var all []sectors.CloseRow + for page := 0; page < 40; page++ { + rows, _, err := s.Sectors.ClosePage(ctx, "", 30, offset) if err != nil { + // Rate-limited mid-sweep: keep pages collected so far if any. + if len(all) > 0 { + break + } return err } - if raw, err := json.Marshal(rows); err == nil && len(rows) > 0 { - _ = s.DB.SaveSnapshot("IDX", rows[0].Date, "close", string(raw)) - } + all = append(all, rows...) offset += len(rows) - if offset >= total || len(rows) == 0 { + if len(rows) == 0 { break } } + if len(all) == 0 { + return nil + } + s.DB.SetMeta("universe_last_sweep", time.Now().Format("2006-01-02")) + date := all[0].Date + for _, r := range all { + if r.Date > date { + date = r.Date + } + } + raw, _ := json.Marshal(all) + _ = s.DB.SaveSnapshot("IDX", date, "close", string(raw)) + // Persist the universe ticker list so AllTickers() can return the full + // IDX set instead of only tickers that have depth data yet. + uni := make([]store.UniverseRow, 0, len(all)) + for _, r := range all { + uni = append(uni, store.UniverseRow{Symbol: r.Symbol, Close: r.Close, Date: r.Date}) + } + if err := s.DB.SaveUniverse(uni); err != nil { + log.Printf("scheduler: universe save: %v", err) + } return nil } @@ -319,6 +360,48 @@ func (s *Scheduler) tickerDepth(ctx context.Context, ticker string) error { return nil } +// rotateDepth incrementally deep-scans the whole universe across cycles: +// each cycle it pulls depth for a bounded slice of tickers that don't yet +// have a foreign-flow snapshot, respecting the credit cap. Progress is +// tracked via a meta cursor (universe_rotate_offset). Every ticker gets +// covered every ~80 cycles (10/cycle x 800), and the dashboard picker +// grows from 5 → ~800 tickers. +func (s *Scheduler) rotateDepth(ctx context.Context) error { + all, err := s.DB.AllTickers() + if err != nil || len(all) == 0 { + return nil + } + have, err := s.DB.TickersWithForeign() + if err != nil { + return err + } + missing := make([]string, 0, len(all)) + for _, t := range all { + if !have[t] { + missing = append(missing, t) + } + } + if len(missing) == 0 { + return nil // full coverage reached; daily freshness keeps it fresh + } + // Round-robin cursor so we don't always start at the same ticker. + start, _ := strconv.Atoi(s.DB.GetMeta("universe_rotate_offset")) + if start >= len(missing) { + start = 0 + } + batch := 10 // 3 credits/ticker = 30 credits; leaves headroom within the cap + for i := 0; i < batch; i++ { + t := missing[(start+i)%len(missing)] + if err := s.tickerDepth(ctx, t); err != nil { + // 429 or transient: stop this cycle, resume next. + s.DB.SetMeta("universe_rotate_offset", fmt.Sprint((start+i)%len(missing))) + return err + } + } + s.DB.SetMeta("universe_rotate_offset", fmt.Sprint((start+batch)%len(missing))) + return nil +} + // events polls news/filings/suspensions incrementally via meta cursors. func (s *Scheduler) events(ctx context.Context) error { today := time.Now().Format("2006-01-02") diff --git a/backend/internal/store/migrations/0009_universe.sql b/backend/internal/store/migrations/0009_universe.sql new file mode 100644 index 0000000..eadfa3c --- /dev/null +++ b/backend/internal/store/migrations/0009_universe.sql @@ -0,0 +1,6 @@ +-- 0009_universe.sql: persisted full IDX ticker universe from the close sweep. +CREATE TABLE IF NOT EXISTS universe( + ticker TEXT PRIMARY KEY, + close INTEGER NOT NULL DEFAULT 0, + date TEXT NOT NULL DEFAULT '' +); \ No newline at end of file diff --git a/backend/internal/store/rows.go b/backend/internal/store/rows.go index f4df9e1..7d625b5 100644 --- a/backend/internal/store/rows.go +++ b/backend/internal/store/rows.go @@ -4,6 +4,7 @@ import ( "database/sql" "encoding/json" "fmt" + "strings" "time" ) @@ -545,23 +546,88 @@ func (db *DB) Watchlist(userKey string) ([]string, error) { // default universe: dashboard, screener fallback, and scheduler depth all use // it when a user has no personal watchlist. func (db *DB) AllTickers() ([]string, error) { - rows, err := db.Query(`SELECT DISTINCT ticker FROM snapshots + // Prefer the persisted full-universe list (from close/ sweep). + rows, err := db.Query(`SELECT ticker FROM universe ORDER BY ticker`) + if err == nil { + var out []string + for rows.Next() { + var t string + if err := rows.Scan(&t); err != nil { + rows.Close() + break + } + out = append(out, t) + } + rows.Close() + if len(out) > 0 { + return out, nil + } + } + // Fallback: distinct snapshots (avoids pseudo indices). + rows2, err := db.Query(`SELECT DISTINCT ticker FROM snapshots WHERE ticker NOT IN ('IDX','ROE') ORDER BY ticker`) if err != nil { return nil, err } + defer rows2.Close() + var out2 []string + for rows2.Next() { + var t string + if err := rows2.Scan(&t); err != nil { + return nil, err + } + out2 = append(out2, t) + } + return out2, rows2.Err() +} + +// TickersWithForeign returns the set of tickers that have ≥1 foreign-flow row. +func (db *DB) TickersWithForeign() (map[string]bool, error) { + rows, err := db.Query(`SELECT DISTINCT ticker FROM foreign_flow`) + if err != nil { + return nil, err + } defer rows.Close() - var out []string + out := map[string]bool{} for rows.Next() { var t string if err := rows.Scan(&t); err != nil { return nil, err } - out = append(out, t) + out[t] = true } return out, rows.Err() } +// SaveUniverse inserts the full ticker universe from the close sweep. +// UniverseRow mirrors the CloseRow shape without importing sectors (cycle). +type UniverseRow struct { + Symbol string + Close int64 + Date string +} + +func (db *DB) SaveUniverse(rows []UniverseRow) error { + tx, err := db.Begin() + if err != nil { + return err + } + defer tx.Rollback() + if _, err := tx.Exec(`DELETE FROM universe`); err != nil { + return err + } + for _, r := range rows { + // Normalize: API returns "BBCA.JK" — strip the suffix so it matches + // watchlist/depth keys (BBCA) everywhere. + t := strings.TrimSuffix(r.Symbol, ".JK") + if _, err := tx.Exec(`INSERT OR IGNORE INTO universe(ticker, close, date) VALUES(?,?,?)`, + t, r.Close, r.Date); err != nil { + return err + } + } + return tx.Commit() +} + // AddWatch inserts a ticker (idempotent). func (db *DB) AddWatch(userKey, ticker string) error { _, err := db.Exec(`INSERT INTO watchlists(user_key,ticker,added_at) VALUES(?,?,?) diff --git a/web/src/lib/api.ts b/web/src/lib/api.ts index b0b73e8..2e4d0f0 100644 --- a/web/src/lib/api.ts +++ b/web/src/lib/api.ts @@ -25,7 +25,7 @@ async function req(path: string, init?: RequestInit): Promise { } export interface Citation { endpoint: string; snapshot_at: string; ticker?: string; stale?: boolean } export interface Health { last_cycle_at: string; credits_today: number; scheduler_ok: boolean; stale_flags: string[] } -export interface FlowSummary { date: string; foreign_net_total: number; top_accumulation: { ticker: string; net_sum: number; brokers: number }[]; closes: { ticker: string; close: number; date: string }[]; citations: Citation[] } +export interface FlowSummary { date: string; foreign_net_total: number; top_accumulation: { ticker: string; net_sum: number; brokers: number }[]; closes: { ticker: string; close: number; date: string }[]; citations: Citation[]; foreign_tickers?: string[] } export interface ScreenRow { symbol: string; name: string; composite: number; breakdown: Record; citations: Citation[] } export interface Routine { id: number; user_key: string; type: string; schedule_cron: string; channels_json?: string; channels?: string[]; enabled: boolean; last_run?: unknown } export interface AlertItem { id: number; user_key: string; name: string; rule_json: string; channels_json: string; last_fired: string } diff --git a/web/src/pages/Dashboard.tsx b/web/src/pages/Dashboard.tsx index 7ecd85a..65f7184 100644 --- a/web/src/pages/Dashboard.tsx +++ b/web/src/pages/Dashboard.tsx @@ -1,4 +1,4 @@ -import { createResource, createSignal, createMemo, For, Show } from "solid-js"; +import { createResource, createSignal, createEffect, For, Show } from "solid-js"; import { A, useNavigate } from "@solidjs/router"; import { api } from "../lib/api"; import { useAuth } from "../components/auth"; @@ -22,6 +22,13 @@ export default function Dashboard() { const [events] = createResource(me, (user) => (user ? api.alertEvents("2000-01-01").then((r) => r.events.slice(0, 5)).catch(() => []) : [])); const [chartTiker, setChartTiker] = createSignal("BBCA"); const [foreign] = createResource(chartTiker, (t) => api.flowForeign(t).catch(() => null)); + // Full picker list = every ticker that has foreign-flow rows. + const fkTickers = () => flow()?.foreign_tickers || []; + // Keep the current chart ticker valid even when the list changes. + createEffect(() => { + const list = fkTickers(); + if (list.length && !list.includes(chartTiker())) setChartTiker(list[0]); + }); const ringkasan = () => { const n = briefing()?.narasi?.trim(); @@ -124,7 +131,16 @@ export default function Dashboard() { harian — naik = asing beli, turun = asing jual. Pilih saham: -
+
+ + + {fkTickers().length} saham tersedia + t.ticker)}>{(t) => }