diff --git a/server/internal/grpc/server.go b/server/internal/grpc/server.go index c279b6d..6b7e8d5 100644 --- a/server/internal/grpc/server.go +++ b/server/internal/grpc/server.go @@ -129,6 +129,9 @@ func (s *vantageServer) CommandStream(stream pb.Vantage_CommandStreamServer) err r := m.Result log.Printf("agent %s cmd %s: success=%v %s", srv.ServerID, r.CommandId, r.Success, r.Message) } + if m.StepResult != nil { + services.StepResults.Deliver(m.StepResult) + } } }() diff --git a/server/internal/services/stepresults.go b/server/internal/services/stepresults.go new file mode 100644 index 0000000..f89b329 --- /dev/null +++ b/server/internal/services/stepresults.go @@ -0,0 +1,49 @@ +package services + +import ( + "sync" + + "github.com/mrhid6/vantage/server/internal/grpc/pb" +) + +type stepResultRegistry struct { + mu sync.Mutex + pending map[string]chan *pb.StepResult +} + +// StepResults correlates agent StepResult replies back to the workflow runner +// goroutine that dispatched the matching RunStepCmd, keyed by command_id. +var StepResults = &stepResultRegistry{pending: make(map[string]chan *pb.StepResult)} + +// Await registers interest in a command's result BEFORE the command is +// dispatched, and returns a buffered channel that receives the single result. +func (r *stepResultRegistry) Await(commandID string) <-chan *pb.StepResult { + ch := make(chan *pb.StepResult, 1) + r.mu.Lock() + r.pending[commandID] = ch + r.mu.Unlock() + return ch +} + +// Cancel removes a pending waiter (call on timeout to avoid leaks). +func (r *stepResultRegistry) Cancel(commandID string) { + r.mu.Lock() + delete(r.pending, commandID) + r.mu.Unlock() +} + +// Deliver routes an incoming StepResult to its waiter, if any. +func (r *stepResultRegistry) Deliver(res *pb.StepResult) { + if res == nil { + return + } + r.mu.Lock() + ch, ok := r.pending[res.CommandId] + if ok { + delete(r.pending, res.CommandId) + } + r.mu.Unlock() + if ok { + ch <- res + } +}