ci / go (push) Waiting to run
ci / go-db (agent) (push) Waiting to run
ci / go-db (config) (push) Waiting to run
ci / go-db (db) (push) Waiting to run
ci / go-db (evidence) (push) Waiting to run
ci / go-db (llmrec) (push) Waiting to run
ci / go-db (server) (push) Waiting to run
detections / detections (push) Waiting to run
web / web (push) Waiting to run
docs / links (push) Canceled after 0s
2022 lines
74 KiB
Go
2022 lines
74 KiB
Go
package db
|
|
|
|
import (
|
|
"bytes"
|
|
"database/sql"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// ErrIntentStateConflict marks a failed compare-and-set transition. Callers use
|
|
// it to distinguish an expected control race (HTTP 409) from a storage failure.
|
|
var ErrIntentStateConflict = errors.New("intent state changed concurrently")
|
|
|
|
// utf8Clean makes a string safe for a PostgreSQL text column. It (1) replaces
|
|
// invalid UTF-8 byte sequences with U+FFFD and (2) strips NUL (0x00) bytes.
|
|
// Tool output (nmap/curl/… stdout) can carry raw/truncated bytes that PostgreSQL's
|
|
// UTF8 encoding rejects ("invalid byte sequence for encoding UTF8"); without this
|
|
// the INSERT fails and the activity record is silently lost. NUL is *valid* UTF-8
|
|
// (U+0000) so ToValidUTF8 leaves it in place, yet PostgreSQL text still rejects it
|
|
// (SQLSTATE 22021) — so it must be removed separately. JSONB columns need the
|
|
// separate jsonbClean below: json.Marshal encodes a NUL as the escape backslash-u-0000,
|
|
// which the text-type json accepts but jsonb rejects (SQLSTATE 22P05).
|
|
func utf8Clean(s string) string {
|
|
if strings.IndexByte(s, 0) >= 0 {
|
|
s = strings.ReplaceAll(s, "\x00", "")
|
|
}
|
|
return strings.ToValidUTF8(s, "�")
|
|
}
|
|
|
|
// jsonbClean makes marshaled JSON safe for a PostgreSQL jsonb column. json.Marshal
|
|
// faithfully encodes a NUL byte (U+0000) as the 6-byte escape sequence \u0000; the
|
|
// json type stores it, but jsonb rejects it with "unsupported Unicode escape
|
|
// sequence" (SQLSTATE 22P05). Captured HTTP/tool bytes can carry NULs, so drop the
|
|
// escape — this also covers NULs nested inside json.RawMessage fields, which are
|
|
// copied verbatim into the output. Only a *real* escape is stripped: a \u0000 is
|
|
// genuine when preceded by an even number of backslashes, so an escaped-backslash
|
|
// run such as \\u0000 (the literal text u0000) is left intact.
|
|
func jsonbClean(b []byte) []byte {
|
|
if !bytes.Contains(b, []byte("\\u0000")) {
|
|
return b
|
|
}
|
|
out := make([]byte, 0, len(b))
|
|
bs := 0 // consecutive backslashes already emitted before position i
|
|
for i := 0; i < len(b); i++ {
|
|
if b[i] == '\\' && bs%2 == 0 && i+5 < len(b) &&
|
|
b[i+1] == 'u' && b[i+2] == '0' && b[i+3] == '0' && b[i+4] == '0' && b[i+5] == '0' {
|
|
i += 5 // skip the whole \u0000
|
|
bs = 0
|
|
continue
|
|
}
|
|
if b[i] == '\\' {
|
|
bs++
|
|
} else {
|
|
bs = 0
|
|
}
|
|
out = append(out, b[i])
|
|
}
|
|
return out
|
|
}
|
|
|
|
// Node is a typed reasoning node (= old task_nodes). kind ∈ goal|intent|finding|hint.
|
|
type Node struct {
|
|
FindingID int64 `json:"finding_id,omitempty"` // populated by finding-aware reads; never inferred from ID
|
|
FindingNodeID int64 `json:"finding_node_id,omitempty"`
|
|
TrafficCount int `json:"traffic_count,omitempty"`
|
|
ID int64 `json:"id"`
|
|
Kind string `json:"kind"`
|
|
Payload json.RawMessage `json:"payload"`
|
|
Priority int `json:"priority"`
|
|
State string `json:"state"`
|
|
Origin string `json:"origin,omitempty"`
|
|
Owner string `json:"owner,omitempty"`
|
|
BlockedReason string `json:"blocked_reason,omitempty"`
|
|
DeleteReason string `json:"delete_reason,omitempty"` // 仅意图假删除(state='deleted')时非空
|
|
Anchors []int64 `json:"anchors,omitempty"`
|
|
CreatedAt time.Time `json:"created_at"`
|
|
SourceTaskID int64 `json:"source_task_id,omitempty"`
|
|
Inherited bool `json:"inherited,omitempty"`
|
|
}
|
|
|
|
// Activity is one worker execution step (= old task_activity).
|
|
type Activity struct {
|
|
ID int64 `json:"id"`
|
|
NodeID *int64 `json:"node_id,omitempty"`
|
|
Worker string `json:"worker,omitempty"`
|
|
Kind string `json:"kind,omitempty"`
|
|
Tool string `json:"tool,omitempty"`
|
|
ToolUseID string `json:"tool_use_id,omitempty"`
|
|
IsError bool `json:"is_error"`
|
|
Summary string `json:"summary,omitempty"`
|
|
Detail string `json:"-"` // full blob; lazy-loaded, omitted from list payloads
|
|
Metadata json.RawMessage `json:"metadata,omitempty"`
|
|
CreatedAt time.Time `json:"created_at"`
|
|
// Token usage — set only on the terminal kind='result' record; nil elsewhere.
|
|
InputTokens *int `json:"input_tokens,omitempty"`
|
|
OutputTokens *int `json:"output_tokens,omitempty"`
|
|
CacheReadTokens *int `json:"cache_read_tokens,omitempty"`
|
|
CacheWriteTokens *int `json:"cache_write_tokens,omitempty"`
|
|
SourceTaskID int64 `json:"source_task_id,omitempty"`
|
|
Inherited bool `json:"inherited,omitempty"`
|
|
// MainSeg is the main-agent conversation segment (nil for non-mainagent rows;
|
|
// nil/0 == the original session). Lets the UI route a mainagent row to its segment.
|
|
MainSeg *int `json:"main_seg,omitempty"`
|
|
}
|
|
|
|
// TokenUsage is a per-worker token aggregate (TokenStatsByWorker).
|
|
type TokenUsage struct {
|
|
Worker string `json:"worker"`
|
|
InputTokens int `json:"input_tokens"`
|
|
OutputTokens int `json:"output_tokens"`
|
|
CacheReadTokens int `json:"cache_read_tokens"`
|
|
CacheWriteTokens int `json:"cache_write_tokens"`
|
|
}
|
|
|
|
// SessionTokenUsage is the authoritative completed-run total for one task UI
|
|
// session. Worker sessions are keyed by intent rather than the reusable work#N
|
|
// executor name, so retries and reassignment remain attached to the same row.
|
|
type SessionTokenUsage struct {
|
|
Session string `json:"session"`
|
|
InputTokens int `json:"input_tokens"`
|
|
OutputTokens int `json:"output_tokens"`
|
|
CacheReadTokens int `json:"cache_read_tokens"`
|
|
CacheWriteTokens int `json:"cache_write_tokens"`
|
|
}
|
|
|
|
// DailyTokenBucket is one day's global token aggregate across all tasks.
|
|
type DailyTokenBucket struct {
|
|
Day string `json:"day"` // "YYYY-MM-DD"
|
|
InputTokens int `json:"input_tokens"`
|
|
OutputTokens int `json:"output_tokens"`
|
|
CacheReadTokens int `json:"cache_read_tokens"`
|
|
}
|
|
|
|
// TokenDailyAll aggregates token consumption by calendar day (UTC) for the
|
|
// past `days` days across all explorations. Only kind='result' rows carry token
|
|
// counts (the terminal per-run summary), so there is no double-counting.
|
|
func (d *DB) TokenDailyAll(days int) ([]DailyTokenBucket, error) {
|
|
if days <= 0 {
|
|
days = 30
|
|
}
|
|
rows, err := d.Query(`
|
|
SELECT TO_CHAR(DATE_TRUNC('day', created_at AT TIME ZONE 'UTC'), 'YYYY-MM-DD') AS day,
|
|
COALESCE(SUM(input_tokens), 0),
|
|
COALESCE(SUM(output_tokens), 0),
|
|
COALESCE(SUM(cache_read_tokens), 0)
|
|
FROM activity
|
|
WHERE kind = 'result'
|
|
AND created_at >= NOW() - ($1 * INTERVAL '1 day')
|
|
GROUP BY day
|
|
ORDER BY day`, days)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []DailyTokenBucket
|
|
for rows.Next() {
|
|
var b DailyTokenBucket
|
|
if err := rows.Scan(&b.Day, &b.InputTokens, &b.OutputTokens, &b.CacheReadTokens); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, b)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// ExplorationStore is a handle onto one exploration's reasoning graph + activity.
|
|
type ExplorationStore struct {
|
|
db *DB
|
|
expID int64
|
|
}
|
|
|
|
// CreateExploration creates a new exploration root and returns its id.
|
|
func (d *DB) CreateExploration(description, goal string) (int64, error) {
|
|
var id int64
|
|
err := d.QueryRow(`INSERT INTO explorations(description, goal) VALUES ($1, $2) RETURNING id`, description, goal).Scan(&id)
|
|
return id, err
|
|
}
|
|
|
|
// Exploration returns a handle bound to an exploration id.
|
|
func (d *DB) Exploration(id int64) *ExplorationStore { return &ExplorationStore{db: d, expID: id} }
|
|
|
|
func (s *ExplorationStore) ID() int64 { return s.expID }
|
|
|
|
// Root returns the exploration's description and goal.
|
|
func (s *ExplorationStore) Root() (description, goal string, err error) {
|
|
var d sql.NullString
|
|
err = s.db.QueryRow(`SELECT description, goal FROM explorations WHERE id=$1`, s.expID).Scan(&d, &goal)
|
|
return d.String, goal, err
|
|
}
|
|
|
|
// AddNode writes a typed reasoning node with optional anchors to asset ids.
|
|
func (s *ExplorationStore) AddNode(kind string, payload map[string]any, priority int, state, origin string, anchors []int64) (int64, error) {
|
|
raw, _ := json.Marshal(payload)
|
|
tx, err := s.db.Begin()
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
defer tx.Rollback()
|
|
var id int64
|
|
if err := tx.QueryRow(`
|
|
INSERT INTO exploration_nodes(exploration_id, kind, payload, priority, state, origin)
|
|
VALUES ($1, $2, $3, $4, $5, $6) RETURNING id`,
|
|
s.expID, kind, string(raw), priority, state, origin).Scan(&id); err != nil {
|
|
return 0, err
|
|
}
|
|
for _, a := range anchors {
|
|
if _, err := tx.Exec(`INSERT INTO exploration_anchors(node_id, asset_id) VALUES ($1,$2) ON CONFLICT DO NOTHING`, id, a); err != nil {
|
|
return 0, err
|
|
}
|
|
}
|
|
return id, tx.Commit()
|
|
}
|
|
|
|
// AddIntent is a convenience: an open intent.
|
|
func (s *ExplorationStore) AddIntent(payload map[string]any, priority int, anchors []int64, origin string) (int64, error) {
|
|
if origin == "" {
|
|
origin = "planner"
|
|
}
|
|
return s.AddNode("intent", payload, priority, "open", origin, anchors)
|
|
}
|
|
|
|
// AddGoal writes a goal node (state open).
|
|
func (s *ExplorationStore) AddGoal(payload map[string]any, origin string) (int64, error) {
|
|
return s.AddNode("goal", payload, 0, "open", origin, nil)
|
|
}
|
|
|
|
// OriginFactID returns this exploration's root fact id — the KindFact node with
|
|
// state='origin' seeded at task creation (the task root). Every intent traces
|
|
// back to it. Returns 0 (no error) when none exists (legacy explorations created
|
|
// under the old 'begin' root, or none yet).
|
|
func (s *ExplorationStore) OriginFactID() (int64, error) {
|
|
var id int64
|
|
err := s.db.QueryRow(`SELECT id FROM exploration_nodes WHERE exploration_id=$1 AND kind='fact' AND state='origin' ORDER BY id LIMIT 1`, s.expID).Scan(&id)
|
|
if err == sql.ErrNoRows {
|
|
return 0, nil
|
|
}
|
|
return id, err
|
|
}
|
|
|
|
// Anchor records a node→asset reference edge (many-to-many): it captures which
|
|
// global assets an exploration node touched, as lineage/provenance. The asset
|
|
// graph itself is global and shared — anchoring no longer gates visibility (any
|
|
// task may read any asset via search); it only preserves the reasoning trail.
|
|
// Idempotent.
|
|
func (s *ExplorationStore) Anchor(nodeID, assetID int64) error {
|
|
if nodeID <= 0 || assetID <= 0 {
|
|
return nil
|
|
}
|
|
_, err := s.db.Exec(`INSERT INTO exploration_anchors(node_id, asset_id) VALUES ($1,$2) ON CONFLICT DO NOTHING`, nodeID, assetID)
|
|
return err
|
|
}
|
|
|
|
// Link adds a typed exploration edge (idempotent).
|
|
func (s *ExplorationStore) Link(from int64, rel string, to int64) error {
|
|
_, err := s.db.Exec(`
|
|
INSERT INTO exploration_edges(exploration_id, src_id, rel, dst_id) VALUES ($1,$2,$3,$4)
|
|
ON CONFLICT (exploration_id, src_id, rel, dst_id) DO NOTHING`, s.expID, from, rel, to)
|
|
return err
|
|
}
|
|
|
|
// SetNodeState updates any node's state (never deletes). content_version bumps
|
|
// so a folded member's state flip invalidates its digest's cached body (§5.3).
|
|
func (s *ExplorationStore) SetNodeState(id int64, state string) error {
|
|
_, err := s.db.Exec(`UPDATE exploration_nodes SET state=$1, blocked_reason=NULL, content_version=content_version+1 WHERE id=$2 AND exploration_id=$3`, state, id, s.expID)
|
|
return err
|
|
}
|
|
|
|
// UpdateGoalPayload rewrites a goal node's payload text (and optional vulnclass);
|
|
// scoped to kind='goal' so it can never mutate an intent/fact by id. Returns an
|
|
// error if no such goal exists in this exploration.
|
|
func (s *ExplorationStore) UpdateGoalPayload(id int64, text, vulnclass string) error {
|
|
payload := map[string]any{"text": text}
|
|
if vulnclass != "" {
|
|
payload["vulnclass"] = vulnclass
|
|
}
|
|
raw, _ := json.Marshal(payload)
|
|
res, err := s.db.Exec(`UPDATE exploration_nodes SET payload=$1 WHERE id=$2 AND exploration_id=$3 AND kind='goal'`, string(raw), id, s.expID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if n, _ := res.RowsAffected(); n == 0 {
|
|
return fmt.Errorf("目标不存在")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// DeleteGoal hard-deletes a goal node; scoped to kind='goal' so it can never drop
|
|
// an intent/fact by id. The edges (spawns) and anchors reference it with ON DELETE
|
|
// CASCADE, so they go with it; activity.node_id is ON DELETE SET NULL. Returns an
|
|
// error if no such goal exists in this exploration.
|
|
func (s *ExplorationStore) DeleteGoal(id int64) error {
|
|
res, err := s.db.Exec(`DELETE FROM exploration_nodes WHERE id=$1 AND exploration_id=$2 AND kind='goal'`, id, s.expID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if n, _ := res.RowsAffected(); n == 0 {
|
|
return fmt.Errorf("目标不存在")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SetIntentState updates an intent node's state.
|
|
// ResetRunningIntents returns any intent left in 'running' back to 'open' so it
|
|
// is re-claimed. Called on startup: a 'running' intent with no live worker (a
|
|
// backend restart or crashed worker goroutine left it stuck) would otherwise spin
|
|
// forever in the UI. Workers resume from their transcript, so reopening is safe.
|
|
func (s *ExplorationStore) ResetRunningIntents() (int64, error) {
|
|
res, err := s.db.Exec(`UPDATE exploration_nodes SET state='open', completed_at=NULL, blocked_reason=NULL
|
|
WHERE exploration_id=$1 AND kind='intent' AND state='running'`, s.expID)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
n, _ := res.RowsAffected()
|
|
return n, nil
|
|
}
|
|
|
|
// ReopenIntent flips ONE not-successfully-finished intent (blocked/exhausted/stopped)
|
|
// back to 'open' so a worker re-claims it (graph writes it already produced stay).
|
|
// The worker resumes from its prior LLM transcript instead of restarting from scratch.
|
|
// done/open/running are left untouched. Returns whether a row changed. Used by "重跑".
|
|
func (s *ExplorationStore) ReopenIntent(id int64) (bool, error) {
|
|
res, err := s.db.Exec(`UPDATE exploration_nodes
|
|
SET state='open', completed_at=NULL, blocked_reason=NULL
|
|
WHERE id=$1 AND exploration_id=$2 AND kind='intent'
|
|
AND state IN ('blocked','exhausted','stopped')`, id, s.expID)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
n, _ := res.RowsAffected()
|
|
return n > 0, nil
|
|
}
|
|
|
|
// ReopenBlockedIntents flips EVERY 'blocked' intent in this exploration back to 'open'
|
|
// (batch rerun after e.g. an LLM/network outage that blocked many at once). Returns the
|
|
// number reopened. Workers resume from their transcripts; kept graph writes remain.
|
|
func (s *ExplorationStore) ReopenBlockedIntents() (int64, error) {
|
|
res, err := s.db.Exec(`UPDATE exploration_nodes SET state='open', completed_at=NULL, blocked_reason=NULL
|
|
WHERE exploration_id=$1 AND kind='intent' AND state='blocked'`, s.expID)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
n, _ := res.RowsAffected()
|
|
return n, nil
|
|
}
|
|
|
|
func (s *ExplorationStore) SetIntentState(id int64, state string) error {
|
|
// terminal states stamp completed_at; reopening (back to open/running) clears it.
|
|
terminal := state == "done" || state == "blocked" || state == "exhausted" || state == "stopped"
|
|
_, err := s.db.Exec(`UPDATE exploration_nodes
|
|
SET state=$1, blocked_reason=NULL, content_version=content_version+1, completed_at = CASE WHEN $4 THEN now() ELSE NULL END
|
|
WHERE id=$2 AND exploration_id=$3 AND kind='intent'`, state, id, s.expID, terminal)
|
|
return err
|
|
}
|
|
|
|
// CompareAndSetIntentState transitions one local intent only when its current
|
|
// state still matches expected. It is the state boundary used by pause/resume and
|
|
// worker settlement, so a stale API read cannot overwrite a concurrent finish.
|
|
func (s *ExplorationStore) CompareAndSetIntentState(id int64, expected, state string) (bool, error) {
|
|
terminal := state == "done" || state == "blocked" || state == "exhausted" || state == "stopped"
|
|
res, err := s.db.Exec(`UPDATE exploration_nodes
|
|
SET state=$1, blocked_reason=NULL, content_version=content_version+1, completed_at = CASE WHEN $5 THEN now() ELSE NULL END
|
|
WHERE id=$2 AND exploration_id=$3 AND kind='intent' AND state=$4`, state, id, s.expID, expected, terminal)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
n, err := res.RowsAffected()
|
|
return n == 1, err
|
|
}
|
|
|
|
// IntentCleanup summarizes the blackboard records removed when a user cancels
|
|
// one worker intent. Global assets and traffic are deliberately not part of this
|
|
// operation; exploration anchors disappear only because their owning nodes do.
|
|
type IntentCleanup struct {
|
|
Intents int64 `json:"intents"`
|
|
Facts int64 `json:"facts"`
|
|
Findings int64 `json:"findings"`
|
|
Activities int64 `json:"activities"`
|
|
}
|
|
|
|
// SoftDeleteIntent 假删除一个待领/运行中/已暂停的意图:置 state='deleted' 并把用户填写的
|
|
// 删除原因记入 delete_reason 字段,保留意图节点及其全部产出/血缘(不再像旧实现那样在图上
|
|
// 另挂一条 fact)。返回删除前意图的 summary,供 planner 通知使用。副会话随删除态一并清理。
|
|
// 调用方须先停掉运行中的 worker,避免其后续写入。
|
|
func (s *ExplorationStore) SoftDeleteIntent(id int64, reason string) (string, error) {
|
|
tx, err := s.db.Begin()
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
var state string
|
|
var rawPayload []byte
|
|
if err := tx.QueryRow(`SELECT state, payload FROM exploration_nodes
|
|
WHERE id=$1 AND exploration_id=$2 AND kind='intent' FOR UPDATE`, id, s.expID).Scan(&state, &rawPayload); err != nil {
|
|
if err == sql.ErrNoRows {
|
|
return "", fmt.Errorf("intent not found")
|
|
}
|
|
return "", err
|
|
}
|
|
if state != "running" && state != "paused" && state != "open" {
|
|
return "", fmt.Errorf("%w: intent state %s cannot be deleted", ErrIntentStateConflict, state)
|
|
}
|
|
if _, err := tx.Exec(`UPDATE exploration_nodes
|
|
SET state='deleted', delete_reason=$3, blocked_reason=NULL,
|
|
content_version=content_version+1, completed_at=now()
|
|
WHERE id=$1 AND exploration_id=$2`, id, s.expID, reason); err != nil {
|
|
return "", err
|
|
}
|
|
// 主意图审计轨迹保留,但副问答会话随删除态原子清理(延迟快照会拒绝它)。
|
|
if _, err := tx.Exec(`DELETE FROM side_question_sessions WHERE intent_id=$1`, id); err != nil {
|
|
return "", err
|
|
}
|
|
var summary string
|
|
var p map[string]any
|
|
if json.Unmarshal(rawPayload, &p) == nil {
|
|
if sm, ok := p["summary"].(string); ok {
|
|
summary = sm
|
|
}
|
|
}
|
|
return summary, tx.Commit()
|
|
}
|
|
|
|
// CancelIntent 物理删除一个意图,以及"仅由该意图支撑"的全部独占子孙节点——从该意图沿
|
|
// yields/derived_from 向下可达、且所有父节点(所有指向它的边的源)都落在删除集内的节点。
|
|
// goal 与 origin fact 永不删除;还被删除集之外的意图/digest 引用的共享节点也保留,以免破坏
|
|
// 其它分支、产生断链。被删的每个 intent 先做 token rollup(保留不可逆计量)并清理其 activity
|
|
// 与副会话;被删的 finding 先删 findings 表行(node_id FK 为 ON DELETE SET NULL,否则留孤儿)。
|
|
// 边随节点 CASCADE 清理,整个清理保持在一个事务内。调用方须先停掉运行中的 worker 防止其后续写入。
|
|
func (s *ExplorationStore) CancelIntent(id int64) (IntentCleanup, error) {
|
|
var out IntentCleanup
|
|
tx, err := s.db.Begin()
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
// 锁定意图行确认存在(幂等:已删则 not found)。状态不校验——真删除对任何状态成立,
|
|
// running 的 worker 由调用方先停。
|
|
if err := tx.QueryRow(`SELECT 1 FROM exploration_nodes
|
|
WHERE id=$1 AND exploration_id=$2 AND kind='intent' FOR UPDATE`, id, s.expID).Scan(new(int)); err != nil {
|
|
if err == sql.ErrNoRows {
|
|
return out, fmt.Errorf("intent not found")
|
|
}
|
|
return out, err
|
|
}
|
|
|
|
// 载入全图节点(判 protected)与边(算向下可达 + 父集)。图规模对真实任务很小。
|
|
kind := map[int64]string{}
|
|
protected := map[int64]bool{}
|
|
nrows, err := tx.Query(`SELECT id, kind, state FROM exploration_nodes WHERE exploration_id=$1`, s.expID)
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
for nrows.Next() {
|
|
var nid int64
|
|
var k, st string
|
|
if err := nrows.Scan(&nid, &k, &st); err != nil {
|
|
nrows.Close()
|
|
return out, err
|
|
}
|
|
kind[nid] = k
|
|
if k == KindGoal || (k == KindFact && st == StateOrigin) {
|
|
protected[nid] = true // 目标与任务根事实永不随意图删除
|
|
}
|
|
}
|
|
nrows.Close()
|
|
if err := nrows.Err(); err != nil {
|
|
return out, err
|
|
}
|
|
|
|
parentsOf := map[int64][]int64{} // dst -> 所有指向它的边的 src(任意 rel,含 covers:被 digest 覆盖的成员因此有 digest 父而被保留)
|
|
downOf := map[int64][]int64{} // src -> 沿 yields/derived_from 的向下邻居
|
|
erows, err := tx.Query(`SELECT src_id, rel, dst_id FROM exploration_edges WHERE exploration_id=$1`, s.expID)
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
for erows.Next() {
|
|
var src, dst int64
|
|
var rel string
|
|
if err := erows.Scan(&src, &rel, &dst); err != nil {
|
|
erows.Close()
|
|
return out, err
|
|
}
|
|
parentsOf[dst] = append(parentsOf[dst], src)
|
|
if rel == RelYields || rel == RelDerivedFrom {
|
|
downOf[src] = append(downOf[src], dst)
|
|
}
|
|
}
|
|
erows.Close()
|
|
if err := erows.Err(); err != nil {
|
|
return out, err
|
|
}
|
|
|
|
// 独占级联:从意图向下扩展,一个节点入删除集当且仅当它未被保护、且它的每个父都已在集内
|
|
// (即除了经过被删节点外再无来路)。迭代到不动点。
|
|
del := map[int64]bool{id: true}
|
|
for changed := true; changed; {
|
|
changed = false
|
|
for src := range del {
|
|
for _, dst := range downOf[src] {
|
|
if del[dst] || protected[dst] {
|
|
continue
|
|
}
|
|
exclusive := true
|
|
for _, p := range parentsOf[dst] {
|
|
if !del[p] {
|
|
exclusive = false
|
|
break
|
|
}
|
|
}
|
|
if exclusive {
|
|
del[dst] = true
|
|
changed = true
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// 按类型分桶。
|
|
var ids, intentIDs, findingIDs []int64
|
|
for nid := range del {
|
|
ids = append(ids, nid)
|
|
switch kind[nid] {
|
|
case KindIntent:
|
|
intentIDs = append(intentIDs, nid)
|
|
out.Intents++
|
|
case KindFact:
|
|
out.Facts++
|
|
case KindFinding:
|
|
findingIDs = append(findingIDs, nid)
|
|
out.Findings++
|
|
}
|
|
}
|
|
|
|
// 每个被删意图:token rollup(保留不可逆计量)后清 activity。rollup 行 node_id 为 NULL,
|
|
// 不会被下面按 node_id 的删除命中。
|
|
for _, iid := range intentIDs {
|
|
tokenBuckets, err := intentTokenRollup(tx, s.expID, iid)
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
res, err := tx.Exec(`DELETE FROM activity WHERE exploration_id=$1 AND node_id=$2`, s.expID, iid)
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
if n, _ := res.RowsAffected(); n > 0 {
|
|
out.Activities += n
|
|
}
|
|
for _, bucket := range tokenBuckets {
|
|
metadata, _ := json.Marshal(map[string]any{
|
|
"cancelled_intent_id": iid,
|
|
"token_day": bucket.Day.Format(time.DateOnly),
|
|
"token_rollup": true,
|
|
})
|
|
if _, err := tx.Exec(`INSERT INTO activity(
|
|
exploration_id, worker, kind, summary, metadata,
|
|
input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, created_at)
|
|
VALUES ($1,'token-ledger','result',$2,$3,$4,$5,$6,$7,$8)`,
|
|
s.expID, fmt.Sprintf("已取消意图 #%d 的 Token 计量", iid), metadata,
|
|
bucket.Usage.InputTokens, bucket.Usage.OutputTokens,
|
|
bucket.Usage.CacheReadTokens, bucket.Usage.CacheWriteTokens, bucket.Day); err != nil {
|
|
return out, err
|
|
}
|
|
}
|
|
if _, err := tx.Exec(`DELETE FROM side_question_sessions WHERE intent_id=$1`, iid); err != nil {
|
|
return out, err
|
|
}
|
|
}
|
|
|
|
// finding 节点删除前先删 findings 表行(否则留 node_id=NULL 的孤儿)。
|
|
for _, fid := range findingIDs {
|
|
if _, err := tx.Exec(`DELETE FROM findings WHERE node_id=$1`, fid); err != nil {
|
|
return out, err
|
|
}
|
|
}
|
|
|
|
// 删节点(边随 CASCADE 清理)。逐个删并核对行数,防并发改动。
|
|
var removed int64
|
|
for _, nid := range ids {
|
|
res, err := tx.Exec(`DELETE FROM exploration_nodes WHERE id=$1 AND exploration_id=$2`, nid, s.expID)
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
if n, _ := res.RowsAffected(); n > 0 {
|
|
removed += n
|
|
}
|
|
}
|
|
if removed != int64(len(ids)) {
|
|
return out, fmt.Errorf("intent cleanup changed concurrently")
|
|
}
|
|
return out, tx.Commit()
|
|
}
|
|
|
|
type tokenUsageBucket struct {
|
|
Day time.Time
|
|
Usage TokenUsage
|
|
}
|
|
|
|
// intentTokenRollup returns exactly-once usage grouped by its original UTC day.
|
|
// Each result terminates one run. A result carrying usage is authoritative; an
|
|
// older error/cancel result without usage falls back to the latest cumulative
|
|
// usage frame since the previous result. A trailing frame represents a run that
|
|
// was interrupted before its terminal result could be persisted.
|
|
func intentTokenRollup(tx *sql.Tx, explorationID, intentID int64) ([]tokenUsageBucket, error) {
|
|
rows, err := tx.Query(`SELECT kind, input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, created_at
|
|
FROM activity
|
|
WHERE exploration_id=$1 AND node_id=$2 AND kind IN ('usage','result')
|
|
ORDER BY id`, explorationID, intentID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
byDay := make(map[time.Time]TokenUsage)
|
|
add := func(at time.Time, usage TokenUsage) {
|
|
if usage.InputTokens == 0 && usage.OutputTokens == 0 &&
|
|
usage.CacheReadTokens == 0 && usage.CacheWriteTokens == 0 {
|
|
return
|
|
}
|
|
utc := at.UTC()
|
|
day := time.Date(utc.Year(), utc.Month(), utc.Day(), 0, 0, 0, 0, time.UTC)
|
|
total := byDay[day]
|
|
total.add(usage)
|
|
byDay[day] = total
|
|
}
|
|
|
|
var pending TokenUsage
|
|
var pendingAt time.Time
|
|
hasPending := false
|
|
for rows.Next() {
|
|
var kind string
|
|
var input, output, read, write sql.NullInt64
|
|
var createdAt time.Time
|
|
if err := rows.Scan(&kind, &input, &output, &read, &write, &createdAt); err != nil {
|
|
return nil, err
|
|
}
|
|
current := TokenUsage{
|
|
InputTokens: int(input.Int64),
|
|
OutputTokens: int(output.Int64),
|
|
CacheReadTokens: int(read.Int64),
|
|
CacheWriteTokens: int(write.Int64),
|
|
}
|
|
hasCurrent := input.Valid || output.Valid || read.Valid || write.Valid
|
|
if kind == "usage" {
|
|
if hasCurrent {
|
|
pending, pendingAt, hasPending = current, createdAt, true
|
|
}
|
|
continue
|
|
}
|
|
|
|
if hasCurrent {
|
|
// Preserve a partially populated legacy terminal row by filling only its
|
|
// missing dimensions from the latest cumulative usage frame.
|
|
if hasPending {
|
|
if !input.Valid {
|
|
current.InputTokens = pending.InputTokens
|
|
}
|
|
if !output.Valid {
|
|
current.OutputTokens = pending.OutputTokens
|
|
}
|
|
if !read.Valid {
|
|
current.CacheReadTokens = pending.CacheReadTokens
|
|
}
|
|
if !write.Valid {
|
|
current.CacheWriteTokens = pending.CacheWriteTokens
|
|
}
|
|
}
|
|
add(createdAt, current)
|
|
} else if hasPending {
|
|
add(createdAt, pending)
|
|
}
|
|
pending, pendingAt, hasPending = TokenUsage{}, time.Time{}, false
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
if hasPending {
|
|
add(pendingAt, pending)
|
|
}
|
|
|
|
buckets := make([]tokenUsageBucket, 0, len(byDay))
|
|
for day, usage := range byDay {
|
|
buckets = append(buckets, tokenUsageBucket{Day: day, Usage: usage})
|
|
}
|
|
sort.Slice(buckets, func(i, j int) bool { return buckets[i].Day.Before(buckets[j].Day) })
|
|
return buckets, nil
|
|
}
|
|
|
|
func (u *TokenUsage) add(other TokenUsage) {
|
|
u.InputTokens += other.InputTokens
|
|
u.OutputTokens += other.OutputTokens
|
|
u.CacheReadTokens += other.CacheReadTokens
|
|
u.CacheWriteTokens += other.CacheWriteTokens
|
|
}
|
|
|
|
const nodeCols = `id, kind, payload, priority, state, COALESCE(origin,''), COALESCE(owner,''), COALESCE(blocked_reason,''), COALESCE(delete_reason,''), created_at`
|
|
|
|
func scanNode(sc interface{ Scan(...any) error }) (*Node, error) {
|
|
var n Node
|
|
var payload []byte
|
|
if err := sc.Scan(&n.ID, &n.Kind, &payload, &n.Priority, &n.State, &n.Origin, &n.Owner, &n.BlockedReason, &n.DeleteReason, &n.CreatedAt); err != nil {
|
|
return nil, err
|
|
}
|
|
n.Payload = json.RawMessage(payload)
|
|
return &n, nil
|
|
}
|
|
|
|
func scanNodes(rows *sql.Rows) ([]*Node, error) {
|
|
var out []*Node
|
|
for rows.Next() {
|
|
n, err := scanNode(rows)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, n)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// ListByKind lists nodes of a kind (newest first).
|
|
func (s *ExplorationStore) ListByKind(kind string, limit int) ([]*Node, error) {
|
|
if limit <= 0 {
|
|
limit = 500
|
|
}
|
|
rows, err := s.db.Query(`SELECT `+nodeCols+` FROM exploration_nodes
|
|
WHERE exploration_id=$1 AND kind=$2 ORDER BY id DESC LIMIT $3`, s.expID, kind, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanNodes(rows)
|
|
}
|
|
|
|
// ListByKindPage returns one newest-first page of nodes of a kind for reverse
|
|
// pagination: up to `limit` nodes with id < `before` (before<=0 = the newest page).
|
|
// hasMore reports whether still-older nodes exist, so a client (e.g. the worker
|
|
// session list) can page past the old fixed cap instead of losing older intents.
|
|
func (s *ExplorationStore) ListByKindPage(kind string, before int64, limit int) ([]*Node, bool, error) {
|
|
if limit <= 0 {
|
|
limit = 300
|
|
}
|
|
rows, err := s.db.Query(`SELECT `+nodeCols+` FROM exploration_nodes
|
|
WHERE exploration_id=$1 AND kind=$2 AND ($3 <= 0 OR id < $3)
|
|
ORDER BY id DESC LIMIT $4`, s.expID, kind, before, limit+1)
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
defer rows.Close()
|
|
nodes, err := scanNodes(rows)
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
hasMore := len(nodes) > limit
|
|
if hasMore {
|
|
nodes = nodes[:limit]
|
|
}
|
|
return nodes, hasMore, nil
|
|
}
|
|
|
|
// listByKindPageFiltered is ListByKindPage with an optional summary keyword
|
|
// filter (payload->>'summary' ILIKE %q%). Fetches one extra row so the caller can
|
|
// probe hasMore. Empty q = no filter. Newest-first by id.
|
|
func (s *ExplorationStore) listByKindPageFiltered(kind string, before int64, limit int, q string) ([]*Node, error) {
|
|
where := `exploration_id=$1 AND kind=$2 AND ($3 <= 0 OR id < $3)`
|
|
args := []any{s.expID, kind, before}
|
|
if q != "" {
|
|
args = append(args, "%"+q+"%")
|
|
where += fmt.Sprintf(` AND payload->>'summary' ILIKE $%d`, len(args))
|
|
}
|
|
args = append(args, limit+1)
|
|
rows, err := s.db.Query(`SELECT `+nodeCols+` FROM exploration_nodes
|
|
WHERE `+where+` ORDER BY id DESC LIMIT $`+fmt.Sprintf("%d", len(args)), args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanNodes(rows)
|
|
}
|
|
|
|
// countByKindFiltered counts nodes of a kind under the same summary filter, for a
|
|
// page's total. Empty q = count all of that kind.
|
|
func (s *ExplorationStore) countByKindFiltered(kind, q string) (int, error) {
|
|
where := `exploration_id=$1 AND kind=$2`
|
|
args := []any{s.expID, kind}
|
|
if q != "" {
|
|
args = append(args, "%"+q+"%")
|
|
where += fmt.Sprintf(` AND payload->>'summary' ILIKE $%d`, len(args))
|
|
}
|
|
var n int
|
|
err := s.db.QueryRow(`SELECT COUNT(*) FROM exploration_nodes WHERE `+where, args...).Scan(&n)
|
|
return n, err
|
|
}
|
|
|
|
// CountFinishedIntents counts this exploration's intents in a terminal explored
|
|
// state (done/blocked/exhausted) — graph_overview's done_intents_total, so the planner
|
|
// knows recent_done_intents (capped at a recent window) is a truncated view and
|
|
// stays cautious about "already tried" dedup. Excludes 'stopped' (killed/deleted),
|
|
// matching exactly what recent_done_intents surfaces.
|
|
func (s *ExplorationStore) CountFinishedIntents() (int, error) {
|
|
var n int
|
|
err := s.db.QueryRow(`SELECT COUNT(*) FROM exploration_nodes
|
|
WHERE exploration_id=$1 AND kind='intent' AND state IN ('done','blocked','exhausted')`, s.expID).Scan(&n)
|
|
return n, err
|
|
}
|
|
|
|
// CountOpenIntents counts this exploration's open intents — graph_overview's
|
|
// frontier_open, so the planner knows open_intents (capped at the top-N by
|
|
// priority) is a truncated view and there may be more claimable work.
|
|
func (s *ExplorationStore) CountOpenIntents() (int, error) {
|
|
var n int
|
|
err := s.db.QueryRow(`SELECT COUNT(*) FROM exploration_nodes
|
|
WHERE exploration_id=$1 AND kind='intent' AND state='open'`, s.expID).Scan(&n)
|
|
return n, err
|
|
}
|
|
|
|
// GetNode returns one node of this exploration by id (nil, nil if not found).
|
|
func (s *ExplorationStore) GetNode(id int64) (*Node, error) {
|
|
n, err := scanNode(s.db.QueryRow(`SELECT `+nodeCols+` FROM exploration_nodes WHERE id=$1 AND exploration_id=$2`, id, s.expID))
|
|
if err == sql.ErrNoRows {
|
|
return nil, nil
|
|
}
|
|
return n, err
|
|
}
|
|
|
|
// Nodes lists all nodes (for graph viz).
|
|
func (s *ExplorationStore) Nodes(limit int) ([]*Node, error) {
|
|
if limit <= 0 {
|
|
limit = 2000
|
|
}
|
|
rows, err := s.db.Query(`SELECT `+nodeCols+` FROM exploration_nodes WHERE exploration_id=$1 ORDER BY id LIMIT $2`, s.expID, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanNodes(rows)
|
|
}
|
|
|
|
// NodeFilter narrows a NodesPage query. Empty fields include every value.
|
|
type NodeFilter struct {
|
|
Kinds []string // exploration_nodes.kind
|
|
States []string // exploration_nodes.state
|
|
Query string // case-insensitive substring over payload text / origin
|
|
Asc bool // true = oldest first (replay); default newest first (broadcast)
|
|
}
|
|
|
|
// NodesPage returns one 1-based page of this exploration's nodes plus the total
|
|
// matching count. The 播报板 reads the graph as a time series, so it pages in SQL
|
|
// rather than pulling the whole graph like Nodes does. Ordering is by id, which
|
|
// is BIGSERIAL and therefore creation order — stable when several nodes share a
|
|
// created_at second.
|
|
func (s *ExplorationStore) NodesPage(f NodeFilter, page, size int) ([]*Node, int, error) {
|
|
if page < 1 {
|
|
page = 1
|
|
}
|
|
if size <= 0 {
|
|
size = 20
|
|
}
|
|
if size > 200 {
|
|
size = 200
|
|
}
|
|
offset := (page - 1) * size
|
|
|
|
conds := []string{"exploration_id=$1"}
|
|
args := []any{s.expID}
|
|
addIn := func(column string, values []string) {
|
|
if len(values) == 0 {
|
|
return
|
|
}
|
|
marks := make([]string, 0, len(values))
|
|
for _, v := range values {
|
|
args = append(args, v)
|
|
marks = append(marks, "$"+fmt.Sprint(len(args)))
|
|
}
|
|
conds = append(conds, column+" IN ("+strings.Join(marks, ",")+")")
|
|
}
|
|
addIn("kind", f.Kinds)
|
|
addIn("state", f.States)
|
|
if q := strings.TrimSpace(f.Query); q != "" {
|
|
args = append(args, "%"+q+"%")
|
|
mark := "$" + fmt.Sprint(len(args))
|
|
ors := []string{"payload::text ILIKE " + mark, "COALESCE(origin,'') ILIKE " + mark}
|
|
// 纯数字(或 UI 里带 # 前缀的形式,如「#41」)当作节点 id 精确匹配,方便直接定位某个节点。
|
|
if idStr := strings.TrimPrefix(q, "#"); idStr != "" {
|
|
if id, err := strconv.ParseInt(idStr, 10, 64); err == nil {
|
|
args = append(args, id)
|
|
ors = append(ors, "id = $"+fmt.Sprint(len(args)))
|
|
}
|
|
}
|
|
conds = append(conds, "("+strings.Join(ors, " OR ")+")")
|
|
}
|
|
where := " WHERE " + strings.Join(conds, " AND ")
|
|
|
|
var total int
|
|
if err := s.db.QueryRow(`SELECT COUNT(*) FROM exploration_nodes`+where, args...).Scan(&total); err != nil {
|
|
return nil, 0, err
|
|
}
|
|
order := "DESC"
|
|
if f.Asc {
|
|
order = "ASC"
|
|
}
|
|
args = append(args, size, offset)
|
|
rows, err := s.db.Query(`SELECT `+nodeCols+` FROM exploration_nodes`+where+
|
|
` ORDER BY id `+order+
|
|
` LIMIT $`+fmt.Sprint(len(args)-1)+` OFFSET $`+fmt.Sprint(len(args)), args...)
|
|
if err != nil {
|
|
return nil, 0, err
|
|
}
|
|
defer rows.Close()
|
|
nodes, err := scanNodes(rows)
|
|
if err != nil {
|
|
return nil, 0, err
|
|
}
|
|
return nodes, total, nil
|
|
}
|
|
|
|
// NodesByIDs loads the given nodes of this exploration in id order. Used to
|
|
// resolve the neighbours of a 播报板 page without fetching the whole graph.
|
|
func (s *ExplorationStore) NodesByIDs(ids []int64) ([]*Node, error) {
|
|
if len(ids) == 0 {
|
|
return nil, nil
|
|
}
|
|
args := make([]any, 0, len(ids)+1)
|
|
args = append(args, s.expID)
|
|
marks := make([]string, 0, len(ids))
|
|
for _, id := range ids {
|
|
args = append(args, id)
|
|
marks = append(marks, "$"+fmt.Sprint(len(args)))
|
|
}
|
|
rows, err := s.db.Query(`SELECT `+nodeCols+` FROM exploration_nodes
|
|
WHERE exploration_id=$1 AND id IN (`+strings.Join(marks, ",")+`) ORDER BY id`, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanNodes(rows)
|
|
}
|
|
|
|
// EdgesTouching returns every edge with one end among ids — the 播报板 uses it to
|
|
// spell out where a node came from and what it produced.
|
|
func (s *ExplorationStore) EdgesTouching(ids []int64) ([]Edge, error) {
|
|
if len(ids) == 0 {
|
|
return nil, nil
|
|
}
|
|
args := make([]any, 0, len(ids)+1)
|
|
args = append(args, s.expID)
|
|
marks := make([]string, 0, len(ids))
|
|
for _, id := range ids {
|
|
args = append(args, id)
|
|
marks = append(marks, "$"+fmt.Sprint(len(args)))
|
|
}
|
|
list := strings.Join(marks, ",")
|
|
rows, err := s.db.Query(`SELECT src_id, rel, dst_id FROM exploration_edges
|
|
WHERE exploration_id=$1 AND (src_id IN (`+list+`) OR dst_id IN (`+list+`))`, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []Edge
|
|
for rows.Next() {
|
|
var e Edge
|
|
if err := rows.Scan(&e.From, &e.Rel, &e.To); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, e)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// Edge is one row from exploration_edges.
|
|
type Edge struct {
|
|
From int64
|
|
Rel string
|
|
To int64
|
|
}
|
|
|
|
// Edges lists exploration edges (for graph viz).
|
|
func (s *ExplorationStore) Edges(limit int) ([]Edge, error) {
|
|
if limit <= 0 {
|
|
limit = 5000
|
|
}
|
|
rows, err := s.db.Query(`SELECT src_id, rel, dst_id FROM exploration_edges WHERE exploration_id=$1 LIMIT $2`, s.expID, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []Edge
|
|
for rows.Next() {
|
|
var e Edge
|
|
if err := rows.Scan(&e.From, &e.Rel, &e.To); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, e)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// FactsYielded returns the ids of fact nodes an intent produced (intent --yields-->
|
|
// fact), oldest first. Used to spell out the round's incremental facts to the planner
|
|
// alongside the finished worker's output. Best-effort: returns nil for an intent with
|
|
// no facts (or a missing one).
|
|
func (s *ExplorationStore) FactsYielded(intentID int64) ([]int64, error) {
|
|
rows, err := s.db.Query(`SELECT n.id
|
|
FROM exploration_edges e
|
|
JOIN exploration_nodes n ON n.id=e.dst_id AND n.exploration_id=e.exploration_id
|
|
WHERE e.exploration_id=$1 AND e.src_id=$2 AND e.rel=$3 AND n.kind=$4
|
|
ORDER BY n.id`, s.expID, intentID, RelYields, KindFact)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var ids []int64
|
|
for rows.Next() {
|
|
var id int64
|
|
if err := rows.Scan(&id); err != nil {
|
|
return nil, err
|
|
}
|
|
ids = append(ids, id)
|
|
}
|
|
return ids, rows.Err()
|
|
}
|
|
|
|
// FindingLineage returns the sub-DAG that leads FROM the exploration root TO the
|
|
// given node: the node itself plus all its ancestors (nodes reverse-reachable by
|
|
// following edges backward), and every edge whose both endpoints are in that set.
|
|
// Used to render "how this finding was reached" — origin → … → finding. Returns
|
|
// empty (not error) for a nonexistent node.
|
|
func (s *ExplorationStore) FindingLineage(nodeID int64) ([]*Node, []Edge, error) {
|
|
// anc = the node + everything that can reach it via edges (walk src<-dst).
|
|
nodeRows, err := s.db.Query(`
|
|
WITH RECURSIVE anc(id) AS (
|
|
SELECT $2::bigint
|
|
UNION
|
|
SELECT e.src_id
|
|
FROM exploration_edges e
|
|
JOIN anc ON e.dst_id = anc.id
|
|
WHERE e.exploration_id = $1
|
|
)
|
|
SELECT `+nodeCols+`
|
|
FROM exploration_nodes
|
|
WHERE exploration_id = $1 AND id IN (SELECT id FROM anc)
|
|
ORDER BY id`, s.expID, nodeID)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
nodes, err := scanNodes(nodeRows)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
if len(nodes) == 0 {
|
|
return nil, nil, nil
|
|
}
|
|
// edges among the ancestor set only (the lineage's internal relations).
|
|
edgeRows, err := s.db.Query(`
|
|
WITH RECURSIVE anc(id) AS (
|
|
SELECT $2::bigint
|
|
UNION
|
|
SELECT e.src_id
|
|
FROM exploration_edges e
|
|
JOIN anc ON e.dst_id = anc.id
|
|
WHERE e.exploration_id = $1
|
|
)
|
|
SELECT src_id, rel, dst_id
|
|
FROM exploration_edges
|
|
WHERE exploration_id = $1
|
|
AND src_id IN (SELECT id FROM anc)
|
|
AND dst_id IN (SELECT id FROM anc)`, s.expID, nodeID)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
defer edgeRows.Close()
|
|
var edges []Edge
|
|
for edgeRows.Next() {
|
|
var e Edge
|
|
if err := edgeRows.Scan(&e.From, &e.Rel, &e.To); err != nil {
|
|
return nil, nil, err
|
|
}
|
|
edges = append(edges, e)
|
|
}
|
|
return nodes, edges, edgeRows.Err()
|
|
}
|
|
|
|
// FindingIntents maps each finding id to the intent that produced it (the
|
|
// intent --yields--> finding edge; report_finding links it). Findings with no
|
|
// producing intent are absent. Precise (JOIN, no edge-scan limit).
|
|
func (s *ExplorationStore) FindingIntents() (map[int64]int64, error) {
|
|
rows, err := s.db.Query(`SELECT e.dst_id, e.src_id
|
|
FROM exploration_edges e
|
|
JOIN exploration_nodes n ON n.id = e.dst_id AND n.exploration_id = e.exploration_id
|
|
WHERE e.exploration_id=$1 AND e.rel=$2 AND n.kind='finding'`, s.expID, RelYields)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := map[int64]int64{}
|
|
for rows.Next() {
|
|
var findingID, intentID int64
|
|
if err := rows.Scan(&findingID, &intentID); err != nil {
|
|
return nil, err
|
|
}
|
|
out[findingID] = intentID
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// FindingIntentsTerminal is the inherited-history variant: it returns finding
|
|
// lineage only when the producing intent has reached an immutable terminal
|
|
// state. This keeps a source task's live worker topology private.
|
|
func (s *ExplorationStore) FindingIntentsTerminal() (map[int64]int64, error) {
|
|
rows, err := s.db.Query(`SELECT e.dst_id, e.src_id
|
|
FROM exploration_edges e
|
|
JOIN exploration_nodes finding ON finding.id=e.dst_id AND finding.exploration_id=e.exploration_id
|
|
JOIN exploration_nodes intent ON intent.id=e.src_id AND intent.exploration_id=e.exploration_id
|
|
WHERE e.exploration_id=$1 AND e.rel=$2 AND finding.kind='finding'
|
|
AND intent.kind='intent' AND intent.state IN ('done','blocked','exhausted','stopped')`, s.expID, RelYields)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := map[int64]int64{}
|
|
for rows.Next() {
|
|
var findingID, intentID int64
|
|
if err := rows.Scan(&findingID, &intentID); err != nil {
|
|
return nil, err
|
|
}
|
|
out[findingID] = intentID
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// Frontier returns open intents ordered by priority desc then id (FIFO tiebreak).
|
|
func (s *ExplorationStore) Frontier(limit int) ([]*Node, error) {
|
|
if limit <= 0 {
|
|
limit = 50
|
|
}
|
|
rows, err := s.db.Query(`SELECT `+nodeCols+` FROM exploration_nodes
|
|
WHERE exploration_id=$1 AND kind='intent' AND state='open'
|
|
ORDER BY priority DESC, id ASC LIMIT $2`, s.expID, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanNodes(rows)
|
|
}
|
|
|
|
// HasActiveIntent reports whether this exploration has any intent still in play
|
|
// (state open OR running). Used to decide whether to kick the FIRST planner round:
|
|
// a seeded task's intent may already have been claimed (open→running) by a worker
|
|
// before the check, so counting only 'open' (Frontier) races with worker claim and
|
|
// wrongly kicks a redundant first round. Counting open+running is race-proof.
|
|
func (s *ExplorationStore) HasActiveIntent() (bool, error) {
|
|
var exists bool
|
|
err := s.db.QueryRow(`SELECT EXISTS(SELECT 1 FROM exploration_nodes
|
|
WHERE exploration_id=$1 AND kind='intent' AND state IN ('open','running'))`, s.expID).Scan(&exists)
|
|
return exists, err
|
|
}
|
|
|
|
// HasOpenGoal reports whether this exploration still has any goal in state 'open'.
|
|
// false ⇒ 所有目标已 met/abandoned(或本任务无目标)⇒ 进入 goalless(人工直投)分支:
|
|
// planner 停跑,任务是否结束改由 frontier 是否抽干决定。met 与 abandoned 都算"已了结"。
|
|
func (s *ExplorationStore) HasOpenGoal() (bool, error) {
|
|
var exists bool
|
|
err := s.db.QueryRow(`SELECT EXISTS(SELECT 1 FROM exploration_nodes
|
|
WHERE exploration_id=$1 AND kind='goal' AND state='open')`, s.expID).Scan(&exists)
|
|
return exists, err
|
|
}
|
|
|
|
// ClaimIntent atomically moves an open intent to running. Returns true if claimed.
|
|
func (s *ExplorationStore) ClaimIntent(id int64, owner string) (bool, error) {
|
|
res, err := s.db.Exec(`UPDATE exploration_nodes SET state='running', owner=$1
|
|
WHERE id=$2 AND exploration_id=$3 AND kind='intent' AND state='open'`, owner, id, s.expID)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
n, _ := res.RowsAffected()
|
|
return n == 1, nil
|
|
}
|
|
|
|
// Stats returns node counts grouped by kind (for dashboard).
|
|
func (s *ExplorationStore) Stats() (map[string]int, error) {
|
|
rows, err := s.db.Query(`SELECT kind, count(*) FROM exploration_nodes WHERE exploration_id=$1 GROUP BY kind`, s.expID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := map[string]int{}
|
|
for rows.Next() {
|
|
var k string
|
|
var c int
|
|
if err := rows.Scan(&k, &c); err != nil {
|
|
return nil, err
|
|
}
|
|
out[k] = c
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// --- activity (poll by global id cursor) ---
|
|
|
|
// AppendActivity records one worker step and returns its id.
|
|
func (s *ExplorationStore) AppendActivity(a Activity) (int64, error) {
|
|
metadata := a.Metadata
|
|
if len(metadata) == 0 {
|
|
metadata = json.RawMessage(`{}`)
|
|
}
|
|
var id int64
|
|
err := s.db.QueryRow(`
|
|
INSERT INTO activity(exploration_id, node_id, worker, kind, tool, tool_use_id, is_error, summary, detail, metadata, input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, main_seg)
|
|
VALUES ($1,$2,NULLIF($3,''),NULLIF($4,''),NULLIF($5,''),NULLIF($6,''),$7,NULLIF($8,''),NULLIF($9,''),$10,$11,$12,$13,$14,$15)
|
|
RETURNING id`, s.expID, a.NodeID, utf8Clean(a.Worker), utf8Clean(a.Kind), utf8Clean(a.Tool), utf8Clean(a.ToolUseID), a.IsError, utf8Clean(a.Summary), utf8Clean(a.Detail),
|
|
metadata, a.InputTokens, a.OutputTokens, a.CacheReadTokens, a.CacheWriteTokens, a.MainSeg).Scan(&id)
|
|
return id, err
|
|
}
|
|
|
|
// ExplorationDiag probes, at the moment of an activity FK violation (23503), why the
|
|
// parent row is unreachable. It reports whether the exploration row still exists, how
|
|
// many task rows reference it (a live task should keep exactly 1 — and RESTRICT on
|
|
// tasks.exploration_id means the exploration CANNOT be deleted while that row lives),
|
|
// and the current MAX(explorations.id). Together these tell apart the failure modes:
|
|
// - expExists=false, taskRefs=0 → the whole row set is gone (DB reset / wrong DB).
|
|
// - expExists=false, taskRefs=1 → impossible under RESTRICT; would mean a broken FK.
|
|
// - expExists=true → the write's expID is NOT this exploration (stale/
|
|
// wrong id in the in-memory store); compare against Store.ID().
|
|
func (d *DB) ExplorationDiag(expID int64) (expExists bool, taskRefs int, maxExpID int64, err error) {
|
|
err = d.QueryRow(`SELECT
|
|
EXISTS(SELECT 1 FROM explorations WHERE id=$1),
|
|
(SELECT COUNT(*) FROM tasks WHERE exploration_id=$1),
|
|
COALESCE((SELECT MAX(id) FROM explorations),0)`, expID).Scan(&expExists, &taskRefs, &maxExpID)
|
|
return
|
|
}
|
|
|
|
// TokenTotal sums token usage across ALL workers for this exploration (whole-task
|
|
// total). Same source as TokenStatsByWorker (kind='result' rows), just ungrouped.
|
|
func (s *ExplorationStore) TokenTotal() (TokenUsage, error) {
|
|
var u TokenUsage
|
|
err := s.db.QueryRow(`SELECT
|
|
COALESCE(SUM(input_tokens),0), COALESCE(SUM(output_tokens),0),
|
|
COALESCE(SUM(cache_read_tokens),0), COALESCE(SUM(cache_write_tokens),0)
|
|
FROM activity WHERE exploration_id=$1 AND kind='result'`, s.expID).
|
|
Scan(&u.InputTokens, &u.OutputTokens, &u.CacheReadTokens, &u.CacheWriteTokens)
|
|
return u, err
|
|
}
|
|
|
|
// TokenTotalsAll returns the whole-task token total for every exploration in one
|
|
// query (exploration_id → total), so the task list can show per-task consumption
|
|
// without a query per row.
|
|
func (d *DB) TokenTotalsAll() (map[int64]TokenUsage, error) {
|
|
rows, err := d.Query(`SELECT exploration_id,
|
|
COALESCE(SUM(input_tokens),0), COALESCE(SUM(output_tokens),0),
|
|
COALESCE(SUM(cache_read_tokens),0), COALESCE(SUM(cache_write_tokens),0)
|
|
FROM activity WHERE kind='result' GROUP BY exploration_id`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := map[int64]TokenUsage{}
|
|
for rows.Next() {
|
|
var eid int64
|
|
var u TokenUsage
|
|
if err := rows.Scan(&eid, &u.InputTokens, &u.OutputTokens, &u.CacheReadTokens, &u.CacheWriteTokens); err != nil {
|
|
return nil, err
|
|
}
|
|
out[eid] = u
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// LastActivityAll returns the unix time of the most recent activity per
|
|
// exploration (exploration_id → max created_at epoch), one query for all tasks.
|
|
// Persisted (unlike Engine.LastActivity's in-memory map), so it survives restarts
|
|
// and gives终态任务 a stable "ran until" time for computing run duration.
|
|
func (d *DB) LastActivityAll() (map[int64]int64, error) {
|
|
rows, err := d.Query(`SELECT exploration_id, EXTRACT(EPOCH FROM MAX(created_at))::bigint FROM activity GROUP BY exploration_id`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := map[int64]int64{}
|
|
for rows.Next() {
|
|
var eid, ts int64
|
|
if err := rows.Scan(&eid, &ts); err != nil {
|
|
return nil, err
|
|
}
|
|
out[eid] = ts
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// GoalCounts is the goal summary for one exploration.
|
|
type GoalCounts struct{ Total, Met int }
|
|
|
|
// FindingSeverityCounts breaks a task's findings down by severity for the task
|
|
// list (severity 白名单外/为空的记录不计入任一档)。
|
|
type FindingSeverityCounts struct{ Critical, High, Medium, Low int }
|
|
|
|
// TaskListMetrics contains the aggregates rendered in task lists.
|
|
type TaskListMetrics struct {
|
|
Tokens TokenUsage
|
|
LastActivity int64
|
|
Goals GoalCounts
|
|
RunningIntents int // kind='intent' 且 state='running' 的条数,即运行中 Worker 数
|
|
Findings FindingSeverityCounts // findings 表里该任务的漏洞数(按严重度分档)
|
|
}
|
|
|
|
// TaskListMetricsAll returns list aggregates for every live task in one query.
|
|
// The lateral lookups use the per-exploration indexes instead of grouping the
|
|
// complete activity history on every poll.
|
|
func (d *DB) TaskListMetricsAll() (map[int64]TaskListMetrics, error) {
|
|
rows, err := d.Query(`
|
|
SELECT task.exploration_id,
|
|
COALESCE(token_metrics.input_tokens,0),
|
|
COALESCE(token_metrics.output_tokens,0),
|
|
COALESCE(token_metrics.cache_read_tokens,0),
|
|
COALESCE(token_metrics.cache_write_tokens,0),
|
|
COALESCE(latest_activity.created_at,0),
|
|
COALESCE(goal_metrics.total,0),
|
|
COALESCE(goal_metrics.met,0),
|
|
COALESCE(intent_metrics.running,0),
|
|
COALESCE(finding_metrics.critical,0),
|
|
COALESCE(finding_metrics.high,0),
|
|
COALESCE(finding_metrics.medium,0),
|
|
COALESCE(finding_metrics.low,0)
|
|
FROM tasks task
|
|
LEFT JOIN LATERAL (
|
|
SELECT SUM(input_tokens) AS input_tokens,
|
|
SUM(output_tokens) AS output_tokens,
|
|
SUM(cache_read_tokens) AS cache_read_tokens,
|
|
SUM(cache_write_tokens) AS cache_write_tokens
|
|
FROM activity
|
|
WHERE exploration_id=task.exploration_id AND kind='result'
|
|
) token_metrics ON true
|
|
LEFT JOIN LATERAL (
|
|
SELECT EXTRACT(EPOCH FROM created_at)::bigint AS created_at
|
|
FROM activity
|
|
WHERE exploration_id=task.exploration_id
|
|
ORDER BY created_at DESC
|
|
LIMIT 1
|
|
) latest_activity ON true
|
|
LEFT JOIN LATERAL (
|
|
SELECT COUNT(*) AS total,
|
|
COUNT(*) FILTER (WHERE state='met') AS met
|
|
FROM exploration_nodes
|
|
WHERE exploration_id=task.exploration_id AND kind='goal'
|
|
) goal_metrics ON true
|
|
LEFT JOIN LATERAL (
|
|
SELECT COUNT(*) AS running
|
|
FROM exploration_nodes
|
|
WHERE exploration_id=task.exploration_id AND kind='intent' AND state='running'
|
|
) intent_metrics ON true
|
|
LEFT JOIN LATERAL (
|
|
SELECT COUNT(*) FILTER (WHERE severity='critical') AS critical,
|
|
COUNT(*) FILTER (WHERE severity='high') AS high,
|
|
COUNT(*) FILTER (WHERE severity='medium') AS medium,
|
|
COUNT(*) FILTER (WHERE severity='low') AS low
|
|
FROM findings
|
|
WHERE task_id=task.id
|
|
) finding_metrics ON true
|
|
WHERE task.deleted_at IS NULL`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := map[int64]TaskListMetrics{}
|
|
for rows.Next() {
|
|
var explorationID int64
|
|
var metrics TaskListMetrics
|
|
if err := rows.Scan(
|
|
&explorationID,
|
|
&metrics.Tokens.InputTokens,
|
|
&metrics.Tokens.OutputTokens,
|
|
&metrics.Tokens.CacheReadTokens,
|
|
&metrics.Tokens.CacheWriteTokens,
|
|
&metrics.LastActivity,
|
|
&metrics.Goals.Total,
|
|
&metrics.Goals.Met,
|
|
&metrics.RunningIntents,
|
|
&metrics.Findings.Critical,
|
|
&metrics.Findings.High,
|
|
&metrics.Findings.Medium,
|
|
&metrics.Findings.Low,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
out[explorationID] = metrics
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// GoalCountsAll returns goal totals (total / met) per exploration id, one query
|
|
// for all tasks — used by listTasks so the task list shows progress for every task.
|
|
func (d *DB) GoalCountsAll() (map[int64]GoalCounts, error) {
|
|
rows, err := d.Query(`
|
|
SELECT exploration_id,
|
|
COUNT(*) AS total,
|
|
COUNT(*) FILTER (WHERE state = 'met') AS met
|
|
FROM exploration_nodes
|
|
WHERE kind = 'goal'
|
|
GROUP BY exploration_id`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := map[int64]GoalCounts{}
|
|
for rows.Next() {
|
|
var eid int64
|
|
var gc GoalCounts
|
|
if err := rows.Scan(&eid, &gc.Total, &gc.Met); err != nil {
|
|
return nil, err
|
|
}
|
|
out[eid] = gc
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// TokenStatsByWorker sums token usage per worker for this exploration (from the
|
|
// kind='result' records that carry usage). Used by the per-agent token display.
|
|
func (s *ExplorationStore) TokenStatsByWorker() ([]TokenUsage, error) {
|
|
rows, err := s.db.Query(`SELECT COALESCE(worker,''),
|
|
COALESCE(SUM(input_tokens),0), COALESCE(SUM(output_tokens),0),
|
|
COALESCE(SUM(cache_read_tokens),0), COALESCE(SUM(cache_write_tokens),0)
|
|
FROM activity WHERE exploration_id=$1 AND kind='result'
|
|
GROUP BY worker ORDER BY worker`, s.expID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := []TokenUsage{}
|
|
for rows.Next() {
|
|
var u TokenUsage
|
|
if err := rows.Scan(&u.Worker, &u.InputTokens, &u.OutputTokens, &u.CacheReadTokens, &u.CacheWriteTokens); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, u)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// TokenStatsBySession aggregates every persisted completed run, independent of
|
|
// activity-history pagination. Main/planner use fixed keys; Worker executions use
|
|
// their intent node id, even when a different work#N process handled a retry.
|
|
func (s *ExplorationStore) TokenStatsBySession() ([]SessionTokenUsage, error) {
|
|
rows, err := s.db.Query(`SELECT
|
|
CASE
|
|
WHEN COALESCE(a.worker,'')='mainagent' THEN 'main:' || COALESCE(a.main_seg,0)::text
|
|
WHEN COALESCE(a.worker,'')='planner' THEN 'plan'
|
|
WHEN n.id IS NOT NULL THEN 'intent:' || n.id::text
|
|
ELSE ''
|
|
END AS session_key,
|
|
COALESCE(SUM(a.input_tokens),0), COALESCE(SUM(a.output_tokens),0),
|
|
COALESCE(SUM(a.cache_read_tokens),0), COALESCE(SUM(a.cache_write_tokens),0)
|
|
FROM activity a
|
|
LEFT JOIN exploration_nodes n
|
|
ON n.id=a.node_id AND n.exploration_id=a.exploration_id AND n.kind='intent'
|
|
WHERE a.exploration_id=$1 AND a.kind='result'
|
|
AND (COALESCE(a.worker,'') IN ('mainagent','planner') OR n.id IS NOT NULL)
|
|
GROUP BY session_key
|
|
ORDER BY session_key`, s.expID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := []SessionTokenUsage{}
|
|
for rows.Next() {
|
|
var u SessionTokenUsage
|
|
if err := rows.Scan(&u.Session, &u.InputTokens, &u.OutputTokens, &u.CacheReadTokens, &u.CacheWriteTokens); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, u)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// ActivityList returns steps after sinceID (exclusive), optionally filtered by node.
|
|
// Returns items and the new cursor (max id seen).
|
|
func (s *ExplorationStore) ActivityList(nodeID *int64, sinceID int64, limit int) ([]Activity, int64, error) {
|
|
if limit <= 0 {
|
|
limit = 300
|
|
}
|
|
var rows *sql.Rows
|
|
var err error
|
|
const cols = `id, node_id, COALESCE(worker,''), COALESCE(kind,''), COALESCE(tool,''), COALESCE(tool_use_id,''), is_error, COALESCE(summary,''), metadata, created_at, input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, main_seg`
|
|
if nodeID != nil {
|
|
rows, err = s.db.Query(`SELECT `+cols+`
|
|
FROM activity WHERE exploration_id=$1 AND node_id=$2 AND id>$3 ORDER BY id LIMIT $4`, s.expID, *nodeID, sinceID, limit)
|
|
} else {
|
|
rows, err = s.db.Query(`SELECT `+cols+`
|
|
FROM activity WHERE exploration_id=$1 AND id>$2 ORDER BY id LIMIT $3`, s.expID, sinceID, limit)
|
|
}
|
|
if err != nil {
|
|
return nil, sinceID, err
|
|
}
|
|
defer rows.Close()
|
|
out := []Activity{}
|
|
cursor := sinceID
|
|
for rows.Next() {
|
|
var a Activity
|
|
if err := rows.Scan(&a.ID, &a.NodeID, &a.Worker, &a.Kind, &a.Tool, &a.ToolUseID, &a.IsError, &a.Summary, &a.Metadata, &a.CreatedAt,
|
|
&a.InputTokens, &a.OutputTokens, &a.CacheReadTokens, &a.CacheWriteTokens, &a.MainSeg); err != nil {
|
|
return nil, sinceID, err
|
|
}
|
|
if a.ID > cursor {
|
|
cursor = a.ID
|
|
}
|
|
out = append(out, a)
|
|
}
|
|
return out, cursor, rows.Err()
|
|
}
|
|
|
|
// ActivityListForTerminalIntent is the inherited, incremental-read boundary.
|
|
// The intent state check and activity read happen in one statement so a source
|
|
// intent reopened between API-level checks cannot expose its new in-flight trace.
|
|
// Model reasoning and accounting rows are never inherited.
|
|
func (s *ExplorationStore) ActivityListForTerminalIntent(nodeID, sinceID int64, limit int) ([]Activity, int64, error) {
|
|
if limit <= 0 {
|
|
limit = 300
|
|
}
|
|
rows, err := s.db.Query(`SELECT a.id, a.node_id, COALESCE(a.worker,''), COALESCE(a.kind,''),
|
|
COALESCE(a.tool,''), COALESCE(a.tool_use_id,''), a.is_error, COALESCE(a.summary,''),
|
|
a.created_at, a.input_tokens, a.output_tokens, a.cache_read_tokens, a.cache_write_tokens
|
|
FROM activity a
|
|
JOIN exploration_nodes n ON n.id=a.node_id AND n.exploration_id=a.exploration_id
|
|
WHERE a.exploration_id=$1 AND a.node_id=$2 AND a.id>$3
|
|
AND n.kind='intent' AND n.state IN ('done','blocked','exhausted','stopped')
|
|
AND a.kind NOT IN ('thinking','usage')
|
|
ORDER BY a.id LIMIT $4`, s.expID, nodeID, sinceID, limit)
|
|
if err != nil {
|
|
return nil, sinceID, err
|
|
}
|
|
defer rows.Close()
|
|
out := []Activity{}
|
|
cursor := sinceID
|
|
for rows.Next() {
|
|
var a Activity
|
|
if err := rows.Scan(&a.ID, &a.NodeID, &a.Worker, &a.Kind, &a.Tool, &a.ToolUseID, &a.IsError, &a.Summary, &a.CreatedAt,
|
|
&a.InputTokens, &a.OutputTokens, &a.CacheReadTokens, &a.CacheWriteTokens); err != nil {
|
|
return nil, sinceID, err
|
|
}
|
|
if a.ID > cursor {
|
|
cursor = a.ID
|
|
}
|
|
out = append(out, a)
|
|
}
|
|
return out, cursor, rows.Err()
|
|
}
|
|
|
|
// ActivitySessionFilter selects one UI "session" within a task's activity stream.
|
|
// Exactly one field is meaningful:
|
|
// - Main == true → the main-agent session (worker="mainagent") for one segment
|
|
// (MainSeg; nil/0 = the original segment, which also matches legacy NULL rows).
|
|
// - Worker != "" → filter by worker name (Plan = "planner").
|
|
// Goal Agent 的第 0 轮拆解也以 worker="planner" 落库,故 Plan 会话完整覆盖 Goal+Planner。
|
|
// - NodeID != nil → a Worker session, filtered by node_id (= intent id).
|
|
//
|
|
// A zero value (all empty) matches the whole task (no session filter).
|
|
type ActivitySessionFilter struct {
|
|
Worker string
|
|
NodeID *int64
|
|
Main bool // main-agent session
|
|
MainSeg *int // segment for the main session (nil == current, resolved by caller; 0 == original)
|
|
}
|
|
|
|
func (f ActivitySessionFilter) cond(argStart int) (string, []any) {
|
|
switch {
|
|
case f.NodeID != nil:
|
|
return fmt.Sprintf(" AND node_id=$%d", argStart), []any{*f.NodeID}
|
|
case f.Main:
|
|
seg := 0
|
|
if f.MainSeg != nil {
|
|
seg = *f.MainSeg
|
|
}
|
|
if seg == 0 { // original segment also owns legacy rows with NULL main_seg
|
|
return " AND worker='mainagent' AND COALESCE(main_seg,0)=0", nil
|
|
}
|
|
return fmt.Sprintf(" AND worker='mainagent' AND main_seg=$%d", argStart), []any{seg}
|
|
case f.Worker != "":
|
|
return fmt.Sprintf(" AND worker=$%d", argStart), []any{f.Worker}
|
|
default:
|
|
return "", nil
|
|
}
|
|
}
|
|
|
|
// ActivityPage returns one page for reverse (newest-first) pagination of a session:
|
|
// up to `limit` steps ending before id `before` (exclusive; before<=0 = the latest
|
|
// page), returned in ASCENDING id order for direct display. hasMore reports whether
|
|
// still-older steps exist before the returned window (so the client can stop loading
|
|
// on scroll-up). Summary-only columns — detail is lazy-loaded via ActivityDetail.
|
|
func (s *ExplorationStore) ActivityPage(f ActivitySessionFilter, before int64, limit int) ([]Activity, bool, error) {
|
|
if limit <= 0 {
|
|
limit = 200
|
|
}
|
|
const cols = `id, node_id, COALESCE(worker,''), COALESCE(kind,''), COALESCE(tool,''), COALESCE(tool_use_id,''), is_error, COALESCE(summary,''), metadata, created_at, input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, main_seg`
|
|
args := []any{s.expID}
|
|
cond, cargs := f.cond(len(args) + 1)
|
|
args = append(args, cargs...)
|
|
beforeArg := len(args) + 1
|
|
args = append(args, before)
|
|
limitArg := len(args) + 1
|
|
args = append(args, limit+1) // one extra row to detect older history
|
|
q := fmt.Sprintf(`SELECT `+cols+`
|
|
FROM activity WHERE exploration_id=$1%s AND ($%d <= 0 OR id < $%d)
|
|
ORDER BY id DESC LIMIT $%d`, cond, beforeArg, beforeArg, limitArg)
|
|
rows, err := s.db.Query(q, args...)
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
defer rows.Close()
|
|
desc := []Activity{}
|
|
for rows.Next() {
|
|
var a Activity
|
|
if err := rows.Scan(&a.ID, &a.NodeID, &a.Worker, &a.Kind, &a.Tool, &a.ToolUseID, &a.IsError, &a.Summary, &a.Metadata, &a.CreatedAt,
|
|
&a.InputTokens, &a.OutputTokens, &a.CacheReadTokens, &a.CacheWriteTokens, &a.MainSeg); err != nil {
|
|
return nil, false, err
|
|
}
|
|
desc = append(desc, a)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, false, err
|
|
}
|
|
hasMore := len(desc) > limit
|
|
if hasMore {
|
|
desc = desc[:limit]
|
|
}
|
|
// reverse the newest-first window into ascending id order for display.
|
|
out := make([]Activity, len(desc))
|
|
for i, a := range desc {
|
|
out[len(desc)-1-i] = a
|
|
}
|
|
return out, hasMore, nil
|
|
}
|
|
|
|
// ActivityPageForTerminalIntent is the inherited session-history boundary. It
|
|
// combines ownership, terminal-state and row-kind checks in one query, closing
|
|
// the race where a previously completed source intent is reopened before paging.
|
|
func (s *ExplorationStore) ActivityPageForTerminalIntent(nodeID, before int64, limit int) ([]Activity, bool, error) {
|
|
if limit <= 0 {
|
|
limit = 200
|
|
}
|
|
rows, err := s.db.Query(`SELECT a.id, a.node_id, COALESCE(a.worker,''), COALESCE(a.kind,''),
|
|
COALESCE(a.tool,''), COALESCE(a.tool_use_id,''), a.is_error, COALESCE(a.summary,''),
|
|
a.created_at, a.input_tokens, a.output_tokens, a.cache_read_tokens, a.cache_write_tokens
|
|
FROM activity a
|
|
JOIN exploration_nodes n ON n.id=a.node_id AND n.exploration_id=a.exploration_id
|
|
WHERE a.exploration_id=$1 AND a.node_id=$2
|
|
AND n.kind='intent' AND n.state IN ('done','blocked','exhausted','stopped')
|
|
AND a.kind NOT IN ('thinking','usage')
|
|
AND ($3 <= 0 OR a.id < $3)
|
|
ORDER BY a.id DESC LIMIT $4`, s.expID, nodeID, before, limit+1)
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
defer rows.Close()
|
|
desc := []Activity{}
|
|
for rows.Next() {
|
|
var a Activity
|
|
if err := rows.Scan(&a.ID, &a.NodeID, &a.Worker, &a.Kind, &a.Tool, &a.ToolUseID, &a.IsError, &a.Summary, &a.CreatedAt,
|
|
&a.InputTokens, &a.OutputTokens, &a.CacheReadTokens, &a.CacheWriteTokens); err != nil {
|
|
return nil, false, err
|
|
}
|
|
desc = append(desc, a)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, false, err
|
|
}
|
|
hasMore := len(desc) > limit
|
|
if hasMore {
|
|
desc = desc[:limit]
|
|
}
|
|
out := make([]Activity, len(desc))
|
|
for i, a := range desc {
|
|
out[len(desc)-1-i] = a
|
|
}
|
|
return out, hasMore, nil
|
|
}
|
|
|
|
// ActivityMaxID returns the task's current maximum activity id (0 when empty). It is
|
|
// the task-level snapshot cursor handed to the SSE stream so history (id<=cursor) and
|
|
// the live tail (id>cursor) meet with no gap — deliberately task-level, not per
|
|
// session, since one SSE covers the whole task.
|
|
func (s *ExplorationStore) ActivityMaxID() (int64, error) {
|
|
var max sql.NullInt64
|
|
err := s.db.QueryRow(`SELECT MAX(id) FROM activity WHERE exploration_id=$1`, s.expID).Scan(&max)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return max.Int64, nil
|
|
}
|
|
|
|
// MainSession is one resettable main-agent conversation segment of a task. Segment 0
|
|
// is the original session (implicit, never stored); further segments are created by
|
|
// "新建会话" to start the main agent on a clean transcript while the task's graph,
|
|
// assets and goal stay shared.
|
|
type MainSession struct {
|
|
Seq int `json:"seq"`
|
|
CreatedAt time.Time `json:"created_at"`
|
|
}
|
|
|
|
// CurrentMainSeg returns the highest main-session segment (0 when none created yet).
|
|
func (s *ExplorationStore) CurrentMainSeg() (int, error) {
|
|
var seg int
|
|
err := s.db.QueryRow(`SELECT COALESCE(MAX(seq),0) FROM main_sessions WHERE exploration_id=$1`, s.expID).Scan(&seg)
|
|
return seg, err
|
|
}
|
|
|
|
// ListMainSessions returns every main-session segment newest-first, always including
|
|
// the implicit original segment 0. Segment 0 carries no stored timestamp.
|
|
func (s *ExplorationStore) ListMainSessions() ([]MainSession, error) {
|
|
rows, err := s.db.Query(`SELECT seq, created_at FROM main_sessions WHERE exploration_id=$1 ORDER BY seq DESC`, s.expID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := []MainSession{}
|
|
for rows.Next() {
|
|
var m MainSession
|
|
if err := rows.Scan(&m.Seq, &m.CreatedAt); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, m)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, MainSession{Seq: 0}) // implicit original session (oldest)
|
|
return out, nil
|
|
}
|
|
|
|
// NewMainSession creates the next main-session segment (seq = current+1) and returns
|
|
// it. It touches only main_sessions — the task's exploration graph/assets/goal are
|
|
// untouched, so the new session starts with a clean transcript over the same task.
|
|
func (s *ExplorationStore) NewMainSession() (MainSession, error) {
|
|
var m MainSession
|
|
err := s.db.QueryRow(`
|
|
INSERT INTO main_sessions(exploration_id, seq)
|
|
VALUES ($1, COALESCE((SELECT MAX(seq) FROM main_sessions WHERE exploration_id=$1),0)+1)
|
|
RETURNING seq, created_at`, s.expID).Scan(&m.Seq, &m.CreatedAt)
|
|
return m, err
|
|
}
|
|
|
|
// ActivityDetail lazily returns the full detail blob for one step.
|
|
func (s *ExplorationStore) ActivityDetail(id int64) (string, error) {
|
|
var d sql.NullString
|
|
err := s.db.QueryRow(`SELECT detail FROM activity WHERE id=$1 AND exploration_id=$2`, id, s.expID).Scan(&d)
|
|
if err == sql.ErrNoRows {
|
|
return "", nil
|
|
}
|
|
return d.String, err
|
|
}
|
|
|
|
// scanTrace reads the summary-only column set shared by the worker-trace tools
|
|
// (ActivityTrace / ActivityTraceSearch). Detail is never selected here — it is
|
|
// pulled separately, on demand, via ActivityByIDs.
|
|
func scanTrace(rows *sql.Rows) ([]Activity, error) {
|
|
out := []Activity{}
|
|
for rows.Next() {
|
|
var a Activity
|
|
if err := rows.Scan(&a.ID, &a.NodeID, &a.Worker, &a.Kind, &a.Tool, &a.IsError, &a.Summary); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, a)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
const traceCols = `id, node_id, COALESCE(worker,''), COALESCE(kind,''), COALESCE(tool,''), is_error, COALESCE(summary,'')`
|
|
|
|
// ActivityTrace lists one work's steps (by intent/node id) for the trace tools:
|
|
// summaries only (detail is lazy — fetch via ActivityByIDs). thinking/usage rows
|
|
// are dropped so the caller sees the work's actions + observations, not the
|
|
// model's internal reasoning or token-accounting noise.
|
|
func (s *ExplorationStore) ActivityTrace(nodeID int64, limit int) ([]Activity, error) {
|
|
if limit <= 0 {
|
|
limit = 500
|
|
}
|
|
rows, err := s.db.Query(`SELECT `+traceCols+`
|
|
FROM activity WHERE exploration_id=$1 AND node_id=$2 AND kind NOT IN ('thinking','usage')
|
|
ORDER BY id LIMIT $3`, s.expID, nodeID, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanTrace(rows)
|
|
}
|
|
|
|
// ActivityTraceForTerminalIntent returns one inherited work trace only while its
|
|
// source intent is still terminal at query time.
|
|
func (s *ExplorationStore) ActivityTraceForTerminalIntent(nodeID int64, limit int) ([]Activity, error) {
|
|
if limit <= 0 {
|
|
limit = 500
|
|
}
|
|
rows, err := s.db.Query(`SELECT a.id, a.node_id, COALESCE(a.worker,''), COALESCE(a.kind,''), COALESCE(a.tool,''), a.is_error, COALESCE(a.summary,'')
|
|
FROM activity a
|
|
JOIN exploration_nodes n ON n.id=a.node_id AND n.exploration_id=a.exploration_id
|
|
WHERE a.exploration_id=$1 AND a.node_id=$2 AND n.kind='intent'
|
|
AND n.state IN ('done','blocked','exhausted','stopped')
|
|
AND a.kind NOT IN ('thinking','usage')
|
|
ORDER BY a.id LIMIT $3`, s.expID, nodeID, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanTrace(rows)
|
|
}
|
|
|
|
// ActivityTraceSearch finds steps whose summary OR detail matches q
|
|
// (case-insensitive), returning summaries only. nodeID != nil scopes to one work;
|
|
// nil searches across ALL works (node_id IS NOT NULL — planner/main steps have no
|
|
// node_id, so this cleanly means "worker traces"). thinking/usage rows dropped.
|
|
func (s *ExplorationStore) ActivityTraceSearch(nodeID *int64, q string, limit int) ([]Activity, error) {
|
|
if limit <= 0 {
|
|
limit = 100
|
|
}
|
|
like := "%" + q + "%"
|
|
var rows *sql.Rows
|
|
var err error
|
|
if nodeID != nil {
|
|
rows, err = s.db.Query(`SELECT `+traceCols+`
|
|
FROM activity WHERE exploration_id=$1 AND node_id=$2 AND kind NOT IN ('thinking','usage')
|
|
AND (summary ILIKE $3 OR detail ILIKE $3) ORDER BY id LIMIT $4`, s.expID, *nodeID, like, limit)
|
|
} else {
|
|
rows, err = s.db.Query(`SELECT `+traceCols+`
|
|
FROM activity WHERE exploration_id=$1 AND node_id IS NOT NULL AND kind NOT IN ('thinking','usage')
|
|
AND (summary ILIKE $2 OR detail ILIKE $2) ORDER BY id LIMIT $3`, s.expID, like, limit)
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanTrace(rows)
|
|
}
|
|
|
|
// ActivityTraceSearchForTerminalIntent is the scoped inherited search variant;
|
|
// terminal eligibility is evaluated by the same statement that reads the rows.
|
|
func (s *ExplorationStore) ActivityTraceSearchForTerminalIntent(nodeID int64, q string, limit int) ([]Activity, error) {
|
|
if limit <= 0 {
|
|
limit = 100
|
|
}
|
|
like := "%" + q + "%"
|
|
rows, err := s.db.Query(`SELECT a.id, a.node_id, COALESCE(a.worker,''), COALESCE(a.kind,''), COALESCE(a.tool,''), a.is_error, COALESCE(a.summary,'')
|
|
FROM activity a
|
|
JOIN exploration_nodes n ON n.id=a.node_id AND n.exploration_id=a.exploration_id
|
|
WHERE a.exploration_id=$1 AND a.node_id=$2 AND n.kind='intent'
|
|
AND n.state IN ('done','blocked','exhausted','stopped')
|
|
AND a.kind NOT IN ('thinking','usage')
|
|
AND (a.summary ILIKE $3 OR a.detail ILIKE $3)
|
|
ORDER BY a.id LIMIT $4`, s.expID, nodeID, like, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanTrace(rows)
|
|
}
|
|
|
|
// ActivityTraceSearchTerminalIntents searches inherited worker history without
|
|
// exposing planner/main rows or work that is still open/running in the source.
|
|
func (s *ExplorationStore) ActivityTraceSearchTerminalIntents(q string, limit int) ([]Activity, error) {
|
|
if limit <= 0 {
|
|
limit = 100
|
|
}
|
|
like := "%" + q + "%"
|
|
rows, err := s.db.Query(`SELECT a.id, a.node_id, COALESCE(a.worker,''), COALESCE(a.kind,''), COALESCE(a.tool,''), a.is_error, COALESCE(a.summary,'')
|
|
FROM activity a
|
|
JOIN exploration_nodes n ON n.id=a.node_id AND n.exploration_id=a.exploration_id
|
|
WHERE a.exploration_id=$1 AND n.kind='intent'
|
|
AND n.state IN ('done','blocked','exhausted','stopped')
|
|
AND a.kind NOT IN ('thinking','usage')
|
|
AND (a.summary ILIKE $2 OR a.detail ILIKE $2)
|
|
ORDER BY a.id LIMIT $3`, s.expID, like, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanTrace(rows)
|
|
}
|
|
|
|
// AssetRef is a compact exploration node (intent/fact/finding) that references an
|
|
// asset via an anchor — for the coverage graph's node drawer.
|
|
type AssetRef struct {
|
|
ID int64 `json:"id"`
|
|
Kind string `json:"kind"`
|
|
State string `json:"state"`
|
|
Summary string `json:"summary"`
|
|
SourceTaskID int64 `json:"source_task_id,omitempty"`
|
|
Inherited bool `json:"inherited,omitempty"`
|
|
}
|
|
|
|
// AssetRefs returns the intents / facts / findings in THIS exploration anchored to
|
|
// the given asset id (newest first) — what this task tested / concluded about it.
|
|
func (s *ExplorationStore) AssetRefs(assetID int64) ([]AssetRef, error) {
|
|
if assetID <= 0 {
|
|
return []AssetRef{}, nil
|
|
}
|
|
rows, err := s.db.Query(`
|
|
SELECT en.id, en.kind, en.state, en.payload
|
|
FROM exploration_anchors ea
|
|
JOIN exploration_nodes en ON en.id = ea.node_id
|
|
WHERE ea.asset_id = $1 AND en.exploration_id = $2 AND en.kind IN ('intent','fact','finding')
|
|
ORDER BY en.id DESC`, assetID, s.expID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := []AssetRef{}
|
|
for rows.Next() {
|
|
var r AssetRef
|
|
var payload []byte
|
|
if err := rows.Scan(&r.ID, &r.Kind, &r.State, &payload); err != nil {
|
|
return nil, err
|
|
}
|
|
var p map[string]any
|
|
_ = json.Unmarshal(payload, &p)
|
|
for _, k := range []string{"summary", "text"} {
|
|
if v, ok := p[k].(string); ok && strings.TrimSpace(v) != "" {
|
|
r.Summary = v
|
|
break
|
|
}
|
|
}
|
|
out = append(out, r)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// ActivityTraceSearchExcluding keyword-searches every node's steps in this
|
|
// exploration EXCEPT one node's own (excludeNodeID) — so a worker's
|
|
// search_all_worker_traces doesn't return its own in-progress trace (which is
|
|
// already in its context). excludeNodeID<=0 → no exclusion (searches all nodes).
|
|
func (s *ExplorationStore) ActivityTraceSearchExcluding(excludeNodeID int64, q string, limit int) ([]Activity, error) {
|
|
if excludeNodeID <= 0 {
|
|
return s.ActivityTraceSearch(nil, q, limit)
|
|
}
|
|
if limit <= 0 {
|
|
limit = 100
|
|
}
|
|
like := "%" + q + "%"
|
|
rows, err := s.db.Query(`SELECT `+traceCols+`
|
|
FROM activity WHERE exploration_id=$1 AND node_id IS NOT NULL AND node_id <> $2
|
|
AND kind NOT IN ('thinking','usage')
|
|
AND (summary ILIKE $3 OR detail ILIKE $3) ORDER BY id LIMIT $4`, s.expID, excludeNodeID, like, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanTrace(rows)
|
|
}
|
|
|
|
// ActivityByIDs loads full detail for specific step ids (the trace drill-down),
|
|
// scoped to this exploration. thinking rows are skipped (never returned, even if
|
|
// their id is asked for). Order is ascending id, not the input order.
|
|
func (s *ExplorationStore) ActivityByIDs(ids []int64) ([]Activity, error) {
|
|
if len(ids) == 0 {
|
|
return nil, nil
|
|
}
|
|
ph := make([]string, len(ids))
|
|
args := make([]any, len(ids)+1)
|
|
args[0] = s.expID
|
|
for i, id := range ids {
|
|
ph[i] = fmt.Sprintf("$%d", i+2)
|
|
args[i+1] = id
|
|
}
|
|
rows, err := s.db.Query(`SELECT id, node_id, COALESCE(kind,''), COALESCE(tool,''), is_error, COALESCE(detail,'')
|
|
FROM activity WHERE exploration_id=$1 AND kind<>'thinking' AND id IN (`+strings.Join(ph, ",")+`) ORDER BY id`, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := []Activity{}
|
|
for rows.Next() {
|
|
var a Activity
|
|
if err := rows.Scan(&a.ID, &a.NodeID, &a.Kind, &a.Tool, &a.IsError, &a.Detail); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, a)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// ActivityByIDsForTerminalIntents is the inherited-detail read boundary. It
|
|
// deliberately excludes node-less planner/main rows, non-intent activities, live
|
|
// source work, and thinking. Local callers that need legacy behavior use
|
|
// ActivityByIDs instead.
|
|
func (s *ExplorationStore) ActivityByIDsForTerminalIntents(ids []int64) ([]Activity, error) {
|
|
if len(ids) == 0 {
|
|
return nil, nil
|
|
}
|
|
ph := make([]string, len(ids))
|
|
args := make([]any, len(ids)+1)
|
|
args[0] = s.expID
|
|
for i, id := range ids {
|
|
ph[i] = fmt.Sprintf("$%d", i+2)
|
|
args[i+1] = id
|
|
}
|
|
rows, err := s.db.Query(`SELECT a.id, a.node_id, COALESCE(a.kind,''), COALESCE(a.tool,''), a.is_error, COALESCE(a.detail,'')
|
|
FROM activity a
|
|
JOIN exploration_nodes n ON n.id=a.node_id AND n.exploration_id=a.exploration_id
|
|
WHERE a.exploration_id=$1 AND a.kind NOT IN ('thinking','usage') AND n.kind='intent'
|
|
AND n.state IN ('done','blocked','exhausted','stopped')
|
|
AND a.id IN (`+strings.Join(ph, ",")+`) ORDER BY a.id`, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := []Activity{}
|
|
for rows.Next() {
|
|
var a Activity
|
|
if err := rows.Scan(&a.ID, &a.NodeID, &a.Kind, &a.Tool, &a.IsError, &a.Detail); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, a)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// AddStandaloneFinding writes a finding to the standalone findings table, which
|
|
// persists across task deletion (task_id / node_id become NULL when the task or
|
|
// exploration node is deleted). taskID and nodeID may be 0 (stored as NULL).
|
|
func (s *ExplorationStore) AddStandaloneFinding(taskID, nodeID int64, vulnclass, name, severity, summary, evidence, worker string, assetIDs []int64) (int64, error) {
|
|
return s.db.AddFinding(taskID, nodeID, vulnclass, name, severity, summary, evidence, worker, assetIDs)
|
|
}
|