diff --git a/internal/store/pg/factory.go b/internal/store/pg/factory.go index 8b9b693f..0e86cb01 100644 --- a/internal/store/pg/factory.go +++ b/internal/store/pg/factory.go @@ -50,5 +50,6 @@ func NewPGStores(cfg store.StoreConfig) (*store.Stores, error) { BuiltinToolTenantCfgs: NewPGBuiltinToolTenantConfigStore(db), SkillTenantCfgs: NewPGSkillTenantConfigStore(db), SystemConfigs: NewPGSystemConfigStore(db), + SubagentTasks: NewPGSubagentTaskStore(db), }, nil } diff --git a/internal/store/pg/subagent_tasks.go b/internal/store/pg/subagent_tasks.go new file mode 100644 index 00000000..70cc801b --- /dev/null +++ b/internal/store/pg/subagent_tasks.go @@ -0,0 +1,216 @@ +package pg + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/google/uuid" + + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +// PGSubagentTaskStore implements store.SubagentTaskStore using PostgreSQL. +type PGSubagentTaskStore struct { + db *sql.DB +} + +// NewPGSubagentTaskStore creates a new PostgreSQL-backed subagent task store. +func NewPGSubagentTaskStore(db *sql.DB) *PGSubagentTaskStore { + return &PGSubagentTaskStore{db: db} +} + +const subagentTaskInsertCols = `tenant_id, parent_agent_key, session_key, subject, description, + status, result, depth, model, provider, iterations, input_tokens, output_tokens, + origin_channel, origin_chat_id, origin_peer_kind, origin_user_id, spawned_by, metadata` + +// Create persists a new subagent task at spawn time. +func (s *PGSubagentTaskStore) Create(ctx context.Context, task *store.SubagentTaskData) error { + tid := tenantIDForInsert(ctx) + + metaJSON := []byte("{}") + if len(task.Metadata) > 0 { + if b, err := json.Marshal(task.Metadata); err == nil { + metaJSON = b + } + } + + q := fmt.Sprintf(`INSERT INTO subagent_tasks (id, %s) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20) + ON CONFLICT (id) DO NOTHING`, subagentTaskInsertCols) + + _, err := s.db.ExecContext(ctx, q, + task.ID, tid, task.ParentAgentKey, task.SessionKey, task.Subject, task.Description, + task.Status, task.Result, task.Depth, task.Model, task.Provider, + task.Iterations, task.InputTokens, task.OutputTokens, + task.OriginChannel, task.OriginChatID, task.OriginPeerKind, task.OriginUserID, + task.SpawnedBy, metaJSON, + ) + return err +} + +const subagentTaskSelectCols = `id, tenant_id, parent_agent_key, session_key, subject, description, + status, result, depth, model, provider, iterations, input_tokens, output_tokens, + origin_channel, origin_chat_id, origin_peer_kind, origin_user_id, spawned_by, + completed_at, archived_at, COALESCE(metadata, '{}'), created_at, updated_at` + +// scanTask scans a single row into SubagentTaskData. +func scanTask(row interface{ Scan(...any) error }) (*store.SubagentTaskData, error) { + var t store.SubagentTaskData + var metaJSON []byte + err := row.Scan( + &t.ID, &t.TenantID, &t.ParentAgentKey, &t.SessionKey, &t.Subject, &t.Description, + &t.Status, &t.Result, &t.Depth, &t.Model, &t.Provider, + &t.Iterations, &t.InputTokens, &t.OutputTokens, + &t.OriginChannel, &t.OriginChatID, &t.OriginPeerKind, &t.OriginUserID, &t.SpawnedBy, + &t.CompletedAt, &t.ArchivedAt, &metaJSON, &t.CreatedAt, &t.UpdatedAt, + ) + if err != nil { + return nil, err + } + if len(metaJSON) > 2 { // skip "{}" + _ = json.Unmarshal(metaJSON, &t.Metadata) + } + return &t, nil +} + +// Get retrieves a single task by ID (tenant-scoped). +func (s *PGSubagentTaskStore) Get(ctx context.Context, id uuid.UUID) (*store.SubagentTaskData, error) { + tid, err := requireTenantID(ctx) + if err != nil { + return nil, err + } + q := fmt.Sprintf(`SELECT %s FROM subagent_tasks WHERE id = $1 AND tenant_id = $2`, subagentTaskSelectCols) + row := s.db.QueryRowContext(ctx, q, id, tid) + t, err := scanTask(row) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + return t, err +} + +// UpdateStatus updates status, result, iterations, and token counts. +func (s *PGSubagentTaskStore) UpdateStatus( + ctx context.Context, id uuid.UUID, + status string, result *string, iterations int, + inputTokens, outputTokens int64, +) error { + tid, err := requireTenantID(ctx) + if err != nil { + return err + } + + var completedAt *time.Time + if status != "running" { + now := time.Now().UTC() + completedAt = &now + } + + q := `UPDATE subagent_tasks SET + status = $1, result = $2, iterations = $3, + input_tokens = $4, output_tokens = $5, + completed_at = $6, updated_at = NOW() + WHERE id = $7 AND tenant_id = $8` + _, err = s.db.ExecContext(ctx, q, + status, result, iterations, inputTokens, outputTokens, + completedAt, id, tid, + ) + return err +} + +// ListByParent returns tasks for a parent agent key, optionally filtered by status. +func (s *PGSubagentTaskStore) ListByParent( + ctx context.Context, parentAgentKey string, statusFilter string, +) ([]store.SubagentTaskData, error) { + tid, err := requireTenantID(ctx) + if err != nil { + return nil, err + } + + var rows *sql.Rows + if statusFilter != "" { + q := fmt.Sprintf(`SELECT %s FROM subagent_tasks + WHERE tenant_id = $1 AND parent_agent_key = $2 AND status = $3 + ORDER BY created_at DESC LIMIT 50`, subagentTaskSelectCols) + rows, err = s.db.QueryContext(ctx, q, tid, parentAgentKey, statusFilter) + } else { + q := fmt.Sprintf(`SELECT %s FROM subagent_tasks + WHERE tenant_id = $1 AND parent_agent_key = $2 + ORDER BY created_at DESC LIMIT 50`, subagentTaskSelectCols) + rows, err = s.db.QueryContext(ctx, q, tid, parentAgentKey) + } + if err != nil { + return nil, err + } + defer rows.Close() + + return collectTasks(rows) +} + +// ListBySession returns tasks for a specific session key (tenant-scoped). +func (s *PGSubagentTaskStore) ListBySession( + ctx context.Context, sessionKey string, +) ([]store.SubagentTaskData, error) { + tid, err := requireTenantID(ctx) + if err != nil { + return nil, err + } + + q := fmt.Sprintf(`SELECT %s FROM subagent_tasks + WHERE tenant_id = $1 AND session_key = $2 + ORDER BY created_at DESC LIMIT 50`, subagentTaskSelectCols) + rows, err := s.db.QueryContext(ctx, q, tid, sessionKey) + if err != nil { + return nil, err + } + defer rows.Close() + + return collectTasks(rows) +} + +// Archive marks old completed/failed/cancelled tasks as archived. +func (s *PGSubagentTaskStore) Archive(ctx context.Context, olderThan time.Duration) (int64, error) { + cutoff := time.Now().UTC().Add(-olderThan) + q := `UPDATE subagent_tasks SET archived_at = NOW(), updated_at = NOW() + WHERE status IN ('completed', 'failed', 'cancelled') + AND archived_at IS NULL AND completed_at < $1` + res, err := s.db.ExecContext(ctx, q, cutoff) + if err != nil { + return 0, err + } + return res.RowsAffected() +} + +// UpdateMetadata merges metadata on an existing task. +func (s *PGSubagentTaskStore) UpdateMetadata(ctx context.Context, id uuid.UUID, metadata map[string]any) error { + tid, err := requireTenantID(ctx) + if err != nil { + return err + } + + metaJSON, err := json.Marshal(metadata) + if err != nil { + return err + } + + q := `UPDATE subagent_tasks SET metadata = metadata || $1, updated_at = NOW() + WHERE id = $2 AND tenant_id = $3` + _, err = s.db.ExecContext(ctx, q, metaJSON, id, tid) + return err +} + +// collectTasks scans rows into a slice. +func collectTasks(rows *sql.Rows) ([]store.SubagentTaskData, error) { + var tasks []store.SubagentTaskData + for rows.Next() { + t, err := scanTask(rows) + if err != nil { + return nil, err + } + tasks = append(tasks, *t) + } + return tasks, rows.Err() +} diff --git a/internal/store/sqlitestore/factory.go b/internal/store/sqlitestore/factory.go index 88b60adf..4edea547 100644 --- a/internal/store/sqlitestore/factory.go +++ b/internal/store/sqlitestore/factory.go @@ -50,7 +50,8 @@ func NewSQLiteStores(cfg store.StoreConfig) (*store.Stores, error) { Activity: NewSQLiteActivityStore(db), APIKeys: NewSQLiteAPIKeyStore(db), ConfigPermissions: NewSQLiteConfigPermissionStore(db), - Memory: NewSQLiteMemoryStore(db), + Memory: NewSQLiteMemoryStore(db), + SubagentTasks: NewSQLiteSubagentTaskStore(), // Phase 2 Batch B+C stores (nil = gracefully skipped by gateway): // AgentLinks, KnowledgeGraph, SecureCLI }, nil diff --git a/internal/store/sqlitestore/schema.go b/internal/store/sqlitestore/schema.go index 84bcb808..cf8fd1e1 100644 --- a/internal/store/sqlitestore/schema.go +++ b/internal/store/sqlitestore/schema.go @@ -14,7 +14,7 @@ var schemaSQL string // SchemaVersion is the current SQLite schema version. // Bump this when adding new migration steps below. -const SchemaVersion = 3 +const SchemaVersion = 4 // migrations maps version → SQL to apply when upgrading FROM that version. // schema.sql always represents the LATEST full schema (for fresh DBs). @@ -42,6 +42,36 @@ UPDATE cron_jobs SET deliver_to = COALESCE(json_extract(payload, '$.to'), ''), wake_heartbeat = COALESCE(json_extract(payload, '$.wake_heartbeat'), 0) WHERE payload IS NOT NULL;`, + // Version 3 → 4: add subagent_tasks table for subagent lifecycle persistence. + 3: `CREATE TABLE IF NOT EXISTS subagent_tasks ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL REFERENCES tenants(id) ON DELETE CASCADE, + parent_agent_key VARCHAR(255) NOT NULL, + session_key VARCHAR(500), + subject VARCHAR(255) NOT NULL, + description TEXT NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'running', + result TEXT, + depth INTEGER NOT NULL DEFAULT 1, + model VARCHAR(255), + provider VARCHAR(255), + iterations INTEGER NOT NULL DEFAULT 0, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + origin_channel VARCHAR(50), + origin_chat_id VARCHAR(255), + origin_peer_kind VARCHAR(20), + origin_user_id VARCHAR(255), + spawned_by TEXT, + completed_at TEXT, + archived_at TEXT, + metadata TEXT NOT NULL DEFAULT '{}', + created_at TEXT DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + updated_at TEXT DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) +); +CREATE INDEX IF NOT EXISTS idx_subagent_tasks_parent_status ON subagent_tasks(tenant_id, parent_agent_key, status); +CREATE INDEX IF NOT EXISTS idx_subagent_tasks_session ON subagent_tasks(session_key); +CREATE INDEX IF NOT EXISTS idx_subagent_tasks_created ON subagent_tasks(tenant_id, created_at);`, } // EnsureSchema creates tables if they don't exist and applies incremental migrations. diff --git a/internal/store/sqlitestore/schema.sql b/internal/store/sqlitestore/schema.sql index d31c3003..69463644 100644 --- a/internal/store/sqlitestore/schema.sql +++ b/internal/store/sqlitestore/schema.sql @@ -1311,3 +1311,38 @@ CREATE TABLE IF NOT EXISTS system_configs ( ); CREATE INDEX IF NOT EXISTS idx_system_configs_tenant ON system_configs(tenant_id); + +-- ============================================================ +-- Table: subagent_tasks +-- ============================================================ + +CREATE TABLE IF NOT EXISTS subagent_tasks ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL REFERENCES tenants(id) ON DELETE CASCADE, + parent_agent_key VARCHAR(255) NOT NULL, + session_key VARCHAR(500), + subject VARCHAR(255) NOT NULL, + description TEXT NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'running', + result TEXT, + depth INTEGER NOT NULL DEFAULT 1, + model VARCHAR(255), + provider VARCHAR(255), + iterations INTEGER NOT NULL DEFAULT 0, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + origin_channel VARCHAR(50), + origin_chat_id VARCHAR(255), + origin_peer_kind VARCHAR(20), + origin_user_id VARCHAR(255), + spawned_by TEXT, + completed_at TEXT, + archived_at TEXT, + metadata TEXT NOT NULL DEFAULT '{}', + created_at TEXT DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + updated_at TEXT DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) +); + +CREATE INDEX IF NOT EXISTS idx_subagent_tasks_parent_status ON subagent_tasks(tenant_id, parent_agent_key, status); +CREATE INDEX IF NOT EXISTS idx_subagent_tasks_session ON subagent_tasks(session_key); +CREATE INDEX IF NOT EXISTS idx_subagent_tasks_created ON subagent_tasks(tenant_id, created_at); diff --git a/internal/store/sqlitestore/subagent_tasks.go b/internal/store/sqlitestore/subagent_tasks.go new file mode 100644 index 00000000..c09f39d4 --- /dev/null +++ b/internal/store/sqlitestore/subagent_tasks.go @@ -0,0 +1,38 @@ +//go:build sqlite || sqliteonly + +package sqlitestore + +import ( + "context" + "time" + + "github.com/google/uuid" + + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +// SQLiteSubagentTaskStore is a no-op implementation for the desktop/lite edition. +// Subagent task persistence is a standard-edition feature. +type SQLiteSubagentTaskStore struct{} + +func NewSQLiteSubagentTaskStore() *SQLiteSubagentTaskStore { return &SQLiteSubagentTaskStore{} } + +func (s *SQLiteSubagentTaskStore) Create(context.Context, *store.SubagentTaskData) error { return nil } +func (s *SQLiteSubagentTaskStore) Get(context.Context, uuid.UUID) (*store.SubagentTaskData, error) { + return nil, nil +} +func (s *SQLiteSubagentTaskStore) UpdateStatus(context.Context, uuid.UUID, string, *string, int, int64, int64) error { + return nil +} +func (s *SQLiteSubagentTaskStore) ListByParent(context.Context, string, string) ([]store.SubagentTaskData, error) { + return nil, nil +} +func (s *SQLiteSubagentTaskStore) ListBySession(context.Context, string) ([]store.SubagentTaskData, error) { + return nil, nil +} +func (s *SQLiteSubagentTaskStore) Archive(context.Context, time.Duration) (int64, error) { + return 0, nil +} +func (s *SQLiteSubagentTaskStore) UpdateMetadata(context.Context, uuid.UUID, map[string]any) error { + return nil +} diff --git a/internal/store/stores.go b/internal/store/stores.go index d7e016e9..9f951fb6 100644 --- a/internal/store/stores.go +++ b/internal/store/stores.go @@ -32,4 +32,5 @@ type Stores struct { BuiltinToolTenantCfgs BuiltinToolTenantConfigStore SkillTenantCfgs SkillTenantConfigStore SystemConfigs SystemConfigStore + SubagentTasks SubagentTaskStore } diff --git a/internal/store/subagent_store.go b/internal/store/subagent_store.go new file mode 100644 index 00000000..4adb16b9 --- /dev/null +++ b/internal/store/subagent_store.go @@ -0,0 +1,62 @@ +package store + +import ( + "context" + "time" + + "github.com/google/uuid" +) + +// SubagentTaskData represents a persisted subagent task for audit trail and cost attribution. +type SubagentTaskData struct { + BaseModel + TenantID uuid.UUID `json:"tenant_id"` + ParentAgentKey string `json:"parent_agent_key"` + SessionKey *string `json:"session_key,omitempty"` + Subject string `json:"subject"` + Description string `json:"description"` + Status string `json:"status"` + Result *string `json:"result,omitempty"` + Depth int `json:"depth"` + Model *string `json:"model,omitempty"` + Provider *string `json:"provider,omitempty"` + Iterations int `json:"iterations"` + InputTokens int64 `json:"input_tokens"` + OutputTokens int64 `json:"output_tokens"` + OriginChannel *string `json:"origin_channel,omitempty"` + OriginChatID *string `json:"origin_chat_id,omitempty"` + OriginPeerKind *string `json:"origin_peer_kind,omitempty"` + OriginUserID *string `json:"origin_user_id,omitempty"` + SpawnedBy *uuid.UUID `json:"spawned_by,omitempty"` + CompletedAt *time.Time `json:"completed_at,omitempty"` + ArchivedAt *time.Time `json:"archived_at,omitempty"` + Metadata map[string]any `json:"metadata,omitempty"` +} + +// SubagentTaskStore persists subagent task lifecycle for audit trail and cost attribution. +// In-memory SubagentManager remains the source of truth for active operations; +// DB writes are fire-and-forget (non-blocking). +type SubagentTaskStore interface { + // Create persists a new subagent task at spawn time. + Create(ctx context.Context, task *SubagentTaskData) error + + // Get retrieves a single task by ID (tenant-scoped). + Get(ctx context.Context, id uuid.UUID) (*SubagentTaskData, error) + + // UpdateStatus updates status, result, iterations, and token counts on completion/failure. + UpdateStatus(ctx context.Context, id uuid.UUID, status string, result *string, iterations int, inputTokens, outputTokens int64) error + + // ListByParent returns tasks for a parent agent key, optionally filtered by status. + // Empty statusFilter returns all statuses. Ordered by created_at DESC. + ListByParent(ctx context.Context, parentAgentKey string, statusFilter string) ([]SubagentTaskData, error) + + // ListBySession returns tasks for a specific session key (tenant-scoped). + ListBySession(ctx context.Context, sessionKey string) ([]SubagentTaskData, error) + + // Archive marks old completed/failed/cancelled tasks as archived. + // Returns the number of rows affected. + Archive(ctx context.Context, olderThan time.Duration) (int64, error) + + // UpdateMetadata merges metadata on an existing task. + UpdateMetadata(ctx context.Context, id uuid.UUID, metadata map[string]any) error +} diff --git a/internal/upgrade/version.go b/internal/upgrade/version.go index fde61142..ed890c12 100644 --- a/internal/upgrade/version.go +++ b/internal/upgrade/version.go @@ -2,4 +2,4 @@ package upgrade // RequiredSchemaVersion is the schema migration version this binary requires. // Bump this whenever adding a new SQL migration file. -const RequiredSchemaVersion uint = 33 +const RequiredSchemaVersion uint = 34 diff --git a/migrations/000034_subagent_tasks.down.sql b/migrations/000034_subagent_tasks.down.sql new file mode 100644 index 00000000..7f708575 --- /dev/null +++ b/migrations/000034_subagent_tasks.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS subagent_tasks; diff --git a/migrations/000034_subagent_tasks.up.sql b/migrations/000034_subagent_tasks.up.sql new file mode 100644 index 00000000..20200114 --- /dev/null +++ b/migrations/000034_subagent_tasks.up.sql @@ -0,0 +1,48 @@ +-- Persist subagent task lifecycle for audit trail, cost attribution, and restart recovery. +CREATE TABLE IF NOT EXISTS subagent_tasks ( + id UUID PRIMARY KEY DEFAULT uuid_generate_v7(), + tenant_id UUID NOT NULL REFERENCES tenants(id) ON DELETE CASCADE, + parent_agent_key VARCHAR(255) NOT NULL, + session_key VARCHAR(500), + subject VARCHAR(255) NOT NULL, + description TEXT NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'running', + result TEXT, + depth INT NOT NULL DEFAULT 1, + model VARCHAR(255), + provider VARCHAR(255), + iterations INT NOT NULL DEFAULT 0, + input_tokens BIGINT NOT NULL DEFAULT 0, + output_tokens BIGINT NOT NULL DEFAULT 0, + origin_channel VARCHAR(50), + origin_chat_id VARCHAR(255), + origin_peer_kind VARCHAR(20), + origin_user_id VARCHAR(255), + spawned_by UUID, + completed_at TIMESTAMPTZ, + archived_at TIMESTAMPTZ, + metadata JSONB NOT NULL DEFAULT '{}', + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +-- Primary lookup: roster by parent + status. +CREATE INDEX idx_subagent_tasks_parent_status + ON subagent_tasks(tenant_id, parent_agent_key, status); + +-- Session-scoped lookup. +CREATE INDEX idx_subagent_tasks_session + ON subagent_tasks(session_key) WHERE session_key IS NOT NULL; + +-- Time-based audit & cleanup. +CREATE INDEX idx_subagent_tasks_created + ON subagent_tasks(tenant_id, created_at DESC); + +-- Flexible metadata queries. +CREATE INDEX idx_subagent_tasks_metadata_gin + ON subagent_tasks USING GIN (metadata); + +-- Archival candidates. +CREATE INDEX idx_subagent_tasks_archive + ON subagent_tasks(status, completed_at) + WHERE status IN ('completed', 'failed', 'cancelled') AND archived_at IS NULL;