Files
artex/db/exploration.go
dela 0335d572de
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
First Commit
2026-10-09 08:38:16 +08:00

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)
}