feat(014): implement reactive agent triggering engine

- Migration 015: extends agents with trigger config, adds reactive_runs table
- Reactor engine: decision chain (mode, depth, budget, cooldown, sequential)
- Reactor store: SQLite persistence for runs with RFC3339 timestamps
- Reactor poller: K8s Job status polling (15s interval)
- Failure notifier: system DM to owner on job failure
- REST API: /api/runs, /api/runs/:id, /api/runs/:id/retry, /api/agents/reactive
- Agent model: trigger_mode, cooldown, budget, depth, k8s_image, pending_work
- K8s runner: GetClientset() for poller
- All 28 test packages pass (8 new reactor tests)

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Algis Dumbris
2026-03-25 17:42:48 +02:00
co-authored by Claude Opus 4.6
parent 68f356b5e3
commit 6afe1853ad
12 changed files with 1691 additions and 16 deletions
+16 -2
View File
@@ -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)
+108 -14
View File
@@ -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 {
+17
View File
@@ -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"`
}
+16
View File
@@ -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)
+165
View File
@@ -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,
})
}
+5
View File
@@ -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
}
+52
View File
@@ -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] + "..."
}
+212
View File
@@ -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)
}
+337
View File
@@ -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)
}
+430
View File
@@ -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
}
+298
View File
@@ -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()
}
@@ -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);