From bee1709f4aebba9826419dd1a1aa40f2f2aaa1c1 Mon Sep 17 00:00:00 2001 From: 0sm0s1z Date: Wed, 22 Apr 2026 20:15:03 -0700 Subject: [PATCH] feat(server): consume engine.commands as defense-in-depth (PR 5) PR 3 made agent.template.sync.jobs the canonical path for "agents need to re-pull template:meta:*". This change adds a defense-in-depth consumer on the legacy engine.commands queue so any producer that still publishes "internal:template upload" / "internal:template delete" there (stale code paths, third-party integrations) gets routed back into the same notify-agents helper instead of being silently dropped. - New EngineCommandQueueProcessor mirrors TemplateSyncQueueProcessor's shape (Listen goroutine, Stop hook, context-scoped processing). - isTemplateNotifyCommand keeps command classification in a tiny pure helper so the routing logic is unit-testable without RabbitMQ. - Unknown commands (engine health pings, scan dispatch, etc.) are ack-and-ignored so the consumer doesn't fight other producers. - Wired into server.NewServer alongside the existing template-sync consumer; failures only log and never abort startup. --- internal/server/engine_commands_consumer.go | 121 ++++++++++++++++++ .../server/engine_commands_consumer_test.go | 42 ++++++ internal/server/server.go | 19 ++- 3 files changed, 180 insertions(+), 2 deletions(-) create mode 100644 internal/server/engine_commands_consumer.go create mode 100644 internal/server/engine_commands_consumer_test.go diff --git a/internal/server/engine_commands_consumer.go b/internal/server/engine_commands_consumer.go new file mode 100644 index 0000000..95c2b8e --- /dev/null +++ b/internal/server/engine_commands_consumer.go @@ -0,0 +1,121 @@ +package server + +import ( + "context" + "encoding/json" + "time" + + "github.com/SiriusScan/go-api/sirius/queue" + "go.uber.org/zap" +) + +const ( + engineCommandsQueueName = "engine.commands" +) + +// EngineCommandMessage is the legacy envelope still used by some +// producers (notably the pre-PR3 sirius-api delete path and any +// third-party integrations that publish directly to engine.commands). +// +// Two commands matter for template-sync purposes: +// - "internal:template upload" - new custom template was written. +// - "internal:template delete" - custom template was removed. +// +// In both cases we just need to nudge connected agents to re-pull the +// template:meta:* namespace; the repository-level sync is unaffected. +// Anything else is acked and ignored so we don't block other producers +// that might land on this queue. +type EngineCommandMessage struct { + Command string `json:"command"` + TemplateID string `json:"template_id,omitempty"` + Timestamp string `json:"timestamp,omitempty"` +} + +// isTemplateNotifyCommand reports whether the legacy command string +// should trigger an agent re-pull. Kept as a small pure helper so the +// classification can be unit-tested without spinning up RabbitMQ. +func isTemplateNotifyCommand(cmd string) bool { + switch cmd { + case "internal:template upload", "internal:template delete": + return true + default: + return false + } +} + +// EngineCommandQueueProcessor is the defense-in-depth consumer that +// catches any producer still publishing to engine.commands. PR 3 made +// the canonical path agent.template.sync.jobs / notify_agents, so this +// consumer is strictly additive: if PR 3's path keeps working this +// consumer logs but never has anything to do. +type EngineCommandQueueProcessor struct { + repositoryMgr *RepositoryManager + logger *zap.Logger + ctx context.Context + cancelFunc context.CancelFunc +} + +// NewEngineCommandQueueProcessor wires the consumer to the repository +// manager (notify-agents path). +func NewEngineCommandQueueProcessor( + repositoryMgr *RepositoryManager, + logger *zap.Logger, +) *EngineCommandQueueProcessor { + ctx, cancel := context.WithCancel(context.Background()) + return &EngineCommandQueueProcessor{ + repositoryMgr: repositoryMgr, + logger: logger, + ctx: ctx, + cancelFunc: cancel, + } +} + +// StartListening spawns the goroutine that pumps messages off +// engine.commands and routes recognized template commands into the +// notify-agents helper. Mirrors the shape of TemplateSyncQueueProcessor. +func (ecp *EngineCommandQueueProcessor) StartListening() error { + ecp.logger.Info("Starting engine.commands queue processor (defense-in-depth)") + + go func() { + processor := func(msg string) { + ecp.logger.Debug("Received engine.commands message", zap.String("message", msg)) + + var cmd EngineCommandMessage + if err := json.Unmarshal([]byte(msg), &cmd); err != nil { + ecp.logger.Warn("Failed to parse engine.commands message", + zap.Error(err), + zap.String("raw", msg)) + return + } + + processCtx, cancel := context.WithTimeout(ecp.ctx, 2*time.Minute) + defer cancel() + + if isTemplateNotifyCommand(cmd.Command) { + ecp.logger.Info("Routing engine.commands template event to notify-agents", + zap.String("command", cmd.Command), + zap.String("template_id", cmd.TemplateID)) + ecp.repositoryMgr.NotifyAgents(processCtx) + } else { + // Other producers (engine health pings, scan dispatch, + // etc.) legitimately share this queue. Ack-and-ignore. + ecp.logger.Debug("Unhandled engine.commands command", + zap.String("command", cmd.Command)) + } + } + + queue.Listen(engineCommandsQueueName, processor) + ecp.logger.Info("engine.commands queue processor stopped") + }() + + return nil +} + +// Stop cancels the consumer's context. The underlying queue.Listen does +// not currently honour cancellation (matches TemplateSyncQueueProcessor), +// but we keep the same lifecycle shape so future plumbing is symmetric. +func (ecp *EngineCommandQueueProcessor) Stop() { + if ecp.cancelFunc != nil { + ecp.cancelFunc() + } +} diff --git a/internal/server/engine_commands_consumer_test.go b/internal/server/engine_commands_consumer_test.go new file mode 100644 index 0000000..28b2c2a --- /dev/null +++ b/internal/server/engine_commands_consumer_test.go @@ -0,0 +1,42 @@ +package server + +import ( + "encoding/json" + "testing" +) + +func TestIsTemplateNotifyCommand(t *testing.T) { + cases := []struct { + cmd string + want bool + }{ + {"internal:template upload", true}, + {"internal:template delete", true}, + {"internal:template-scan --template foo", false}, + {"scan:start", false}, + {"", false}, + {"INTERNAL:template upload", false}, // case-sensitive on purpose; legacy producers always lowercase + } + for _, tc := range cases { + if got := isTemplateNotifyCommand(tc.cmd); got != tc.want { + t.Errorf("isTemplateNotifyCommand(%q) = %v, want %v", tc.cmd, got, tc.want) + } + } +} + +func TestEngineCommandMessage_Unmarshal(t *testing.T) { + const payload = `{"command":"internal:template upload","template_id":"smoke-test","timestamp":"2026-04-22T00:00:00Z"}` + var msg EngineCommandMessage + if err := json.Unmarshal([]byte(payload), &msg); err != nil { + t.Fatalf("unmarshal: %v", err) + } + if msg.Command != "internal:template upload" { + t.Errorf("Command = %q", msg.Command) + } + if msg.TemplateID != "smoke-test" { + t.Errorf("TemplateID = %q", msg.TemplateID) + } + if !isTemplateNotifyCommand(msg.Command) { + t.Error("expected isTemplateNotifyCommand true for parsed upload command") + } +} diff --git a/internal/server/server.go b/internal/server/server.go index 0c0f092..9117acf 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -115,8 +115,9 @@ type Server struct { // Template management templateManager *ServerTemplateManager repositoryManager *RepositoryManager - syncQueueProcessor *TemplateSyncQueueProcessor - valkeyClient valkey.Client + syncQueueProcessor *TemplateSyncQueueProcessor + engineCommandConsumer *EngineCommandQueueProcessor + valkeyClient valkey.Client // KVStore for agent token authentication kvStore goapistore.KVStore @@ -210,6 +211,20 @@ func NewServer(cfg *config.ServerConfig, logger *zap.Logger) (*Server, error) { logger.Info("Template sync queue processor started") } + // Defense-in-depth: also consume the legacy engine.commands queue + // so any producer that still publishes "internal:template upload" + // or "internal:template delete" there gets routed back into the + // notify-agents path. Strictly redundant with PR 3's canonical + // agent.template.sync.jobs flow when everything works; catches + // future drift without noise. + engineCmdConsumer := NewEngineCommandQueueProcessor(repositoryManager, logger) + server.engineCommandConsumer = engineCmdConsumer + if err := engineCmdConsumer.StartListening(); err != nil { + logger.Error("Failed to start engine.commands consumer", zap.Error(err)) + } else { + logger.Info("engine.commands consumer started") + } + // Perform initial sync of all repositories go func() { ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)