feat(store): persist subagent tasks to PostgreSQL (#600)

- Migration 000034: subagent_tasks table with tenant scope, JSONB
  metadata + GIN index, partial index for archival candidates
- SubagentTaskStore interface with Create/Get/UpdateStatus/List/Archive
- PG implementation with parameterized queries and tenant isolation
- SQLite schema v3→4 migration + no-op stub for Lite edition
- Wire into store.Stores and factories
This commit is contained in:
viettranx
2026-03-31 11:45:03 +07:00
parent 7d35ee53c7
commit d8fc97ec63
11 changed files with 436 additions and 3 deletions
+1
View File
@@ -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
}
+216
View File
@@ -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()
}
+2 -1
View File
@@ -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
+31 -1
View File
@@ -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.
+35
View File
@@ -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);
@@ -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
}
+1
View File
@@ -32,4 +32,5 @@ type Stores struct {
BuiltinToolTenantCfgs BuiltinToolTenantConfigStore
SkillTenantCfgs SkillTenantConfigStore
SystemConfigs SystemConfigStore
SubagentTasks SubagentTaskStore
}
+62
View File
@@ -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
}
+1 -1
View File
@@ -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
@@ -0,0 +1 @@
DROP TABLE IF EXISTS subagent_tasks;
+48
View File
@@ -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;