Files
goclaw/internal/store/pg/knowledge_graph_traversal.go
viettranx 2aa3ca31ae refactor(store/pg): error-propagate parseUUID in remaining kg + memory files
Five files, targeted CRITICAL-site migrations:

- knowledge_graph_temporal.go: SupersedeEntity (line 71 — C1 from red-team,
  UPDATE + INSERT in transaction). ListEntitiesTemporal remains on
  mustParseUUID (WARN-acceptable SELECT WHERE).
- memory_admin.go: ListAllDocuments, GetDocumentDetail, ListChunks (3 sites,
  all return errors). Adds fmt import.
- memory_search.go: Search (1 site, agent_id resolves hybrid FTS+vector).
- knowledge_graph_traversal.go: Traverse (2 sites — agent + startEntity).
- vault_documents_enrichment.go: UpdateSummaryAndReembed (2 sites — tenant,
  doc). FindSimilarDocs was migrated earlier with optAgentUUID wrapper fix.

Phase 4 Steps 4b.6, 4b.8, 4b.9, 4b.11, 4b.12 of agent identity hardening.

Still pending:
- vault_links.go (16 sites — next commit)
- knowledge_graph_embedding.go line 124 (SAFE per scout — audit only)
- agents_export_queries.go lines 312,354 (SAFE cursor per scout — audit only)
2026-04-11 21:22:23 +07:00

149 lines
4.5 KiB
Go

package pg
import (
"context"
"fmt"
"github.com/nextlevelbuilder/goclaw/internal/store"
)
// Traverse walks the knowledge graph from startEntityID up to maxDepth hops
// using a recursive CTE. Returns all reachable entities (excluding the start node).
// A 5-second statement timeout is applied for safety.
func (s *PGKnowledgeGraphStore) Traverse(ctx context.Context, agentID, userID, startEntityID string, maxDepth int) ([]store.TraversalResult, error) {
if maxDepth <= 0 {
maxDepth = 3
}
aid, err := parseUUID(agentID)
if err != nil {
return nil, fmt.Errorf("kg traverse: agent: %w", err)
}
startID, err := parseUUID(startEntityID)
if err != nil {
return nil, fmt.Errorf("kg traverse: start: %w", err)
}
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return nil, err
}
defer tx.Rollback() //nolint:errcheck
if _, err := tx.ExecContext(ctx, `SET LOCAL statement_timeout = '5000'`); err != nil {
return nil, err
}
var q string
var args []any
if store.IsSharedKG(ctx) {
// fixed params: $1=startID, $2=aid; tenant at $3 (if needed); maxDepth last
tc, tcArgs, _, tcErr := scopeClause(ctx, 3)
if tcErr != nil {
return nil, tcErr
}
depthN := 3 + len(tcArgs)
q = fmt.Sprintf(`
WITH RECURSIVE paths AS (
SELECT
e.id, e.agent_id, e.user_id, e.external_id,
e.name, e.entity_type, e.description,
e.properties, e.source_id, e.confidence,
e.created_at, e.updated_at,
1 AS depth,
ARRAY[e.id::text] AS path,
''::text AS via
FROM kg_entities e
WHERE e.id = $1 AND e.agent_id = $2 AND e.valid_until IS NULL%s
UNION ALL
SELECT
e.id, e.agent_id, e.user_id, e.external_id,
e.name, e.entity_type, e.description,
e.properties, e.source_id, e.confidence,
e.created_at, e.updated_at,
p.depth + 1,
p.path || e.id::text,
CASE WHEN r.source_entity_id = p.id
THEN r.relation_type
ELSE '~' || r.relation_type
END
FROM paths p
JOIN kg_relations r ON (r.source_entity_id = p.id OR r.target_entity_id = p.id) AND r.agent_id = $2 AND r.valid_until IS NULL
JOIN kg_entities e ON e.id = (CASE WHEN r.source_entity_id = p.id THEN r.target_entity_id ELSE r.source_entity_id END) AND e.agent_id = $2 AND e.valid_until IS NULL
WHERE p.depth < $%d
AND NOT e.id::text = ANY(p.path)
)
SELECT
id, agent_id, user_id, external_id,
name, entity_type, description,
properties, source_id, confidence,
created_at, updated_at,
depth, path, via
FROM paths WHERE depth > 1`, tc, depthN)
args = append([]any{startID, aid}, tcArgs...)
args = append(args, maxDepth)
} else {
// fixed params: $1=startID, $2=aid, $3=userID; tenant at $4 (if needed); maxDepth last
tc, tcArgs, _, tcErr := scopeClause(ctx, 4)
if tcErr != nil {
return nil, tcErr
}
depthN := 4 + len(tcArgs)
q = fmt.Sprintf(`
WITH RECURSIVE paths AS (
SELECT
e.id, e.agent_id, e.user_id, e.external_id,
e.name, e.entity_type, e.description,
e.properties, e.source_id, e.confidence,
e.created_at, e.updated_at,
1 AS depth,
ARRAY[e.id::text] AS path,
''::text AS via
FROM kg_entities e
WHERE e.id = $1 AND e.agent_id = $2 AND e.user_id = $3 AND e.valid_until IS NULL%s
UNION ALL
SELECT
e.id, e.agent_id, e.user_id, e.external_id,
e.name, e.entity_type, e.description,
e.properties, e.source_id, e.confidence,
e.created_at, e.updated_at,
p.depth + 1,
p.path || e.id::text,
CASE WHEN r.source_entity_id = p.id
THEN r.relation_type
ELSE '~' || r.relation_type
END
FROM paths p
JOIN kg_relations r ON (r.source_entity_id = p.id OR r.target_entity_id = p.id) AND r.user_id = $3 AND r.valid_until IS NULL
JOIN kg_entities e ON e.id = (CASE WHEN r.source_entity_id = p.id THEN r.target_entity_id ELSE r.source_entity_id END) AND e.user_id = $3 AND e.valid_until IS NULL
WHERE p.depth < $%d
AND NOT e.id::text = ANY(p.path)
)
SELECT
id, agent_id, user_id, external_id,
name, entity_type, description,
properties, source_id, confidence,
created_at, updated_at,
depth, path, via
FROM paths WHERE depth > 1`, tc, depthN)
args = append([]any{startID, aid, userID}, tcArgs...)
args = append(args, maxDepth)
}
// Use sqlx on the transaction for struct scanning with pq.StringArray support.
txSqlx := sqlxTx(tx)
var tRows []traversalRow
if err = txSqlx.SelectContext(ctx, &tRows, q, args...); err != nil {
return nil, err
}
results := make([]store.TraversalResult, len(tRows))
for i := range tRows {
results[i] = tRows[i].toTraversalResult()
}
return results, tx.Commit()
}