feat: full IDX universe (780 ticker) — universe table + close sweep harian, dashboard foreign-flow picker dari semua ticker, rotateDepth incremental fill 10/cylcle, strip .JK, backfill cmd/universe-sweep
This commit is contained in:
@@ -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)
|
||||
}
|
||||
@@ -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,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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 ''
|
||||
);
|
||||
@@ -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(?,?,?)
|
||||
|
||||
Reference in New Issue
Block a user