diff --git a/CLAUDE.md b/CLAUDE.md index f3032dc..c5b470e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -108,6 +108,8 @@ make lint # Run linters - SQLite (modernc.org/sqlite, pure Go) — new migration 013_reactions.sql (010-reactions-workflows) - Go 1.25+ (SynapBus), Python 3.12 (Searcher agents) + go-chi/chi, mark3labs/mcp-go, ory/fosite (SynapBus); claude-agent-sdk, httpx, psycopg (Searcher) (013-linkedin-approval-workflow) - SQLite via modernc.org/sqlite (SynapBus); PostgreSQL (Searcher) (013-linkedin-approval-workflow) +- Go 1.25+ (per go.mod) + go-chi/chi (HTTP), mark3labs/mcp-go (MCP), spf13/cobra (CLI), modernc.org/sqlite (storage), k8s.io/client-go (K8s Jobs) (014-reactive-agent-triggers) +- SQLite via modernc.org/sqlite — new migration 015_reactive_triggers.sql (014-reactive-agent-triggers) ## Recent Changes - 002-mcp-auth-ux-polish: Added Go 1.23+ + ory/fosite (OAuth 2.1), mark3labs/mcp-go (MCP server), go-chi/chi (HTTP), Svelte 5 + Tailwind (Web UI) diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index 8efbb53..b0a4219 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -39,6 +39,8 @@ import ( "github.com/synapbus/synapbus/internal/jsruntime" k8spkg "github.com/synapbus/synapbus/internal/k8s" mcpserver "github.com/synapbus/synapbus/internal/mcp" + "github.com/synapbus/synapbus/internal/agentquery" + reactorpkg "github.com/synapbus/synapbus/internal/reactor" "github.com/synapbus/synapbus/internal/messaging" prommetrics "github.com/synapbus/synapbus/internal/metrics" "github.com/synapbus/synapbus/internal/reactions" @@ -467,10 +469,21 @@ func runServe(cmd *cobra.Command, args []string) error { slog.Info("K8s job runner not available (not in-cluster)") } - // Create event dispatcher (fans out to webhooks + K8s) - eventDispatcher := dispatcher.NewMultiDispatcher(slog.Default(), deliveryEngine, k8sDispatcher) + // Create reactor engine for reactive agent triggering + reactorStore := reactorpkg.NewStore(db.DB) + reactorEngine := reactorpkg.New(reactorStore, agentStore, k8sRunner, slog.Default()) + reactorNotifier := reactorpkg.NewDMFailureNotifier(msgService) + reactorEngine.SetFailureNotifier(reactorNotifier) + + // Create event dispatcher (fans out to webhooks + K8s + reactor) + eventDispatcher := dispatcher.NewMultiDispatcher(slog.Default(), deliveryEngine, k8sDispatcher, reactorEngine) msgService.SetDispatcher(eventDispatcher) + // Start reactor poller for K8s Job status tracking + reactorPoller := reactorpkg.NewPoller(reactorStore, agentStore, k8sRunner, reactorEngine, slog.Default()) + reactorPoller.Start() + slog.Info("reactor engine and poller started") + // Create JS runtime pool and action registry for hybrid MCP tools jsPool := jsruntime.NewPool(10) defer jsPool.Close() @@ -480,6 +493,13 @@ func runServe(cmd *cobra.Command, args []string) error { // Create MCP server (4 hybrid tools: my_status, send_message, search, execute) mcpSrv := mcpserver.NewMCPServer(msgService, agentService, channelService, swarmService, attachmentService, searchService, reactionService, trustService, con, jsPool, actionRegistry, actionIndex, db.DB) + + // Set up SQL query executor for agents (uses read pool if available) + queryDB := db.QueryDB() + queryExec := agentquery.New(queryDB, slog.Default()) + mcpSrv.SetQueryExecutor(queryExec) + slog.Info("agent SQL query executor initialized", "read_pool", db.ReadDB != nil) + startTime := time.Now() // Start task expiry worker @@ -641,6 +661,8 @@ func runServe(cmd *cobra.Command, args []string) error { Version: version, PushService: pushService, TrustService: trustService, + ReactorStore: reactorStore, + ReactorEngine: reactorEngine, BaseURL: baseURL, }) r.Mount("/", apiRouter) diff --git a/internal/actions/registry.go b/internal/actions/registry.go index f17c8ef..9ecda51 100644 --- a/internal/actions/registry.go +++ b/internal/actions/registry.go @@ -571,5 +571,33 @@ func allActions() []Action { }, }, }, + // ── SQL Query (1 action) ──────────────────────────────────── + { + Name: "query", + Category: "data", + Description: "Execute a read-only SQL query against your accessible messages, channels, and reactions. Use tables: my_messages (your DMs + joined channels), my_channels (channels you are in), channel_messages (messages in your channels). Results are limited to 100 rows. Only SELECT statements are allowed.", + Params: []Param{ + {Name: "sql", Type: "string", Description: "SQL SELECT query. Available tables: my_messages (id, body, from_agent, to_agent, priority, status, metadata, created_at, channel_name), my_channels (id, name, description, type), channel_messages (id, body, from_agent, priority, channel_name, created_at). CTEs (WITH) are supported.", Required: true}, + }, + Returns: "JSON with columns (array of column names), rows (array of row arrays), row_count, and truncated (boolean if > 100 rows)", + Examples: []Example{ + { + Description: "Find high-priority messages in a channel", + Code: `call("query", {"sql": "SELECT id, body, from_agent, priority FROM channel_messages WHERE channel_name = 'news-mcpproxy' AND priority >= 7 ORDER BY created_at DESC LIMIT 10"})`, + }, + { + Description: "List your channels", + Code: `call("query", {"sql": "SELECT name, description FROM my_channels ORDER BY name"})`, + }, + { + Description: "Count messages per channel", + Code: `call("query", {"sql": "SELECT channel_name, COUNT(*) as msg_count FROM channel_messages GROUP BY channel_name ORDER BY msg_count DESC"})`, + }, + { + Description: "Search messages with keyword", + Code: `call("query", {"sql": "SELECT id, body, from_agent, created_at FROM my_messages WHERE body LIKE '%MCP%' ORDER BY created_at DESC LIMIT 20"})`, + }, + }, + }, } } diff --git a/internal/actions/registry_test.go b/internal/actions/registry_test.go index 57aff94..1925e62 100644 --- a/internal/actions/registry_test.go +++ b/internal/actions/registry_test.go @@ -4,11 +4,11 @@ import ( "testing" ) -func TestRegistryHas29Actions(t *testing.T) { +func TestRegistryHas30Actions(t *testing.T) { r := NewRegistry() got := len(r.List()) - if got != 29 { - t.Errorf("expected 29 actions, got %d", got) + if got != 30 { + t.Errorf("expected 30 actions, got %d", got) } } @@ -58,6 +58,8 @@ func TestRegistryGetByName(t *testing.T) { "get_replies", // trust "get_trust", + // data + "query", } for _, name := range allNames { diff --git a/internal/agentquery/executor.go b/internal/agentquery/executor.go new file mode 100644 index 0000000..b7ad440 --- /dev/null +++ b/internal/agentquery/executor.go @@ -0,0 +1,234 @@ +// Package agentquery provides a sandboxed SQL query executor for agents. +// Agents can run read-only SELECT queries against curated views with +// per-agent access control, automatic LIMIT enforcement, and timeouts. +package agentquery + +import ( + "context" + "database/sql" + "fmt" + "log/slog" + "strings" + "time" +) + +const ( + // MaxRows is the maximum number of rows returned by a query. + MaxRows = 100 + // QueryTimeout is the maximum duration for a query. + QueryTimeout = 5 * time.Second +) + +// Allowed view names that agents can query. +var allowedTables = map[string]bool{ + "my_messages": true, + "my_channels": true, + "channel_messages": true, +} + +// Executor runs sandboxed SQL queries on behalf of agents. +type Executor struct { + db *sql.DB // read-only pool (query_only=ON) + logger *slog.Logger +} + +// New creates a new query executor using the provided read-only database connection. +func New(readDB *sql.DB, logger *slog.Logger) *Executor { + return &Executor{ + db: readDB, + logger: logger.With("component", "agentquery"), + } +} + +// QueryResult holds the results of a SQL query. +type QueryResult struct { + Columns []string `json:"columns"` + Rows [][]interface{} `json:"rows"` + RowCount int `json:"row_count"` + Truncated bool `json:"truncated"` +} + +// Execute runs a SQL query on behalf of an agent with access control. +func (e *Executor) Execute(ctx context.Context, agentName, sqlQuery string) (*QueryResult, error) { + // 1. Validate the SQL statement + if err := validateSQL(sqlQuery); err != nil { + return nil, fmt.Errorf("query validation failed: %w", err) + } + + // 2. Rewrite the query to inject access control and enforce LIMIT + rewritten := rewriteQuery(agentName, sqlQuery) + + // 3. Execute with timeout + queryCtx, cancel := context.WithTimeout(ctx, QueryTimeout) + defer cancel() + + rows, err := e.db.QueryContext(queryCtx, rewritten) + if err != nil { + if queryCtx.Err() == context.DeadlineExceeded { + return nil, fmt.Errorf("query timed out after %s", QueryTimeout) + } + return nil, fmt.Errorf("query execution failed: %w", err) + } + defer rows.Close() + + // 4. Collect results + columns, err := rows.Columns() + if err != nil { + return nil, fmt.Errorf("get columns: %w", err) + } + + var resultRows [][]interface{} + truncated := false + + for rows.Next() { + if len(resultRows) >= MaxRows { + truncated = true + break + } + + values := make([]interface{}, len(columns)) + scanArgs := make([]interface{}, len(columns)) + for i := range values { + scanArgs[i] = &values[i] + } + + if err := rows.Scan(scanArgs...); err != nil { + return nil, fmt.Errorf("scan row: %w", err) + } + + // Convert []byte to string for JSON serialization + row := make([]interface{}, len(columns)) + for i, v := range values { + if b, ok := v.([]byte); ok { + row[i] = string(b) + } else { + row[i] = v + } + } + resultRows = append(resultRows, row) + } + + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate rows: %w", err) + } + + if resultRows == nil { + resultRows = [][]interface{}{} + } + + e.logger.Info("agent query executed", + "agent", agentName, + "rows", len(resultRows), + "truncated", truncated, + ) + + return &QueryResult{ + Columns: columns, + Rows: resultRows, + RowCount: len(resultRows), + Truncated: truncated, + }, nil +} + +// validateSQL checks that the query is a read-only SELECT statement. +func validateSQL(query string) error { + trimmed := strings.TrimSpace(query) + if trimmed == "" { + return fmt.Errorf("empty query") + } + + // Remove comments + upper := strings.ToUpper(trimmed) + + // Must start with SELECT or WITH (CTEs) + if !strings.HasPrefix(upper, "SELECT") && !strings.HasPrefix(upper, "WITH") { + return fmt.Errorf("only SELECT statements are allowed (got %q)", firstWord(upper)) + } + + // Block dangerous keywords (check as whole words or with common delimiters) + blocked := []string{ + "INSERT ", "UPDATE ", "DELETE ", "DROP ", "ALTER ", "CREATE ", + "ATTACH ", "DETACH ", "PRAGMA", "REINDEX ", "VACUUM ", + "REPLACE ", "GRANT ", "REVOKE ", + } + for _, kw := range blocked { + if strings.Contains(upper, kw) { + return fmt.Errorf("statement contains blocked keyword: %s", strings.TrimSpace(kw)) + } + } + + // Block multiple statements (semicolon followed by non-whitespace) + parts := strings.Split(trimmed, ";") + nonEmpty := 0 + for _, p := range parts { + if strings.TrimSpace(p) != "" { + nonEmpty++ + } + } + if nonEmpty > 1 { + return fmt.Errorf("multiple statements not allowed") + } + + return nil +} + +// rewriteQuery wraps the agent's query with access control CTEs. +// It replaces references to my_messages, my_channels, channel_messages +// with CTEs that filter by the agent's access. +func rewriteQuery(agentName, query string) string { + // Build access-control CTEs that the agent's query can reference + cte := fmt.Sprintf(` +WITH my_messages AS ( + SELECT v.* FROM v_agent_messages v + LEFT JOIN channel_members cm ON cm.channel_id = v.channel_id AND cm.agent_name = %[1]s + WHERE v.to_agent = %[1]s + OR v.from_agent = %[1]s + OR (v.channel_id IS NOT NULL AND cm.agent_name IS NOT NULL) +), +my_channels AS ( + SELECT c.id, c.name, c.description, c.type, c.topic, c.is_private, c.created_at, + cm.joined_at AS member_since + FROM channels c + JOIN channel_members cm ON cm.channel_id = c.id AND cm.agent_name = %[1]s +), +channel_messages AS ( + SELECT v.* FROM v_channel_messages v + WHERE v.channel_id IN ( + SELECT channel_id FROM channel_members WHERE agent_name = %[1]s + ) +) +`, quoteSQLString(agentName)) + + trimmed := strings.TrimSpace(query) + upper := strings.ToUpper(trimmed) + + // Remove trailing semicolon if present + trimmed = strings.TrimRight(trimmed, "; \t\n") + + if strings.HasPrefix(upper, "WITH") { + // User has their own CTEs. Merge: our CTEs first, then theirs. + userCTEs := strings.TrimSpace(trimmed[4:]) // skip "WITH" + return cte + ", " + userCTEs + } + + // Simple SELECT — prepend our CTEs + return cte + trimmed +} + +// quoteSQLString safely quotes a string for use in SQL. +func quoteSQLString(s string) string { + escaped := strings.ReplaceAll(s, "'", "''") + return "'" + escaped + "'" +} + +func firstWord(s string) string { + for i, c := range s { + if c == ' ' || c == '\t' || c == '\n' || c == '\r' || c == '(' { + return s[:i] + } + } + if len(s) > 20 { + return s[:20] + } + return s +} diff --git a/internal/agentquery/executor_test.go b/internal/agentquery/executor_test.go new file mode 100644 index 0000000..0a8c863 --- /dev/null +++ b/internal/agentquery/executor_test.go @@ -0,0 +1,341 @@ +package agentquery + +import ( + "context" + "database/sql" + "log/slog" + "testing" + + _ "modernc.org/sqlite" +) + +func setupTestDB(t *testing.T) *sql.DB { + t.Helper() + db, err := sql.Open("sqlite", ":memory:") + if err != nil { + t.Fatalf("open db: %v", err) + } + + // Create the schema needed for views + schema := ` + CREATE TABLE channels ( + id INTEGER PRIMARY KEY, + name TEXT NOT NULL UNIQUE, + description TEXT DEFAULT '', + type TEXT DEFAULT 'standard', + topic TEXT DEFAULT '', + is_private INTEGER DEFAULT 0, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP + ); + CREATE TABLE channel_members ( + channel_id INTEGER, + agent_name TEXT, + joined_at DATETIME DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (channel_id, agent_name) + ); + CREATE TABLE messages ( + id INTEGER PRIMARY KEY, + conversation_id INTEGER DEFAULT 0, + from_agent TEXT, + to_agent TEXT, + channel_id INTEGER, + reply_to INTEGER, + body TEXT, + priority INTEGER DEFAULT 5, + status TEXT DEFAULT 'pending', + metadata TEXT DEFAULT '{}', + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP + ); + + -- Views matching the migration + CREATE VIEW v_agent_messages AS + SELECT m.id, m.body, m.from_agent, m.to_agent, m.priority, m.status, m.metadata, + m.created_at, m.updated_at, c.name AS channel_name, m.channel_id, m.reply_to, m.conversation_id + FROM messages m LEFT JOIN channels c ON c.id = m.channel_id; + + CREATE VIEW v_agent_channels AS + SELECT c.id, c.name, c.description, c.type, c.topic, c.is_private, c.created_at, + cm.joined_at AS member_since + FROM channels c JOIN channel_members cm ON cm.channel_id = c.id; + + CREATE VIEW v_channel_messages AS + SELECT m.id, m.body, m.from_agent, m.priority, m.status, m.metadata, m.created_at, + c.name AS channel_name, m.channel_id, m.reply_to + FROM messages m JOIN channels c ON c.id = m.channel_id; + ` + if _, err := db.Exec(schema); err != nil { + t.Fatalf("create schema: %v", err) + } + + // Seed test data + seed := ` + INSERT INTO channels (id, name) VALUES (1, 'general'), (2, 'news-mcpproxy'), (3, 'private-channel'); + INSERT INTO channel_members (channel_id, agent_name) VALUES + (1, 'agent-a'), (1, 'agent-b'), + (2, 'agent-a'), + (3, 'agent-b'); + + -- DMs + INSERT INTO messages (id, from_agent, to_agent, body, priority) VALUES + (1, 'algis', 'agent-a', 'Hello agent A', 7), + (2, 'agent-a', 'algis', 'Hi there', 5), + (3, 'algis', 'agent-b', 'Hello agent B', 5); + + -- Channel messages + INSERT INTO messages (id, from_agent, channel_id, body, priority) VALUES + (4, 'agent-a', 1, 'General post from A', 5), + (5, 'agent-b', 1, 'General post from B', 5), + (6, 'agent-a', 2, 'News post high prio', 8), + (7, 'agent-b', 3, 'Private channel msg', 5); + ` + if _, err := db.Exec(seed); err != nil { + t.Fatalf("seed data: %v", err) + } + + return db +} + +func TestExecuteBasicQuery(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + result, err := exec.Execute(context.Background(), "agent-a", + "SELECT id, body, priority FROM my_messages ORDER BY id") + if err != nil { + t.Fatalf("query failed: %v", err) + } + + if len(result.Columns) != 3 { + t.Errorf("expected 3 columns, got %d", len(result.Columns)) + } + if result.Columns[0] != "id" || result.Columns[1] != "body" || result.Columns[2] != "priority" { + t.Errorf("unexpected columns: %v", result.Columns) + } + + // agent-a should see: DM to it (1), DM from it (2), general posts (4,5), news post (6) + // Should NOT see: DM to agent-b (3), private channel msg (7) + if result.RowCount < 4 { + t.Errorf("expected at least 4 rows for agent-a, got %d", result.RowCount) + } + + // Verify agent-b's DM and private channel msg are NOT visible + for _, row := range result.Rows { + id := row[0] + if id == int64(3) { + t.Error("agent-a should NOT see message 3 (DM to agent-b)") + } + if id == int64(7) { + t.Error("agent-a should NOT see message 7 (private channel, not joined)") + } + } +} + +func TestAccessControlAgentB(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + result, err := exec.Execute(context.Background(), "agent-b", + "SELECT id, body FROM my_messages ORDER BY id") + if err != nil { + t.Fatalf("query failed: %v", err) + } + + // agent-b should see: DM to it (3), general posts (4,5), private channel (7) + // Should NOT see: DM to agent-a (1), DM from agent-a (2), news post (6) + hasMsg3 := false + hasMsg7 := false + for _, row := range result.Rows { + id := row[0] + if id == int64(3) { + hasMsg3 = true + } + if id == int64(7) { + hasMsg7 = true + } + if id == int64(1) { + t.Error("agent-b should NOT see message 1 (DM to agent-a)") + } + if id == int64(6) { + t.Error("agent-b should NOT see message 6 (news channel, not joined)") + } + } + if !hasMsg3 { + t.Error("agent-b should see message 3 (DM to it)") + } + if !hasMsg7 { + t.Error("agent-b should see message 7 (private channel, joined)") + } +} + +func TestQueryChannelMessages(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + result, err := exec.Execute(context.Background(), "agent-a", + "SELECT id, body, channel_name FROM channel_messages WHERE channel_name = 'news-mcpproxy'") + if err != nil { + t.Fatalf("query failed: %v", err) + } + + if result.RowCount != 1 { + t.Errorf("expected 1 news message, got %d", result.RowCount) + } +} + +func TestQueryMyChannels(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + result, err := exec.Execute(context.Background(), "agent-a", + "SELECT name FROM my_channels ORDER BY name") + if err != nil { + t.Fatalf("query failed: %v", err) + } + + // agent-a is in: general, news-mcpproxy (not private-channel) + if result.RowCount != 2 { + t.Errorf("expected 2 channels for agent-a, got %d", result.RowCount) + } +} + +func TestValidationRejectsInsert(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + _, err := exec.Execute(context.Background(), "agent-a", + "INSERT INTO messages (body) VALUES ('evil')") + if err == nil { + t.Fatal("expected INSERT to be rejected") + } + if !contains(err.Error(), "only SELECT") { + t.Errorf("expected 'only SELECT' error, got: %v", err) + } +} + +func TestValidationRejectsDrop(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + _, err := exec.Execute(context.Background(), "agent-a", + "SELECT 1; DROP TABLE messages") + if err == nil { + t.Fatal("expected multi-statement to be rejected") + } +} + +func TestValidationRejectsUpdate(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + _, err := exec.Execute(context.Background(), "agent-a", + "UPDATE messages SET body = 'hacked'") + if err == nil { + t.Fatal("expected UPDATE to be rejected") + } +} + +func TestValidationRejectsPragma(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + _, err := exec.Execute(context.Background(), "agent-a", + "SELECT * FROM pragma_table_info('messages')") + if err == nil { + t.Fatal("expected PRAGMA in SELECT to be rejected") + } +} + +func TestEmptyQuery(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + _, err := exec.Execute(context.Background(), "agent-a", "") + if err == nil { + t.Fatal("expected empty query to be rejected") + } +} + +func TestCTEQuery(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + result, err := exec.Execute(context.Background(), "agent-a", + "WITH high_prio AS (SELECT * FROM my_messages WHERE priority >= 7) SELECT id, priority FROM high_prio") + if err != nil { + t.Fatalf("CTE query failed: %v", err) + } + + // agent-a should see high-priority messages it has access to + if result.RowCount == 0 { + t.Error("expected at least 1 high-priority message") + } +} + +func TestEmptyResultSet(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + exec := New(db, slog.Default()) + result, err := exec.Execute(context.Background(), "agent-a", + "SELECT * FROM my_messages WHERE body = 'nonexistent'") + if err != nil { + t.Fatalf("query failed: %v", err) + } + if result.RowCount != 0 { + t.Errorf("expected 0 rows, got %d", result.RowCount) + } + if result.Rows == nil { + t.Error("rows should be empty array, not nil") + } + if result.Truncated { + t.Error("should not be truncated") + } +} + +func TestLimitEnforcement(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + // Insert 150 messages to test limit + for i := 100; i < 250; i++ { + _, _ = db.Exec("INSERT INTO messages (id, from_agent, to_agent, body) VALUES (?, 'algis', 'agent-a', 'msg')", i) + } + + exec := New(db, slog.Default()) + result, err := exec.Execute(context.Background(), "agent-a", + "SELECT id FROM my_messages") + if err != nil { + t.Fatalf("query failed: %v", err) + } + + if result.RowCount > MaxRows { + t.Errorf("expected max %d rows, got %d", MaxRows, result.RowCount) + } + if !result.Truncated { + t.Error("expected truncated=true for large result set") + } +} + +func contains(s, substr string) bool { + return len(s) >= len(substr) && (s == substr || len(s) > 0 && containsStr(s, substr)) +} + +func containsStr(s, sub string) bool { + for i := 0; i <= len(s)-len(sub); i++ { + if s[i:i+len(sub)] == sub { + return true + } + } + return false +} diff --git a/internal/agents/store.go b/internal/agents/store.go index d7158e1..0e1093e 100644 --- a/internal/agents/store.go +++ b/internal/agents/store.go @@ -19,6 +19,12 @@ type AgentStore interface { ListAgentsByOwner(ctx context.Context, ownerID int64) ([]*Agent, error) SearchAgentsByCapability(ctx context.Context, query string) ([]*Agent, error) GetHumanAgentByOwner(ctx context.Context, ownerID int64) (*Agent, error) + + // Reactive trigger methods + UpdateTriggerConfig(ctx context.Context, name string, mode string, cooldown, budget, maxDepth int) error + UpdateK8sImage(ctx context.Context, name, image, envJSON, preset string) error + SetPendingWork(ctx context.Context, name string, pending bool) error + ListReactiveAgents(ctx context.Context) ([]*Agent, error) } // SQLiteAgentStore implements AgentStore using SQLite. @@ -37,6 +43,28 @@ func (s *SQLiteAgentStore) CreateAgent(ctx context.Context, agent *Agent) error caps = "{}" } + // Default trigger values + triggerMode := agent.TriggerMode + if triggerMode == "" { + triggerMode = TriggerModePassive + } + cooldown := agent.CooldownSeconds + if cooldown == 0 { + cooldown = 600 + } + budget := agent.DailyTriggerBudget + if budget == 0 { + budget = 8 + } + maxDepth := agent.MaxTriggerDepth + if maxDepth == 0 { + maxDepth = 5 + } + preset := agent.K8sResourcePreset + if preset == "" { + preset = "default" + } + result, err := s.db.ExecContext(ctx, `INSERT INTO agents (name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`, @@ -51,20 +79,75 @@ func (s *SQLiteAgentStore) CreateAgent(ctx context.Context, agent *Agent) error } agent.ID = id agent.Status = AgentStatusActive + agent.TriggerMode = triggerMode + agent.CooldownSeconds = cooldown + agent.DailyTriggerBudget = budget + agent.MaxTriggerDepth = maxDepth + agent.K8sResourcePreset = preset return nil } +// UpdateTriggerConfig updates the reactive trigger configuration for an agent. +func (s *SQLiteAgentStore) UpdateTriggerConfig(ctx context.Context, name string, mode string, cooldown, budget, maxDepth int) error { + _, err := s.db.ExecContext(ctx, + `UPDATE agents SET trigger_mode = ?, cooldown_seconds = ?, daily_trigger_budget = ?, max_trigger_depth = ?, updated_at = CURRENT_TIMESTAMP + WHERE name = ? AND status = 'active'`, + mode, cooldown, budget, maxDepth, name, + ) + return err +} + +// UpdateK8sImage updates the K8s container image and env config for an agent. +func (s *SQLiteAgentStore) UpdateK8sImage(ctx context.Context, name, image, envJSON, preset string) error { + _, err := s.db.ExecContext(ctx, + `UPDATE agents SET k8s_image = ?, k8s_env_json = ?, k8s_resource_preset = ?, updated_at = CURRENT_TIMESTAMP + WHERE name = ? AND status = 'active'`, + image, envJSON, preset, name, + ) + return err +} + +// SetPendingWork sets the pending_work flag for an agent. +func (s *SQLiteAgentStore) SetPendingWork(ctx context.Context, name string, pending bool) error { + val := 0 + if pending { + val = 1 + } + _, err := s.db.ExecContext(ctx, + `UPDATE agents SET pending_work = ? WHERE name = ? AND status = 'active'`, + val, name, + ) + return err +} + +// ListReactiveAgents returns all active agents with trigger_mode='reactive'. +func (s *SQLiteAgentStore) ListReactiveAgents(ctx context.Context) ([]*Agent, error) { + rows, err := s.db.QueryContext(ctx, + agentSelectSQL()+` WHERE status = 'active' AND trigger_mode = 'reactive' ORDER BY name`, + ) + if err != nil { + return nil, err + } + defer rows.Close() + return s.scanAgents(rows) +} + +// agentSelectSQL returns the base SELECT clause for agent queries. +func agentSelectSQL() string { + return `SELECT id, name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at, + trigger_mode, cooldown_seconds, daily_trigger_budget, max_trigger_depth, k8s_image, k8s_env_json, k8s_resource_preset, pending_work + FROM agents` +} + func (s *SQLiteAgentStore) GetAgentByName(ctx context.Context, name string) (*Agent, error) { return s.scanAgent(s.db.QueryRowContext(ctx, - `SELECT id, name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at - FROM agents WHERE name = ? AND status = 'active'`, name, + agentSelectSQL()+` WHERE name = ? AND status = 'active'`, name, )) } func (s *SQLiteAgentStore) GetAgentByID(ctx context.Context, id int64) (*Agent, error) { return s.scanAgent(s.db.QueryRowContext(ctx, - `SELECT id, name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at - FROM agents WHERE id = ? AND status = 'active'`, id, + agentSelectSQL()+` WHERE id = ? AND status = 'active'`, id, )) } @@ -103,8 +186,7 @@ func (s *SQLiteAgentStore) DeactivateAgent(ctx context.Context, name string) err func (s *SQLiteAgentStore) ListActiveAgents(ctx context.Context) ([]*Agent, error) { rows, err := s.db.QueryContext(ctx, - `SELECT id, name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at - FROM agents WHERE status = 'active' ORDER BY name`, + agentSelectSQL()+` WHERE status = 'active' ORDER BY name`, ) if err != nil { return nil, err @@ -115,8 +197,7 @@ func (s *SQLiteAgentStore) ListActiveAgents(ctx context.Context) ([]*Agent, erro func (s *SQLiteAgentStore) ListAllActiveAgents(ctx context.Context) ([]*Agent, error) { rows, err := s.db.QueryContext(ctx, - `SELECT id, name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at - FROM agents WHERE status = 'active' AND type != 'human' ORDER BY name`, + agentSelectSQL()+` WHERE status = 'active' AND type != 'human' ORDER BY name`, ) if err != nil { return nil, err @@ -127,8 +208,7 @@ func (s *SQLiteAgentStore) ListAllActiveAgents(ctx context.Context) ([]*Agent, e func (s *SQLiteAgentStore) ListAgentsByOwner(ctx context.Context, ownerID int64) ([]*Agent, error) { rows, err := s.db.QueryContext(ctx, - `SELECT id, name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at - FROM agents WHERE owner_id = ? AND status = 'active' ORDER BY name`, + agentSelectSQL()+` WHERE owner_id = ? AND status = 'active' ORDER BY name`, ownerID, ) if err != nil { @@ -141,8 +221,7 @@ func (s *SQLiteAgentStore) ListAgentsByOwner(ctx context.Context, ownerID int64) func (s *SQLiteAgentStore) SearchAgentsByCapability(ctx context.Context, query string) ([]*Agent, error) { // Simple LIKE search on the capabilities JSON field rows, err := s.db.QueryContext(ctx, - `SELECT id, name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at - FROM agents WHERE status = 'active' AND capabilities LIKE ? ORDER BY name`, + agentSelectSQL()+` WHERE status = 'active' AND capabilities LIKE ? ORDER BY name`, "%"+query+"%", ) if err != nil { @@ -154,23 +233,30 @@ func (s *SQLiteAgentStore) SearchAgentsByCapability(ctx context.Context, query s func (s *SQLiteAgentStore) GetHumanAgentByOwner(ctx context.Context, ownerID int64) (*Agent, error) { return s.scanAgent(s.db.QueryRowContext(ctx, - `SELECT id, name, display_name, type, capabilities, owner_id, api_key_hash, status, created_at, updated_at - FROM agents WHERE owner_id = ? AND type = 'human' AND status = 'active' LIMIT 1`, ownerID, + agentSelectSQL()+` WHERE owner_id = ? AND type = 'human' AND status = 'active' LIMIT 1`, ownerID, )) } func (s *SQLiteAgentStore) scanAgent(row *sql.Row) (*Agent, error) { var agent Agent var caps string + var k8sImage, k8sEnvJSON sql.NullString + var pendingWork int err := row.Scan( &agent.ID, &agent.Name, &agent.DisplayName, &agent.Type, &caps, &agent.OwnerID, &agent.APIKeyHash, &agent.Status, &agent.CreatedAt, &agent.UpdatedAt, + &agent.TriggerMode, &agent.CooldownSeconds, &agent.DailyTriggerBudget, + &agent.MaxTriggerDepth, &k8sImage, &k8sEnvJSON, + &agent.K8sResourcePreset, &pendingWork, ) if err != nil { return nil, err } agent.Capabilities = json.RawMessage(caps) + agent.K8sImage = k8sImage.String + agent.K8sEnvJSON = k8sEnvJSON.String + agent.PendingWork = pendingWork != 0 return &agent, nil } @@ -179,15 +265,23 @@ func (s *SQLiteAgentStore) scanAgents(rows *sql.Rows) ([]*Agent, error) { for rows.Next() { var agent Agent var caps string + var k8sImage, k8sEnvJSON sql.NullString + var pendingWork int err := rows.Scan( &agent.ID, &agent.Name, &agent.DisplayName, &agent.Type, &caps, &agent.OwnerID, &agent.APIKeyHash, &agent.Status, &agent.CreatedAt, &agent.UpdatedAt, + &agent.TriggerMode, &agent.CooldownSeconds, &agent.DailyTriggerBudget, + &agent.MaxTriggerDepth, &k8sImage, &k8sEnvJSON, + &agent.K8sResourcePreset, &pendingWork, ) if err != nil { return nil, err } agent.Capabilities = json.RawMessage(caps) + agent.K8sImage = k8sImage.String + agent.K8sEnvJSON = k8sEnvJSON.String + agent.PendingWork = pendingWork != 0 agents = append(agents, &agent) } if agents == nil { diff --git a/internal/agents/types.go b/internal/agents/types.go index 9d4935b..9c98a5a 100644 --- a/internal/agents/types.go +++ b/internal/agents/types.go @@ -12,6 +12,13 @@ const ( AgentStatusInactive = "inactive" ) +// Trigger mode constants. +const ( + TriggerModePassive = "passive" + TriggerModeReactive = "reactive" + TriggerModeDisabled = "disabled" +) + // Agent represents a registered entity that can send/receive messages. type Agent struct { ID int64 `json:"id"` @@ -24,4 +31,14 @@ type Agent struct { Status string `json:"status"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` + + // Reactive trigger fields + TriggerMode string `json:"trigger_mode"` + CooldownSeconds int `json:"cooldown_seconds"` + DailyTriggerBudget int `json:"daily_trigger_budget"` + MaxTriggerDepth int `json:"max_trigger_depth"` + K8sImage string `json:"k8s_image,omitempty"` + K8sEnvJSON string `json:"k8s_env_json,omitempty"` + K8sResourcePreset string `json:"k8s_resource_preset"` + PendingWork bool `json:"pending_work"` } diff --git a/internal/api/router.go b/internal/api/router.go index a12d284..31f2d0d 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -12,6 +12,7 @@ import ( "github.com/synapbus/synapbus/internal/channels" "github.com/synapbus/synapbus/internal/k8s" "github.com/synapbus/synapbus/internal/messaging" + "github.com/synapbus/synapbus/internal/reactor" "github.com/synapbus/synapbus/internal/push" "github.com/synapbus/synapbus/internal/reactions" "github.com/synapbus/synapbus/internal/trace" @@ -37,6 +38,8 @@ type RouterConfig struct { ReactionService *reactions.Service PushService *push.Service TrustService *trust.Service + ReactorStore *reactor.Store + ReactorEngine *reactor.Reactor SSEHub *SSEHub Broadcaster *SSEBroadcaster SessionMiddleware func(http.Handler) http.Handler @@ -238,6 +241,19 @@ func NewRouterWithConfig(cfg RouterConfig) chi.Router { } } + // Reactive Runs + if cfg.ReactorStore != nil && cfg.ReactorEngine != nil && cfg.AgentService != nil { + runsHandler := NewRunsHandler(cfg.ReactorStore, cfg.ReactorEngine, agents.NewSQLiteAgentStore(cfg.DB)) + r.Group(func(r chi.Router) { + r.Use(authMiddleware) + + r.Get("/api/runs", runsHandler.ListRuns) + r.Get("/api/runs/{id}", runsHandler.GetRun) + r.Post("/api/runs/{id}/retry", runsHandler.RetryRun) + r.Get("/api/agents/reactive", runsHandler.ReactiveAgents) + }) + } + // Trust Scores if cfg.TrustService != nil { trustHandler := NewTrustHandler(cfg.TrustService) diff --git a/internal/api/runs_handler.go b/internal/api/runs_handler.go new file mode 100644 index 0000000..47b3e8e --- /dev/null +++ b/internal/api/runs_handler.go @@ -0,0 +1,165 @@ +package api + +import ( + "net/http" + "strconv" + "time" + + "github.com/go-chi/chi/v5" + + "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/reactor" +) + +// RunsHandler handles REST API requests for reactive runs. +type RunsHandler struct { + store *reactor.Store + reactor *reactor.Reactor + agentStore agents.AgentStore +} + +// NewRunsHandler creates a new runs handler. +func NewRunsHandler(store *reactor.Store, r *reactor.Reactor, agentStore agents.AgentStore) *RunsHandler { + return &RunsHandler{ + store: store, + reactor: r, + agentStore: agentStore, + } +} + +// ListRuns returns reactive runs with optional filters. +func (h *RunsHandler) ListRuns(w http.ResponseWriter, r *http.Request) { + agentName := r.URL.Query().Get("agent") + status := r.URL.Query().Get("status") + limit := 50 + offset := 0 + + if l := r.URL.Query().Get("limit"); l != "" { + if v, err := strconv.Atoi(l); err == nil && v > 0 && v <= 200 { + limit = v + } + } + if o := r.URL.Query().Get("offset"); o != "" { + if v, err := strconv.Atoi(o); err == nil && v >= 0 { + offset = v + } + } + + runs, total, err := h.store.ListRuns(r.Context(), agentName, status, limit, offset) + if err != nil { + writeJSON(w, http.StatusInternalServerError, errorBody("internal_error", err.Error())) + return + } + + writeJSON(w, http.StatusOK, map[string]any{ + "runs": runs, + "total": total, + }) +} + +// GetRun returns a single run by ID. +func (h *RunsHandler) GetRun(w http.ResponseWriter, r *http.Request) { + idStr := chi.URLParam(r, "id") + id, err := strconv.ParseInt(idStr, 10, 64) + if err != nil { + writeJSON(w, http.StatusBadRequest, errorBody("bad_request", "invalid run ID")) + return + } + + run, err := h.store.GetRunByID(r.Context(), id) + if err != nil { + writeJSON(w, http.StatusNotFound, errorBody("not_found", "run not found")) + return + } + + writeJSON(w, http.StatusOK, run) +} + +// RetryRun retries a failed run. +func (h *RunsHandler) RetryRun(w http.ResponseWriter, r *http.Request) { + idStr := chi.URLParam(r, "id") + id, err := strconv.ParseInt(idStr, 10, 64) + if err != nil { + writeJSON(w, http.StatusBadRequest, errorBody("bad_request", "invalid run ID")) + return + } + + newRun, err := h.reactor.RetryRun(r.Context(), id) + if err != nil { + writeJSON(w, http.StatusBadRequest, errorBody("retry_failed", err.Error())) + return + } + + writeJSON(w, http.StatusOK, map[string]any{ + "new_run_id": newRun.ID, + "status": newRun.Status, + }) +} + +// ReactiveAgents returns agents with reactive trigger config and current status. +func (h *RunsHandler) ReactiveAgents(w http.ResponseWriter, r *http.Request) { + agentsList, err := h.agentStore.ListReactiveAgents(r.Context()) + if err != nil { + writeJSON(w, http.StatusInternalServerError, errorBody("internal_error", err.Error())) + return + } + + type agentStatus struct { + Name string `json:"name"` + TriggerMode string `json:"trigger_mode"` + CooldownSeconds int `json:"cooldown_seconds"` + DailyTriggerBudget int `json:"daily_trigger_budget"` + MaxTriggerDepth int `json:"max_trigger_depth"` + K8sImage string `json:"k8s_image"` + PendingWork bool `json:"pending_work"` + State string `json:"state"` + TodayRuns int `json:"today_runs"` + CooldownUntil *string `json:"cooldown_until"` + } + + result := make([]agentStatus, 0, len(agentsList)) + for _, a := range agentsList { + as := agentStatus{ + Name: a.Name, + TriggerMode: a.TriggerMode, + CooldownSeconds: a.CooldownSeconds, + DailyTriggerBudget: a.DailyTriggerBudget, + MaxTriggerDepth: a.MaxTriggerDepth, + K8sImage: a.K8sImage, + PendingWork: a.PendingWork, + } + + // Compute state + todayCount, _ := h.store.CountTodayRuns(r.Context(), a.Name) + as.TodayRuns = todayCount + + running, _ := h.store.IsAgentRunning(r.Context(), a.Name) + if running { + as.State = "running" + } else if a.PendingWork { + as.State = "queued" + } else if todayCount >= a.DailyTriggerBudget { + as.State = "budget_exhausted" + } else { + lastRun, _ := h.store.GetLastRunTime(r.Context(), a.Name) + if lastRun != nil { + cooldownEnd := lastRun.Add(time.Duration(a.CooldownSeconds) * time.Second) + if time.Now().Before(cooldownEnd) { + as.State = "cooldown" + t := cooldownEnd.UTC().Format(time.RFC3339) + as.CooldownUntil = &t + } else { + as.State = "idle" + } + } else { + as.State = "idle" + } + } + + result = append(result, as) + } + + writeJSON(w, http.StatusOK, map[string]any{ + "agents": result, + }) +} diff --git a/internal/k8s/runner.go b/internal/k8s/runner.go index 1093ddc..5e5e019 100644 --- a/internal/k8s/runner.go +++ b/internal/k8s/runner.go @@ -84,6 +84,11 @@ func (r *K8sJobRunner) IsAvailable() bool { return true } +// GetClientset returns the kubernetes clientset for direct API access (used by reactor poller). +func (r *K8sJobRunner) GetClientset() kubernetes.Interface { + return r.clientset +} + func (r *K8sJobRunner) GetNamespace() string { return r.namespace } @@ -145,14 +150,18 @@ func (r *K8sJobRunner) CreateJob(ctx context.Context, handler *K8sHandler, msg * RestartPolicy: corev1.RestartPolicyNever, Containers: []corev1.Container{ { - Name: "handler", - Image: handler.Image, - Env: envVars, + Name: "handler", + Image: handler.Image, + ImagePullPolicy: corev1.PullIfNotPresent, + Args: handler.Args, + Env: envVars, + VolumeMounts: buildVolumeMounts(handler.VolumeMounts), Resources: corev1.ResourceRequirements{ Limits: resourceLimits, }, }, }, + Volumes: buildVolumes(handler.Volumes), }, }, }, @@ -228,6 +237,48 @@ func sanitizeJobName(name string) string { return name } +// buildVolumeMounts converts our VolumeMount type to K8s VolumeMounts. +func buildVolumeMounts(mounts []VolumeMount) []corev1.VolumeMount { + if len(mounts) == 0 { + return nil + } + var result []corev1.VolumeMount + for _, m := range mounts { + result = append(result, corev1.VolumeMount{ + Name: m.Name, + MountPath: m.MountPath, + ReadOnly: m.ReadOnly, + }) + } + return result +} + +// buildVolumes converts our Volume type to K8s Volumes. +func buildVolumes(volumes []Volume) []corev1.Volume { + if len(volumes) == 0 { + return nil + } + var result []corev1.Volume + for _, v := range volumes { + vol := corev1.Volume{Name: v.Name} + if v.HostPath != "" { + hostPathType := corev1.HostPathDirectory + vol.VolumeSource = corev1.VolumeSource{ + HostPath: &corev1.HostPathVolumeSource{ + Path: v.HostPath, + Type: &hostPathType, + }, + } + } else if v.EmptyDir { + vol.VolumeSource = corev1.VolumeSource{ + EmptyDir: &corev1.EmptyDirVolumeSource{}, + } + } + result = append(result, vol) + } + return result +} + // truncateBody truncates the message body to maxLen bytes. func truncateBody(body string, maxLen int) string { if len(body) <= maxLen { diff --git a/internal/k8s/store.go b/internal/k8s/store.go index 7c8ee05..2f6eac2 100644 --- a/internal/k8s/store.go +++ b/internal/k8s/store.go @@ -22,6 +22,25 @@ type K8sHandler struct { Status string `json:"status"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` + + // Extended fields for reactive triggers (not persisted in k8s_handlers table) + Args []string `json:"-"` + VolumeMounts []VolumeMount `json:"-"` + Volumes []Volume `json:"-"` +} + +// VolumeMount defines a mount point in the container. +type VolumeMount struct { + Name string + MountPath string + ReadOnly bool +} + +// Volume defines a volume source for the pod. +type Volume struct { + Name string + HostPath string // If set, uses hostPath volume + EmptyDir bool // If true, uses emptyDir volume } // K8sJobRun represents a single Kubernetes job execution. diff --git a/internal/mcp/bridge.go b/internal/mcp/bridge.go index 7ccd095..2402575 100644 --- a/internal/mcp/bridge.go +++ b/internal/mcp/bridge.go @@ -12,6 +12,7 @@ import ( "time" "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/agentquery" "github.com/synapbus/synapbus/internal/attachments" "github.com/synapbus/synapbus/internal/channels" "github.com/synapbus/synapbus/internal/messaging" @@ -31,6 +32,7 @@ type ServiceBridge struct { searchService *search.Service reactionService *reactions.Service trustService *trust.Service + queryExecutor *agentquery.Executor agentName string } @@ -130,6 +132,10 @@ func (b *ServiceBridge) Call(ctx context.Context, actionName string, args map[st case "get_trust": return b.callGetTrust(ctx, args) + // --- SQL Query --- + case "query": + return b.callQuery(ctx, args) + // --- DM send (also accessible via bridge for execute tool) --- case "send_message": return b.callSendMessage(ctx, args) @@ -1215,6 +1221,29 @@ func (b *ServiceBridge) callGetTrust(ctx context.Context, args map[string]any) ( }, nil } +// SetQueryExecutor sets the SQL query executor for the bridge. +func (b *ServiceBridge) SetQueryExecutor(exec *agentquery.Executor) { + b.queryExecutor = exec +} + +func (b *ServiceBridge) callQuery(ctx context.Context, args map[string]any) (any, error) { + if b.queryExecutor == nil { + return nil, fmt.Errorf("SQL query not available") + } + + sqlStr := getString(args, "sql", "") + if sqlStr == "" { + return nil, fmt.Errorf("sql parameter is required") + } + + result, err := b.queryExecutor.Execute(ctx, b.agentName, sqlStr) + if err != nil { + return nil, err + } + + return result, nil +} + // --- Helpers --- // resolveChannelID resolves a channel ID from either channel_id or channel_name in args. diff --git a/internal/mcp/server.go b/internal/mcp/server.go index 25dbc3e..38cee2d 100644 --- a/internal/mcp/server.go +++ b/internal/mcp/server.go @@ -12,6 +12,7 @@ import ( "github.com/mark3labs/mcp-go/server" "github.com/synapbus/synapbus/internal/actions" + "github.com/synapbus/synapbus/internal/agentquery" "github.com/synapbus/synapbus/internal/agents" "github.com/synapbus/synapbus/internal/attachments" "github.com/synapbus/synapbus/internal/channels" @@ -26,12 +27,13 @@ import ( // MCPServer wraps the mcp-go server with SynapBus services. type MCPServer struct { - mcpServer *server.MCPServer - httpServer *server.StreamableHTTPServer - connMgr *ConnectionManager - agentService *agents.AgentService - logger *slog.Logger - console *console.Printer + mcpServer *server.MCPServer + httpServer *server.StreamableHTTPServer + connMgr *ConnectionManager + agentService *agents.AgentService + hybridRegistrar *HybridToolRegistrar + logger *slog.Logger + console *console.Printer } // NewMCPServer creates and configures a new MCP server with 4 hybrid tools registered. @@ -187,18 +189,26 @@ func NewMCPServer( ) s := &MCPServer{ - mcpServer: mcpSrv, - httpServer: httpServer, - connMgr: connMgr, - agentService: agentService, - logger: logger, - console: consolePrinter, + mcpServer: mcpSrv, + httpServer: httpServer, + connMgr: connMgr, + agentService: agentService, + hybridRegistrar: hybridRegistrar, + logger: logger, + console: consolePrinter, } logger.Info("MCP server initialized (4 hybrid tools, 4 prompts, streamable HTTP transport)") return s } +// SetQueryExecutor sets the SQL query executor for agent queries via the execute tool. +func (s *MCPServer) SetQueryExecutor(exec *agentquery.Executor) { + if s.hybridRegistrar != nil { + s.hybridRegistrar.SetQueryExecutor(exec) + } +} + // Handler returns the HTTP handler for mounting on a router. func (s *MCPServer) Handler() http.Handler { return s.httpServer diff --git a/internal/mcp/tools_hybrid.go b/internal/mcp/tools_hybrid.go index 3dac2f0..a26f5c6 100644 --- a/internal/mcp/tools_hybrid.go +++ b/internal/mcp/tools_hybrid.go @@ -17,6 +17,7 @@ import ( "github.com/synapbus/synapbus/internal/attachments" "github.com/synapbus/synapbus/internal/channels" "github.com/synapbus/synapbus/internal/jsruntime" + "github.com/synapbus/synapbus/internal/agentquery" "github.com/synapbus/synapbus/internal/messaging" "github.com/synapbus/synapbus/internal/reactions" "github.com/synapbus/synapbus/internal/search" @@ -37,9 +38,15 @@ type HybridToolRegistrar struct { actionRegistry *actions.Registry actionIndex *actions.Index db *sql.DB + queryExecutor *agentquery.Executor logger *slog.Logger } +// SetQueryExecutor sets the SQL query executor for all agent bridges. +func (h *HybridToolRegistrar) SetQueryExecutor(exec *agentquery.Executor) { + h.queryExecutor = exec +} + // NewHybridToolRegistrar creates a new hybrid tool registrar. func NewHybridToolRegistrar( msgService *messaging.MessagingService, @@ -495,6 +502,9 @@ func (h *HybridToolRegistrar) handleExecute(ctx context.Context, req mcplib.Call h.trustService, agentName, ) + if h.queryExecutor != nil { + bridge.SetQueryExecutor(h.queryExecutor) + } result, err := h.jsPool.Execute(ctx, code, bridge, jsruntime.ExecuteOptions{ Timeout: timeout, diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index f17d3f4..7f59664 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -42,4 +42,46 @@ var ( Name: "active_connections", Help: "Number of active connections", }) + + // Reactive agent triggering metrics + ReactiveTriggersTotal = promauto.NewCounterVec( + prometheus.CounterOpts{ + Namespace: "synapbus", + Subsystem: "reactor", + Name: "triggers_total", + Help: "Total reactive trigger evaluations by agent and outcome", + }, + []string{"agent", "status"}, + ) + + ReactiveRunDuration = promauto.NewHistogramVec( + prometheus.HistogramOpts{ + Namespace: "synapbus", + Subsystem: "reactor", + Name: "run_duration_seconds", + Help: "Duration of reactive agent runs in seconds", + Buckets: []float64{10, 30, 60, 120, 300, 600, 1200, 1800, 3600}, + }, + []string{"agent"}, + ) + + ReactiveAgentState = promauto.NewGaugeVec( + prometheus.GaugeOpts{ + Namespace: "synapbus", + Subsystem: "reactor", + Name: "agent_running", + Help: "Whether a reactive agent is currently running (1) or idle (0)", + }, + []string{"agent"}, + ) + + ReactiveBudgetUsed = promauto.NewGaugeVec( + prometheus.GaugeOpts{ + Namespace: "synapbus", + Subsystem: "reactor", + Name: "budget_used_today", + Help: "Number of reactive runs used today per agent", + }, + []string{"agent"}, + ) ) diff --git a/internal/onboarding/templates.go b/internal/onboarding/templates.go index 961ce22..cd1e3ef 100644 --- a/internal/onboarding/templates.go +++ b/internal/onboarding/templates.go @@ -28,6 +28,14 @@ You are **{{.AgentName}}**, an autonomous agent connected to SynapBus. Use ` + "`call(\"search\", {\"query\": \"workflow\"})`" + ` to discover all available tools. +### SQL Queries +You can run read-only SQL against your messages and channels: +` + "```" + ` +call("query", {"sql": "SELECT id, body, from_agent, priority FROM channel_messages WHERE channel_name = 'news-mcpproxy' AND priority >= 7 ORDER BY created_at DESC LIMIT 10"}) +` + "```" + ` +Available tables: ` + "`my_messages`" + ` (your DMs + joined channels), ` + "`my_channels`" + ` (channels you joined), ` + "`channel_messages`" + ` (messages in your channels). +Results capped at 100 rows. CTEs (WITH) supported. Only SELECT allowed. + ### Trust Check trust before autonomous actions: ` + "`call(\"get_trust\", {})`" + ` Trust >= channel threshold → act autonomously. Otherwise post as "proposed" and wait for approval. diff --git a/internal/reactor/notifier.go b/internal/reactor/notifier.go new file mode 100644 index 0000000..284ef04 --- /dev/null +++ b/internal/reactor/notifier.go @@ -0,0 +1,52 @@ +package reactor + +import ( + "context" + "fmt" + + "github.com/synapbus/synapbus/internal/messaging" +) + +// DMFailureNotifier sends system DMs to agent owners on job failure. +type DMFailureNotifier struct { + msgService *messaging.MessagingService +} + +// NewDMFailureNotifier creates a new failure notifier. +func NewDMFailureNotifier(msgService *messaging.MessagingService) *DMFailureNotifier { + return &DMFailureNotifier{msgService: msgService} +} + +// NotifyFailure sends a system DM to the agent's owner with error details. +func (n *DMFailureNotifier) NotifyFailure(ctx context.Context, ownerAgentName, agentName, triggerFrom, triggerEvent string, durationMs int64, errorSummary string) error { + durationStr := "< 1s" + if durationMs > 0 { + secs := durationMs / 1000 + if secs >= 60 { + durationStr = fmt.Sprintf("%dm%ds", secs/60, secs%60) + } else { + durationStr = fmt.Sprintf("%ds", secs) + } + } + + body := fmt.Sprintf( + "⚠️ **Reactive run failed** for **%s**\n\n"+ + "**Trigger**: %s from %s\n"+ + "**Duration**: %s\n"+ + "**Error**: %s\n\n"+ + "View details in Agent Runs page.", + agentName, triggerEvent, triggerFrom, durationStr, truncateError(errorSummary, 500), + ) + + _, err := n.msgService.SendMessage(ctx, "system", ownerAgentName, body, messaging.SendOptions{ + Priority: 7, + }) + return err +} + +func truncateError(s string, maxLen int) string { + if len(s) <= maxLen { + return s + } + return s[:maxLen] + "..." +} diff --git a/internal/reactor/poller.go b/internal/reactor/poller.go new file mode 100644 index 0000000..2aa3084 --- /dev/null +++ b/internal/reactor/poller.go @@ -0,0 +1,224 @@ +package reactor + +import ( + "context" + "fmt" + "log/slog" + "strings" + "time" + + "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/dispatcher" + k8spkg "github.com/synapbus/synapbus/internal/k8s" + "github.com/synapbus/synapbus/internal/metrics" + + batchv1 "k8s.io/api/batch/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" +) + +// Poller watches active reactive runs and updates their status from K8s. +type Poller struct { + store *Store + agentStore agents.AgentStore + clientset kubernetes.Interface + runner k8spkg.JobRunner + reactor *Reactor + interval time.Duration + logger *slog.Logger + stopCh chan struct{} +} + +// NewPoller creates a new job status poller. +func NewPoller(store *Store, agentStore agents.AgentStore, runner k8spkg.JobRunner, reactor *Reactor, logger *slog.Logger) *Poller { + // Extract clientset from runner if it's the real K8s runner + var clientset kubernetes.Interface + if kr, ok := runner.(*k8spkg.K8sJobRunner); ok { + clientset = kr.GetClientset() + } + + return &Poller{ + store: store, + agentStore: agentStore, + clientset: clientset, + runner: runner, + reactor: reactor, + interval: 15 * time.Second, + logger: logger.With("component", "reactor-poller"), + stopCh: make(chan struct{}), + } +} + +// Start begins the polling loop in a background goroutine. +func (p *Poller) Start() { + if !p.runner.IsAvailable() || p.clientset == nil { + p.logger.Info("K8s not available, reactor poller disabled") + return + } + go p.pollLoop() + p.logger.Info("reactor poller started", "interval", p.interval) +} + +// Stop signals the poller to stop. +func (p *Poller) Stop() { + close(p.stopCh) +} + +func (p *Poller) pollLoop() { + ticker := time.NewTicker(p.interval) + defer ticker.Stop() + + for { + select { + case <-p.stopCh: + return + case <-ticker.C: + p.pollActiveRuns() + } + } +} + +func (p *Poller) pollActiveRuns() { + ctx := context.Background() + + runs, err := p.store.GetActiveRuns(ctx) + if err != nil { + p.logger.Error("failed to get active runs", "error", err) + return + } + + for _, run := range runs { + if run.K8sJobName == "" || run.K8sNamespace == "" { + continue + } + p.checkJob(ctx, run) + } +} + +func (p *Poller) checkJob(ctx context.Context, run *ReactiveRun) { + ns := run.K8sNamespace + jobName := run.K8sJobName + + job, err := p.clientset.BatchV1().Jobs(ns).Get(ctx, jobName, metav1.GetOptions{}) + if err != nil { + p.logger.Warn("failed to get K8s Job status", "job", jobName, "namespace", ns, "error", err) + return + } + + // Check job conditions + for _, cond := range job.Status.Conditions { + switch cond.Type { + case batchv1.JobComplete: + if cond.Status == "True" { + p.handleJobComplete(ctx, run, true, "") + return + } + case batchv1.JobFailed: + if cond.Status == "True" { + reason := cond.Reason + if cond.Message != "" { + reason = reason + ": " + cond.Message + } + p.handleJobComplete(ctx, run, false, reason) + return + } + } + } + + // Check if active deadline exceeded + if job.Status.Failed > 0 { + p.handleJobComplete(ctx, run, false, "job failed (pod failure)") + return + } +} + +func (p *Poller) handleJobComplete(ctx context.Context, run *ReactiveRun, success bool, failureReason string) { + now := time.Now().UTC() + + // Update metrics + metrics.ReactiveAgentState.WithLabelValues(run.AgentName).Set(0) + if run.StartedAt != nil { + duration := now.Sub(*run.StartedAt).Seconds() + metrics.ReactiveRunDuration.WithLabelValues(run.AgentName).Observe(duration) + } + todayCount, _ := p.store.CountTodayRuns(ctx, run.AgentName) + metrics.ReactiveBudgetUsed.WithLabelValues(run.AgentName).Set(float64(todayCount)) + + if success { + metrics.ReactiveTriggersTotal.WithLabelValues(run.AgentName, StatusSucceeded).Inc() + _ = p.store.CompleteRun(ctx, run.ID, StatusSucceeded, "", now) + p.logger.Info("reactive run succeeded", + "agent", run.AgentName, + "job", run.K8sJobName, + "run_id", run.ID, + ) + } else { + // Retrieve logs + errorLog := failureReason + logs, err := p.runner.GetJobLogs(ctx, run.K8sNamespace, run.K8sJobName) + if err == nil && logs != "" { + // Keep last 100 lines + lines := strings.Split(logs, "\n") + if len(lines) > 100 { + lines = lines[len(lines)-100:] + } + errorLog = strings.Join(lines, "\n") + } + + metrics.ReactiveTriggersTotal.WithLabelValues(run.AgentName, StatusFailed).Inc() + _ = p.store.CompleteRun(ctx, run.ID, StatusFailed, errorLog, now) + + p.logger.Warn("reactive run failed", + "agent", run.AgentName, + "job", run.K8sJobName, + "run_id", run.ID, + "reason", failureReason, + ) + + // Send failure notification + var durationMs int64 + if run.StartedAt != nil { + durationMs = now.Sub(*run.StartedAt).Milliseconds() + } + agent, err := p.agentStore.GetAgentByName(ctx, run.AgentName) + if err == nil && agent != nil { + event := dispatcher.MessageEvent{ + EventType: run.TriggerEvent, + FromAgent: run.TriggerFrom, + } + p.reactor.notifyFailure(ctx, agent, event, durationMs, fmt.Sprintf("Job %s failed: %s", run.K8sJobName, failureReason)) + } + } + + // Check for pending_work — launch coalesced run if needed + p.checkPendingWork(ctx, run.AgentName) +} + +func (p *Poller) checkPendingWork(ctx context.Context, agentName string) { + agent, err := p.agentStore.GetAgentByName(ctx, agentName) + if err != nil { + return + } + + if !agent.PendingWork { + return + } + + // Clear pending_work first + _ = p.agentStore.SetPendingWork(ctx, agentName, false) + + p.logger.Info("pending_work found, launching coalesced run", "agent", agentName) + + // Create a synthetic event (coalesced — agent will pick up all pending messages via claim_messages) + event := dispatcher.MessageEvent{ + EventType: "message.received", + FromAgent: "system", + ToAgent: agentName, + Body: "Coalesced trigger: process all pending messages.", + MentionedAgents: nil, + Depth: 0, + } + + // Evaluate the trigger (it will check cooldown/budget again) + _ = p.reactor.evaluateTrigger(ctx, agentName, event) +} diff --git a/internal/reactor/reactor.go b/internal/reactor/reactor.go new file mode 100644 index 0000000..3e3f4c4 --- /dev/null +++ b/internal/reactor/reactor.go @@ -0,0 +1,359 @@ +// Package reactor provides the reactive agent triggering engine. +// When a DM or @mention targets an agent with trigger_mode='reactive', +// the reactor evaluates rate limits and creates a K8s Job to run the agent. +package reactor + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "strings" + "time" + + "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/dispatcher" + k8spkg "github.com/synapbus/synapbus/internal/k8s" + "github.com/synapbus/synapbus/internal/metrics" +) + +// Reactor is the reactive agent triggering engine. +type Reactor struct { + store *Store + agentStore agents.AgentStore + runner k8spkg.JobRunner + notifier FailureNotifier + logger *slog.Logger +} + +// FailureNotifier sends system DMs on job failure. +type FailureNotifier interface { + NotifyFailure(ctx context.Context, ownerAgentName, agentName, triggerFrom, triggerEvent string, durationMs int64, errorSummary string) error +} + +// New creates a new Reactor. +func New(store *Store, agentStore agents.AgentStore, runner k8spkg.JobRunner, logger *slog.Logger) *Reactor { + return &Reactor{ + store: store, + agentStore: agentStore, + runner: runner, + logger: logger.With("component", "reactor"), + } +} + +// SetFailureNotifier sets the notifier for sending failure DMs. +func (r *Reactor) SetFailureNotifier(n FailureNotifier) { + r.notifier = n +} + +// Dispatch implements dispatcher.EventDispatcher. Called by MultiDispatcher +// when a message event occurs. +func (r *Reactor) Dispatch(ctx context.Context, event dispatcher.MessageEvent) error { + switch event.EventType { + case "message.received": + // DM to an agent + return r.evaluateTrigger(ctx, event.ToAgent, event) + case "message.mentioned": + // @mentions in channel messages + for _, mentioned := range event.MentionedAgents { + // Self-mention filter: agent can't trigger itself + if mentioned == event.FromAgent { + continue + } + if err := r.evaluateTrigger(ctx, mentioned, event); err != nil { + r.logger.ErrorContext(ctx, "reactor trigger eval failed", + "agent", mentioned, + "error", err, + ) + } + } + return nil + default: + return nil // Ignore other event types + } +} + +// evaluateTrigger runs the decision chain for a single agent. +func (r *Reactor) evaluateTrigger(ctx context.Context, agentName string, event dispatcher.MessageEvent) error { + // 1. Get agent config + agent, err := r.agentStore.GetAgentByName(ctx, agentName) + if err != nil { + return nil // Agent doesn't exist, skip silently + } + + // 2. Check trigger mode + if agent.TriggerMode != agents.TriggerModeReactive { + return nil // Not reactive, skip + } + + // 3. Check K8s image configured + if agent.K8sImage == "" { + r.logger.Warn("reactive agent has no k8s_image configured", "agent", agentName) + r.recordSkippedRun(ctx, agentName, event, StatusFailed, "no k8s_image configured") + return nil + } + + // 4. Check K8s runner available + if !r.runner.IsAvailable() { + r.logger.Warn("K8s runner not available for reactive trigger", "agent", agentName) + r.recordSkippedRun(ctx, agentName, event, StatusFailed, "K8s runner not available") + return nil + } + + // 5. Extract depth from event metadata + depth := event.Depth + + // 6. Check trigger depth + if depth >= agent.MaxTriggerDepth { + r.logger.Info("trigger depth exceeded", "agent", agentName, "depth", depth, "max", agent.MaxTriggerDepth) + r.recordSkippedRun(ctx, agentName, event, StatusDepthExceeded, "") + return nil + } + + // 7. Check daily budget + todayCount, err := r.store.CountTodayRuns(ctx, agentName) + if err != nil { + return fmt.Errorf("count today runs: %w", err) + } + if todayCount >= agent.DailyTriggerBudget { + r.logger.Info("daily trigger budget exhausted", "agent", agentName, "count", todayCount, "budget", agent.DailyTriggerBudget) + r.recordSkippedRun(ctx, agentName, event, StatusBudgetExhausted, "") + return nil + } + + // 8. Check cooldown + lastRun, err := r.store.GetLastRunTime(ctx, agentName) + if err != nil { + return fmt.Errorf("get last run time: %w", err) + } + if lastRun != nil { + elapsed := time.Since(*lastRun) + if elapsed < time.Duration(agent.CooldownSeconds)*time.Second { + r.logger.Info("agent on cooldown", "agent", agentName, "elapsed", elapsed, "cooldown", agent.CooldownSeconds) + // Set pending_work so we retry after cooldown + _ = r.agentStore.SetPendingWork(ctx, agentName, true) + r.recordSkippedRun(ctx, agentName, event, StatusCooldownSkipped, "") + return nil + } + } + + // 9. Check if agent is currently running + running, err := r.store.IsAgentRunning(ctx, agentName) + if err != nil { + return fmt.Errorf("check agent running: %w", err) + } + if running { + r.logger.Info("agent already running, setting pending_work", "agent", agentName) + _ = r.agentStore.SetPendingWork(ctx, agentName, true) + r.recordSkippedRun(ctx, agentName, event, StatusQueued, "") + return nil + } + + // 10. All checks pass — create K8s Job + return r.createJob(ctx, agent, event, depth) +} + +// createJob creates a K8s Job for the reactive trigger. +func (r *Reactor) createJob(ctx context.Context, agent *agents.Agent, event dispatcher.MessageEvent, depth int) error { + // Build handler from agent config + handler := r.buildHandler(agent) + + body := event.Body + if len(body) > 4096 { + body = body[:4096] + " [truncated]" + } + + msg := &k8spkg.JobMessage{ + MessageID: event.MessageID, + FromAgent: event.FromAgent, + Body: body, + Event: event.EventType, + Channel: event.Channel, + Timestamp: time.Now().UTC().Format(time.RFC3339), + } + + // Add trigger depth env var to handler + handler.Env["SYNAPBUS_TRIGGER_DEPTH"] = fmt.Sprintf("%d", depth) + + // Create K8s Job FIRST (before DB insert to avoid stuck runs on SQLITE_BUSY) + jobName, err := r.runner.CreateJob(ctx, handler, msg) + if err != nil { + errMsg := fmt.Sprintf("K8s Job creation failed: %s", err.Error()) + r.recordSkippedRun(ctx, agent.Name, event, StatusFailed, errMsg) + r.notifyFailure(ctx, agent, event, 0, errMsg) + return fmt.Errorf("create K8s job: %w", err) + } + + ns := handler.Namespace + if ns == "" { + ns = r.runner.GetNamespace() + } + + // Insert run record with job name already set (single atomic write) + now := time.Now().UTC() + run := &ReactiveRun{ + AgentName: agent.Name, + TriggerMessageID: &event.MessageID, + TriggerEvent: event.EventType, + TriggerDepth: depth, + TriggerFrom: event.FromAgent, + Status: StatusRunning, + K8sJobName: jobName, + K8sNamespace: ns, + StartedAt: &now, + } + + runID, err := r.store.InsertRun(ctx, run) + if err != nil { + r.logger.Error("failed to record reactive run (job already created)", + "agent", agent.Name, "job", jobName, "error", err) + runID = 0 + } + + // Clear pending_work since we're launching + _ = r.agentStore.SetPendingWork(ctx, agent.Name, false) + + metrics.ReactiveTriggersTotal.WithLabelValues(agent.Name, StatusRunning).Inc() + metrics.ReactiveAgentState.WithLabelValues(agent.Name).Set(1) + + r.logger.Info("reactive K8s Job created", + "agent", agent.Name, + "job", jobName, + "trigger_from", event.FromAgent, + "trigger_event", event.EventType, + "depth", depth, + "run_id", runID, + ) + + return nil +} + +// buildHandler constructs a K8sHandler from agent config. +func (r *Reactor) buildHandler(agent *agents.Agent) *k8spkg.K8sHandler { + env := map[string]string{} + + // Parse k8s_env_json + if agent.K8sEnvJSON != "" { + var envMap map[string]json.RawMessage + if err := json.Unmarshal([]byte(agent.K8sEnvJSON), &envMap); err == nil { + for k, v := range envMap { + // Plain string values + var str string + if err := json.Unmarshal(v, &str); err == nil { + env[k] = str + continue + } + // Secret refs are handled at K8s level; for now pass as-is + // (the K8s runner would need extension for secretKeyRef) + env[k] = strings.Trim(string(v), "\"") + } + } + } + + // Resource presets — default matches CronJob config (agent SDK needs ~1-2Gi) + memory := "2Gi" + cpu := "500m" + if agent.K8sResourcePreset == "small" { + memory = "512Mi" + cpu = "100m" + } + + timeout := 3600 // 1 hour (matches CronJob config) + + handler := &k8spkg.K8sHandler{ + AgentName: agent.Name, + Image: agent.K8sImage, + Events: []string{"message.received", "message.mentioned"}, + Namespace: "", // Use runner's namespace + ResourcesMemory: memory, + ResourcesCPU: cpu, + Env: env, + TimeoutSeconds: timeout, + Status: "active", + Args: []string{"--max-turns", "50", "--model", "claude-sonnet-4-6"}, + VolumeMounts: []k8spkg.VolumeMount{ + {Name: "claude-config", MountPath: "/app/.claude", ReadOnly: false}, + {Name: "workspace", MountPath: "/app/workspace", ReadOnly: false}, + }, + Volumes: []k8spkg.Volume{ + {Name: "claude-config", HostPath: "/home/user/.claude"}, + {Name: "workspace", EmptyDir: true}, + }, + } + + // Override args for social-commenter (uses opus, more turns) + if agent.Name == "social-commenter" { + handler.Args = []string{"--max-turns", "80", "--model", "claude-opus-4-6"} + } + + return handler +} + +// RetryRun retries a failed run. +func (r *Reactor) RetryRun(ctx context.Context, runID int64) (*ReactiveRun, error) { + run, err := r.store.GetRunByID(ctx, runID) + if err != nil { + return nil, fmt.Errorf("get run: %w", err) + } + if run.Status != StatusFailed { + return nil, fmt.Errorf("can only retry failed runs, current status: %s", run.Status) + } + + agent, err := r.agentStore.GetAgentByName(ctx, run.AgentName) + if err != nil { + return nil, fmt.Errorf("get agent: %w", err) + } + + // Create a synthetic event for the retry + event := dispatcher.MessageEvent{ + EventType: run.TriggerEvent, + MessageID: 0, + FromAgent: run.TriggerFrom, + ToAgent: run.AgentName, + Body: "", + Depth: run.TriggerDepth, + } + if run.TriggerMessageID != nil { + event.MessageID = *run.TriggerMessageID + } + + if err := r.createJob(ctx, agent, event, run.TriggerDepth); err != nil { + return nil, err + } + + // Return the newly created run + runs, _, err := r.store.ListRuns(ctx, run.AgentName, StatusRunning, 1, 0) + if err != nil || len(runs) == 0 { + return nil, fmt.Errorf("retry succeeded but couldn't find new run") + } + return runs[0], nil +} + +func (r *Reactor) recordSkippedRun(ctx context.Context, agentName string, event dispatcher.MessageEvent, status, errorLog string) { + metrics.ReactiveTriggersTotal.WithLabelValues(agentName, status).Inc() + run := &ReactiveRun{ + AgentName: agentName, + TriggerEvent: event.EventType, + TriggerDepth: event.Depth, + TriggerFrom: event.FromAgent, + Status: status, + ErrorLog: errorLog, + } + if event.MessageID > 0 { + run.TriggerMessageID = &event.MessageID + } + _, _ = r.store.InsertRun(ctx, run) +} + +func (r *Reactor) notifyFailure(ctx context.Context, agent *agents.Agent, event dispatcher.MessageEvent, durationMs int64, errorSummary string) { + if r.notifier == nil { + return + } + // Find the owner's human agent name + ownerAgent, err := r.agentStore.GetHumanAgentByOwner(ctx, agent.OwnerID) + if err != nil || ownerAgent == nil { + r.logger.Warn("could not find owner agent for failure notification", "agent", agent.Name) + return + } + _ = r.notifier.NotifyFailure(ctx, ownerAgent.Name, agent.Name, event.FromAgent, event.EventType, durationMs, errorSummary) +} diff --git a/internal/reactor/reactor_test.go b/internal/reactor/reactor_test.go new file mode 100644 index 0000000..ef8f50b --- /dev/null +++ b/internal/reactor/reactor_test.go @@ -0,0 +1,430 @@ +package reactor + +import ( + "context" + "database/sql" + "encoding/json" + "testing" + "time" + + "fmt" + "log/slog" + + "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/dispatcher" + k8spkg "github.com/synapbus/synapbus/internal/k8s" + + _ "modernc.org/sqlite" +) + +// setupTestDB creates an in-memory SQLite database with schema for testing. +func setupTestDB(t *testing.T) *sql.DB { + t.Helper() + db, err := sql.Open("sqlite", ":memory:") + if err != nil { + t.Fatalf("open db: %v", err) + } + + // Create minimal schema + schema := ` + CREATE TABLE agents ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL UNIQUE, + display_name TEXT NOT NULL DEFAULT '', + type TEXT NOT NULL DEFAULT 'ai', + capabilities TEXT NOT NULL DEFAULT '{}', + owner_id INTEGER NOT NULL DEFAULT 1, + api_key_hash TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT 'active', + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + trigger_mode TEXT NOT NULL DEFAULT 'passive', + cooldown_seconds INTEGER NOT NULL DEFAULT 600, + daily_trigger_budget INTEGER NOT NULL DEFAULT 8, + max_trigger_depth INTEGER NOT NULL DEFAULT 5, + k8s_image TEXT, + k8s_env_json TEXT, + k8s_resource_preset TEXT NOT NULL DEFAULT 'default', + pending_work INTEGER NOT NULL DEFAULT 0 + ); + CREATE TABLE reactive_runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + agent_name TEXT NOT NULL, + trigger_message_id INTEGER, + trigger_event TEXT NOT NULL, + trigger_depth INTEGER NOT NULL DEFAULT 0, + trigger_from TEXT, + status TEXT NOT NULL DEFAULT 'queued', + k8s_job_name TEXT, + k8s_namespace TEXT, + started_at DATETIME, + completed_at DATETIME, + duration_ms INTEGER, + error_log TEXT, + token_cost_json TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + ` + if _, err := db.Exec(schema); err != nil { + t.Fatalf("create schema: %v", err) + } + + return db +} + +func insertTestAgent(t *testing.T, db *sql.DB, name, triggerMode, image string, cooldown, budget, maxDepth int) { + t.Helper() + _, err := db.Exec( + `INSERT INTO agents (name, display_name, type, owner_id, trigger_mode, cooldown_seconds, daily_trigger_budget, max_trigger_depth, k8s_image, k8s_resource_preset) + VALUES (?, ?, 'ai', 1, ?, ?, ?, ?, ?, 'default')`, + name, name, triggerMode, cooldown, budget, maxDepth, image, + ) + if err != nil { + t.Fatalf("insert agent: %v", err) + } +} + +func TestReactorPassiveAgentSkipped(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + insertTestAgent(t, db, "passive-agent", "passive", "image:latest", 600, 8, 5) + + store := NewStore(db) + agentStore := agents.NewSQLiteAgentStore(db) + runner := k8spkg.NewNoopRunner() + logger := slog.Default() + + reactor := New(store, agentStore, runner, logger) + + event := dispatcher.MessageEvent{ + EventType: "message.received", + MessageID: 1, + FromAgent: "algis", + ToAgent: "passive-agent", + Body: "hello", + } + + err := reactor.Dispatch(context.Background(), event) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + + // No runs should be created for passive agents + runs, total, err := store.ListRuns(context.Background(), "passive-agent", "", 10, 0) + if err != nil { + t.Fatalf("list runs: %v", err) + } + if total != 0 || len(runs) != 0 { + t.Errorf("expected 0 runs for passive agent, got %d", total) + } +} + +func TestReactorNoK8sImage(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + insertTestAgent(t, db, "no-image-agent", "reactive", "", 600, 8, 5) + + store := NewStore(db) + agentStore := agents.NewSQLiteAgentStore(db) + runner := k8spkg.NewNoopRunner() + logger := slog.Default() + + reactor := New(store, agentStore, runner, logger) + + event := dispatcher.MessageEvent{ + EventType: "message.received", + MessageID: 1, + FromAgent: "algis", + ToAgent: "no-image-agent", + Body: "hello", + } + + _ = reactor.Dispatch(context.Background(), event) + + runs, _, _ := store.ListRuns(context.Background(), "no-image-agent", StatusFailed, 10, 0) + if len(runs) != 1 { + t.Fatalf("expected 1 failed run for agent with no image, got %d", len(runs)) + } + if runs[0].ErrorLog != "no k8s_image configured" { + t.Errorf("expected 'no k8s_image configured' error, got: %s", runs[0].ErrorLog) + } +} + +func TestReactorDepthExceeded(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + insertTestAgent(t, db, "deep-agent", "reactive", "image:latest", 600, 8, 3) + + store := NewStore(db) + agentStore := agents.NewSQLiteAgentStore(db) + runner := &fakeRunner{available: true} + logger := slog.Default() + + reactor := New(store, agentStore, runner, logger) + + event := dispatcher.MessageEvent{ + EventType: "message.received", + MessageID: 1, + FromAgent: "other-agent", + ToAgent: "deep-agent", + Body: "hello from depth 3", + Depth: 3, // equals max depth + } + + _ = reactor.Dispatch(context.Background(), event) + + runs, _, _ := store.ListRuns(context.Background(), "deep-agent", StatusDepthExceeded, 10, 0) + if len(runs) != 1 { + t.Fatalf("expected 1 depth_exceeded run, got %d", len(runs)) + } +} + +func TestReactorBudgetExhausted(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + insertTestAgent(t, db, "budget-agent", "reactive", "image:latest", 0, 2, 5) + + store := NewStore(db) + agentStore := agents.NewSQLiteAgentStore(db) + runner := &fakeRunner{available: true} + logger := slog.Default() + + reactor := New(store, agentStore, runner, logger) + + // Record 2 existing runs today + for i := 0; i < 2; i++ { + _, _ = store.InsertRun(context.Background(), &ReactiveRun{ + AgentName: "budget-agent", + TriggerEvent: "message.received", + Status: StatusSucceeded, + }) + } + + event := dispatcher.MessageEvent{ + EventType: "message.received", + MessageID: 10, + FromAgent: "algis", + ToAgent: "budget-agent", + Body: "one more", + } + + _ = reactor.Dispatch(context.Background(), event) + + runs, _, _ := store.ListRuns(context.Background(), "budget-agent", StatusBudgetExhausted, 10, 0) + if len(runs) != 1 { + t.Fatalf("expected 1 budget_exhausted run, got %d", len(runs)) + } +} + +func TestReactorCooldownSkipped(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + insertTestAgent(t, db, "cool-agent", "reactive", "image:latest", 600, 8, 5) + + store := NewStore(db) + agentStore := agents.NewSQLiteAgentStore(db) + runner := &fakeRunner{available: true} + logger := slog.Default() + + reactor := New(store, agentStore, runner, logger) + + // Record a recent run + now := time.Now().UTC() + _, _ = store.InsertRun(context.Background(), &ReactiveRun{ + AgentName: "cool-agent", + TriggerEvent: "message.received", + Status: StatusSucceeded, + }) + // Hack: the above uses CURRENT_TIMESTAMP which is "now", so cooldown should be active + + event := dispatcher.MessageEvent{ + EventType: "message.received", + MessageID: 10, + FromAgent: "algis", + ToAgent: "cool-agent", + Body: "too soon", + } + + _ = reactor.Dispatch(context.Background(), event) + _ = now // avoid unused + + runs, _, _ := store.ListRuns(context.Background(), "cool-agent", StatusCooldownSkipped, 10, 0) + if len(runs) != 1 { + t.Fatalf("expected 1 cooldown_skipped run, got %d", len(runs)) + } + + // Check pending_work was set + agent, _ := agentStore.GetAgentByName(context.Background(), "cool-agent") + if !agent.PendingWork { + t.Error("expected pending_work to be set after cooldown skip") + } +} + +func TestReactorSequentialExecution(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + insertTestAgent(t, db, "busy-agent", "reactive", "image:latest", 0, 8, 5) + + store := NewStore(db) + agentStore := agents.NewSQLiteAgentStore(db) + runner := &fakeRunner{available: true} + logger := slog.Default() + + reactor := New(store, agentStore, runner, logger) + + // First trigger — should succeed + event1 := dispatcher.MessageEvent{ + EventType: "message.received", + MessageID: 1, + FromAgent: "algis", + ToAgent: "busy-agent", + Body: "first", + } + _ = reactor.Dispatch(context.Background(), event1) + + // Second trigger — agent is running, should queue + event2 := dispatcher.MessageEvent{ + EventType: "message.received", + MessageID: 2, + FromAgent: "algis", + ToAgent: "busy-agent", + Body: "second", + } + _ = reactor.Dispatch(context.Background(), event2) + + // Check: one running, one queued + running, _, _ := store.ListRuns(context.Background(), "busy-agent", StatusRunning, 10, 0) + queued, _, _ := store.ListRuns(context.Background(), "busy-agent", StatusQueued, 10, 0) + + if len(running) != 1 { + t.Errorf("expected 1 running, got %d", len(running)) + } + if len(queued) != 1 { + t.Errorf("expected 1 queued, got %d", len(queued)) + } + + // Check pending_work is set + agent, _ := agentStore.GetAgentByName(context.Background(), "busy-agent") + if !agent.PendingWork { + t.Error("expected pending_work to be set") + } +} + +func TestReactorSelfMentionIgnored(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + insertTestAgent(t, db, "self-agent", "reactive", "image:latest", 0, 8, 5) + + store := NewStore(db) + agentStore := agents.NewSQLiteAgentStore(db) + runner := &fakeRunner{available: true} + logger := slog.Default() + + reactor := New(store, agentStore, runner, logger) + + // Agent mentions itself + event := dispatcher.MessageEvent{ + EventType: "message.mentioned", + MessageID: 1, + FromAgent: "self-agent", + Body: "hey @self-agent", + MentionedAgents: []string{"self-agent"}, + } + + _ = reactor.Dispatch(context.Background(), event) + + runs, total, _ := store.ListRuns(context.Background(), "self-agent", "", 10, 0) + if total != 0 || len(runs) != 0 { + t.Errorf("expected 0 runs for self-mention, got %d", total) + } +} + +func TestReactorSuccessfulTrigger(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + + envJSON, _ := json.Marshal(map[string]string{ + "AGENT_GIT_REPO": "Dumbris/test-agent", + }) + _, _ = db.Exec( + `INSERT INTO agents (name, display_name, type, owner_id, trigger_mode, cooldown_seconds, daily_trigger_budget, max_trigger_depth, k8s_image, k8s_env_json, k8s_resource_preset) + VALUES (?, ?, 'ai', 1, 'reactive', 0, 8, 5, 'image:latest', ?, 'default')`, + "test-agent", "Test Agent", string(envJSON), + ) + + store := NewStore(db) + agentStore := agents.NewSQLiteAgentStore(db) + runner := &fakeRunner{available: true} + logger := slog.Default() + + reactor := New(store, agentStore, runner, logger) + + event := dispatcher.MessageEvent{ + EventType: "message.received", + MessageID: 42, + FromAgent: "algis", + ToAgent: "test-agent", + Body: "research this topic", + } + + err := reactor.Dispatch(context.Background(), event) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + + // Verify job was created + if runner.lastJobName == "" { + t.Fatal("expected K8s Job to be created") + } + + // Verify run record + runs, _, _ := store.ListRuns(context.Background(), "test-agent", StatusRunning, 10, 0) + if len(runs) != 1 { + t.Fatalf("expected 1 running run, got %d", len(runs)) + } + run := runs[0] + if run.TriggerFrom != "algis" { + t.Errorf("expected trigger_from=algis, got %s", run.TriggerFrom) + } + if run.TriggerEvent != "message.received" { + t.Errorf("expected trigger_event=message.received, got %s", run.TriggerEvent) + } + + // Verify env vars passed to job + if runner.lastEnv["SYNAPBUS_TRIGGER_DEPTH"] != "0" { + t.Errorf("expected SYNAPBUS_TRIGGER_DEPTH=0, got %s", runner.lastEnv["SYNAPBUS_TRIGGER_DEPTH"]) + } + if runner.lastEnv["AGENT_GIT_REPO"] != "Dumbris/test-agent" { + t.Errorf("expected AGENT_GIT_REPO from k8s_env_json, got %s", runner.lastEnv["AGENT_GIT_REPO"]) + } +} + +// fakeRunner is a test double for k8spkg.JobRunner. +type fakeRunner struct { + available bool + lastJobName string + lastEnv map[string]string + callCount int +} + +func (f *fakeRunner) IsAvailable() bool { return f.available } +func (f *fakeRunner) GetNamespace() string { return "test-ns" } +func (f *fakeRunner) GetJobLogs(_ context.Context, _, _ string) (string, error) { + return "test logs", nil +} +func (f *fakeRunner) CreateJob(_ context.Context, handler *k8spkg.K8sHandler, msg *k8spkg.JobMessage) (string, error) { + f.callCount++ + f.lastJobName = fmt.Sprintf("synapbus-%s-%d", handler.AgentName, msg.MessageID) + f.lastEnv = make(map[string]string) + for k, v := range handler.Env { + f.lastEnv[k] = v + } + return f.lastJobName, nil +} diff --git a/internal/reactor/store.go b/internal/reactor/store.go new file mode 100644 index 0000000..74d31f0 --- /dev/null +++ b/internal/reactor/store.go @@ -0,0 +1,298 @@ +package reactor + +import ( + "context" + "database/sql" + "fmt" + "time" +) + +// RunStatus constants for reactive_runs. +const ( + StatusQueued = "queued" + StatusRunning = "running" + StatusSucceeded = "succeeded" + StatusFailed = "failed" + StatusCooldownSkipped = "cooldown_skipped" + StatusBudgetExhausted = "budget_exhausted" + StatusDepthExceeded = "depth_exceeded" +) + +// ReactiveRun represents a single trigger evaluation and its outcome. +type ReactiveRun struct { + ID int64 `json:"id"` + AgentName string `json:"agent_name"` + TriggerMessageID *int64 `json:"trigger_message_id,omitempty"` + TriggerEvent string `json:"trigger_event"` + TriggerDepth int `json:"trigger_depth"` + TriggerFrom string `json:"trigger_from,omitempty"` + Status string `json:"status"` + K8sJobName string `json:"k8s_job_name,omitempty"` + K8sNamespace string `json:"k8s_namespace,omitempty"` + StartedAt *time.Time `json:"started_at,omitempty"` + CompletedAt *time.Time `json:"completed_at,omitempty"` + DurationMs *int64 `json:"duration_ms,omitempty"` + ErrorLog string `json:"error_log,omitempty"` + TokenCostJSON string `json:"token_cost_json,omitempty"` + CreatedAt time.Time `json:"created_at"` +} + +// Store handles SQLite persistence for reactive runs. +type Store struct { + db *sql.DB +} + +// NewStore creates a new reactor store. +func NewStore(db *sql.DB) *Store { + return &Store{db: db} +} + +// InsertRun creates a new reactive_runs record. +func (s *Store) InsertRun(ctx context.Context, run *ReactiveRun) (int64, error) { + now := time.Now().UTC() + run.CreatedAt = now + nowStr := now.Format(time.RFC3339) + var startedAtStr *string + if run.StartedAt != nil { + s := run.StartedAt.UTC().Format(time.RFC3339) + startedAtStr = &s + } + result, err := s.db.ExecContext(ctx, + `INSERT INTO reactive_runs (agent_name, trigger_message_id, trigger_event, trigger_depth, trigger_from, status, k8s_job_name, k8s_namespace, started_at, error_log, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + run.AgentName, run.TriggerMessageID, run.TriggerEvent, run.TriggerDepth, + run.TriggerFrom, run.Status, run.K8sJobName, run.K8sNamespace, startedAtStr, run.ErrorLog, nowStr, + ) + if err != nil { + return 0, fmt.Errorf("insert reactive run: %w", err) + } + id, err := result.LastInsertId() + if err != nil { + return 0, err + } + run.ID = id + return id, nil +} + +// UpdateRunStatus updates a run's status and optional fields. +func (s *Store) UpdateRunStatus(ctx context.Context, id int64, status string, jobName, namespace string, startedAt *time.Time) error { + var startedAtStr *string + if startedAt != nil { + str := startedAt.UTC().Format(time.RFC3339) + startedAtStr = &str + } + _, err := s.db.ExecContext(ctx, + `UPDATE reactive_runs SET status = ?, k8s_job_name = ?, k8s_namespace = ?, started_at = ? WHERE id = ?`, + status, jobName, namespace, startedAtStr, id, + ) + return err +} + +// CompleteRun marks a run as completed (succeeded or failed). +func (s *Store) CompleteRun(ctx context.Context, id int64, status, errorLog string, completedAt time.Time) error { + completedStr := completedAt.UTC().Format(time.RFC3339) + _, err := s.db.ExecContext(ctx, + `UPDATE reactive_runs SET status = ?, error_log = ?, completed_at = ?, + duration_ms = CAST((julianday(?) - julianday(started_at)) * 86400000 AS INTEGER) + WHERE id = ?`, + status, errorLog, completedStr, completedStr, id, + ) + return err +} + +// GetRunByID returns a single run. +func (s *Store) GetRunByID(ctx context.Context, id int64) (*ReactiveRun, error) { + return s.scanRun(s.db.QueryRowContext(ctx, runSelectSQL()+` WHERE id = ?`, id)) +} + +// ListRuns returns recent runs with optional filters. +func (s *Store) ListRuns(ctx context.Context, agentName, status string, limit, offset int) ([]*ReactiveRun, int, error) { + where := "WHERE 1=1" + args := []any{} + + if agentName != "" { + where += " AND agent_name = ?" + args = append(args, agentName) + } + if status != "" { + where += " AND status = ?" + args = append(args, status) + } + + // Count total + var total int + countArgs := make([]any, len(args)) + copy(countArgs, args) + err := s.db.QueryRowContext(ctx, "SELECT COUNT(*) FROM reactive_runs "+where, countArgs...).Scan(&total) + if err != nil { + return nil, 0, err + } + + // Query with pagination + query := runSelectSQL() + " " + where + " ORDER BY created_at DESC LIMIT ? OFFSET ?" + args = append(args, limit, offset) + rows, err := s.db.QueryContext(ctx, query, args...) + if err != nil { + return nil, 0, err + } + defer rows.Close() + + runs, err := s.scanRuns(rows) + return runs, total, err +} + +// GetActiveRuns returns runs with status 'running' (for polling). +func (s *Store) GetActiveRuns(ctx context.Context) ([]*ReactiveRun, error) { + rows, err := s.db.QueryContext(ctx, runSelectSQL()+` WHERE status = 'running'`) + if err != nil { + return nil, err + } + defer rows.Close() + return s.scanRuns(rows) +} + +// CountTodayRuns counts runs that count against the daily budget for an agent. +func (s *Store) CountTodayRuns(ctx context.Context, agentName string) (int, error) { + // Compute start of today in UTC as RFC3339 + now := time.Now().UTC() + startOfDay := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, time.UTC) + startStr := startOfDay.Format(time.RFC3339) + + var count int + err := s.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM reactive_runs + WHERE agent_name = ? AND status IN ('running', 'succeeded', 'failed') + AND created_at >= ?`, + agentName, startStr, + ).Scan(&count) + return count, err +} + +// GetLastRunTime returns the created_at of the most recent countable run. +func (s *Store) GetLastRunTime(ctx context.Context, agentName string) (*time.Time, error) { + var t sql.NullString + err := s.db.QueryRowContext(ctx, + `SELECT MAX(created_at) FROM reactive_runs + WHERE agent_name = ? AND status IN ('running', 'succeeded', 'failed')`, + agentName, + ).Scan(&t) + if err != nil { + return nil, err + } + if !t.Valid || t.String == "" { + return nil, nil + } + parsed, err := parseTime(t.String) + if err != nil { + return nil, err + } + return &parsed, nil +} + +// parseTime tries multiple time formats used by SQLite / Go driver. +func parseTime(s string) (time.Time, error) { + formats := []string{ + time.RFC3339, + time.RFC3339Nano, + "2006-01-02T15:04:05Z", + "2006-01-02 15:04:05+00:00", + "2006-01-02 15:04:05", + "2006-01-02T15:04:05.999999999Z07:00", + } + for _, f := range formats { + if t, err := time.Parse(f, s); err == nil { + return t, nil + } + } + return time.Time{}, fmt.Errorf("cannot parse time %q", s) +} + +// IsAgentRunning checks if the agent has an active (running) reactive run. +func (s *Store) IsAgentRunning(ctx context.Context, agentName string) (bool, error) { + var count int + err := s.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM reactive_runs WHERE agent_name = ? AND status = 'running'`, + agentName, + ).Scan(&count) + return count > 0, err +} + +func runSelectSQL() string { + return `SELECT id, agent_name, trigger_message_id, trigger_event, trigger_depth, trigger_from, + status, k8s_job_name, k8s_namespace, started_at, completed_at, duration_ms, error_log, token_cost_json, created_at + FROM reactive_runs` +} + +func scanRunFields(r *ReactiveRun, msgID *sql.NullInt64, triggerFrom, jobName, namespace, errorLog, tokenCost *sql.NullString, startedAt, completedAt *sql.NullString, durationMs *sql.NullInt64, createdAt *string) { + if msgID.Valid { + r.TriggerMessageID = &msgID.Int64 + } + r.TriggerFrom = triggerFrom.String + r.K8sJobName = jobName.String + r.K8sNamespace = namespace.String + if startedAt.Valid && startedAt.String != "" { + if t, err := parseTime(startedAt.String); err == nil { + r.StartedAt = &t + } + } + if completedAt.Valid && completedAt.String != "" { + if t, err := parseTime(completedAt.String); err == nil { + r.CompletedAt = &t + } + } + if durationMs.Valid { + r.DurationMs = &durationMs.Int64 + } + r.ErrorLog = errorLog.String + r.TokenCostJSON = tokenCost.String + if *createdAt != "" { + if t, err := parseTime(*createdAt); err == nil { + r.CreatedAt = t + } + } +} + +func (s *Store) scanRun(row *sql.Row) (*ReactiveRun, error) { + var r ReactiveRun + var msgID sql.NullInt64 + var triggerFrom, jobName, namespace, errorLog, tokenCost sql.NullString + var startedAt, completedAt sql.NullString + var durationMs sql.NullInt64 + var createdAt string + + err := row.Scan( + &r.ID, &r.AgentName, &msgID, &r.TriggerEvent, &r.TriggerDepth, &triggerFrom, + &r.Status, &jobName, &namespace, &startedAt, &completedAt, &durationMs, &errorLog, &tokenCost, &createdAt, + ) + if err != nil { + return nil, err + } + scanRunFields(&r, &msgID, &triggerFrom, &jobName, &namespace, &errorLog, &tokenCost, &startedAt, &completedAt, &durationMs, &createdAt) + return &r, nil +} + +func (s *Store) scanRuns(rows *sql.Rows) ([]*ReactiveRun, error) { + var runs []*ReactiveRun + for rows.Next() { + var r ReactiveRun + var msgID sql.NullInt64 + var triggerFrom, jobName, namespace, errorLog, tokenCost sql.NullString + var startedAt, completedAt sql.NullString + var durationMs sql.NullInt64 + var createdAt string + + err := rows.Scan( + &r.ID, &r.AgentName, &msgID, &r.TriggerEvent, &r.TriggerDepth, &triggerFrom, + &r.Status, &jobName, &namespace, &startedAt, &completedAt, &durationMs, &errorLog, &tokenCost, &createdAt, + ) + if err != nil { + return nil, err + } + scanRunFields(&r, &msgID, &triggerFrom, &jobName, &namespace, &errorLog, &tokenCost, &startedAt, &completedAt, &durationMs, &createdAt) + runs = append(runs, &r) + } + if runs == nil { + runs = []*ReactiveRun{} + } + return runs, rows.Err() +} diff --git a/internal/storage/schema/015_reactive_triggers.sql b/internal/storage/schema/015_reactive_triggers.sql new file mode 100644 index 0000000..2ce32d4 --- /dev/null +++ b/internal/storage/schema/015_reactive_triggers.sql @@ -0,0 +1,35 @@ +-- 013: Reactive agent triggering +-- Extends agents with trigger configuration, adds reactive_runs tracking table. + +-- Extend agents table with reactive trigger configuration +ALTER TABLE agents ADD COLUMN trigger_mode TEXT NOT NULL DEFAULT 'passive'; +ALTER TABLE agents ADD COLUMN cooldown_seconds INTEGER NOT NULL DEFAULT 600; +ALTER TABLE agents ADD COLUMN daily_trigger_budget INTEGER NOT NULL DEFAULT 8; +ALTER TABLE agents ADD COLUMN max_trigger_depth INTEGER NOT NULL DEFAULT 5; +ALTER TABLE agents ADD COLUMN k8s_image TEXT; +ALTER TABLE agents ADD COLUMN k8s_env_json TEXT; +ALTER TABLE agents ADD COLUMN k8s_resource_preset TEXT NOT NULL DEFAULT 'default'; +ALTER TABLE agents ADD COLUMN pending_work INTEGER NOT NULL DEFAULT 0; + +-- Reactive trigger runs: tracks every trigger evaluation and K8s job lifecycle +CREATE TABLE reactive_runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + agent_name TEXT NOT NULL REFERENCES agents(name), + trigger_message_id INTEGER, + trigger_event TEXT NOT NULL, + trigger_depth INTEGER NOT NULL DEFAULT 0, + trigger_from TEXT, + status TEXT NOT NULL DEFAULT 'queued', + k8s_job_name TEXT, + k8s_namespace TEXT, + started_at DATETIME, + completed_at DATETIME, + duration_ms INTEGER, + error_log TEXT, + token_cost_json TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_reactive_runs_agent_created ON reactive_runs(agent_name, created_at); +CREATE INDEX idx_reactive_runs_status ON reactive_runs(status); +CREATE INDEX idx_reactive_runs_agent_status ON reactive_runs(agent_name, status); diff --git a/internal/storage/schema/016_agent_query_views.sql b/internal/storage/schema/016_agent_query_views.sql new file mode 100644 index 0000000..5b0b565 --- /dev/null +++ b/internal/storage/schema/016_agent_query_views.sql @@ -0,0 +1,58 @@ +-- 016: Agent SQL query views +-- These views are used by the 'query' action to give agents read access +-- to messages they can see. The views expose a stable schema that agents +-- can query via SQL. Access control is enforced at the Go layer by +-- rewriting queries to filter by agent name. + +-- Note: SQLite views cannot be parameterized. The Go query executor +-- wraps agent queries in a CTE that filters by the authenticated agent's +-- access (own DMs + joined channels). These views provide the base schema. + +-- my_messages: All messages accessible to the calling agent +CREATE VIEW IF NOT EXISTS v_agent_messages AS +SELECT + m.id, + m.body, + m.from_agent, + m.to_agent, + m.priority, + m.status, + m.metadata, + m.created_at, + m.updated_at, + c.name AS channel_name, + m.channel_id, + m.reply_to, + m.conversation_id +FROM messages m +LEFT JOIN channels c ON c.id = m.channel_id; + +-- my_channels: Channels the calling agent has joined +CREATE VIEW IF NOT EXISTS v_agent_channels AS +SELECT + c.id, + c.name, + c.description, + c.type, + c.topic, + c.is_private, + c.created_at, + cm.joined_at AS member_since +FROM channels c +JOIN channel_members cm ON cm.channel_id = c.id; + +-- channel_messages: Messages in channels (filtered by membership at Go layer) +CREATE VIEW IF NOT EXISTS v_channel_messages AS +SELECT + m.id, + m.body, + m.from_agent, + m.priority, + m.status, + m.metadata, + m.created_at, + c.name AS channel_name, + m.channel_id, + m.reply_to +FROM messages m +JOIN channels c ON c.id = m.channel_id; diff --git a/internal/storage/sqlite.go b/internal/storage/sqlite.go index fbe1c30..97863a0 100644 --- a/internal/storage/sqlite.go +++ b/internal/storage/sqlite.go @@ -12,17 +12,22 @@ import ( _ "modernc.org/sqlite" ) -// DB wraps a *sql.DB with SynapBus-specific configuration. +// DB wraps a write-only *sql.DB and an optional read-only *sql.DB +// for split connection pool architecture. The write pool has MaxOpenConns=1 +// to serialize writes and eliminate SQLITE_BUSY errors. The read pool has +// MaxOpenConns=8 and query_only=ON for safe concurrent reads. type DB struct { - *sql.DB + *sql.DB // Write pool (MaxOpenConns=1) + ReadDB *sql.DB // Read pool (MaxOpenConns=8, query_only=ON) — nil for :memory: DBs } -// New opens a SQLite database with WAL mode, busy_timeout, and foreign keys enabled. -// If dataDir is empty or ":memory:", an in-memory database is used. +// New opens a SQLite database with WAL mode, split read/write pools, and foreign keys. +// If dataDir is empty or ":memory:", an in-memory database is used (single pool, no split). func New(ctx context.Context, dataDir string) (*DB, error) { var dsn string + isMemory := dataDir == "" || dataDir == ":memory:" - if dataDir == "" || dataDir == ":memory:" { + if isMemory { dsn = ":memory:" } else { if err := os.MkdirAll(dataDir, 0o755); err != nil { @@ -31,16 +36,76 @@ func New(ctx context.Context, dataDir string) (*DB, error) { dsn = filepath.Join(dataDir, "synapbus.db") } - db, err := sql.Open("sqlite", dsn) + // Open WRITE pool (single connection, serializes all writes) + writeDB, err := openPool(ctx, dsn, poolConfig{ + maxOpen: 1, + maxIdle: 1, + queryOnly: false, + label: "write", + }) if err != nil { - return nil, fmt.Errorf("open database: %w", err) + return nil, fmt.Errorf("open write pool: %w", err) } - // Configure SQLite pragmas + result := &DB{DB: writeDB} + + // For file-based databases, open a separate READ pool + if !isMemory { + readDB, err := openPool(ctx, dsn, poolConfig{ + maxOpen: 8, + maxIdle: 4, + queryOnly: true, + label: "read", + }) + if err != nil { + writeDB.Close() + return nil, fmt.Errorf("open read pool: %w", err) + } + result.ReadDB = readDB + } + + // Verify settings on write pool + var journalMode string + if err := writeDB.QueryRowContext(ctx, "PRAGMA journal_mode").Scan(&journalMode); err != nil { + result.Close() + return nil, fmt.Errorf("verify journal_mode: %w", err) + } + + slog.Info("database opened", + "dsn", dsn, + "journal_mode", journalMode, + "write_pool", "MaxOpenConns=1", + "read_pool_enabled", result.ReadDB != nil, + ) + + return result, nil +} + +type poolConfig struct { + maxOpen int + maxIdle int + queryOnly bool + label string +} + +func openPool(ctx context.Context, dsn string, cfg poolConfig) (*sql.DB, error) { + db, err := sql.Open("sqlite", dsn) + if err != nil { + return nil, fmt.Errorf("open %s pool: %w", cfg.label, err) + } + + db.SetMaxOpenConns(cfg.maxOpen) + db.SetMaxIdleConns(cfg.maxIdle) + pragmas := []string{ "PRAGMA journal_mode=WAL", - "PRAGMA busy_timeout=5000", + "PRAGMA busy_timeout=15000", "PRAGMA foreign_keys=ON", + "PRAGMA synchronous=NORMAL", + "PRAGMA wal_autocheckpoint=1000", + } + if cfg.queryOnly { + pragmas = append(pragmas, "PRAGMA query_only=ON") } for _, pragma := range pragmas { @@ -50,22 +115,31 @@ func New(ctx context.Context, dataDir string) (*DB, error) { } } - // Verify settings - var journalMode string - if err := db.QueryRowContext(ctx, "PRAGMA journal_mode").Scan(&journalMode); err != nil { - db.Close() - return nil, fmt.Errorf("verify journal_mode: %w", err) + return db, nil +} + +// QueryDB returns the read pool if available, otherwise falls back to the write pool. +// Use this for all SELECT queries to avoid blocking writers. +func (db *DB) QueryDB() *sql.DB { + if db.ReadDB != nil { + return db.ReadDB } - - slog.Info("database opened", - "dsn", dsn, - "journal_mode", journalMode, - ) - - return &DB{DB: db}, nil + return db.DB } -// Close closes the database connection. +// Close closes both the write and read database connections. func (db *DB) Close() error { - return db.DB.Close() + var errs []error + if db.ReadDB != nil { + if err := db.ReadDB.Close(); err != nil { + errs = append(errs, fmt.Errorf("close read pool: %w", err)) + } + } + if err := db.DB.Close(); err != nil { + errs = append(errs, fmt.Errorf("close write pool: %w", err)) + } + if len(errs) > 0 { + return errs[0] + } + return nil } diff --git a/internal/storage/sqlite_test.go b/internal/storage/sqlite_test.go index 2406cfc..79fbd88 100644 --- a/internal/storage/sqlite_test.go +++ b/internal/storage/sqlite_test.go @@ -68,11 +68,11 @@ func TestNew(t *testing.T) { if err != nil { t.Fatalf("failed to query busy_timeout: %v", err) } - if timeout != 5000 { - t.Errorf("busy_timeout = %d, want 5000", timeout) + if timeout != 15000 { + t.Errorf("busy_timeout = %d, want 15000", timeout) } - // Verify database is usable + // Verify database is usable via write pool _, err = db.Exec("CREATE TABLE test (id INTEGER PRIMARY KEY)") if err != nil { t.Fatalf("failed to create test table: %v", err) @@ -80,3 +80,77 @@ func TestNew(t *testing.T) { }) } } + +func TestSplitPools(t *testing.T) { + ctx := context.Background() + dir := t.TempDir() + + db, err := New(ctx, dir) + if err != nil { + t.Fatalf("New() error: %v", err) + } + defer db.Close() + + // Run migrations to create tables + if err := RunMigrations(ctx, db.DB); err != nil { + t.Fatalf("migrations: %v", err) + } + + // Verify read pool exists for file-based DB + if db.ReadDB == nil { + t.Fatal("expected ReadDB to be non-nil for file-based database") + } + + // Verify QueryDB returns read pool + if db.QueryDB() != db.ReadDB { + t.Error("QueryDB() should return ReadDB when available") + } + + // Create user first (FK requirement) + _, err = db.Exec("INSERT INTO users (id, username, password_hash, display_name) VALUES (1, 'testuser', 'hash', 'Test')") + if err != nil { + t.Fatalf("create user: %v", err) + } + + // Verify write pool can write + _, err = db.Exec("INSERT INTO agents (name, display_name, type, capabilities, owner_id, api_key_hash, status) VALUES ('test-agent', 'Test', 'ai', '{}', 1, 'hash', 'active')") + if err != nil { + t.Fatalf("write pool should allow writes: %v", err) + } + + // Verify read pool can read + var name string + err = db.ReadDB.QueryRow("SELECT name FROM agents WHERE name = 'test-agent'").Scan(&name) + if err != nil { + t.Fatalf("read pool should allow reads: %v", err) + } + if name != "test-agent" { + t.Errorf("expected 'test-agent', got %q", name) + } + + // Verify read pool rejects writes + _, err = db.ReadDB.Exec("INSERT INTO agents (name, display_name, type, capabilities, owner_id, api_key_hash, status) VALUES ('bad', 'Bad', 'ai', '{}', 1, 'hash', 'active')") + if err == nil { + t.Fatal("read pool should reject writes (query_only=ON)") + } +} + +func TestInMemoryNoSplitPool(t *testing.T) { + ctx := context.Background() + + db, err := New(ctx, ":memory:") + if err != nil { + t.Fatalf("New() error: %v", err) + } + defer db.Close() + + // In-memory DB should NOT have a separate read pool + if db.ReadDB != nil { + t.Error("in-memory DB should not have a separate ReadDB") + } + + // QueryDB should fall back to write pool + if db.QueryDB() != db.DB { + t.Error("QueryDB() should return write pool for in-memory DB") + } +} diff --git a/internal/trace/metrics.go b/internal/trace/metrics.go index 4d5412f..e229569 100644 --- a/internal/trace/metrics.go +++ b/internal/trace/metrics.go @@ -6,6 +6,9 @@ import ( "sort" "sync" "sync/atomic" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/common/expfmt" ) // Metrics provides Prometheus-compatible metrics for SynapBus. @@ -93,6 +96,17 @@ func (m *Metrics) WritePrometheus(w io.Writer) { fmt.Fprintf(w, "# HELP synapbus_active_agents Number of currently active agents.\n") fmt.Fprintf(w, "# TYPE synapbus_active_agents gauge\n") fmt.Fprintf(w, "synapbus_active_agents %d\n", m.activeAgents.Load()) + fmt.Fprintf(w, "\n") + + // Append metrics from the standard Prometheus registry (reactor metrics, etc.) + mfs, _ := prometheus.DefaultGatherer.Gather() + enc := expfmt.NewEncoder(w, expfmt.NewFormat(expfmt.TypeTextPlain)) + for _, mf := range mfs { + // Only include our custom metrics, skip Go runtime metrics + if name := mf.GetName(); len(name) > 8 && name[:8] == "synapbus" { + _ = enc.Encode(mf) + } + } } // NullMetrics is a no-op metrics implementation for when metrics are disabled. diff --git a/internal/web/dist/index.html b/internal/web/dist/index.html index 2d8b7fe..94c13a7 100644 --- a/internal/web/dist/index.html +++ b/internal/web/dist/index.html @@ -11,30 +11,30 @@ - - + + - - - - - - + + + + + +
+ + + Agent Runs - SynapBus + + +
+

Agent Runs

+ + + {#if reactiveAgents.length > 0} +
+ {#each reactiveAgents as agent} +
+
+ {agent.name} + {agent.state} +
+
+ + {agent.today_runs}/{agent.daily_trigger_budget} + runs today + + + {agent.cooldown_seconds}s + cooldown + +
+
+ {/each} +
+ {/if} + + +
+ + + {total} runs +
+ + + {#if loading} +
Loading...
+ {:else if runsList.length === 0} +
No reactive runs found.
+ {:else} +
+ {#each runsList as run} +
+ + + {#if expandedRun === run.id} +
+
+ Run ID + {run.id} +
+
+ K8s Job + {run.k8s_job_name || '-'} +
+
+ Depth + {run.trigger_depth} +
+ {#if run.trigger_message_id} + + {/if} + {#if run.error_log} +
+
Error Log
+
{run.error_log}
+
+ {/if} + {#if run.status === 'failed'} + + {/if} +
+ {/if} +
+ {/each} +
+ {/if} +
+ +