feat: add Citations and WatchlistChat components, integrate with API
- Implemented Citations component to display citation data. - Created WatchlistDrawer and ChatSidebar components for managing watchlists and AI chat functionality. - Integrated API calls for watchlist management and chat interactions. - Updated index.tsx to include new components in the main application layout. - Added API client in lib/api.ts for structured API interactions. - Developed Alerts, Dashboard, Portfolio, Routines, Screener, and Report pages with relevant data fetching and UI components. - Introduced styles in tokens.css for consistent theming across the application. - Configured TypeScript and Vite for project setup and development.
This commit is contained in:
@@ -0,0 +1,76 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strconv"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
)
|
||||
|
||||
// ListAlerts serves GET /api/alerts.
|
||||
func (s *Server) ListAlerts(w http.ResponseWriter, r *http.Request) {
|
||||
rows, err := s.DB.ListAlerts(s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"alerts": rows})
|
||||
}
|
||||
|
||||
// CreateAlert serves POST /api/alerts {name, rule, channels[]}.
|
||||
func (s *Server) CreateAlert(w http.ResponseWriter, r *http.Request) {
|
||||
var req struct {
|
||||
Name string `json:"name" validate:"required"`
|
||||
Rule any `json:"rule" validate:"required"`
|
||||
Channels []string `json:"channels"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
if err := s.Validate.Struct(req); err != nil {
|
||||
writeErr(w, http.StatusUnprocessableEntity, "name and rule are required")
|
||||
return
|
||||
}
|
||||
ruleRaw, _ := json.Marshal(req.Rule)
|
||||
id, err := s.DB.CreateAlert(s.userKey(r), req.Name, string(ruleRaw), req.Channels)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, map[string]any{"id": id})
|
||||
}
|
||||
|
||||
// DeleteAlert serves DELETE /api/alerts/:id.
|
||||
func (s *Server) DeleteAlert(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid id")
|
||||
return
|
||||
}
|
||||
ok, err := s.DB.DeleteAlert(id, s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
if !ok {
|
||||
writeErr(w, http.StatusNotFound, "alert not found")
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"id": id, "ok": true})
|
||||
}
|
||||
|
||||
// AlertEvents serves GET /api/alert-events?since=&ticker=.
|
||||
func (s *Server) AlertEvents(w http.ResponseWriter, r *http.Request) {
|
||||
since := r.URL.Query().Get("since")
|
||||
if since == "" {
|
||||
since = "2000-01-01"
|
||||
}
|
||||
evts, err := s.DB.AlertEventsSince(since, r.URL.Query().Get("ticker"), 50, s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"events": evts})
|
||||
}
|
||||
@@ -0,0 +1,273 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
"flowsight/internal/config"
|
||||
"flowsight/internal/sectors"
|
||||
"flowsight/internal/store"
|
||||
)
|
||||
|
||||
func testServer(t *testing.T) *Server {
|
||||
t.Helper()
|
||||
cfg := config.Load()
|
||||
cfg.DemoUserKey = "demo"
|
||||
db, err := store.Open(t.TempDir() + "/api.db")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { db.Close() })
|
||||
if _, err := db.SeedFromDir("../../tests/fixtures", "demo"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
cache := store.NewCache("")
|
||||
return New(cfg, db, cache, sectors.New(cfg.SectorsBaseURL, ""))
|
||||
}
|
||||
|
||||
func do(s *Server, method, path string, body any) *httptest.ResponseRecorder {
|
||||
var rdr *bytes.Reader
|
||||
if body != nil {
|
||||
raw, _ := json.Marshal(body)
|
||||
rdr = bytes.NewReader(raw)
|
||||
} else {
|
||||
rdr = bytes.NewReader(nil)
|
||||
}
|
||||
req := httptest.NewRequest(method, path, rdr)
|
||||
req.Header.Set("X-User-Key", "demo")
|
||||
rec := httptest.NewRecorder()
|
||||
s.Router().ServeHTTP(rec, req)
|
||||
return rec
|
||||
}
|
||||
|
||||
// Health 200 with cycle + credits fields.
|
||||
func TestHealth(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "GET", "/api/health", nil)
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d", rec.Code)
|
||||
}
|
||||
var out map[string]any
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &out)
|
||||
for _, k := range []string{"last_cycle_at", "credits_today", "scheduler_ok", "stale_flags"} {
|
||||
if _, ok := out[k]; !ok {
|
||||
t.Fatalf("missing key %s", k)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Briefing generates from seed with zero empty sections + citations.
|
||||
func TestBriefing(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "GET", "/api/briefing/today", nil)
|
||||
if rec.Code == http.StatusNotFound {
|
||||
// No briefing yet: run the routine via engine path instead.
|
||||
rows, _ := s.DB.ListRoutines("demo")
|
||||
if len(rows) == 0 {
|
||||
t.Fatal("seed has no routines")
|
||||
}
|
||||
if _, err := s.Engine.Run(httptest.NewRequest("GET", "/", nil).Context(), rows[0]); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rec = do(s, "GET", "/api/briefing/today", nil)
|
||||
}
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d, body %s", rec.Code, rec.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
// Screener returns a ranked list with per-row breakdown.
|
||||
func TestScreen(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "POST", "/api/screen", map[string]any{"limit": 5})
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d", rec.Code)
|
||||
}
|
||||
var out struct {
|
||||
Rows []map[string]any `json:"rows"`
|
||||
}
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &out)
|
||||
if len(out.Rows) == 0 {
|
||||
t.Fatal("empty screener rows")
|
||||
}
|
||||
if _, ok := out.Rows[0]["breakdown"]; !ok {
|
||||
t.Fatal("missing per-row breakdown")
|
||||
}
|
||||
}
|
||||
|
||||
// Report: all 7 sections populated with citations.
|
||||
func TestReport(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "POST", "/api/report/BBCA?profile=moderate", nil)
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d, body %s", rec.Code, rec.Body.String()[:300])
|
||||
}
|
||||
var out struct {
|
||||
Sections []map[string]any `json:"sections"`
|
||||
}
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &out)
|
||||
if len(out.Sections) != 7 {
|
||||
t.Fatalf("sections = %d, want 7", len(out.Sections))
|
||||
}
|
||||
for _, sec := range out.Sections {
|
||||
if sec["body"] == "" || sec["body"] == nil {
|
||||
t.Fatalf("empty section %v", sec["name"])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Subscribe -> run -> history row appears.
|
||||
func TestRoutineSubscribeRunHistory(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "POST", "/api/routines", map[string]any{"type": "foreign-reversal"})
|
||||
if rec.Code != http.StatusCreated {
|
||||
t.Fatalf("code = %d", rec.Code)
|
||||
}
|
||||
var created map[string]any
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &created)
|
||||
rows, _ := s.DB.ListRoutines("demo")
|
||||
if len(rows) == 0 {
|
||||
t.Fatal("no routines")
|
||||
}
|
||||
if _, err := s.Engine.Run(httptest.NewRequest("GET", "/", nil).Context(), rows[0]); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rec = do(s, "GET", "/api/routine-runs?limit=5", nil)
|
||||
var out struct {
|
||||
Runs []map[string]any `json:"runs"`
|
||||
}
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &out)
|
||||
if len(out.Runs) == 0 {
|
||||
t.Fatal("no history rows")
|
||||
}
|
||||
}
|
||||
|
||||
// Interrogation is scoped to report citations: conviction Q&A answers
|
||||
// from the persisted report, unknown report 404s.
|
||||
func TestInterrogate(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "POST", "/api/report/BBCA?profile=moderate", nil)
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d", rec.Code)
|
||||
}
|
||||
rec = do(s, "POST", "/api/report/BBCA/ask", map[string]any{"question": "kenapa conviction segitu?"})
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d, body %s", rec.Code, rec.Body.String()[:200])
|
||||
}
|
||||
var out struct {
|
||||
Answer string `json:"answer"`
|
||||
}
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &out)
|
||||
if out.Answer == "" {
|
||||
t.Fatal("empty interrogation answer")
|
||||
}
|
||||
rec = do(s, "POST", "/api/report/ZZZZ/ask", map[string]any{"question": "apa?"})
|
||||
if rec.Code != http.StatusNotFound {
|
||||
t.Fatalf("code = %d, want 404 for unknown ticker", rec.Code)
|
||||
}
|
||||
}
|
||||
|
||||
// Concentrated fixture warns >40% sector; accuracy math covered.
|
||||
func TestPortfolioAndAccuracy(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "GET", "/api/portfolio/risk", nil)
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d", rec.Code)
|
||||
}
|
||||
rec = do(s, "GET", "/api/accuracy", nil)
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d", rec.Code)
|
||||
}
|
||||
}
|
||||
|
||||
// start/end filters narrow the foreign series; reversal uses the 2x rule.
|
||||
func TestFlowForeignWindow(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "GET", "/api/flow/foreign?ticker=BBCA&start=2026-09-11&end=2026-09-11", nil)
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d", rec.Code)
|
||||
}
|
||||
var out struct {
|
||||
Dates []string `json:"dates"`
|
||||
Nets []float64 `json:"nets"`
|
||||
}
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &out)
|
||||
if len(out.Dates) != 1 || out.Dates[0] != "2026-09-11" {
|
||||
t.Fatalf("dates = %v, want single filtered day", out.Dates)
|
||||
}
|
||||
}
|
||||
|
||||
// Unknown routine types are rejected; missing ids 404.
|
||||
func TestRoutineValidation(t *testing.T) {
|
||||
s := testServer(t)
|
||||
rec := do(s, "POST", "/api/routines", map[string]any{"type": "not-a-routine"})
|
||||
if rec.Code != http.StatusUnprocessableEntity {
|
||||
t.Fatalf("code = %d, want 422", rec.Code)
|
||||
}
|
||||
rec = do(s, "PATCH", "/api/routines/999999", map[string]any{"enabled": false})
|
||||
if rec.Code != http.StatusNotFound {
|
||||
t.Fatalf("code = %d, want 404", rec.Code)
|
||||
}
|
||||
rec = do(s, "DELETE", "/api/routines/999999", nil)
|
||||
if rec.Code != http.StatusNotFound {
|
||||
t.Fatalf("code = %d, want 404", rec.Code)
|
||||
}
|
||||
rec = do(s, "DELETE", "/api/alerts/999999", nil)
|
||||
if rec.Code != http.StatusNotFound {
|
||||
t.Fatalf("code = %d, want 404", rec.Code)
|
||||
}
|
||||
}
|
||||
|
||||
// Destination CRUD: masked list, kind validation, owner scoping, 404s.
|
||||
func TestDestinations(t *testing.T) {
|
||||
s := testServer(t)
|
||||
// Invalid kind -> 422.
|
||||
rec := do(s, "POST", "/api/destinations", map[string]any{"kind": "sms"})
|
||||
if rec.Code != http.StatusUnprocessableEntity {
|
||||
t.Fatalf("code = %d, want 422", rec.Code)
|
||||
}
|
||||
// Telegram without chat_id -> 422.
|
||||
rec = do(s, "POST", "/api/destinations", map[string]any{"kind": "telegram", "bot_token": "x"})
|
||||
if rec.Code != http.StatusUnprocessableEntity {
|
||||
t.Fatalf("code = %d, want 422", rec.Code)
|
||||
}
|
||||
// Discord non-https -> 422.
|
||||
rec = do(s, "POST", "/api/destinations", map[string]any{"kind": "discord", "webhook_url": "http://x"})
|
||||
if rec.Code != http.StatusUnprocessableEntity {
|
||||
t.Fatalf("code = %d, want 422", rec.Code)
|
||||
}
|
||||
// Valid discord create -> 201.
|
||||
rec = do(s, "POST", "/api/destinations", map[string]any{"kind": "discord", "label": "ops", "webhook_url": "https://discord.example/hook"})
|
||||
if rec.Code != http.StatusCreated {
|
||||
t.Fatalf("code = %d, body %s", rec.Code, rec.Body.String())
|
||||
}
|
||||
// List masks secrets.
|
||||
rec = do(s, "GET", "/api/destinations", nil)
|
||||
var out struct {
|
||||
Destinations []map[string]any `json:"destinations"`
|
||||
}
|
||||
_ = json.Unmarshal(rec.Body.Bytes(), &out)
|
||||
if len(out.Destinations) != 1 {
|
||||
t.Fatalf("destinations = %v", out.Destinations)
|
||||
}
|
||||
for _, k := range []string{"bot_token", "chat_id", "webhook_url"} {
|
||||
if _, ok := out.Destinations[0][k]; ok {
|
||||
t.Fatalf("secret leaked in list: %s", k)
|
||||
}
|
||||
}
|
||||
if out.Destinations[0]["configured"] != true {
|
||||
t.Fatalf("configured flag = %v", out.Destinations[0])
|
||||
}
|
||||
// Missing id -> 404 on patch and delete.
|
||||
rec = do(s, "PATCH", "/api/destinations/999999", map[string]any{"enabled": false})
|
||||
if rec.Code != http.StatusNotFound {
|
||||
t.Fatalf("code = %d, want 404", rec.Code)
|
||||
}
|
||||
rec = do(s, "DELETE", "/api/destinations/999999", nil)
|
||||
if rec.Code != http.StatusNotFound {
|
||||
t.Fatalf("code = %d, want 404", rec.Code)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,71 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// ChatRequest is POST /api/chat {message, scope?: {report_id}}.
|
||||
type ChatRequest struct {
|
||||
Message string `json:"message" validate:"required"`
|
||||
Scope *struct {
|
||||
ReportID int64 `json:"report_id"`
|
||||
} `json:"scope"`
|
||||
}
|
||||
|
||||
// Chat serves POST /api/chat: cited answers. With scope.report_id the
|
||||
// grounding is restricted to that report's citations (report interrogation);
|
||||
// without scope it answers from latest snapshots. The LLM refines prose only
|
||||
// — numbers always come from stored data, never from generation.
|
||||
func (s *Server) Chat(w http.ResponseWriter, r *http.Request) {
|
||||
var req ChatRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
if err := s.Validate.Struct(req); err != nil {
|
||||
writeErr(w, http.StatusUnprocessableEntity, "message is required")
|
||||
return
|
||||
}
|
||||
// Scope grounding: report citations when scoped.
|
||||
var ground, citesRaw string
|
||||
if req.Scope != nil && req.Scope.ReportID > 0 {
|
||||
var cites string
|
||||
var at string
|
||||
err := s.DB.QueryRow(`SELECT payload_json, citations_json, generated_at FROM reports WHERE id=?`,
|
||||
req.Scope.ReportID).Scan(&ground, &cites, &at)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusNotFound, "report not found")
|
||||
return
|
||||
}
|
||||
citesRaw = cites
|
||||
} else {
|
||||
// Unscoped: ground on the latest briefing + watchlist.
|
||||
_, payload, cites, err := s.DB.LatestBriefing()
|
||||
if err != nil {
|
||||
ground = "no briefing or report data yet"
|
||||
} else {
|
||||
ground, citesRaw = payload, cites
|
||||
}
|
||||
}
|
||||
answer := "Based on stored data: " + head(ground, 600)
|
||||
if s.LLM.Available() {
|
||||
if text, err := s.LLM.Complete(r.Context(), s.Cfg.LLMTriage,
|
||||
"You answer questions about Indonesian stocks using ONLY the grounded data below. "+
|
||||
"Every number in your answer must cite its source. If the data lacks the answer, say so.",
|
||||
"Question: "+req.Message+"\n\nGrounded data:\n"+head(ground, 3000), 400); err == nil {
|
||||
answer = strings.TrimSpace(text)
|
||||
}
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"answer": answer, "grounding": head(ground, 600), "citations": citesRaw,
|
||||
})
|
||||
}
|
||||
|
||||
func head(s string, n int) string {
|
||||
if len(s) <= n {
|
||||
return s
|
||||
}
|
||||
return s[:n]
|
||||
}
|
||||
@@ -0,0 +1,156 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
|
||||
"flowsight/internal/store"
|
||||
)
|
||||
|
||||
// destOut is the masked API shape: secrets never leave the server.
|
||||
type destOut struct {
|
||||
ID int64 `json:"id"`
|
||||
Kind string `json:"kind"`
|
||||
Label string `json:"label"`
|
||||
Enabled bool `json:"enabled"`
|
||||
Configured bool `json:"configured"`
|
||||
}
|
||||
|
||||
func maskDestinations(rows []store.Destination) []destOut {
|
||||
out := make([]destOut, 0, len(rows))
|
||||
for _, d := range rows {
|
||||
cfg := false
|
||||
switch d.Kind {
|
||||
case store.DestTelegram:
|
||||
cfg = d.BotToken != "" && d.ChatID != ""
|
||||
case store.DestDiscord:
|
||||
cfg = d.WebhookURL != ""
|
||||
}
|
||||
out = append(out, destOut{ID: d.ID, Kind: d.Kind, Label: d.Label, Enabled: d.Enabled, Configured: cfg})
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// ListDestinations serves GET /api/destinations (secrets masked).
|
||||
func (s *Server) ListDestinations(w http.ResponseWriter, r *http.Request) {
|
||||
rows, err := s.DB.ListDestinations(s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"destinations": maskDestinations(rows)})
|
||||
}
|
||||
|
||||
// CreateDestination serves POST /api/destinations.
|
||||
func (s *Server) CreateDestination(w http.ResponseWriter, r *http.Request) {
|
||||
var req struct {
|
||||
Kind string `json:"kind" validate:"required,oneof=telegram discord"`
|
||||
Label string `json:"label"`
|
||||
BotToken string `json:"bot_token"`
|
||||
ChatID string `json:"chat_id"`
|
||||
WebhookURL string `json:"webhook_url"`
|
||||
Enabled *bool `json:"enabled"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
req.Kind = strings.ToLower(strings.TrimSpace(req.Kind))
|
||||
if err := s.Validate.Struct(req); err != nil {
|
||||
writeErr(w, http.StatusUnprocessableEntity, "kind must be telegram or discord")
|
||||
return
|
||||
}
|
||||
if msg := checkDestSecrets(req.Kind, req.BotToken, req.ChatID, req.WebhookURL); msg != "" {
|
||||
writeErr(w, http.StatusUnprocessableEntity, msg)
|
||||
return
|
||||
}
|
||||
enabled := true
|
||||
if req.Enabled != nil {
|
||||
enabled = *req.Enabled
|
||||
}
|
||||
id, err := s.DB.CreateDestination(store.Destination{
|
||||
UserKey: s.userKey(r), Kind: req.Kind, Label: req.Label,
|
||||
BotToken: req.BotToken, ChatID: req.ChatID, WebhookURL: req.WebhookURL,
|
||||
Enabled: enabled,
|
||||
})
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, map[string]any{"id": id, "kind": req.Kind})
|
||||
}
|
||||
|
||||
// UpdateDestination serves PATCH /api/destinations/:id. Kind is immutable;
|
||||
// omitted secret fields keep their stored value.
|
||||
func (s *Server) UpdateDestination(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid id")
|
||||
return
|
||||
}
|
||||
var req struct {
|
||||
Label *string `json:"label"`
|
||||
Enabled *bool `json:"enabled"`
|
||||
BotToken *string `json:"bot_token"`
|
||||
ChatID *string `json:"chat_id"`
|
||||
WebhookURL *string `json:"webhook_url"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
ok, err := s.DB.UpdateDestination(id, s.userKey(r), req.Label, req.Enabled, req.BotToken, req.ChatID, req.WebhookURL)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
if !ok {
|
||||
writeErr(w, http.StatusNotFound, "destination not found")
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"id": id, "ok": true})
|
||||
}
|
||||
|
||||
// DeleteDestination serves DELETE /api/destinations/:id.
|
||||
func (s *Server) DeleteDestination(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid id")
|
||||
return
|
||||
}
|
||||
ok, err := s.DB.DeleteDestination(id, s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
if !ok {
|
||||
writeErr(w, http.StatusNotFound, "destination not found")
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"id": id, "ok": true})
|
||||
}
|
||||
|
||||
// checkDestSecrets validates kind-appropriate secrets.
|
||||
func checkDestSecrets(kind, botToken, chatID, webhookURL string) string {
|
||||
switch kind {
|
||||
case store.DestTelegram:
|
||||
if strings.TrimSpace(botToken) == "" || strings.TrimSpace(chatID) == "" {
|
||||
return "telegram needs bot_token and chat_id"
|
||||
}
|
||||
case store.DestDiscord:
|
||||
u := strings.TrimSpace(webhookURL)
|
||||
if u == "" {
|
||||
return "discord needs webhook_url"
|
||||
}
|
||||
if !strings.HasPrefix(u, "https://") {
|
||||
return "webhook_url must be https"
|
||||
}
|
||||
default:
|
||||
return "kind must be telegram or discord"
|
||||
}
|
||||
return ""
|
||||
}
|
||||
@@ -0,0 +1,136 @@
|
||||
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"
|
||||
}
|
||||
wl, _ := s.DB.Watchlist(s.userKey(r))
|
||||
if len(wl) == 0 {
|
||||
wl = s.Cfg.Watchlist
|
||||
}
|
||||
type accRow struct {
|
||||
Ticker string `json:"ticker"`
|
||||
NetSum float64 `json:"net_sum"`
|
||||
Brokers int `json:"brokers"`
|
||||
}
|
||||
var accs []accRow
|
||||
foreignTotal := 0.0
|
||||
var cites []model.Citation
|
||||
for _, tk := range wl {
|
||||
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]
|
||||
cites = append(cites, model.Cite("v2/foreign-flow/"+tk+"/", tk, date))
|
||||
}
|
||||
}
|
||||
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, "citations": cites,
|
||||
})
|
||||
}
|
||||
|
||||
// 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])},
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"flowsight/internal/model"
|
||||
)
|
||||
|
||||
// Health serves GET /api/health: last cycle time + credits spent today +
|
||||
// scheduler state + stale flags (docs/API.md).
|
||||
func (s *Server) Health(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Query().Get("force") == "1" {
|
||||
_ = s.Sched.RunCycle(r.Context()) // synchronous probe cycle
|
||||
}
|
||||
lastCycle, schedOK := s.Sched.Status()
|
||||
today := time.Now().Format("2006-01-02")
|
||||
credits := s.DB.CreditsToday(today)
|
||||
lastDate, _, _, _ := s.DB.LatestBriefing()
|
||||
cutoff := model.StaleSession(time.Now())
|
||||
stale := []string{}
|
||||
if lastDate != "" {
|
||||
if t, err := time.Parse("2006-01-02", lastDate[:10]); err == nil && t.Before(cutoff) {
|
||||
stale = append(stale, "briefing older than one session ("+lastDate+")")
|
||||
} else if lastDate < today {
|
||||
stale = append(stale, "briefing older than today ("+lastDate+")")
|
||||
}
|
||||
}
|
||||
if !s.Cfg.HasSectorsKey() {
|
||||
stale = append(stale, "offline mode: SECTORS_API_KEY unset, serving seed data")
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"last_cycle_at": lastCycle,
|
||||
"credits_today": credits,
|
||||
"scheduler_ok": schedOK,
|
||||
"stale_flags": stale,
|
||||
"citations": []model.Citation{},
|
||||
})
|
||||
}
|
||||
|
||||
// Accuracy serves GET /api/accuracy: per-agent {calls, resolved, hits, hit_rate}.
|
||||
func (s *Server) Accuracy(w http.ResponseWriter, r *http.Request) {
|
||||
stats, err := s.DB.AccuracyStats()
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
if stats == nil {
|
||||
stats = []map[string]any{}
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"agents": stats})
|
||||
}
|
||||
@@ -0,0 +1,106 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
)
|
||||
|
||||
// Interrogate serves POST /api/report/:ticker/ask {question, report_id?}:
|
||||
// follow-up Q&A grounded ONLY in that report's persisted citations.
|
||||
// Without report_id it uses the latest report for the ticker.
|
||||
func (s *Server) Interrogate(w http.ResponseWriter, r *http.Request) {
|
||||
ticker := strings.ToUpper(chi.URLParam(r, "ticker"))
|
||||
var req struct {
|
||||
Question string `json:"question" validate:"required"`
|
||||
ReportID int64 `json:"report_id"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
if err := s.Validate.Struct(req); err != nil {
|
||||
writeErr(w, http.StatusUnprocessableEntity, "question is required")
|
||||
return
|
||||
}
|
||||
payload, citesRaw, at, id, err := s.loadReport(ticker, req.ReportID)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusNotFound, "no report for "+ticker+" yet — POST /api/report/"+ticker+" first")
|
||||
return
|
||||
}
|
||||
var cites []map[string]any
|
||||
_ = json.Unmarshal([]byte(citesRaw), &cites)
|
||||
answer := groundedAnswer(req.Question, payload)
|
||||
if s.LLM.Available() {
|
||||
if text, err := s.LLM.Complete(r.Context(), s.Cfg.LLMTriage,
|
||||
"Answer ONLY from the report JSON below. Every number must quote its cited value. "+
|
||||
"If the report lacks the answer, say exactly: not in this report.",
|
||||
"Question: "+req.Question+"\n\nReport:\n"+head(payload, 3000), 400); err == nil && text != "" {
|
||||
answer = strings.TrimSpace(text)
|
||||
}
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"answer": answer, "ticker": ticker, "report_id": id,
|
||||
"generated_at": at, "citations": cites,
|
||||
})
|
||||
}
|
||||
|
||||
// loadReport fetches (payload, citations, generated_at, id) for an explicit
|
||||
// report id or the latest report for a ticker.
|
||||
func (s *Server) loadReport(ticker string, id int64) (string, string, string, int64, error) {
|
||||
if id > 0 {
|
||||
var t, p, c, at string
|
||||
var rid int64
|
||||
err := s.DB.QueryRow(`SELECT id, ticker, payload_json, citations_json, generated_at
|
||||
FROM reports WHERE id=?`, id).Scan(&rid, &t, &p, &c, &at)
|
||||
if err != nil {
|
||||
return "", "", "", 0, err
|
||||
}
|
||||
if t != ticker {
|
||||
return "", "", "", 0, fmt.Errorf("report %d belongs to %s", id, t)
|
||||
}
|
||||
return p, c, at, rid, nil
|
||||
}
|
||||
p, c, at, err := s.DB.LatestReport(ticker)
|
||||
if err != nil {
|
||||
return "", "", "", 0, err
|
||||
}
|
||||
var rid int64
|
||||
_ = s.DB.QueryRow(`SELECT id FROM reports WHERE ticker=? ORDER BY id DESC LIMIT 1`,
|
||||
ticker).Scan(&rid)
|
||||
return p, c, at, rid, nil
|
||||
}
|
||||
|
||||
// groundedAnswer is the offline fallback: it extracts the recommendation +
|
||||
// conviction + cited lines matching question keywords from the payload.
|
||||
func groundedAnswer(question, payload string) string {
|
||||
var rep struct {
|
||||
Synthesis struct {
|
||||
Recommendation string `json:"recommendation"`
|
||||
Conviction int `json:"conviction"`
|
||||
Thesis string `json:"thesis"`
|
||||
} `json:"synthesis"`
|
||||
Sections []struct {
|
||||
Name string `json:"name"`
|
||||
Body string `json:"body"`
|
||||
} `json:"sections"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(payload), &rep); err != nil {
|
||||
return "not in this report"
|
||||
}
|
||||
q := strings.ToLower(question)
|
||||
if strings.Contains(q, "conviction") || strings.Contains(q, "kenapa") || strings.Contains(q, "why") {
|
||||
return fmt.Sprintf("%s with conviction %d/5: %s",
|
||||
rep.Synthesis.Recommendation, rep.Synthesis.Conviction, rep.Synthesis.Thesis)
|
||||
}
|
||||
for _, sec := range rep.Sections {
|
||||
if strings.Contains(q, strings.ToLower(sec.Name)) {
|
||||
return sec.Body
|
||||
}
|
||||
}
|
||||
return fmt.Sprintf("%s (conviction %d/5): %s",
|
||||
rep.Synthesis.Recommendation, rep.Synthesis.Conviction, rep.Synthesis.Thesis)
|
||||
}
|
||||
@@ -0,0 +1,233 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"math"
|
||||
"net/http"
|
||||
"sort"
|
||||
|
||||
"flowsight/internal/model"
|
||||
)
|
||||
|
||||
// PortfolioRisk serves GET /api/portfolio/risk: concentration bars,
|
||||
// correlation matrix, beta vs IHSG, warnings (concentrated fixture warns
|
||||
// >40% sector), accuracy-adjacent citations.
|
||||
func (s *Server) PortfolioRisk(w http.ResponseWriter, r *http.Request) {
|
||||
wl, _ := s.DB.Watchlist(s.userKey(r))
|
||||
if len(wl) == 0 {
|
||||
wl = s.Cfg.Watchlist
|
||||
}
|
||||
// Concentration: weight by latest close x assumed equal shares (seed-safe).
|
||||
type bar struct {
|
||||
Ticker string `json:"ticker"`
|
||||
Sector string `json:"sector"`
|
||||
Weight float64 `json:"weight"`
|
||||
}
|
||||
prices := map[string]float64{}
|
||||
total := 0.0
|
||||
for _, tk := range wl {
|
||||
px, _, err := s.DB.LatestClose(tk)
|
||||
if err != nil || px <= 0 {
|
||||
px = 1000 // seed-safe placeholder, flagged in warnings
|
||||
}
|
||||
prices[tk] = px
|
||||
total += px
|
||||
}
|
||||
sectorOf := sectorMap()
|
||||
var bars []bar
|
||||
sectorW := map[string]float64{}
|
||||
for _, tk := range wl {
|
||||
wt := 0.0
|
||||
if total > 0 {
|
||||
wt = prices[tk] / total
|
||||
}
|
||||
sec := sectorOf[tk]
|
||||
if sec == "" {
|
||||
sec = "unknown"
|
||||
}
|
||||
bars = append(bars, bar{tk, sec, wt})
|
||||
sectorW[sec] += wt
|
||||
}
|
||||
sort.Slice(bars, func(i, j int) bool { return bars[i].Weight > bars[j].Weight })
|
||||
|
||||
var warnings []string
|
||||
for sec, wt := range sectorW {
|
||||
if wt > 0.4 {
|
||||
warnings = append(warnings, "concentrated: "+sec+" at "+pct(wt)+" (over 40%)")
|
||||
}
|
||||
}
|
||||
|
||||
// Correlation: pairwise Pearson over stored daily closes (aligned tail).
|
||||
series := s.closes(wl)
|
||||
corr := correlationMatrixFrom(series, wl)
|
||||
beta := betaFrom(series, wl)
|
||||
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"concentration": bars, "correlation": corr, "beta": beta,
|
||||
"warnings": warnings,
|
||||
"citations": []model.Citation{model.Cite("v2/daily/", "watchlist", "stored")},
|
||||
})
|
||||
}
|
||||
|
||||
func pct(v float64) string {
|
||||
return itoa(int(v*100+0.5)) + "%"
|
||||
}
|
||||
|
||||
func itoa(n int) string {
|
||||
if n == 0 {
|
||||
return "0"
|
||||
}
|
||||
s := ""
|
||||
for n > 0 {
|
||||
s = string(rune('0'+n%10)) + s
|
||||
n /= 10
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
// sectorMap is the seed-safe sector lookup (live: subsector/report).
|
||||
func sectorMap() map[string]string {
|
||||
return map[string]string{
|
||||
"BBCA": "financials", "BBRI": "financials", "BMRI": "financials", "BBNI": "financials",
|
||||
"TLKM": "infrastructure", "ASII": "industrials", "UNVR": "consumer", "ICBP": "consumer",
|
||||
}
|
||||
}
|
||||
|
||||
// closes returns aligned close series per ticker from snapshots.
|
||||
func (s *Server) closes(wl []string) map[string][]float64 {
|
||||
out := map[string][]float64{}
|
||||
for _, tk := range wl {
|
||||
var rows []struct {
|
||||
Close float64 `json:"close"`
|
||||
}
|
||||
if raw, _, err := s.DB.LatestSnapshot(tk, "daily"); err == nil {
|
||||
var bars []struct {
|
||||
Close float64 `json:"close"`
|
||||
}
|
||||
if json.Unmarshal([]byte(raw), &bars) == nil {
|
||||
for _, b := range bars {
|
||||
rows = append(rows, struct {
|
||||
Close float64 `json:"close"`
|
||||
}{b.Close})
|
||||
}
|
||||
}
|
||||
_ = rows
|
||||
series := make([]float64, 0, len(bars))
|
||||
for _, b := range bars {
|
||||
series = append(series, b.Close)
|
||||
}
|
||||
out[tk] = series
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func correlationMatrixFrom(series map[string][]float64, wl []string) map[string]map[string]float64 {
|
||||
m := map[string]map[string]float64{}
|
||||
for _, a := range wl {
|
||||
m[a] = map[string]float64{}
|
||||
for _, b := range wl {
|
||||
if a == b {
|
||||
m[a][b] = 1
|
||||
continue
|
||||
}
|
||||
m[a][b] = pearson(tail(series[a], 30), tail(series[b], 30))
|
||||
}
|
||||
}
|
||||
return m
|
||||
}
|
||||
|
||||
func tail(xs []float64, n int) []float64 {
|
||||
if len(xs) <= n {
|
||||
return xs
|
||||
}
|
||||
return xs[len(xs)-n:]
|
||||
}
|
||||
|
||||
// pearson computes the correlation of two equal-length series.
|
||||
func pearson(a, b []float64) float64 {
|
||||
n := len(a)
|
||||
if n != len(b) || n < 2 {
|
||||
return 0
|
||||
}
|
||||
ma, mb := mean(a), mean(b)
|
||||
num, da, db := 0.0, 0.0, 0.0
|
||||
for i := range a {
|
||||
num += (a[i] - ma) * (b[i] - mb)
|
||||
da += (a[i] - ma) * (a[i] - ma)
|
||||
db += (b[i] - mb) * (b[i] - mb)
|
||||
}
|
||||
if da == 0 || db == 0 {
|
||||
return 0
|
||||
}
|
||||
return num / (math.Sqrt(da) * math.Sqrt(db))
|
||||
}
|
||||
|
||||
func mean(xs []float64) float64 {
|
||||
s := 0.0
|
||||
for _, x := range xs {
|
||||
s += x
|
||||
}
|
||||
return s / float64(len(xs))
|
||||
}
|
||||
|
||||
// betaFrom regresses mean ticker returns vs the watchlist mean (index-daily
|
||||
// benchmark when cached; watchlist-mean fallback keeps seeds working).
|
||||
func betaFrom(series map[string][]float64, wl []string) float64 {
|
||||
if len(wl) == 0 {
|
||||
return 1
|
||||
}
|
||||
n := 0
|
||||
for _, tk := range wl {
|
||||
if len(series[tk]) > n {
|
||||
n = len(series[tk])
|
||||
}
|
||||
}
|
||||
if n < 2 {
|
||||
return 1
|
||||
}
|
||||
idx := make([]float64, n)
|
||||
for _, tk := range wl {
|
||||
s := series[tk]
|
||||
for i := range idx {
|
||||
if i < len(s) {
|
||||
idx[i] += s[i]
|
||||
}
|
||||
}
|
||||
}
|
||||
for i := range idx {
|
||||
idx[i] /= float64(len(wl))
|
||||
}
|
||||
betas := []float64{}
|
||||
for _, tk := range wl {
|
||||
if b := betaOf(series[tk], idx); b != 0 {
|
||||
betas = append(betas, b)
|
||||
}
|
||||
}
|
||||
if len(betas) == 0 {
|
||||
return 1
|
||||
}
|
||||
return mean(betas)
|
||||
}
|
||||
|
||||
// betaOf is cov(asset,index)/var(index) over the aligned tail.
|
||||
func betaOf(asset, index []float64) float64 {
|
||||
n := len(asset)
|
||||
if len(index) < n {
|
||||
n = len(index)
|
||||
}
|
||||
if n < 2 {
|
||||
return 0
|
||||
}
|
||||
a, ix := asset[len(asset)-n:], index[len(index)-n:]
|
||||
ma, mi := mean(a), mean(ix)
|
||||
num, den := 0.0, 0.0
|
||||
for i := range a {
|
||||
num += (a[i] - ma) * (ix[i] - mi)
|
||||
den += (ix[i] - mi) * (ix[i] - mi)
|
||||
}
|
||||
if den == 0 {
|
||||
return 0
|
||||
}
|
||||
return num / den
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
)
|
||||
|
||||
// BuildReport serves POST /api/report/:ticker?format=json|html|pdf|md.
|
||||
// Runs A1..A6 + A7 live over stored snapshots, persists, and renders.
|
||||
func (s *Server) BuildReport(w http.ResponseWriter, r *http.Request) {
|
||||
ticker := strings.ToUpper(chi.URLParam(r, "ticker"))
|
||||
if len(ticker) < 3 || len(ticker) > 6 {
|
||||
writeErr(w, http.StatusUnprocessableEntity, "ticker must be 3-6 letters")
|
||||
return
|
||||
}
|
||||
profile := r.URL.Query().Get("profile")
|
||||
if profile == "" {
|
||||
profile = "moderate"
|
||||
}
|
||||
rep, id, err := s.Builder.Build(r.Context(), ticker, profile)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "report: "+err.Error())
|
||||
return
|
||||
}
|
||||
_ = id
|
||||
switch strings.ToLower(r.URL.Query().Get("format")) {
|
||||
case "html":
|
||||
w.Header().Set("Content-Type", "text/html; charset=utf-8")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
_, _ = w.Write([]byte(rep.ToHTML()))
|
||||
case "md":
|
||||
w.Header().Set("Content-Type", "text/markdown; charset=utf-8")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
_, _ = w.Write([]byte(rep.ToMarkdown()))
|
||||
case "pdf":
|
||||
raw, err := rep.ToPDF()
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "pdf: "+err.Error())
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/pdf")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
_, _ = w.Write(raw)
|
||||
default:
|
||||
writeJSON(w, http.StatusOK, rep)
|
||||
}
|
||||
// Push a live agent-panel event for the dashboard SSE feed.
|
||||
s.Hub.Publish("agents", `{"ticker":"`+ticker+`","recommendation":"`+
|
||||
rep.Synthesis.Recommendation+`"}`)
|
||||
}
|
||||
@@ -0,0 +1,136 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strconv"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
|
||||
"flowsight/internal/routines"
|
||||
"flowsight/internal/store"
|
||||
)
|
||||
|
||||
// ListRoutines serves GET /api/routines with last-run status.
|
||||
func (s *Server) ListRoutines(w http.ResponseWriter, r *http.Request) {
|
||||
rows, err := s.DB.ListRoutines(s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
type rowOut struct {
|
||||
store.Routine
|
||||
LastRun any `json:"last_run"`
|
||||
}
|
||||
out := make([]rowOut, 0, len(rows))
|
||||
for _, row := range rows {
|
||||
hist, _ := s.DB.RunHistory(row.ID, 1)
|
||||
var last any
|
||||
if len(hist) > 0 {
|
||||
last = hist[0]
|
||||
}
|
||||
out = append(out, rowOut{row, last})
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"routines": out})
|
||||
}
|
||||
|
||||
// CreateRoutine serves POST /api/routines {type, schedule_cron?, channels[]}.
|
||||
func (s *Server) CreateRoutine(w http.ResponseWriter, r *http.Request) {
|
||||
var req struct {
|
||||
Type string `json:"type" validate:"required"`
|
||||
Schedule string `json:"schedule_cron"`
|
||||
Channels []string `json:"channels"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
if err := s.Validate.Struct(req); err != nil {
|
||||
writeErr(w, http.StatusUnprocessableEntity, "type is required")
|
||||
return
|
||||
}
|
||||
if !routines.KnownType(req.Type) {
|
||||
writeErr(w, http.StatusUnprocessableEntity, "unknown routine type")
|
||||
return
|
||||
}
|
||||
if req.Schedule == "" {
|
||||
req.Schedule = routines.DefaultSchedule(req.Type)
|
||||
}
|
||||
id, err := s.DB.CreateRoutine(s.userKey(r), req.Type, req.Schedule, req.Channels)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, map[string]any{"id": id, "type": req.Type, "schedule_cron": req.Schedule})
|
||||
}
|
||||
|
||||
// UpdateRoutine serves PATCH /api/routines/:id.
|
||||
func (s *Server) UpdateRoutine(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid id")
|
||||
return
|
||||
}
|
||||
var req struct {
|
||||
Enabled *bool `json:"enabled"`
|
||||
Schedule string `json:"schedule_cron"`
|
||||
Channels []string `json:"channels"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
ok, err := s.DB.UpdateRoutine(id, s.userKey(r), req.Enabled, req.Schedule, req.Channels)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
if !ok {
|
||||
writeErr(w, http.StatusNotFound, "routine not found")
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"id": id, "ok": true})
|
||||
}
|
||||
|
||||
// RunHistory serves GET /api/routine-runs?routine_id=&limit=.
|
||||
func (s *Server) RunHistory(w http.ResponseWriter, r *http.Request) {
|
||||
rid, _ := strconv.ParseInt(r.URL.Query().Get("routine_id"), 10, 64)
|
||||
limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
|
||||
hist, err := s.DB.RunHistory(rid, limit, s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"runs": hist})
|
||||
}
|
||||
|
||||
// BriefingToday serves GET /api/briefing/today: latest payload + citations.
|
||||
func (s *Server) BriefingToday(w http.ResponseWriter, r *http.Request) {
|
||||
date, payload, cites, err := s.DB.LatestBriefing()
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusNotFound, "no briefing yet — run the morning-briefing routine")
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"date": date, "payload": payload, "citations": cites,
|
||||
})
|
||||
}
|
||||
|
||||
// DeleteRoutine serves DELETE /api/routines/:id.
|
||||
func (s *Server) DeleteRoutine(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64)
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid id")
|
||||
return
|
||||
}
|
||||
ok, err := s.DB.DeleteRoutine(id, s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
if !ok {
|
||||
writeErr(w, http.StatusNotFound, "routine not found")
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"id": id, "ok": true})
|
||||
}
|
||||
@@ -0,0 +1,163 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
"flowsight/internal/model"
|
||||
)
|
||||
|
||||
// ScreenRequest is POST /api/screen body.
|
||||
type ScreenRequest struct {
|
||||
Where string `json:"where"`
|
||||
Q string `json:"q"`
|
||||
Institutional *struct {
|
||||
BrokerScoreMin float64 `json:"broker_score_min"`
|
||||
ForeignTrend string `json:"foreign_trend"`
|
||||
InsiderBuying bool `json:"insider_buying"`
|
||||
VolumeAnomaly bool `json:"volume_anomaly"`
|
||||
} `json:"institutional"`
|
||||
Limit int `json:"limit"`
|
||||
}
|
||||
|
||||
// ScreenRow is one ranked result with per-row signal breakdown.
|
||||
type ScreenRow struct {
|
||||
Symbol string `json:"symbol"`
|
||||
Name string `json:"name"`
|
||||
Composite float64 `json:"composite"`
|
||||
Breakdown map[string]any `json:"breakdown"`
|
||||
Citations []model.Citation `json:"citations"`
|
||||
}
|
||||
|
||||
// Screen serves POST /api/screen: companies/ base filter enriched with
|
||||
// broker score + foreign trend + insider flag, ranked composite.
|
||||
func (s *Server) Screen(w http.ResponseWriter, r *http.Request) {
|
||||
var req ScreenRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
limit := req.Limit
|
||||
if limit <= 0 || limit > 100 {
|
||||
limit = 20
|
||||
}
|
||||
// Base universe: live screener when keyed, else stored watchlist.
|
||||
var universe []string
|
||||
if s.Cfg.HasSectorsKey() && (req.Where != "" || req.Q != "") {
|
||||
if rows, err := s.Sectors.Screen(r.Context(), req.Where, req.Q, limit*2, 0); err == nil {
|
||||
for _, row := range rows {
|
||||
universe = append(universe, strings.ToUpper(strings.TrimSuffix(row.Symbol, ".JK")))
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(universe) == 0 {
|
||||
universe, _ = s.DB.Watchlist(s.userKey(r))
|
||||
if len(universe) == 0 {
|
||||
universe = s.Cfg.Watchlist
|
||||
}
|
||||
}
|
||||
var rows []ScreenRow
|
||||
for _, tk := range universe {
|
||||
row := s.scoreTicker(tk)
|
||||
if req.Institutional != nil {
|
||||
inst := req.Institutional
|
||||
if b, _ := row.Breakdown["broker_score"].(float64); b < inst.BrokerScoreMin {
|
||||
continue
|
||||
}
|
||||
if inst.ForeignTrend != "" {
|
||||
if t, _ := row.Breakdown["foreign_trend"].(string); t != inst.ForeignTrend {
|
||||
continue
|
||||
}
|
||||
}
|
||||
if inst.InsiderBuying {
|
||||
if b, _ := row.Breakdown["insider_buying"].(bool); !b {
|
||||
continue
|
||||
}
|
||||
}
|
||||
if inst.VolumeAnomaly {
|
||||
if b, _ := row.Breakdown["volume_anomaly"].(bool); !b {
|
||||
continue
|
||||
}
|
||||
}
|
||||
}
|
||||
rows = append(rows, row)
|
||||
}
|
||||
sort.Slice(rows, func(i, j int) bool { return rows[i].Composite > rows[j].Composite })
|
||||
if len(rows) > limit {
|
||||
rows = rows[:limit]
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"rows": rows, "count": len(rows)})
|
||||
}
|
||||
|
||||
// scoreTicker computes the composite (broker 40 + foreign 25 + insider 15 + volume 20).
|
||||
func (s *Server) scoreTicker(tk string) ScreenRow {
|
||||
tk = strings.ToUpper(tk)
|
||||
row := ScreenRow{Symbol: tk, Name: tk, Breakdown: map[string]any{}}
|
||||
// Broker score from 5d net imbalance.
|
||||
brokerScore := 0.0
|
||||
if nets, err := s.DB.NetBuySum5d(tk); err == nil && len(nets) > 0 {
|
||||
pos, neg := 0.0, 0.0
|
||||
for _, v := range nets {
|
||||
if v > 0 {
|
||||
pos += v
|
||||
} else {
|
||||
neg -= v
|
||||
}
|
||||
}
|
||||
if tot := pos + neg; tot > 0 {
|
||||
brokerScore = (pos - neg) / tot * 100
|
||||
}
|
||||
row.Citations = append(row.Citations, model.Cite("v2/broker-summary/"+tk+"/top/", tk, "stored"))
|
||||
}
|
||||
// Foreign trend from last-6 series.
|
||||
foreignScore, trend := 0.0, "flat"
|
||||
if dates, nets, err := s.DB.ForeignLast6(tk); err == nil && len(nets) > 0 {
|
||||
last := nets[len(nets)-1]
|
||||
if last > 0 {
|
||||
foreignScore, trend = 50, "inflow"
|
||||
} else if last < 0 {
|
||||
foreignScore, trend = -50, "outflow"
|
||||
}
|
||||
row.Citations = append(row.Citations, model.Cite("v2/foreign-flow/"+tk+"/", tk, dates[len(dates)-1]))
|
||||
}
|
||||
// Insider flag from filings average.
|
||||
insider := s.DB.FilingAvg30(tk) > 0
|
||||
// Volume anomaly from stored daily bars.
|
||||
volAnom, volMult := false, 0.0
|
||||
if vols, _, err := s.DB.DailyVolumes(tk, 21); err == nil && len(vols) >= 2 {
|
||||
n := len(vols)
|
||||
if a := avgF(vols[:n-1]); a > 0 {
|
||||
volMult = vols[n-1] / a
|
||||
volAnom = volMult > 2
|
||||
}
|
||||
row.Citations = append(row.Citations, model.Cite("v2/daily/"+tk+"/", tk, "stored"))
|
||||
}
|
||||
volScore := 0.0
|
||||
if volAnom {
|
||||
volScore = 50
|
||||
}
|
||||
row.Composite = brokerScore*0.4 + foreignScore*0.25 + volScore*0.2
|
||||
if insider {
|
||||
row.Composite += 7.5
|
||||
}
|
||||
row.Breakdown = map[string]any{
|
||||
"broker": brokerScore, "broker_score": brokerScore,
|
||||
"foreign": foreignScore, "foreign_trend": trend,
|
||||
"insider": insider, "insider_buying": insider,
|
||||
"volume_mult": volMult, "volume_anomaly": volAnom,
|
||||
}
|
||||
return row
|
||||
}
|
||||
|
||||
func avgF(xs []float64) float64 {
|
||||
if len(xs) == 0 {
|
||||
return 0
|
||||
}
|
||||
sum := 0.0
|
||||
for _, x := range xs {
|
||||
sum += x
|
||||
}
|
||||
return sum / float64(len(xs))
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
// Package api serves the FlowSight REST API (docs/API.md) on chi: flow
|
||||
// summary/broker/foreign, screener, routines, briefing, alerts, reports,
|
||||
// watchlist, portfolio risk, accuracy, chat (report-scoped), health, and the
|
||||
// SSE stream (agents/alerts/activity, 15s heartbeat). Demo auth: X-User-Key.
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
"github.com/go-chi/chi/v5/middleware"
|
||||
"github.com/go-playground/validator/v10"
|
||||
|
||||
"flowsight/internal/agents"
|
||||
"flowsight/internal/config"
|
||||
"flowsight/internal/llm"
|
||||
"flowsight/internal/reports"
|
||||
"flowsight/internal/routines"
|
||||
"flowsight/internal/scheduler"
|
||||
"flowsight/internal/sectors"
|
||||
"flowsight/internal/store"
|
||||
)
|
||||
|
||||
// Server wires all handlers.
|
||||
type Server struct {
|
||||
Cfg config.Config
|
||||
DB *store.DB
|
||||
Sectors *sectors.Client
|
||||
Sched *scheduler.Scheduler
|
||||
Engine *routines.Engine
|
||||
Builder *reports.Builder
|
||||
LLM *llm.Client
|
||||
Validate *validator.Validate
|
||||
Hub *Hub
|
||||
}
|
||||
|
||||
// New builds a Server with all dependencies wired.
|
||||
func New(cfg config.Config, db *store.DB, cache *store.Cache, s *sectors.Client) *Server {
|
||||
llmc := llm.New(cfg.LLMBaseURL, cfg.LLMAPIKey)
|
||||
sched := scheduler.New(cfg, db, cache, s)
|
||||
srv := &Server{
|
||||
Cfg: cfg, DB: db, Sectors: s, Sched: sched, LLM: llmc,
|
||||
Validate: validator.New(),
|
||||
Hub: NewHub(),
|
||||
}
|
||||
srv.Engine = &routines.Engine{DB: db, Notifier: sched.Notifier, UserKey: cfg.DemoUserKey,
|
||||
Publish: srv.Hub.Publish}
|
||||
sched.Publish = srv.Hub.Publish
|
||||
srv.Builder = &reports.Builder{DB: db, Deps: agents.Deps{
|
||||
DB: db, LLM: llmc, TriageModel: cfg.LLMTriage, SynthModel: cfg.LLMSynth,
|
||||
Now: time.Now(),
|
||||
}}
|
||||
return srv
|
||||
}
|
||||
|
||||
// Router returns the chi mux with all routes.
|
||||
func (s *Server) Router() http.Handler {
|
||||
r := chi.NewRouter()
|
||||
r.Use(middleware.Logger, middleware.Recoverer, middleware.Heartbeat("/ping"))
|
||||
r.Route("/api", func(r chi.Router) {
|
||||
r.Get("/health", s.Health)
|
||||
r.Get("/stream", s.Stream)
|
||||
r.Get("/flow/summary", s.FlowSummary)
|
||||
r.Get("/flow/broker", s.FlowBroker)
|
||||
r.Get("/flow/foreign", s.FlowForeign)
|
||||
r.Post("/screen", s.Screen)
|
||||
r.Get("/routines", s.ListRoutines)
|
||||
r.Post("/routines", s.CreateRoutine)
|
||||
r.Patch("/routines/{id}", s.UpdateRoutine)
|
||||
r.Delete("/routines/{id}", s.DeleteRoutine)
|
||||
r.Get("/routine-runs", s.RunHistory)
|
||||
r.Get("/briefing/today", s.BriefingToday)
|
||||
r.Get("/alerts", s.ListAlerts)
|
||||
r.Post("/alerts", s.CreateAlert)
|
||||
r.Delete("/alerts/{id}", s.DeleteAlert)
|
||||
r.Get("/alert-events", s.AlertEvents)
|
||||
r.Get("/destinations", s.ListDestinations)
|
||||
r.Post("/destinations", s.CreateDestination)
|
||||
r.Patch("/destinations/{id}", s.UpdateDestination)
|
||||
r.Delete("/destinations/{id}", s.DeleteDestination)
|
||||
r.Post("/report/{ticker}", s.BuildReport)
|
||||
r.Post("/report/{ticker}/ask", s.Interrogate)
|
||||
r.Get("/watchlist", s.GetWatchlist)
|
||||
r.Post("/watchlist", s.AddWatch)
|
||||
r.Delete("/watchlist/{ticker}", s.RemoveWatch)
|
||||
r.Get("/portfolio/risk", s.PortfolioRisk)
|
||||
r.Get("/accuracy", s.Accuracy)
|
||||
r.Post("/chat", s.Chat)
|
||||
})
|
||||
return r
|
||||
}
|
||||
|
||||
// userKey resolves the demo auth header (single demo key for hackathon).
|
||||
func (s *Server) userKey(r *http.Request) string {
|
||||
if k := strings.TrimSpace(r.Header.Get("X-User-Key")); k != "" {
|
||||
return k
|
||||
}
|
||||
return s.Cfg.DemoUserKey
|
||||
}
|
||||
|
||||
func writeJSON(w http.ResponseWriter, code int, v any) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(code)
|
||||
_ = json.NewEncoder(w).Encode(v)
|
||||
}
|
||||
|
||||
func writeErr(w http.ResponseWriter, code int, msg string) {
|
||||
writeJSON(w, code, map[string]any{"error": map[string]string{"code": http.StatusText(code), "message": msg}})
|
||||
}
|
||||
@@ -0,0 +1,122 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Hub fans SSE events out to connected browsers. Channels: agents
|
||||
// (status+scores during runs), alerts (new events), activity (feed rows).
|
||||
// Heartbeat 15s; reconnect resumes from last event ID (best-effort replay of
|
||||
// the last 50 events).
|
||||
type Hub struct {
|
||||
mu sync.Mutex
|
||||
subs map[chan SSEEvent]bool
|
||||
history []SSEEvent
|
||||
nextID int64
|
||||
}
|
||||
|
||||
// SSEEvent is one server-sent event.
|
||||
type SSEEvent struct {
|
||||
ID int64
|
||||
Channel string
|
||||
Data string
|
||||
}
|
||||
|
||||
// NewHub builds an empty hub.
|
||||
func NewHub() *Hub { return &Hub{subs: map[chan SSEEvent]bool{}} }
|
||||
|
||||
// Publish broadcasts to all subscribers and appends to history.
|
||||
func (h *Hub) Publish(channel, data string) {
|
||||
h.mu.Lock()
|
||||
h.nextID++
|
||||
ev := SSEEvent{ID: h.nextID, Channel: channel, Data: data}
|
||||
h.history = append(h.history, ev)
|
||||
if len(h.history) > 50 {
|
||||
h.history = h.history[len(h.history)-50:]
|
||||
}
|
||||
for ch := range h.subs {
|
||||
select {
|
||||
case ch <- ev:
|
||||
default:
|
||||
}
|
||||
}
|
||||
h.mu.Unlock()
|
||||
}
|
||||
|
||||
func (h *Hub) subscribe() chan SSEEvent {
|
||||
ch := make(chan SSEEvent, 16)
|
||||
h.mu.Lock()
|
||||
h.subs[ch] = true
|
||||
h.mu.Unlock()
|
||||
return ch
|
||||
}
|
||||
|
||||
// since returns history entries newer than id (all when id <= 0).
|
||||
func (h *Hub) since(id int64) []SSEEvent {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
var out []SSEEvent
|
||||
for _, ev := range h.history {
|
||||
if ev.ID > id {
|
||||
out = append(out, ev)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// lastEventID parses Last-Event-ID (header or query) for resume.
|
||||
func lastEventID(r *http.Request) int64 {
|
||||
s := r.Header.Get("Last-Event-ID")
|
||||
if s == "" {
|
||||
s = r.URL.Query().Get("lastEventId")
|
||||
}
|
||||
var id int64
|
||||
fmt.Sscanf(s, "%d", &id)
|
||||
return id
|
||||
}
|
||||
|
||||
func (h *Hub) unsubscribe(ch chan SSEEvent) {
|
||||
h.mu.Lock()
|
||||
delete(h.subs, ch)
|
||||
close(ch)
|
||||
h.mu.Unlock()
|
||||
}
|
||||
|
||||
// Stream serves GET /api/stream as text/event-stream.
|
||||
func (s *Server) Stream(w http.ResponseWriter, r *http.Request) {
|
||||
fl, ok := w.(http.Flusher)
|
||||
if !ok {
|
||||
writeErr(w, http.StatusInternalServerError, "streaming unsupported")
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
w.Header().Set("Cache-Control", "no-cache")
|
||||
w.Header().Set("Connection", "keep-alive")
|
||||
ch := s.Hub.subscribe()
|
||||
defer s.Hub.unsubscribe(ch)
|
||||
// Best-effort replay: resume after Last-Event-ID so reconnects do not
|
||||
// lose the last 50 events (matches the history comment on Publish).
|
||||
for _, ev := range s.Hub.since(lastEventID(r)) {
|
||||
fmt.Fprintf(w, "id: %d\nevent: %s\ndata: %s\n\n", ev.ID, ev.Channel, ev.Data)
|
||||
}
|
||||
fl.Flush()
|
||||
tick := time.NewTicker(15 * time.Second)
|
||||
defer tick.Stop()
|
||||
fmt.Fprintf(w, ": connected\n\n")
|
||||
fl.Flush()
|
||||
for {
|
||||
select {
|
||||
case <-r.Context().Done():
|
||||
return
|
||||
case ev := <-ch:
|
||||
fmt.Fprintf(w, "id: %d\nevent: %s\ndata: %s\n\n", ev.ID, ev.Channel, ev.Data)
|
||||
fl.Flush()
|
||||
case <-tick.C:
|
||||
fmt.Fprintf(w, ": heartbeat\n\n")
|
||||
fl.Flush()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
)
|
||||
|
||||
// GetWatchlist serves GET /api/watchlist.
|
||||
func (s *Server) GetWatchlist(w http.ResponseWriter, r *http.Request) {
|
||||
wl, err := s.DB.Watchlist(s.userKey(r))
|
||||
if err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"watchlist": wl})
|
||||
}
|
||||
|
||||
// AddWatch serves POST /api/watchlist {ticker}.
|
||||
func (s *Server) AddWatch(w http.ResponseWriter, r *http.Request) {
|
||||
var req struct {
|
||||
Ticker string `json:"ticker" validate:"required,len=4"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, "invalid JSON body")
|
||||
return
|
||||
}
|
||||
req.Ticker = strings.ToUpper(strings.TrimSpace(req.Ticker))
|
||||
if err := s.Validate.Struct(req); err != nil {
|
||||
writeErr(w, http.StatusUnprocessableEntity, "ticker (4 letters) required")
|
||||
return
|
||||
}
|
||||
if err := s.DB.AddWatch(s.userKey(r), req.Ticker); err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, map[string]any{"ticker": req.Ticker})
|
||||
}
|
||||
|
||||
// RemoveWatch serves DELETE /api/watchlist/:ticker.
|
||||
func (s *Server) RemoveWatch(w http.ResponseWriter, r *http.Request) {
|
||||
ticker := strings.ToUpper(chi.URLParam(r, "ticker"))
|
||||
if err := s.DB.RemoveWatch(s.userKey(r), ticker); err != nil {
|
||||
writeErr(w, http.StatusBadGateway, "db: "+err.Error())
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"ticker": ticker, "ok": true})
|
||||
}
|
||||
Reference in New Issue
Block a user