package app import ( "crypto/rand" "encoding/hex" "fmt" "strings" "time" "github.com/correx/tui-go/internal/protocol" ) func nowMillis() int64 { return time.Now().UnixMilli() } func newSessionID() string { b := make([]byte, 16) _, _ = rand.Read(b) return hex.EncodeToString(b) } func containsFold(haystack, needle string) bool { return strings.Contains(strings.ToLower(haystack), strings.ToLower(needle)) } func formatTime(epochMillis int64) string { return time.UnixMilli(epochMillis).UTC().Format("15:04:05") } // inferCategory maps an event type string to a display category, matching the // Kotlin inferCategory heuristic. func inferCategory(t string) string { lt := strings.ToLower(t) switch { case strings.Contains(lt, "approval"): return "Approval" case strings.Contains(lt, "infer"): return "Inference" case strings.Contains(lt, "tool"): return "Tool" case strings.Contains(lt, "context"): return "Context" case strings.Contains(lt, "stage"), strings.Contains(lt, "lifecycle"): return "Lifecycle" case strings.Contains(lt, "session"), strings.Contains(lt, "artifact"): return "Domain" default: return "Lifecycle" } } // narrationLine maps a progress frame to a dim router-feed line, or "" if the // frame is not a narration trigger. func narrationLine(msg protocol.ServerMessage) string { switch msg.Type { case protocol.TypeStageStarted: return "→ " + msg.StageID case protocol.TypeStageCompleted: return "✓ " + msg.StageID case protocol.TypeStageFailed: if msg.Reason != "" { return "⚠ " + msg.StageID + ": " + msg.Reason } return "⚠ " + msg.StageID case protocol.TypeSessionPaused: return "⏸ paused" case protocol.TypeSessionResumed: return "▶ resumed" case protocol.TypeSessionCompleted: return "✓ workflow complete" case protocol.TypeSessionFailed: if msg.Reason != "" { return "✗ workflow failed: " + msg.Reason } return "✗ workflow failed" default: return "" } } // applyServer mutates the model for a single (non-buffered) server message. func (m *Model) applyServer(msg protocol.ServerMessage) { debugLog("SRV type=%s session=%s reason=%s", msg.Type, msg.SessionID, msg.Reason) // Auto-vivify: any event for an unknown session creates it, so session // existence is derived from the event stream, not a dedicated control frame. if msg.IsEventBearing() && msg.SessionID != "" { m.ensureSession(msg.SessionID) } // Inject dim narration lines for progress frames. Runs in addition to the // per-type handling below (which feeds the events panel separately). if line := narrationLine(msg); line != "" { m.appendRouter(msg.SessionID, RouterEntry{Role: "narration", Content: line}) } switch msg.Type { case protocol.TypeSessionAnnounced: m.onSessionAnnounced(msg) case protocol.TypeRouterNarration: entry := RouterEntry{Role: "narration_llm", Content: msg.Content} if msg.LatencyMs != nil && msg.TotalTokens != nil { entry.Metrics = &TurnMetrics{LatencyMs: *msg.LatencyMs, TotalTokens: *msg.TotalTokens} } m.appendRouter(msg.SessionID, entry) case protocol.TypeChatTurn: m.routerConnected = true role := "router" if msg.Role == "USER" { role = "user" } entry := RouterEntry{Role: role, Content: msg.Content} if msg.LatencyMs != nil && msg.TotalTokens != nil { entry.Metrics = &TurnMetrics{LatencyMs: *msg.LatencyMs, TotalTokens: *msg.TotalTokens} } m.appendRouter(msg.SessionID, entry) if s := m.session(msg.SessionID); s != nil { s.LastEventAt = nowMillis() } case protocol.TypeSessionPaused: label := "PAUSED" if msg.Reason == "APPROVAL_PENDING" { label = "PAUSED awaiting approval" } m.touch(msg.SessionID, label) if s := m.session(msg.SessionID); s != nil { s.Active = false } case protocol.TypeSessionResumed: if s := m.session(msg.SessionID); s != nil { s.Status = "ACTIVE" s.Pending = nil s.LastEventAt = nowMillis() } case protocol.TypeSessionCompleted: m.touch(msg.SessionID, "COMPLETED") if s := m.session(msg.SessionID); s != nil { s.Active = false } case protocol.TypeSessionFailed: m.touch(msg.SessionID, "FAILED") if s := m.session(msg.SessionID); s != nil { s.Active = false } case protocol.TypeStageStarted: if s := m.session(msg.SessionID); s != nil { s.CurrentStage = msg.StageID s.Tools = nil s.addEvent(msg.OccurredAt, "StageStarted", msg.StageID) } case protocol.TypeStageCompleted: if s := m.session(msg.SessionID); s != nil { s.CurrentStage = "" s.addEvent(msg.OccurredAt, "StageCompleted", msg.StageID) } case protocol.TypeStageFailed: if s := m.session(msg.SessionID); s != nil { s.CurrentStage = "" s.addEvent(msg.OccurredAt, "StageFailed", msg.StageID) } case protocol.TypeInferenceStarted: if s := m.session(msg.SessionID); s != nil { s.Active = true s.addEvent(nowMillis(), "InferenceStarted", msg.StageID) } case protocol.TypeInferenceDone: if s := m.session(msg.SessionID); s != nil { s.Active = false s.LastOutput = msg.Summary if msg.Response != "" { s.LastResponse = msg.Response } s.addEvent(msg.OccurredAt, "InferenceCompleted", msg.StageID) } case protocol.TypeInferenceTimeout: if s := m.session(msg.SessionID); s != nil { s.Active = false s.addEvent(nowMillis(), "InferenceTimedOut", msg.StageID) } case protocol.TypeInferenceFailed: if s := m.session(msg.SessionID); s != nil { s.Active = false s.addEvent(nowMillis(), "InferenceFailed", msg.StageID+": "+msg.Reason) } case protocol.TypeInferenceRetry: if s := m.session(msg.SessionID); s != nil { s.addEvent(nowMillis(), "RetryAttempted", fmt.Sprintf("%s (%d/%d): %s", msg.StageID, msg.AttemptNumber, msg.MaxAttempts, msg.FailureReason)) } case protocol.TypeToolStarted: if s := m.session(msg.SessionID); s != nil { s.Active = true s.Tools = append(s.Tools, ToolRecord{Name: msg.ToolName, Status: ToolStarted}) if len(s.Tools) > 8 { s.Tools = s.Tools[len(s.Tools)-8:] } s.LastEventAt = nowMillis() } case protocol.TypeToolCompleted: if s := m.session(msg.SessionID); s != nil { s.Active = false s.markTool(msg.ToolName, ToolCompleted) s.LastOutput = msg.ToolName + ": " + msg.Summary s.addEvent(msg.OccurredAt, "ToolCompleted", msg.ToolName) } if msg.Diff != nil && *msg.Diff != "" { m.appendRouter(msg.SessionID, RouterEntry{Role: "tool", Content: *msg.Diff}) } case protocol.TypeToolAssessed: if s := m.session(msg.SessionID); s != nil { if msg.Disposition != "PROCEED" || len(msg.AssessedIssues) > 0 { detail := msg.Disposition if len(msg.AssessedIssues) > 0 { detail += " [" + msg.AssessedIssues[0].Code + "]" } s.LastOutput = msg.ToolName + " assessed: " + detail s.addEvent(msg.OccurredAt, "ToolAssessed", detail) } } case protocol.TypeToolFailed: if s := m.session(msg.SessionID); s != nil { s.markTool(msg.ToolName, ToolFailed) s.LastEventAt = msg.OccurredAt } case protocol.TypeToolRejected: if s := m.session(msg.SessionID); s != nil { s.markTool(msg.ToolName, ToolRejected) s.LastEventAt = nowMillis() } case protocol.TypeArtifactCreated: if s := m.session(msg.SessionID); s != nil { s.addEvent(nowMillis(), "ArtifactCreated", msg.ArtifactID) } case protocol.TypeWorkspaceBound: if s := m.session(msg.SessionID); s != nil { s.WorkspaceRoot = msg.WorkspaceRoot } case protocol.TypeApprovalRequired: m.onApprovalRequired(msg) case protocol.TypeApprovalResolved: if s := m.session(msg.SessionID); s != nil { s.Pending = nil detail := msg.Outcome if msg.Reason != "" { detail += " — " + msg.Reason } s.addEvent(msg.OccurredAt, "ApprovalResolved", detail) s.LastEventAt = nowMillis() } case protocol.TypeSessionSnapshot: m.onSnapshot(msg) case protocol.TypeStageManifest: if s := m.session(msg.SessionID); s != nil { byStage := map[string][]ManifestTool{} for _, st := range msg.Stages { tools := make([]ManifestTool, 0, len(st.Tools)) for _, td := range st.Tools { tools = append(tools, ManifestTool{Name: td.Name, Tier: td.Tier}) } byStage[st.StageID] = tools } s.ToolsByStage = byStage } case protocol.TypeProviderStatus: m.currentModel = msg.ProviderID if containsFold(msg.ProviderID, "llama") || containsFold(msg.ProviderID, "local") { m.providerType = "LOCAL" } else { m.providerType = "REMOTE" } case protocol.TypeWorkflowList: m.workflows = m.workflows[:0] for _, w := range msg.Workflows { m.workflows = append(m.workflows, Workflow{ID: w.WorkflowID, Description: w.Description}) } case protocol.TypeModelChanged: if msg.Loaded { m.currentModel = msg.ModelID m.providerType = "LOCAL" } case protocol.TypeModelList: m.availableModels = append(m.availableModels[:0], msg.Models...) if msg.Current != "" { m.currentModel = msg.Current m.providerType = "LOCAL" } case protocol.TypeResourceStatus: m.gpuUsedMB = msg.GpuMemoryUsedMb m.gpuTotalMB = msg.GpuMemoryTotalMb m.gpuUtil = msg.GpuUtilizationPct m.ramMB = msg.ProcessRssMb m.sysRamUsedMB = msg.SystemRamUsedMb m.sysRamTotalMB = msg.SystemRamTotalMb case protocol.TypeArtifactList: m.artifacts = msg.Artifacts m.artifactsFor = msg.SessionID m.artifactsLoading = false m.artifactScroll = 0 if m.artifactsIndex >= len(m.artifacts) { m.artifactsIndex = 0 } case protocol.TypeConfigSnapshot: m.configFields = msg.ConfigFields m.configRestart = msg.ConfigRestartRequired m.configLoading = false if msg.ConfigError != nil { m.configError = *msg.ConfigError } else { // A clean snapshot means the staged edits were accepted (or this is a fresh fetch). m.configError = "" m.configStaged = map[string]string{} } if m.configIndex >= len(m.configFields) { m.configIndex = 0 } } // Background-update badge for non-selected sessions. if sid := sessionIDOf(msg); sid != "" && sid != m.selectedID { m.bgUpdates++ } } func sessionIDOf(msg protocol.ServerMessage) string { switch msg.Type { case protocol.TypeStageManifest, protocol.TypeSnapshotComplete, protocol.TypeProtocolError, protocol.TypeProviderStatus, protocol.TypeWorkflowList, protocol.TypeRouterResponse, protocol.TypeModelChanged, protocol.TypeModelList, protocol.TypeResourceStatus, protocol.TypeArtifactList, protocol.TypeConfigSnapshot: return "" default: return msg.SessionID } } // onSessionAnnounced fills in a session's workflow identity (the announce is the // only event carrying workflowId) and applies auto-focus. The session entry itself // was already created by the auto-vivify path in applyServer. func (m *Model) onSessionAnnounced(msg protocol.ServerMessage) { s := m.ensureSession(msg.SessionID) s.Status = "ACTIVE" if msg.WorkflowID != "" { s.WorkflowID = msg.WorkflowID s.Name = msg.WorkflowID } s.LastEventAt = nowMillis() if m.pendingWorkflowFocus { m.selectedID = msg.SessionID m.sessionEntered = true m.pendingWorkflowFocus = false } else if m.selectedID == "" { m.selectedID = msg.SessionID } } func (m *Model) onApprovalRequired(msg protocol.ServerMessage) { risk := "unknown" var rationale []string if msg.RiskSummary != nil { risk = msg.RiskSummary.Level rationale = msg.RiskSummary.Rationale } info := &Approval{ RequestID: msg.RequestID, SessionID: msg.SessionID, Tier: msg.Tier, Risk: risk, ToolName: deref(&msg.ToolName), Preview: derefp(msg.Preview), Rationale: rationale, } if s := m.session(msg.SessionID); s != nil { s.Pending = info } } func (m *Model) onSnapshot(msg protocol.ServerMessage) { var pending *Approval if len(msg.PendingAppr) > 0 { a := msg.PendingAppr[0] pending = &Approval{ RequestID: a.RequestID, SessionID: msg.SessionID, Tier: a.Tier, Risk: "unknown", ToolName: derefp(a.ToolName), Preview: derefp(a.Preview), } } status := "running" if msg.State != nil { status = msg.State.Status } if pending != nil { status = "PAUSED awaiting approval" } sess := Session{ ID: msg.SessionID, Status: status, WorkflowID: msg.WorkflowID, Name: msg.WorkflowID, LastEventAt: nowMillis(), Pending: pending, } if msg.State != nil && msg.State.CurrentStageID != nil { sess.CurrentStage = *msg.State.CurrentStageID } for _, e := range msg.RecentEvents { sess.Events = append(sess.Events, EventEntry{formatTime(e.Timestamp), e.Type, e.Detail}) } for _, t := range msg.Tools { sess.Tools = append(sess.Tools, ToolRecord{Name: t.Name, Tier: t.Tier, Status: toolStatusOf(t.Status)}) } // Replace existing session with same id, else append. replaced := false for i := range m.sessions { if m.sessions[i].ID == msg.SessionID { m.sessions[i] = sess replaced = true break } } if !replaced { m.sessions = append(m.sessions, sess) } if m.selectedID == "" { m.selectedID = msg.SessionID } } func (m *Model) touch(id, status string) { if s := m.session(id); s != nil { s.Status = status s.LastEventAt = nowMillis() } } // --- Session helpers --- // maxSessionEvents bounds retained per-session history (effectively "all" for a // session, while keeping memory bounded on very long runs). The EVENTS panel // shows the tail that fits; the `e` inspector scrolls the full retained list. const maxSessionEvents = 1000 func (s *Session) addEvent(epochMillis int64, typ, detail string) { s.Events = append(s.Events, EventEntry{formatTime(epochMillis), typ, detail}) if len(s.Events) > maxSessionEvents { s.Events = s.Events[len(s.Events)-maxSessionEvents:] } s.LastEventAt = epochMillis } func (s *Session) markTool(name string, status ToolStatus) { for i := range s.Tools { if s.Tools[i].Name == name && s.Tools[i].Status == ToolStarted { s.Tools[i].Status = status } } } func toolStatusOf(s string) ToolStatus { switch s { case "COMPLETED": return ToolCompleted case "FAILED": return ToolFailed case "REJECTED": return ToolRejected default: return ToolStarted } } func deref(s *string) string { if s == nil { return "" } return *s } func derefp(s *string) string { return deref(s) }