Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
121 changes: 121 additions & 0 deletions internal/server/engine_commands_consumer.go
Original file line number Diff line number Diff line change
@@ -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))
Comment on lines +100 to +103

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Do not discard non-template engine.commands messages

This new consumer subscribes to engine.commands and explicitly treats unknown commands as "Ack-and-ignore"; on a shared RabbitMQ queue, that means non-template messages can be consumed here and silently dropped instead of reaching the component that actually handles them. The risk is highest in deployments where other command types (e.g. scan/health events noted in the comment) are published to the same queue, because those events may disappear intermittently once this service is running.

Useful? React with 👍 / 👎.

}
}

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()
}
}
42 changes: 42 additions & 0 deletions internal/server/engine_commands_consumer_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
19 changes: 17 additions & 2 deletions internal/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
Loading