diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index 8efbb53..6c8b1c9 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -39,6 +39,7 @@ import ( "github.com/synapbus/synapbus/internal/jsruntime" k8spkg "github.com/synapbus/synapbus/internal/k8s" mcpserver "github.com/synapbus/synapbus/internal/mcp" + 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 +468,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() @@ -641,6 +653,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/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..9aeac47 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 } 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..488bdb8 --- /dev/null +++ b/internal/reactor/poller.go @@ -0,0 +1,212 @@ +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" + + 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() + + if success { + _ = 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") + } + + _ = 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..9134734 --- /dev/null +++ b/internal/reactor/reactor.go @@ -0,0 +1,337 @@ +// 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" +) + +// 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), + } + + now := time.Now().UTC() + run := &ReactiveRun{ + AgentName: agent.Name, + TriggerMessageID: &event.MessageID, + TriggerEvent: event.EventType, + TriggerDepth: depth, + TriggerFrom: event.FromAgent, + Status: StatusRunning, + StartedAt: &now, + } + + // Insert run record first + runID, err := r.store.InsertRun(ctx, run) + if err != nil { + return fmt.Errorf("insert reactive run: %w", err) + } + + // Add trigger depth env var to handler + handler.Env["SYNAPBUS_TRIGGER_DEPTH"] = fmt.Sprintf("%d", depth) + + // Create K8s Job + jobName, err := r.runner.CreateJob(ctx, handler, msg) + if err != nil { + // Record failure + errMsg := fmt.Sprintf("K8s Job creation failed: %s", err.Error()) + _ = r.store.CompleteRun(ctx, runID, StatusFailed, errMsg, time.Now().UTC()) + r.notifyFailure(ctx, agent, event, 0, errMsg) + return fmt.Errorf("create K8s job: %w", err) + } + + // Update run with job name + ns := handler.Namespace + if ns == "" { + ns = r.runner.GetNamespace() + } + _ = r.store.UpdateRunStatus(ctx, runID, StatusRunning, jobName, ns, &now) + + // Clear pending_work since we're launching + _ = r.agentStore.SetPendingWork(ctx, agent.Name, false) + + 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 + memory := "256Mi" + cpu := "100m" + if agent.K8sResourcePreset == "large" { + memory = "2Gi" + cpu = "1000m" + } + + timeout := 600 // 10 minutes default + + return &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", + } +} + +// 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) { + 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);