package opsconsole import ( "context" "fmt" "io" "strings" "sync" "go.uber.org/zap" edgenode "iop/apps/edge/internal/node" eventpkg "iop/packages/go/events" iop "iop/proto/gen/iop" ) // EventRouter renders run and node lifecycle events to the ops console. // Foreground runs tracked via TrackForeground are suppressed from the // asynchronous stream so the SendRun loop owns their output. type EventRouter struct { mu sync.Mutex out io.Writer registry *edgenode.Registry logger *zap.Logger foreground map[string]struct{} streams map[string]*ResponseStream nodeLabels map[string]string aliasLabels map[string]string } func NewEventRouter(out io.Writer, registry *edgenode.Registry, logger *zap.Logger) *EventRouter { return &EventRouter{ out: out, registry: registry, logger: logger, foreground: make(map[string]struct{}), streams: make(map[string]*ResponseStream), nodeLabels: make(map[string]string), aliasLabels: make(map[string]string), } } func (r *EventRouter) TrackForeground(runID string) func() { r.mu.Lock() r.foreground[runID] = struct{}{} r.mu.Unlock() return func() { r.mu.Lock() delete(r.foreground, runID) r.mu.Unlock() } } func (r *EventRouter) Handle(event *iop.RunEvent) { if event == nil || !event.GetBackground() { return } runID := event.GetRunId() r.mu.Lock() if _, ok := r.foreground[runID]; ok { r.mu.Unlock() return } r.printAsyncLocked(event) r.mu.Unlock() } func (r *EventRouter) HandleNodeEvent(event *iop.EdgeNodeEvent) { r.mu.Lock() r.printNodeEventLocked(event) r.mu.Unlock() } func (r *EventRouter) Drain(ctx context.Context, runEvents <-chan *iop.RunEvent, nodeEvents <-chan *iop.EdgeNodeEvent) { for { select { case <-ctx.Done(): return case event, ok := <-runEvents: if !ok { runEvents = nil } else { r.Handle(event) } case event, ok := <-nodeEvents: if !ok { nodeEvents = nil } else { r.HandleNodeEvent(event) } } if runEvents == nil && nodeEvents == nil { return } } } func (r *EventRouter) NodeLabel(event *iop.RunEvent) string { r.mu.Lock() defer r.mu.Unlock() return r.nodeLabelLocked(event) } func (r *EventRouter) nodeLabelLocked(event *iop.RunEvent) string { nodeID := event.GetNodeId() alias := event.GetNodeAlias() return r.resolveNodeLabelLocked(nodeID, alias) } func (r *EventRouter) edgeNodeLabelLocked(event *iop.EdgeNodeEvent) string { nodeID := event.GetNodeId() alias := event.GetAlias() return r.resolveNodeLabelLocked(nodeID, alias) } func (r *EventRouter) resolveNodeLabelLocked(nodeID, alias string) string { if r.registry != nil { if entry, ok := r.registry.Get(nodeID); ok { return r.cacheEntryLabelLocked(entry) } if alias != "" { if entry, err := r.registry.Resolve(alias); err == nil { return r.cacheEntryLabelLocked(entry) } } } if nodeID != "" { if label := r.nodeLabels[nodeID]; label != "" { if alias != "" { r.aliasLabels[alias] = label } return label } } if alias != "" { if label := r.aliasLabels[alias]; label != "" { if nodeID != "" { r.nodeLabels[nodeID] = label } return label } return alias } if nodeID != "" { return nodeID } return "unknown" } func (r *EventRouter) cacheEntryLabelLocked(entry *edgenode.NodeEntry) string { label := entry.DisplayLabel() if entry.NodeID != "" { r.nodeLabels[entry.NodeID] = label } if entry.Alias != "" { r.aliasLabels[entry.Alias] = label } return label } func (r *EventRouter) printNodeEventLocked(event *iop.EdgeNodeEvent) { if r.out == nil { return } label := r.edgeNodeLabelLocked(event) detail := transportCloseDetail(event.GetMetadata()) switch event.GetType() { case eventpkg.TypeNodeConnected: fmt.Fprintf(r.out, "%s connected reason=%q%s\n", nodeEventPrefix(label), event.GetReason(), detail) case eventpkg.TypeNodeDisconnected: fmt.Fprintf(r.out, "%s disconnected reason=%q%s\n", nodeEventPrefix(label), event.GetReason(), detail) default: fmt.Fprintf(r.out, "%s %s reason=%q%s\n", nodeEventPrefix(label), event.GetType(), event.GetReason(), detail) } } func transportCloseDetail(metadata map[string]string) string { if metadata == nil { return "" } detail := "" if reason := metadata[eventpkg.MetadataTransportCloseReason]; reason != "" { detail += fmt.Sprintf(" transport_close_reason=%q", reason) } if err := metadata[eventpkg.MetadataTransportCloseError]; err != "" { detail += fmt.Sprintf(" transport_close_error=%q", err) } return detail } func (r *EventRouter) printAsyncLocked(event *iop.RunEvent) { if r.out == nil { return } runID := event.GetRunId() label := r.nodeLabelLocked(event) switch event.GetType() { case "start": fmt.Fprintf(r.out, "%s start run_id=%s session=%s background=%v\n", nodeEventPrefix(label), runID, event.GetSessionId(), event.GetBackground()) case "reasoning_delta": writeReasoningDelta(r.out, label, event.GetDelta()) case "delta": r.responseStream(runID, label).Write(event.GetDelta()) case "complete": r.finishStream(runID, label) fmt.Fprintf(r.out, "%s complete run_id=%s detail=%q\n", nodeEventPrefix(label), runID, event.GetMessage()) case "cancelled": r.finishStreamIfStarted(runID) fmt.Fprintf(r.out, "%s cancelled run_id=%s\n", nodeEventPrefix(label), runID) case "error": r.finishStreamIfStarted(runID) fmt.Fprintf(r.out, "%s error run_id=%s detail=%q\n", nodeEventPrefix(label), runID, event.GetError()) default: fmt.Fprintf(r.out, "%s %s run_id=%s detail=%q\n", nodeEventPrefix(label), event.GetType(), runID, event.GetMessage()) } } func (r *EventRouter) responseStream(runID, label string) *ResponseStream { if s, ok := r.streams[runID]; ok { return s } s := NewResponseStream(r.out, nodeMessagePrefix(label)) r.streams[runID] = s return s } func (r *EventRouter) finishStream(runID, label string) { if s, ok := r.streams[runID]; ok { s.Finish() delete(r.streams, runID) return } NewResponseStream(r.out, nodeMessagePrefix(label)).Finish() } func (r *EventRouter) finishStreamIfStarted(runID string) { if s, ok := r.streams[runID]; ok { s.FinishIfStarted() delete(r.streams, runID) } } type ResponseStream struct { out io.Writer prefix string started bool endedWithNewline bool } func NewResponseStream(out io.Writer, prefix string) *ResponseStream { return &ResponseStream{out: out, prefix: prefix, endedWithNewline: true} } func nodeEventPrefix(label string) string { return fmt.Sprintf("[%s-evt]", label) } func nodeMessagePrefix(label string) string { return fmt.Sprintf("[%s-msg] ", label) } func (s *ResponseStream) Write(delta string) { if s.out == nil || delta == "" { return } for delta != "" { if !s.started || s.endedWithNewline { fmt.Fprint(s.out, s.prefix) s.started = true s.endedWithNewline = false } idx := strings.IndexByte(delta, '\n') if idx < 0 { fmt.Fprint(s.out, delta) return } fmt.Fprint(s.out, delta[:idx+1]) s.endedWithNewline = true delta = delta[idx+1:] } } func (s *ResponseStream) Finish() { if s.out == nil { return } if !s.started { fmt.Fprintf(s.out, "%s\n", s.prefix) return } if !s.endedWithNewline { fmt.Fprintln(s.out) } } func (s *ResponseStream) FinishIfStarted() { if s.started { s.Finish() } } func nodeReasoningPrefix(label string) string { return fmt.Sprintf("[%s-reasoning] ", label) } func writeReasoningDelta(out io.Writer, label, delta string) { if delta == "" { return } stream := NewResponseStream(out, nodeReasoningPrefix(label)) stream.Write(delta) stream.FinishIfStarted() }