From 5d0560e523ec803d9190b6fc9b7fd5814359bb18 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alex=20S=C3=A1nchez?= Date: Wed, 26 Aug 2026 16:15:17 -0600 Subject: [PATCH 1/2] i18n[frontend](soar): add WAITING, EXECUTING, DEAD execution statuses --- frontend/src/shared/i18n/locales/de.json | 7 +++++-- frontend/src/shared/i18n/locales/en.json | 7 +++++-- frontend/src/shared/i18n/locales/es.json | 7 +++++-- frontend/src/shared/i18n/locales/fr.json | 7 +++++-- frontend/src/shared/i18n/locales/it.json | 7 +++++-- frontend/src/shared/i18n/locales/pt.json | 7 +++++-- frontend/src/shared/i18n/locales/ru.json | 7 +++++-- 7 files changed, 35 insertions(+), 14 deletions(-) diff --git a/frontend/src/shared/i18n/locales/de.json b/frontend/src/shared/i18n/locales/de.json index bdc58666a..ecddf9629 100644 --- a/frontend/src/shared/i18n/locales/de.json +++ b/frontend/src/shared/i18n/locales/de.json @@ -4795,9 +4795,12 @@ "manual": "Interaktive Konsole" }, "executionStatus": { - "EXECUTED": "Ausgeführt", + "WAITING": "Wartet", "PENDING": "Ausstehend", - "FAILED": "Fehlgeschlagen" + "EXECUTING": "Wird ausgeführt", + "EXECUTED": "Ausgeführt", + "FAILED": "Fehlgeschlagen", + "DEAD": "Nicht erreichbar" }, "nonExecutionCause": { "AGENT_OFFLINE": "Agent offline", diff --git a/frontend/src/shared/i18n/locales/en.json b/frontend/src/shared/i18n/locales/en.json index e26f301e5..79a10f847 100644 --- a/frontend/src/shared/i18n/locales/en.json +++ b/frontend/src/shared/i18n/locales/en.json @@ -5185,9 +5185,12 @@ "manual": "Interactive console" }, "executionStatus": { - "EXECUTED": "Executed", + "WAITING": "Waiting", "PENDING": "Pending", - "FAILED": "Failed" + "EXECUTING": "Executing", + "EXECUTED": "Executed", + "FAILED": "Failed", + "DEAD": "Unreachable" }, "executionOrigin": { "FLOW": "Flow", diff --git a/frontend/src/shared/i18n/locales/es.json b/frontend/src/shared/i18n/locales/es.json index cbaae547e..8464c0961 100644 --- a/frontend/src/shared/i18n/locales/es.json +++ b/frontend/src/shared/i18n/locales/es.json @@ -4921,9 +4921,12 @@ "manual": "Consola interactiva" }, "executionStatus": { - "EXECUTED": "Ejecutado", + "WAITING": "En espera", "PENDING": "Pendiente", - "FAILED": "Fallido" + "EXECUTING": "Ejecutando", + "EXECUTED": "Ejecutado", + "FAILED": "Fallido", + "DEAD": "Inaccesible" }, "executionOrigin": { "FLOW": "Flujo", diff --git a/frontend/src/shared/i18n/locales/fr.json b/frontend/src/shared/i18n/locales/fr.json index 90407326e..4b00ae73a 100644 --- a/frontend/src/shared/i18n/locales/fr.json +++ b/frontend/src/shared/i18n/locales/fr.json @@ -4799,9 +4799,12 @@ "manual": "Console interactive" }, "executionStatus": { + "WAITING": "En attente", + "PENDING": "En file d'attente", + "EXECUTING": "En cours d'exécution", "EXECUTED": "Exécuté", - "PENDING": "En attente", - "FAILED": "Échoué" + "FAILED": "Échoué", + "DEAD": "Inaccessible" }, "executionOrigin": { "FLOW": "Flux", diff --git a/frontend/src/shared/i18n/locales/it.json b/frontend/src/shared/i18n/locales/it.json index 6ff09bb2d..f1e7b6b7f 100644 --- a/frontend/src/shared/i18n/locales/it.json +++ b/frontend/src/shared/i18n/locales/it.json @@ -4799,9 +4799,12 @@ "manual": "Console interattiva" }, "executionStatus": { + "WAITING": "In attesa", + "PENDING": "In coda", + "EXECUTING": "In esecuzione", "EXECUTED": "Eseguito", - "PENDING": "In attesa", - "FAILED": "Fallito" + "FAILED": "Fallito", + "DEAD": "Non raggiungibile" }, "executionOrigin": { "FLOW": "Flusso", diff --git a/frontend/src/shared/i18n/locales/pt.json b/frontend/src/shared/i18n/locales/pt.json index cd7eb7f3f..26f9ba6f1 100644 --- a/frontend/src/shared/i18n/locales/pt.json +++ b/frontend/src/shared/i18n/locales/pt.json @@ -4921,9 +4921,12 @@ "manual": "Consola interativa" }, "executionStatus": { - "EXECUTED": "Executado", + "WAITING": "Aguardando", "PENDING": "Pendente", - "FAILED": "Falhou" + "EXECUTING": "Executando", + "EXECUTED": "Executado", + "FAILED": "Falhou", + "DEAD": "Inacessível" }, "executionOrigin": { "FLOW": "Fluxo", diff --git a/frontend/src/shared/i18n/locales/ru.json b/frontend/src/shared/i18n/locales/ru.json index dcb799281..8613b7396 100644 --- a/frontend/src/shared/i18n/locales/ru.json +++ b/frontend/src/shared/i18n/locales/ru.json @@ -4587,9 +4587,12 @@ "manual": "Интерактивная консоль" }, "executionStatus": { + "WAITING": "В ожидании", + "PENDING": "В очереди", + "EXECUTING": "Выполняется", "EXECUTED": "Выполнено", - "PENDING": "Ожидание", - "FAILED": "Ошибка" + "FAILED": "Ошибка", + "DEAD": "Недостижимо" }, "nonExecutionCause": { "AGENT_OFFLINE": "Агент офлайн", From 826c50e4db8b7ffcbdc00f5547c587abc70b4a01 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alex=20S=C3=A1nchez?= Date: Thu, 27 Aug 2026 15:44:27 -0600 Subject: [PATCH 2/2] fix[backend,agent-manager](soar): stop SOAR runs ghosting the agent stream --- agent-manager/agent/agent_imp.go | 33 +++++++++++++------- agent-manager/agent/agent_imp_test.go | 44 +++++++++++++++++++++++++++ backend/pkg/agentmanager/client.go | 14 +++++++-- 3 files changed, 77 insertions(+), 14 deletions(-) create mode 100644 agent-manager/agent/agent_imp_test.go diff --git a/agent-manager/agent/agent_imp.go b/agent-manager/agent/agent_imp.go index b1fd9f25d..2c54ca89b 100644 --- a/agent-manager/agent/agent_imp.go +++ b/agent-manager/agent/agent_imp.go @@ -296,6 +296,17 @@ func (s *AgentService) GetAgentAuth(ctx context.Context, req *ConnectorAuthReque return &ConnectorAuthResponse{Key: agent.AgentKey, TenantId: agent.TenantID}, nil } +// evictIfOwner deletes the AgentStreamMap entry for agentID only if it still +// points to stream. Prevents a slow-exiting prior AgentStream goroutine from +// clobbering the fresh entry a newly-reconnected agent installed. +func (s *AgentService) evictIfOwner(agentID uint, stream AgentService_AgentStreamServer) { + s.AgentStreamMutex.Lock() + if s.AgentStreamMap[agentID] == stream { + delete(s.AgentStreamMap, agentID) + } + s.AgentStreamMutex.Unlock() +} + func (s *AgentService) AgentStream(stream AgentService_AgentStreamServer) error { id, _, _, err := utils.GetItemsFromContext(stream.Context()) if err != nil { @@ -307,11 +318,12 @@ func (s *AgentService) AgentStream(stream AgentService_AgentStreamServer) error } idUint := uint(idInt) + // Replace any prior entry rather than rejecting the reconnect. A dead + // prior stream's goroutine may still be looping on Recv (see + // utils.WaitForReconnect) and would otherwise block the agent from + // re-registering for minutes. evictIfOwner guards the map so the old + // goroutine's eventual delete does not clobber the fresh entry. s.AgentStreamMutex.Lock() - if _, ok := s.AgentStreamMap[idUint]; ok { - s.AgentStreamMutex.Unlock() - return status.Error(codes.AlreadyExists, "stream already exists") - } s.AgentStreamMap[idUint] = stream s.AgentStreamMutex.Unlock() @@ -324,18 +336,17 @@ func (s *AgentService) AgentStream(stream AgentService_AgentStreamServer) error if err == io.EOF { err = utils.WaitForReconnect(stream.Context(), stream) if err != nil { - s.AgentStreamMutex.Lock() - delete(s.AgentStreamMap, idUint) - s.AgentStreamMutex.Unlock() - + catcher.Info("AgentStream: WaitForReconnect failed, evicting stream", + map[string]any{"agent_id": idUint, "err": err.Error(), "process": "agent-manager"}) + s.evictIfOwner(idUint, stream) return status.Error(codes.Internal, fmt.Sprintf("failed to reconnect: %v", err)) } continue } if err != nil { - s.AgentStreamMutex.Lock() - delete(s.AgentStreamMap, idUint) - s.AgentStreamMutex.Unlock() + catcher.Info("AgentStream: Recv errored, evicting stream", + map[string]any{"agent_id": idUint, "err": err.Error(), "process": "agent-manager"}) + s.evictIfOwner(idUint, stream) return status.Error(codes.Internal, fmt.Sprintf("failed to receive message: %v", err)) } diff --git a/agent-manager/agent/agent_imp_test.go b/agent-manager/agent/agent_imp_test.go new file mode 100644 index 000000000..b63adc5d3 --- /dev/null +++ b/agent-manager/agent/agent_imp_test.go @@ -0,0 +1,44 @@ +package agent + +import ( + "context" + "testing" + + "google.golang.org/grpc/metadata" +) + +// fakeAgentStream is the smallest thing that satisfies +// AgentService_AgentStreamServer for identity-comparison tests. +type fakeAgentStream struct{ id int } + +func (fakeAgentStream) Send(*BidirectionalStream) error { return nil } +func (fakeAgentStream) Recv() (*BidirectionalStream, error) { return nil, nil } +func (fakeAgentStream) SetHeader(metadata.MD) error { return nil } +func (fakeAgentStream) SendHeader(metadata.MD) error { return nil } +func (fakeAgentStream) SetTrailer(metadata.MD) {} +func (fakeAgentStream) Context() context.Context { return context.Background() } +func (fakeAgentStream) SendMsg(any) error { return nil } +func (fakeAgentStream) RecvMsg(any) error { return nil } + +// TestEvictIfOwner_LeavesForeignStream: an old goroutine returning long after +// a fresh reconnect must NOT clobber the fresh entry. +// TestEvictIfOwner_RemovesOwnedStream: the current owner cleans up on exit. +func TestEvictIfOwner(t *testing.T) { + s := &AgentService{AgentStreamMap: map[uint]AgentService_AgentStreamServer{}} + old := &fakeAgentStream{id: 1} + fresh := &fakeAgentStream{id: 2} + + s.AgentStreamMap[42] = fresh + s.evictIfOwner(42, old) + if _, ok := s.AgentStreamMap[42]; !ok { + t.Fatal("evictIfOwner clobbered a fresh stream owned by a different goroutine") + } + if s.AgentStreamMap[42] != fresh { + t.Fatal("evictIfOwner replaced the fresh entry with something else") + } + + s.evictIfOwner(42, fresh) + if _, ok := s.AgentStreamMap[42]; ok { + t.Fatal("evictIfOwner did not remove the owned entry") + } +} diff --git a/backend/pkg/agentmanager/client.go b/backend/pkg/agentmanager/client.go index 7e3f9a69e..a0336c6a9 100644 --- a/backend/pkg/agentmanager/client.go +++ b/backend/pkg/agentmanager/client.go @@ -352,7 +352,18 @@ func (c *AgentManagerClient) GetCollectorIntegrationState(ctx context.Context, c return resp, nil } +// ProcessCommand opens a bidi stream, sends one command, and reads one result. +// It does NOT call CloseSend — parity with ProcessCommandStream / Java. Half- +// closing the panel-side stream races the agent-manager's ProcessCommand +// handler (agent-manager/agent/agent_imp.go) into an EOF path that has been +// observed to leave AgentStreamMap[agentID] empty, after which every +// subsequent panel call (SOAR + console) returns codes.NotFound "agent not +// found or is disconnected". The ctx cancellation on function return is what +// tears the stream down cleanly. func (c *AgentManagerClient) ProcessCommand(ctx context.Context, cmd *agent.UtmCommand) (*agent.CommandResult, error) { + ctx, cancel := context.WithCancel(ctx) + defer cancel() + stream, err := c.panelService.ProcessCommand(ctx) if err != nil { return nil, fmt.Errorf("agentmanager: ProcessCommand open stream: %w", err) @@ -361,9 +372,6 @@ func (c *AgentManagerClient) ProcessCommand(ctx context.Context, cmd *agent.UtmC if err := stream.Send(cmd); err != nil { return nil, fmt.Errorf("agentmanager: ProcessCommand send: %w", err) } - if err := stream.CloseSend(); err != nil { - return nil, fmt.Errorf("agentmanager: ProcessCommand close send: %w", err) - } result, err := stream.Recv() if err != nil && err != io.EOF {