package api import ( "encoding/json" "net/http" "sort" "strings" "flowsight/internal/model" "flowsight/internal/sectors" ) // FlowSummary serves GET /api/flow/summary: foreign net total, top-5 // accumulation rows, rotation signal, mover of the day — all cited. func (s *Server) FlowSummary(w http.ResponseWriter, r *http.Request) { date := r.URL.Query().Get("date") 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 { wl = all } else { wl = s.Cfg.Watchlist } } type accRow struct { Ticker string `json:"ticker"` NetSum float64 `json:"net_sum"` Brokers int `json:"brokers"` } var accs []accRow var fkTickers []string foreignTotal := 0.0 var closes []map[string]any var cites []model.Citation // 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 { if v > 0 { n++ sum += v } } accs = append(accs, accRow{tk, sum, n}) cites = append(cites, model.Cite("v2/broker-summary/"+tk+"/top/", tk, date)) } 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). 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] } writeJSON(w, http.StatusOK, map[string]any{ "date": date, "foreign_net_total": foreignTotal, "top_accumulation": accs, "closes": closes, "citations": cites, "foreign_tickers": fkTickers, }) } // brokerQuery validates ticker/start/end query params. type brokerQuery struct { Ticker string `validate:"required,len=4"` Start string `validate:"omitempty,datetime=2006-01-02"` End string `validate:"omitempty,datetime=2006-01-02"` } // FlowBroker serves GET /api/flow/broker: buyers/sellers + 5d net series. func (s *Server) FlowBroker(w http.ResponseWriter, r *http.Request) { q := brokerQuery{ Ticker: strings.ToUpper(r.URL.Query().Get("ticker")), Start: r.URL.Query().Get("start"), End: r.URL.Query().Get("end"), } if err := s.Validate.Struct(q); err != nil { writeErr(w, http.StatusUnprocessableEntity, "ticker (4 letters) required, dates YYYY-MM-DD") return } var top sectors.BrokerSummaryTop if raw, d, err := s.DB.SnapshotAt(q.Ticker, "broker-summary-top", q.End); err == nil { _ = json.Unmarshal([]byte(raw), &top) writeJSON(w, http.StatusOK, map[string]any{ "ticker": q.Ticker, "buyers": top.TopBuyers, "sellers": top.TopSellers, "citations": []model.Citation{model.Cite("v2/broker-summary/"+q.Ticker+"/top/", q.Ticker, d)}, }) return } writeJSON(w, http.StatusOK, map[string]any{ "ticker": q.Ticker, "buyers": []any{}, "sellers": []any{}, "citations": []model.Citation{}, "note": "no snapshots yet", }) } // FlowForeign serves GET /api/flow/foreign: inflow series + reversal flag. func (s *Server) FlowForeign(w http.ResponseWriter, r *http.Request) { q := brokerQuery{ Ticker: strings.ToUpper(r.URL.Query().Get("ticker")), Start: r.URL.Query().Get("start"), End: r.URL.Query().Get("end"), } if err := s.Validate.Struct(q); err != nil { writeErr(w, http.StatusUnprocessableEntity, "ticker (4 letters) required, dates YYYY-MM-DD") return } dates, nets, err := s.DB.ForeignWindow(q.Ticker, q.Start, q.End, 30) if err != nil || len(nets) == 0 { writeJSON(w, http.StatusOK, map[string]any{ "ticker": q.Ticker, "series": []any{}, "reversal": false, "citations": []model.Citation{}, "note": "no snapshots yet", }) return } // Same 2x-magnitude rule as the alert engine: 5d cumulative one way, // last day the other way at >2x the trailing 5d daily average. reversal := false if len(nets) >= 6 { tail := nets[len(nets)-6:] sum5, absAvg := 0.0, 0.0 for _, v := range tail[:5] { sum5 += v if v < 0 { absAvg -= v } else { absAvg += v } } absAvg /= 5 last := tail[5] if absAvg > 0 && ((sum5 < 0 && last > 0 && last > 2*absAvg) || (sum5 > 0 && last < 0 && -last > 2*absAvg)) { reversal = true } } writeJSON(w, http.StatusOK, map[string]any{ "ticker": q.Ticker, "dates": dates, "nets": nets, "reversal": reversal, "citations": []model.Citation{model.Cite("v2/foreign-flow/"+q.Ticker+"/", q.Ticker, dates[len(dates)-1])}, "start": startOf(dates), "end": dates[len(dates)-1], }) } // startOf returns the first series date ("" when empty). func startOf(dates []string) string { if len(dates) == 0 { return "" } return dates[0] }