From 0fe111747e6e43dfc23006ee7a70462c5a6002be Mon Sep 17 00:00:00 2001 From: lyon Date: Sun, 21 Jun 2026 21:25:14 +0800 Subject: [PATCH] move workbench session list to go runtime --- internal/cloud/server-workbench-http.ts | 29 + internal/cloud/workbench-runtime-client.ts | 1 + internal/workbenchruntime/service.go | 833 +++++++++++++++++---- 3 files changed, 731 insertions(+), 132 deletions(-) diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index 8d7d1c45..47eace1f 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -14,6 +14,7 @@ import { sendJson } from "./server-http-utils.ts"; import { createWorkbenchReadModel } from "./workbench-read-model.ts"; +import { createWorkbenchRuntimeClient } from "./workbench-runtime-client.ts"; import { createWorkbenchTurnProjection, createWorkbenchTurnTimingProjection, durableTraceStatus, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts"; import { emitHttpServerRequestSpan } from "./otel-trace.ts"; @@ -404,6 +405,34 @@ async function authenticateWorkbenchRead(request, response, options) { } async function handleWorkbenchSessionList(response, url, options, actor) { + if (url.searchParams.has("projectId") || url.searchParams.has("workspaceId")) return sendJson(response, 400, workbenchError("workbench_authority_removed", "Workbench session list is keyed by sessionId only.")); + const limit = boundedSessionListLimit(url.searchParams.get("limit")); + const offset = cursorOffset(url.searchParams.get("cursor") ?? url.searchParams.get("after")); + const includeSessionId = safeSessionId(url.searchParams.get("includeSessionId")); + const includeRouteId = includeSessionId ?? safeConversationId(url.searchParams.get("includeSessionId")); + const runtime = options.workbenchRuntime ?? createWorkbenchRuntimeClient({ env: options.env ?? process.env, fetch: options.fetch, logger: options.logger }); + if (typeof runtime?.listWorkbenchSessions !== "function") { + throw runtimeDependencyError("workbench_runtime_unconfigured", "Workbench runtime service is required for session list.", false); + } + const payload = await runtime.listWorkbenchSessions({ + limit, + offset, + cursor: url.searchParams.get("cursor") ?? url.searchParams.get("after") ?? null, + includeSessionId: includeRouteId ?? null, + actor: { id: actor.id, role: actor.role ?? "user", valuesRedacted: true } + }); + return sendJson(response, 200, payload); +} + +function runtimeDependencyError(code, message, retryable = true) { + const error = new Error(message || code); + error.name = "WorkbenchRuntimeDependencyError"; + error.code = code; + error.data = { serviceId: "hwlab-workbench-runtime", retryable, transient: retryable, valuesRedacted: true }; + return error; +} + +async function handleWorkbenchSessionListBunLegacy(response, url, options, actor) { if (url.searchParams.has("projectId") || url.searchParams.has("workspaceId")) return sendJson(response, 400, workbenchError("workbench_authority_removed", "Workbench session list is keyed by sessionId only.")); const perf = options.backendPerformance; const limit = boundedSessionListLimit(url.searchParams.get("limit")); diff --git a/internal/cloud/workbench-runtime-client.ts b/internal/cloud/workbench-runtime-client.ts index 1831e9f4..5f9b4f56 100644 --- a/internal/cloud/workbench-runtime-client.ts +++ b/internal/cloud/workbench-runtime-client.ts @@ -15,6 +15,7 @@ export function createWorkbenchRuntimeClient(options = {}) { serviceId: "hwlab-workbench-runtime", queryWorkbenchFacts: (params = {}) => postRuntime({ fetchFn, baseUrl, timeoutMs, path: "/v1/workbench-runtime/facts", body: params }), queryAgentTraceEvents: (params = {}) => postRuntime({ fetchFn, baseUrl, timeoutMs, path: "/v1/workbench-runtime/trace-events", body: params }), + listWorkbenchSessions: (params = {}) => postRuntime({ fetchFn, baseUrl, timeoutMs, path: "/v1/workbench-runtime/sessions", body: params }), readWorkbenchProjectionOutbox: (params = {}) => postRuntime({ fetchFn, baseUrl, timeoutMs, path: "/v1/workbench-runtime/projection-outbox", body: params }).then((result) => Array.isArray(result?.rows) ? result.rows : []) }; } diff --git a/internal/workbenchruntime/service.go b/internal/workbenchruntime/service.go index b40e4096..fd55a3af 100644 --- a/internal/workbenchruntime/service.go +++ b/internal/workbenchruntime/service.go @@ -35,41 +35,41 @@ type Server struct { } type factQuery struct { - Families []string `json:"families"` - FactFamilies []string `json:"factFamilies"` - Limit int `json:"limit"` - AfterProjectedSeq int64 `json:"afterProjectedSeq"` - Order string `json:"order"` - SessionsOrder string `json:"sessionsOrder"` - MessagesOrder string `json:"messagesOrder"` - PartsOrder string `json:"partsOrder"` - TurnsOrder string `json:"turnsOrder"` - CheckpointsOrder string `json:"checkpointsOrder"` - SessionProjection string `json:"sessionProjection"` - SessionID string `json:"sessionId"` - SessionIDs []string `json:"sessionIds"` - TraceID string `json:"traceId"` - TraceIDs []string `json:"traceIds"` - MessageID string `json:"messageId"` - MessageIDs []string `json:"messageIds"` - TurnID string `json:"turnId"` - TurnIDs []string `json:"turnIds"` - Status string `json:"status"` - Statuses []string `json:"statuses"` - OwnerUserID string `json:"ownerUserId"` - OwnerUserIDs []string `json:"ownerUserIds"` - ConversationID string `json:"conversationId"` - ConversationIDs []string `json:"conversationIds"` - EventType string `json:"eventType"` - EventTypes []string `json:"eventTypes"` - ProjectionStatus string `json:"projectionStatus"` + Families []string `json:"families"` + FactFamilies []string `json:"factFamilies"` + Limit int `json:"limit"` + AfterProjectedSeq int64 `json:"afterProjectedSeq"` + Order string `json:"order"` + SessionsOrder string `json:"sessionsOrder"` + MessagesOrder string `json:"messagesOrder"` + PartsOrder string `json:"partsOrder"` + TurnsOrder string `json:"turnsOrder"` + CheckpointsOrder string `json:"checkpointsOrder"` + SessionProjection string `json:"sessionProjection"` + SessionID string `json:"sessionId"` + SessionIDs []string `json:"sessionIds"` + TraceID string `json:"traceId"` + TraceIDs []string `json:"traceIds"` + MessageID string `json:"messageId"` + MessageIDs []string `json:"messageIds"` + TurnID string `json:"turnId"` + TurnIDs []string `json:"turnIds"` + Status string `json:"status"` + Statuses []string `json:"statuses"` + OwnerUserID string `json:"ownerUserId"` + OwnerUserIDs []string `json:"ownerUserIds"` + ConversationID string `json:"conversationId"` + ConversationIDs []string `json:"conversationIds"` + EventType string `json:"eventType"` + EventTypes []string `json:"eventTypes"` + ProjectionStatus string `json:"projectionStatus"` ProjectionStatuses []string `json:"projectionStatuses"` - ProjectionHealth string `json:"projectionHealth"` - ProjectionHealths []string `json:"projectionHealths"` - RunID string `json:"runId"` - RunIDs []string `json:"runIds"` - CommandID string `json:"commandId"` - CommandIDs []string `json:"commandIds"` + ProjectionHealth string `json:"projectionHealth"` + ProjectionHealths []string `json:"projectionHealths"` + RunID string `json:"runId"` + RunIDs []string `json:"runIds"` + CommandID string `json:"commandId"` + CommandIDs []string `json:"commandIds"` } type traceEventsQuery struct { @@ -87,6 +87,19 @@ type outboxQuery struct { SessionID string `json:"sessionId"` } +type sessionListQuery struct { + Limit int `json:"limit"` + Offset int `json:"offset"` + Cursor string `json:"cursor"` + IncludeSessionID string `json:"includeSessionId"` + Actor sessionActor `json:"actor"` +} + +type sessionActor struct { + ID string `json:"id"` + Role string `json:"role"` +} + type fieldColumn struct { field string values []string @@ -140,6 +153,7 @@ func (s *Server) routes() { s.mux.HandleFunc("/health/ready", s.handleReady) s.mux.HandleFunc("/v1/workbench-runtime/facts", s.handleFacts) s.mux.HandleFunc("/v1/workbench-runtime/trace-events", s.handleTraceEvents) + s.mux.HandleFunc("/v1/workbench-runtime/sessions", s.handleSessions) s.mux.HandleFunc("/v1/workbench-runtime/projection-outbox", s.handleProjectionOutbox) } @@ -179,10 +193,10 @@ func (s *Server) handleFacts(w http.ResponseWriter, r *http.Request) { } writeJSON(w, http.StatusOK, map[string]any{ "contractVersion": "workbench-runtime-facts-v1", - "facts": facts, - "count": count, - "persistence": s.persistenceSummary(), - "valuesRedacted": true, + "facts": facts, + "count": count, + "persistence": s.persistenceSummary(), + "valuesRedacted": true, }) } @@ -204,10 +218,10 @@ func (s *Server) handleTraceEvents(w http.ResponseWriter, r *http.Request) { } writeJSON(w, http.StatusOK, map[string]any{ "contractVersion": "workbench-runtime-trace-events-v1", - "events": events, - "count": len(events), - "persistence": s.persistenceSummary(), - "valuesRedacted": true, + "events": events, + "count": len(events), + "persistence": s.persistenceSummary(), + "valuesRedacted": true, }) } @@ -229,13 +243,119 @@ func (s *Server) handleProjectionOutbox(w http.ResponseWriter, r *http.Request) } writeJSON(w, http.StatusOK, map[string]any{ "contractVersion": "workbench-runtime-projection-outbox-v1", - "rows": rows, - "count": len(rows), - "persistence": s.persistenceSummary(), - "valuesRedacted": true, + "rows": rows, + "count": len(rows), + "persistence": s.persistenceSummary(), + "valuesRedacted": true, }) } +func (s *Server) handleSessions(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + methodNotAllowed(w, "POST") + return + } + var query sessionListQuery + if !decodeJSON(w, r, &query) { + return + } + ctx, cancel := context.WithTimeout(r.Context(), s.config.QueryTimeout) + defer cancel() + payload, err := s.listSessions(ctx, query) + if err != nil { + writeQueryError(w, err) + return + } + writeJSON(w, http.StatusOK, payload) +} + +func (s *Server) listSessions(ctx context.Context, query sessionListQuery) (map[string]any, error) { + limit := boundedLimit(query.Limit, 20, 100) + offset := query.Offset + if offset <= 0 { + offset = cursorOffset(query.Cursor) + } + if offset < 0 { + offset = 0 + } + ownerUserID := "" + if strings.ToLower(strings.TrimSpace(query.Actor.Role)) != "admin" { + ownerUserID = strings.TrimSpace(query.Actor.ID) + } + pageFacts, err := s.queryFacts(ctx, factQuery{OwnerUserID: ownerUserID, Limit: limit + offset + 1, Families: []string{"sessions"}, SessionsOrder: "updated_desc", SessionProjection: "summary"}) + if err != nil { + return nil, err + } + natural := visibleSessions(pageFacts["sessions"], query.Actor) + if offset > 0 { + if offset >= len(natural) { + natural = []map[string]any{} + } else { + natural = natural[offset:] + } + } + if len(natural) > limit+1 { + natural = natural[:limit+1] + } + includeID := strings.TrimSpace(query.IncludeSessionID) + if includeID != "" && !sessionListContainsRoute(natural, includeID) { + includeFacts, err := s.queryIncludedSession(ctx, includeID) + if err != nil { + return nil, err + } + included := visibleSessions(includeFacts["sessions"], query.Actor) + if len(included) > 0 { + pageFacts = mergeFacts(pageFacts, includeFacts) + natural = append([]map[string]any{included[0]}, natural...) + } + } + hasMore := len(natural) > limit + pageSessions := natural + if len(pageSessions) > limit { + pageSessions = pageSessions[:limit] + } + summaryFacts, err := s.querySummaryFacts(ctx, pageSessions) + if err != nil { + return nil, err + } + facts := mergeFacts(pageFacts, summaryFacts) + summaries := []any{} + for _, session := range pageSessions { + if summary := sessionSummary(session, facts); summary != nil { + summaries = append(summaries, summary) + } + } + return map[string]any{"ok": true, "status": "succeeded", "contractVersion": "workbench-sessions-v1", "sessions": summaries, "count": len(summaries), "cursor": cursorOrNil(offset), "hasMore": hasMore, "nextCursor": nextCursor(offset, limit, hasMore), "persistence": s.persistenceSummary(), "servedBy": serviceID, "valuesRedacted": true, "secretMaterialStored": false}, nil +} + +func (s *Server) queryIncludedSession(ctx context.Context, routeID string) (map[string][]any, error) { + if strings.HasPrefix(routeID, "ses_") { + return s.queryFacts(ctx, factQuery{SessionID: routeID, Limit: 1, Families: []string{"sessions"}, SessionProjection: "summary"}) + } + return s.queryFacts(ctx, factQuery{ConversationID: routeID, Limit: 1, Families: []string{"sessions"}, SessionProjection: "summary"}) +} + +func (s *Server) querySummaryFacts(ctx context.Context, sessions []map[string]any) (map[string][]any, error) { + sessionIDs := uniqueStrings(mapSessions(sessions, sessionID)) + traceIDs := uniqueStrings(mapSessions(sessions, lastTraceID)) + facts := emptyFacts() + if len(sessionIDs) > 0 { + bySession, err := s.queryFacts(ctx, factQuery{SessionIDs: sessionIDs, Families: []string{"messages", "turns"}, Limit: 1000, MessagesOrder: "updated_asc", TurnsOrder: "updated_desc"}) + if err != nil { + return nil, err + } + facts = mergeFacts(facts, bySession) + } + if len(traceIDs) > 0 { + byTrace, err := s.queryFacts(ctx, factQuery{TraceIDs: traceIDs, Families: []string{"checkpoints"}, Limit: 1000, CheckpointsOrder: "updated_desc"}) + if err != nil { + return nil, err + } + facts = mergeFacts(facts, byTrace) + } + return facts, nil +} + func (s *Server) queryFacts(ctx context.Context, query factQuery) (map[string][]any, error) { families := factFamilySet(query) facts := emptyFacts() @@ -248,7 +368,9 @@ func (s *Server) queryFacts(ctx context.Context, query factQuery) (map[string][] {"conversationId", query.ConversationIDs, "conversation_id"}, {"status", query.Statuses, "status"}, }) - if err != nil { return nil, err } + if err != nil { + return nil, err + } } if families["messages"] { facts["messages"], err = s.queryFactRows(ctx, "workbench_messages", "message_json", query, []fieldColumn{ @@ -258,7 +380,9 @@ func (s *Server) queryFacts(ctx context.Context, query factQuery) (map[string][] {"turnId", query.TurnIDs, "turn_id"}, {"status", query.Statuses, "status"}, }) - if err != nil { return nil, err } + if err != nil { + return nil, err + } } if families["parts"] { facts["parts"], err = s.queryFactRows(ctx, "workbench_parts", "part_json", query, []fieldColumn{ @@ -268,7 +392,9 @@ func (s *Server) queryFacts(ctx context.Context, query factQuery) (map[string][] {"turnId", query.TurnIDs, "turn_id"}, {"status", query.Statuses, "status"}, }) - if err != nil { return nil, err } + if err != nil { + return nil, err + } } if families["turns"] { facts["turns"], err = s.queryFactRows(ctx, "workbench_turns", "turn_json", query, []fieldColumn{ @@ -278,7 +404,9 @@ func (s *Server) queryFacts(ctx context.Context, query factQuery) (map[string][] {"messageId", query.MessageIDs, "message_id"}, {"status", query.Statuses, "status"}, }) - if err != nil { return nil, err } + if err != nil { + return nil, err + } } if families["traceEvents"] { facts["traceEvents"], err = s.queryFactRows(ctx, "workbench_trace_events", "event_json", query, []fieldColumn{ @@ -288,7 +416,9 @@ func (s *Server) queryFacts(ctx context.Context, query factQuery) (map[string][] {"messageId", query.MessageIDs, "message_id"}, {"eventType", query.EventTypes, "event_type"}, }) - if err != nil { return nil, err } + if err != nil { + return nil, err + } } if families["checkpoints"] { facts["checkpoints"], err = s.queryFactRows(ctx, "workbench_projection_checkpoints", "checkpoint_json", query, []fieldColumn{ @@ -300,7 +430,9 @@ func (s *Server) queryFacts(ctx context.Context, query factQuery) (map[string][] {"projectionStatus", query.ProjectionStatuses, "projection_status"}, {"projectionHealth", query.ProjectionHealths, "projection_health"}, }) - if err != nil { return nil, err } + if err != nil { + return nil, err + } } return facts, nil } @@ -379,22 +511,22 @@ func (s *Server) querySessionSummaries(ctx context.Context, sqlText string, args continue } item := map[string]any{ - "id": sessionID.String, - "sessionId": sessionID.String, - "ownerUserId": nullableString(ownerUserID), - "projectId": nullableString(projectID), + "id": sessionID.String, + "sessionId": sessionID.String, + "ownerUserId": nullableString(ownerUserID), + "projectId": nullableString(projectID), "conversationId": nullableString(conversationID), - "threadId": nullableString(threadID), - "status": first(nullableStringValue(status), "unknown"), - "lastTraceId": nullableString(lastTraceID), - "projectedSeq": nullableInt(projectedSeq), - "sourceSeq": nullableInt(sourceSeq), - "sourceEventId": nullableString(sourceEventID), - "terminal": nullableBool(terminal), - "sealed": nullableBool(sealed), - "createdAt": nullableTime(createdAt), - "updatedAt": first(nullableTime(updatedAt), nullableTime(createdAt)), - "valuesPrinted": false, + "threadId": nullableString(threadID), + "status": first(nullableStringValue(status), "unknown"), + "lastTraceId": nullableString(lastTraceID), + "projectedSeq": nullableInt(projectedSeq), + "sourceSeq": nullableInt(sourceSeq), + "sourceEventId": nullableString(sourceEventID), + "terminal": nullableBool(terminal), + "sealed": nullableBool(sealed), + "createdAt": nullableTime(createdAt), + "updatedAt": first(nullableTime(updatedAt), nullableTime(createdAt)), + "valuesPrinted": false, "valuesRedacted": true, } if providerProfile.Valid && strings.TrimSpace(providerProfile.String) != "" { @@ -449,19 +581,19 @@ func (s *Server) readProjectionOutbox(ctx context.Context, query outboxQuery) ([ _ = json.Unmarshal(raw, &payload) } result = append(result, map[string]any{ - "outboxSeq": nullableInt(outboxSeq), - "traceId": nullableString(traceID), - "sessionId": nullableString(sessionID), - "turnId": nullableString(turnID), - "messageId": nullableString(messageID), - "projectedSeq": nullableInt(projectedSeq), - "sourceSeq": nullableInt(sourceSeq), + "outboxSeq": nullableInt(outboxSeq), + "traceId": nullableString(traceID), + "sessionId": nullableString(sessionID), + "turnId": nullableString(turnID), + "messageId": nullableString(messageID), + "projectedSeq": nullableInt(projectedSeq), + "sourceSeq": nullableInt(sourceSeq), "sourceEventId": nullableString(sourceEventID), - "commitType": nullableString(commitType), - "terminal": nullableBool(terminal), - "sealed": nullableBool(sealed), - "payload": payload, - "createdAt": nullableTime(createdAt), + "commitType": nullableString(commitType), + "terminal": nullableBool(terminal), + "sealed": nullableBool(sealed), + "payload": payload, + "createdAt": nullableTime(createdAt), "valuesPrinted": false, }) } @@ -508,19 +640,32 @@ func addTextListClause(clauses *[]string, args *[]any, column string, values []s func singularField(query factQuery, field string) string { switch field { - case "sessionId": return query.SessionID - case "traceId": return query.TraceID - case "messageId": return query.MessageID - case "turnId": return query.TurnID - case "status": return query.Status - case "ownerUserId": return query.OwnerUserID - case "conversationId": return query.ConversationID - case "eventType": return query.EventType - case "projectionStatus": return query.ProjectionStatus - case "projectionHealth": return query.ProjectionHealth - case "runId": return query.RunID - case "commandId": return query.CommandID - default: return "" + case "sessionId": + return query.SessionID + case "traceId": + return query.TraceID + case "messageId": + return query.MessageID + case "turnId": + return query.TurnID + case "status": + return query.Status + case "ownerUserId": + return query.OwnerUserID + case "conversationId": + return query.ConversationID + case "eventType": + return query.EventType + case "projectionStatus": + return query.ProjectionStatus + case "projectionHealth": + return query.ProjectionHealth + case "runId": + return query.RunID + case "commandId": + return query.CommandID + default: + return "" } } @@ -533,10 +678,10 @@ func whereSQL(clauses []string) string { func emptyFacts() map[string][]any { return map[string][]any{ - "sessions": []any{}, - "messages": []any{}, - "parts": []any{}, - "turns": []any{}, + "sessions": []any{}, + "messages": []any{}, + "parts": []any{}, + "turns": []any{}, "traceEvents": []any{}, "checkpoints": []any{}, } @@ -550,43 +695,75 @@ func factFamilySet(query factQuery) map[string]bool { } result := map[string]bool{} if len(selected) == 0 { - for _, family := range all { result[family] = true } + for _, family := range all { + result[family] = true + } return result } valid := map[string]bool{} - for _, family := range all { valid[family] = true } + for _, family := range all { + valid[family] = true + } for _, family := range selected { family = strings.TrimSpace(family) - if valid[family] { result[family] = true } + if valid[family] { + result[family] = true + } } if len(result) == 0 { - for _, family := range all { result[family] = true } + for _, family := range all { + result[family] = true + } } return result } func factFamilyForTable(table string) string { switch table { - case "workbench_sessions": return "sessions" - case "workbench_messages": return "messages" - case "workbench_parts": return "parts" - case "workbench_turns": return "turns" - case "workbench_trace_events": return "traceEvents" - case "workbench_projection_checkpoints": return "checkpoints" - default: return "unknown" + case "workbench_sessions": + return "sessions" + case "workbench_messages": + return "messages" + case "workbench_parts": + return "parts" + case "workbench_turns": + return "turns" + case "workbench_trace_events": + return "traceEvents" + case "workbench_projection_checkpoints": + return "checkpoints" + default: + return "unknown" } } func factOrder(query factQuery, family string) string { value := query.Order switch family { - case "sessions": if query.SessionsOrder != "" { value = query.SessionsOrder } - case "messages": if query.MessagesOrder != "" { value = query.MessagesOrder } - case "parts": if query.PartsOrder != "" { value = query.PartsOrder } - case "turns": if query.TurnsOrder != "" { value = query.TurnsOrder } - case "checkpoints": if query.CheckpointsOrder != "" { value = query.CheckpointsOrder } + case "sessions": + if query.SessionsOrder != "" { + value = query.SessionsOrder + } + case "messages": + if query.MessagesOrder != "" { + value = query.MessagesOrder + } + case "parts": + if query.PartsOrder != "" { + value = query.PartsOrder + } + case "turns": + if query.TurnsOrder != "" { + value = query.TurnsOrder + } + case "checkpoints": + if query.CheckpointsOrder != "" { + value = query.CheckpointsOrder + } + } + if value == "updated_desc" { + return "updated_desc" } - if value == "updated_desc" { return "updated_desc" } return "updated_asc" } @@ -620,14 +797,388 @@ func boundedLimit(value int, fallback int, max int) int { return value } +func visibleSessions(rows []any, actor sessionActor) []map[string]any { + result := []map[string]any{} + for _, row := range rows { + item := objectMap(row) + if len(item) == 0 { + continue + } + if strings.ToLower(text(item["status"])) == "archived" { + continue + } + if strings.ToLower(strings.TrimSpace(actor.Role)) != "admin" && text(item["ownerUserId"]) != strings.TrimSpace(actor.ID) { + continue + } + result = append(result, item) + } + return result +} + +func sessionSummary(session map[string]any, facts map[string][]any) map[string]any { + sid := sessionID(session) + if sid == "" { + return nil + } + traceID := first(lastTraceID(session), latestTraceIDForSession(facts, sid)) + var projection map[string]any + if traceID != "" { + projection = projectionForTrace(facts, traceID) + } + turn := turnForTrace(facts, traceID) + checkpoint := checkpointForTrace(facts, traceID) + messages := messagesForSession(facts, sid) + checkpointStatus := terminalStatus(text(checkpoint["status"])) + turnStatus := normalizeStatus(turn["status"]) + status := normalizeStatus(first(checkpointStatus, turnStatus, text(session["status"]))) + timingSource := turn + if checkpointStatus != "" { + timingSource = checkpoint + } + timing := timingProjection(timingSource, status) + var turnSummary any + if len(turn) > 0 { + turnSummary = map[string]any{"turnId": first(text(turn["turnId"]), traceID), "traceId": traceID, "status": status, "running": isRunning(status), "terminal": isTerminal(status), "eventCount": projection["lastProjectedSeq"], "projection": projection, "projectionStatus": projection["projectionStatus"], "projectionHealth": projection["projectionHealth"], "staleMs": projection["staleMs"], "blocker": projection["blocker"], "timing": timing, "startedAt": timing["startedAt"], "lastEventAt": timing["lastEventAt"], "finishedAt": timing["finishedAt"], "durationMs": timing["durationMs"], "updatedAt": updatedAt(turn)} + } + return map[string]any{"sessionId": sid, "threadId": nullableText(session["threadId"]), "agentId": first(text(session["agentId"]), "hwlab-code-agent"), "status": status, "running": isRunning(status), "terminal": isTerminal(status), "lastTraceId": nullableText(traceID), "projection": projection, "projectionStatus": projection["projectionStatus"], "projectionHealth": projection["projectionHealth"], "staleMs": projection["staleMs"], "blocker": projection["blocker"], "providerProfile": nullableText(first(text(session["providerProfile"]), text(objectMap(session["sessionJson"])["providerProfile"]))), "messageCount": len(messages), "firstUserMessagePreview": firstUserPreview(messages), "updatedAt": updatedAt(session), "turnSummary": turnSummary, "valuesRedacted": true} +} + +func messagesForSession(facts map[string][]any, sid string) []map[string]any { + result := []map[string]any{} + for _, row := range facts["messages"] { + item := objectMap(row) + if text(item["sessionId"]) == sid { + result = append(result, item) + } + } + return result +} + +func firstUserPreview(messages []map[string]any) any { + for _, message := range messages { + if text(message["role"]) != "user" { + continue + } + if preview := text(message["textPreview"]); preview != "" { + return preview + } + body := first(text(message["text"]), text(message["content"]), text(message["message"])) + if body == "" { + return nil + } + if len(body) > 240 { + return body[:240] + } + return body + } + return nil +} + +func projectionForTrace(facts map[string][]any, traceID string) map[string]any { + checkpoint := checkpointForTrace(facts, traceID) + if len(checkpoint) == 0 { + return map[string]any{"projectionStatus": "unknown", "projectionHealth": "unknown", "lastProjectedSeq": nil, "sourceRunId": nil, "sourceCommandId": nil, "staleMs": nil, "blocker": nil, "updatedAt": nil, "valuesRedacted": true} + } + status := normalizeProjectionStatus(text(checkpoint["projectionStatus"])) + health := normalizeProjectionHealth(text(checkpoint["projectionHealth"]), status) + timing := timingProjection(checkpoint, normalizeStatus(checkpoint["status"])) + return map[string]any{"projectionStatus": status, "projectionHealth": health, "lastProjectedSeq": seq(checkpoint), "sourceRunId": nullableText(first(text(checkpoint["runId"]), text(checkpoint["sourceRunId"]))), "sourceCommandId": nullableText(first(text(checkpoint["commandId"]), text(checkpoint["sourceCommandId"]))), "staleMs": nil, "blocker": firstAny(objectMap(checkpoint["diagnostic"])["blocker"], checkpoint["blocker"]), "timing": timing, "startedAt": timing["startedAt"], "lastEventAt": timing["lastEventAt"], "finishedAt": timing["finishedAt"], "durationMs": timing["durationMs"], "updatedAt": updatedAt(checkpoint), "valuesRedacted": true} +} + +func timingProjection(record map[string]any, status string) map[string]any { + source := objectMap(record["timing"]) + observedAt := time.Now().UTC().Format(time.RFC3339Nano) + startedAt := timestamp(source["startedAt"]) + lastEventAt := timestamp(source["lastEventAt"]) + terminal := boolValue(record["terminal"]) || isTerminal(status) + finishedAt := "" + if terminal { + finishedAt = timestamp(source["finishedAt"]) + } + return map[string]any{"startedAt": nullableText(startedAt), "lastEventAt": nullableText(lastEventAt), "finishedAt": nullableText(finishedAt), "durationMs": elapsedMs(startedAt, first(finishedAt, observedAt)), "observedAt": nilIf(terminal, observedAt), "lastEventAgeMs": nilIf(terminal, elapsedMs(lastEventAt, observedAt)), "valuesRedacted": true} +} + +func checkpointForTrace(facts map[string][]any, traceID string) map[string]any { + if traceID == "" { + return map[string]any{} + } + for _, row := range facts["checkpoints"] { + item := objectMap(row) + if text(item["traceId"]) == traceID { + return item + } + } + return map[string]any{} +} + +func turnForTrace(facts map[string][]any, traceID string) map[string]any { + if traceID == "" { + return map[string]any{} + } + for _, row := range facts["turns"] { + item := objectMap(row) + if text(item["traceId"]) == traceID { + return item + } + } + return map[string]any{} +} + +func latestTraceIDForSession(facts map[string][]any, sid string) string { + for _, row := range facts["turns"] { + item := objectMap(row) + if text(item["sessionId"]) == sid && text(item["traceId"]) != "" { + return text(item["traceId"]) + } + } + return "" +} + +func mergeFacts(sets ...map[string][]any) map[string][]any { + merged := emptyFacts() + keys := map[string]string{"sessions": "sessionId", "messages": "messageId", "parts": "partId", "turns": "turnId", "traceEvents": "id", "checkpoints": "traceId"} + for _, set := range sets { + for family, idKey := range keys { + seen := map[string]bool{} + for _, row := range merged[family] { + if id := text(objectMap(row)[idKey]); id != "" { + seen[id] = true + } + } + for _, row := range set[family] { + id := text(objectMap(row)[idKey]) + if id != "" && seen[id] { + continue + } + merged[family] = append(merged[family], row) + if id != "" { + seen[id] = true + } + } + } + } + return merged +} + +func sessionListContainsRoute(sessions []map[string]any, routeID string) bool { + for _, session := range sessions { + if sessionID(session) == routeID || text(session["conversationId"]) == routeID { + return true + } + } + return false +} + +func sessionID(session map[string]any) string { + return first(text(session["sessionId"]), text(session["id"])) +} +func lastTraceID(session map[string]any) string { + return first(text(session["lastTraceId"]), text(session["traceId"])) +} +func updatedAt(record map[string]any) any { + return nullableText(first(text(record["updatedAt"]), text(record["occurredAt"]), text(record["createdAt"]))) +} +func mapSessions(sessions []map[string]any, fn func(map[string]any) string) []string { + out := []string{} + for _, session := range sessions { + if value := fn(session); value != "" { + out = append(out, value) + } + } + return out +} +func uniqueStrings(values []string) []string { + seen := map[string]bool{} + out := []string{} + for _, value := range values { + if value != "" && !seen[value] { + seen[value] = true + out = append(out, value) + } + } + return out +} +func cursorOffset(value string) int { + text := strings.TrimSpace(value) + if strings.HasPrefix(text, "idx:") { + text = strings.TrimPrefix(text, "idx:") + } + parsed, err := strconv.Atoi(text) + if err != nil || parsed < 0 { + return 0 + } + return parsed +} +func cursorOrNil(offset int) any { + if offset > 0 { + return fmt.Sprintf("idx:%d", offset) + } + return nil +} +func nextCursor(offset int, limit int, hasMore bool) any { + if hasMore { + return fmt.Sprintf("idx:%d", offset+limit) + } + return nil +} +func objectMap(value any) map[string]any { + if item, ok := value.(map[string]any); ok { + return item + } + return map[string]any{} +} +func text(value any) string { + if value == nil { + return "" + } + if stringValue, ok := value.(string); ok { + return strings.TrimSpace(stringValue) + } + result := strings.TrimSpace(fmt.Sprint(value)) + if result == "" { + return "" + } + return result +} +func nullableText(value any) any { + text := text(value) + if text == "" { + return nil + } + return text +} +func boolValue(value any) bool { b, _ := value.(bool); return b } +func nilIf(condition bool, value any) any { + if condition { + return nil + } + return value +} +func firstAny(values ...any) any { + for _, value := range values { + if value != nil { + return value + } + } + return nil +} + +func seq(record map[string]any) any { + for _, key := range []string{"projectedSeq", "sourceSeq", "seq"} { + switch value := record[key].(type) { + case int64: + if value >= 0 { + return value + } + case int: + if value >= 0 { + return value + } + case float64: + if value >= 0 { + return int64(value) + } + } + } + return nil +} + +func timestamp(value any) string { + parsed, err := time.Parse(time.RFC3339Nano, text(value)) + if err == nil { + return parsed.UTC().Format(time.RFC3339Nano) + } + parsed, err = time.Parse(time.RFC3339, text(value)) + if err == nil { + return parsed.UTC().Format(time.RFC3339Nano) + } + return "" +} + +func elapsedMs(startedAt string, endedAt string) any { + start, err := time.Parse(time.RFC3339Nano, startedAt) + if err != nil { + return nil + } + end, err := time.Parse(time.RFC3339Nano, endedAt) + if err != nil || end.Before(start) { + return nil + } + return end.Sub(start).Milliseconds() +} + +func normalizeStatus(value any) string { + status := strings.ReplaceAll(strings.ToLower(text(value)), "_", "-") + if status == "cancelled" { + return "canceled" + } + if status == "" || status == "" { + return "unknown" + } + return status +} +func terminalStatus(value string) string { + status := normalizeStatus(value) + if isTerminal(status) { + return status + } + return "" +} +func isRunning(status string) bool { + switch normalizeStatus(status) { + case "running", "retrying", "pending", "queued", "accepted", "dispatching", "streaming", "processing", "busy": + return true + default: + return false + } +} +func isTerminal(status string) bool { + switch normalizeStatus(status) { + case "completed", "failed", "blocked", "timeout", "canceled", "stale", "thread-resume-failed", "interrupted", "expired": + return true + default: + return false + } +} +func normalizeProjectionStatus(value string) string { + status := normalizeStatus(value) + if status == "caughtup" { + return "caught-up" + } + switch status { + case "caught-up", "projecting", "blocked", "stalled", "unknown": + return status + default: + return "unknown" + } +} +func normalizeProjectionHealth(value string, projectionStatus string) string { + health := normalizeStatus(value) + if health == "healthy" && projectionStatus == "caught-up" { + return "caught-up" + } + if health == "healthy" && projectionStatus == "projecting" { + return "projecting" + } + switch health { + case "caught-up", "projecting", "degraded", "stalled", "unavailable", "unknown": + return health + default: + if projectionStatus == "caught-up" || projectionStatus == "projecting" { + return projectionStatus + } + return "unknown" + } +} + func configFromEnv() Config { port := first(os.Getenv("HWLAB_WORKBENCH_RUNTIME_PORT"), os.Getenv("PORT"), "6671") return Config{ - Addr: ":" + port, - DatabaseURL: first(os.Getenv("HWLAB_WORKBENCH_RUNTIME_DB_URL"), os.Getenv("HWLAB_CLOUD_DB_URL")), - QueryTimeout: time.Duration(envInt("HWLAB_WORKBENCH_RUNTIME_QUERY_TIMEOUT_MS", 10000)) * time.Millisecond, + Addr: ":" + port, + DatabaseURL: first(os.Getenv("HWLAB_WORKBENCH_RUNTIME_DB_URL"), os.Getenv("HWLAB_CLOUD_DB_URL")), + QueryTimeout: time.Duration(envInt("HWLAB_WORKBENCH_RUNTIME_QUERY_TIMEOUT_MS", 10000)) * time.Millisecond, ReadinessDelay: time.Duration(envInt("HWLAB_WORKBENCH_RUNTIME_READY_TIMEOUT_MS", 2000)) * time.Millisecond, - PoolMax: envInt("HWLAB_WORKBENCH_RUNTIME_DB_POOL_MAX", 8), + PoolMax: envInt("HWLAB_WORKBENCH_RUNTIME_DB_POOL_MAX", 8), } } @@ -650,11 +1201,11 @@ func writeQueryError(w http.ResponseWriter, err error) { writeJSON(w, http.StatusServiceUnavailable, map[string]any{ "ok": false, "error": map[string]any{ - "code": "workbench_runtime_query_failed", - "message": "workbench runtime query failed", - "retryable": true, - "transient": true, - "causeCode": sqlErrorCode(err), + "code": "workbench_runtime_query_failed", + "message": "workbench runtime query failed", + "retryable": true, + "transient": true, + "causeCode": sqlErrorCode(err), "valuesRedacted": true, }, "valuesRedacted": true, @@ -678,31 +1229,43 @@ func writeJSON(w http.ResponseWriter, status int, value any) { func envInt(name string, fallback int) int { value := strings.TrimSpace(os.Getenv(name)) - if value == "" { return fallback } + if value == "" { + return fallback + } parsed, err := strconv.Atoi(value) - if err != nil || parsed <= 0 { return fallback } + if err != nil || parsed <= 0 { + return fallback + } return parsed } func first(values ...string) string { for _, value := range values { - if strings.TrimSpace(value) != "" { return value } + if strings.TrimSpace(value) != "" { + return value + } } return "" } func nullableString(value sql.NullString) any { - if value.Valid { return value.String } + if value.Valid { + return value.String + } return nil } func nullableStringValue(value sql.NullString) string { - if value.Valid { return value.String } + if value.Valid { + return value.String + } return "" } func nullableInt(value sql.NullInt64) any { - if value.Valid { return value.Int64 } + if value.Valid { + return value.Int64 + } return int64(0) } @@ -711,15 +1274,21 @@ func nullableBool(value sql.NullBool) bool { } func nullableTime(value sql.NullTime) string { - if value.Valid { return value.Time.UTC().Format(time.RFC3339Nano) } + if value.Valid { + return value.Time.UTC().Format(time.RFC3339Nano) + } return "" } -type sqlStateError interface { SQLState() string } +type sqlStateError interface{ SQLState() string } func sqlErrorCode(err error) string { - if err == nil { return "" } + if err == nil { + return "" + } var stateErr sqlStateError - if errors.As(err, &stateErr) { return stateErr.SQLState() } + if errors.As(err, &stateErr) { + return stateErr.SQLState() + } return "" }