Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions packages/tui/src/observability/local-observability.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,27 @@ export interface TuiEventStreamObservation {
readonly errorKind?: string;
}

/**
* Run-ownership transitions the TUI cannot reconstruct from Runtime history alone.
* Session and Turn ids are opaque identifiers; no message content is recorded.
*/
export interface TuiRunLifecycleObservation {
readonly kind:
| 'runtime-turn-adopted'
| 'queue-turn-started'
| 'terminal-not-owned'
| 'stale-run-reconciled'
| 'stale-in-process-run';
readonly sessionId: string;
readonly turnId?: string;
readonly projectedTurnId?: string;
readonly inProcessTurnId?: string;
readonly reason?: string;
readonly runtimeState?: string;
readonly queuedCount?: number;
readonly stalledForMs?: number;
}

export interface TuiTerminalObservation {
readonly terminalId: string;
readonly platform: NodeJS.Platform;
Expand Down Expand Up @@ -138,6 +159,7 @@ export interface TuiObservability {
recordStartup(observation: TuiStartupObservation): void;
recordAccess(observation: TuiRuntimeAccessObservation): void;
recordEventStream(observation: TuiEventStreamObservation): void;
recordRunLifecycle?(observation: TuiRunLifecycleObservation): void;
recordTerminal(observation: TuiTerminalObservation): void;
recordProcessStop(observation: TuiProcessStopObservation): void;
recordTheme?(observation: TuiThemeObservation): void;
Expand Down Expand Up @@ -268,6 +290,10 @@ class LocalTuiObservability implements TuiObservability {
this.record('event.stream', observation);
}

recordRunLifecycle(observation: TuiRunLifecycleObservation): void {
this.record('run.lifecycle', observation);
}

recordTerminal(observation: TuiTerminalObservation): void {
this.terminal = Object.freeze({ ...observation });
this.record('terminal.capability', observation);
Expand Down
212 changes: 210 additions & 2 deletions packages/tui/src/tui/controller/runtime/runtime-event-flow.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ import {
noopTuiObservability,
type TuiIncidentSink,
type TuiObservability,
type TuiRunLifecycleObservation,
} from '../../../observability/index.js';
import { formatTuiRuntimeFailure, resolveTuiRuntimeFailure } from './runtime-error-presentation.js';
import {
Expand All @@ -47,6 +48,8 @@ const EVENT_BUS_RECONNECT_MAX_DELAY_MS = 10_000;
const INCIDENT_RECONNECT_WINDOW_MS = 30_000;
const INCIDENT_RECONNECT_FAILURE_THRESHOLD = 3;
const TERMINAL_RECONCILIATION_TIMEOUT_MS = 2_000;
// A projected run must read as settled on two consecutive checks before the TUI clears it.
const STALE_RUN_CHECK_INTERVAL_MS = 5_000;

type AppendLocalCell = (content: string, kind?: 'final-summary' | 'warning' | 'error') => void;

Expand Down Expand Up @@ -86,6 +89,8 @@ export interface TuiRuntimeEventFlowOptions {
readonly observability?: TuiObservability;
readonly incidentReporter?: TuiIncidentSink;
readonly terminalReconciliationTimeoutMs?: number;
/** Interval for the stale-run safety net; `0` disables the timer. */
readonly staleRunCheckIntervalMs?: number;
readonly notify?: (kind: TuiTerminalNotificationKind, key: string) => void;
}

Expand All @@ -105,6 +110,13 @@ export class TuiRuntimeEventFlow {
}
| undefined;
private readonly llmRetryCalls = new Map<string, TuiLlmRetryEvent>();
private staleRunTimer: ReturnType<typeof setInterval> | undefined;
private staleRunCheck: Promise<boolean> | undefined;
private staleRunGeneration = 0;
private staleRunCandidate:
| { readonly sessionId: string; readonly turnId: string; readonly firstSeenAtMs: number }
| undefined;
private reportedStaleInProcessTurnId: string | undefined;

constructor(private readonly options: TuiRuntimeEventFlowOptions) {}

Expand Down Expand Up @@ -139,6 +151,7 @@ export class TuiRuntimeEventFlow {

start(): void {
if (this.task || this.options.isStopped()) return;
this.startStaleRunWatchdog();
this.abortController = new AbortController();
const busController = this.abortController;
const task = this.consume(busController);
Expand All @@ -161,6 +174,10 @@ export class TuiRuntimeEventFlow {

stop(): void {
this.clearLlmRetry();
if (this.staleRunTimer) clearInterval(this.staleRunTimer);
this.staleRunTimer = undefined;
this.staleRunGeneration += 1;
this.staleRunCandidate = undefined;
this.abortController?.abort();
this.liveTurn?.controller.abort();
void this.task?.catch(() => undefined);
Expand All @@ -171,6 +188,8 @@ export class TuiRuntimeEventFlow {

restart(): void {
if (this.options.isStopped()) return;
this.staleRunGeneration += 1;
this.staleRunCandidate = undefined;
const previousTask = this.task;
this.abortController?.abort();
if (this.task === previousTask) this.task = undefined;
Expand All @@ -187,6 +206,172 @@ export class TuiRuntimeEventFlow {
this.clearLlmRetry();
}

/**
* Safety net for a run the TUI still projects after Runtime has settled it.
*
* Runtime is authoritative for Turn ownership. When the TUI keeps a Runtime-owned Turn
* (adopted, recovered, or queue-started) while two consecutive reads report no running
* Turn and no Queue handoff, the missed terminal is reconciled locally so the activity
* line, Enter routing, and Queue admission stop treating the Session as busy. A
* foreground in-process submission owns its own completion and is only reported.
*/
reconcileStaleRun(): Promise<boolean> {
if (!this.staleRunCheck) {
const check = this.checkStaleRun()
.catch((error: unknown) => {
this.staleRunCandidate = undefined;
throw error;
})
.finally(() => {
if (this.staleRunCheck === check) this.staleRunCheck = undefined;
});
this.staleRunCheck = check;
}
return this.staleRunCheck;
}

private startStaleRunWatchdog(): void {
const intervalMs = this.options.staleRunCheckIntervalMs ?? STALE_RUN_CHECK_INTERVAL_MS;
if (this.staleRunTimer || intervalMs <= 0) return;
this.staleRunTimer = setInterval(() => {
void this.reconcileStaleRun().catch(() => undefined);
}, intervalMs);
this.staleRunTimer.unref?.();
}

private projectedRunTurnId(): string | undefined {
return (
this.options.controller.snapshot().activeTurnId ??
this.options.runProjection.snapshot().latestRuntimeTurnId ??
this.liveTurn?.turnId
);
}

private staleRunOwnershipBlocksReconciliation(
sessionId: string,
projectedTurnId: string,
): boolean {
const snapshot = this.options.controller.snapshot();
const runProjection = this.options.runProjection.snapshot();
return (
this.options.isStopped() ||
snapshot.session?.sessionId !== sessionId ||
this.projectedRunTurnId() !== projectedTurnId ||
Boolean(snapshot.retiringTurnId) ||
runProjection.queueHandoffPending ||
Boolean(runProjection.stoppingRuntimeTurnId) ||
this.options.interactionFlow.continuesTurn(sessionId, projectedTurnId)
);
}

private async checkStaleRun(): Promise<boolean> {
const snapshot = this.options.controller.snapshot();
const sessionId = snapshot.session?.sessionId;
const projectedTurnId = this.projectedRunTurnId();
if (
!sessionId ||
!projectedTurnId ||
this.staleRunOwnershipBlocksReconciliation(sessionId, projectedTurnId)
) {
this.staleRunCandidate = undefined;
return false;
}
const generation = this.staleRunGeneration;
const [activeRun, queued] = await Promise.all([
this.options.runtime.getActiveRun(sessionId),
this.options.runtime.listQueuedMessages(sessionId),
]);
if (generation !== this.staleRunGeneration) return false;
if (this.staleRunOwnershipBlocksReconciliation(sessionId, projectedTurnId)) {
this.staleRunCandidate = undefined;
return false;
}
const runtimeBusy =
activeRun.state === 'running' ||
activeRun.state === 'decision-blocked' ||
queued.some((item) => item.status === 'accepted' || item.status === 'running');
if (runtimeBusy) {
this.staleRunCandidate = undefined;
return false;
}
const nowMs = Date.now();
const candidate = this.staleRunCandidate;
if (candidate?.sessionId !== sessionId || candidate.turnId !== projectedTurnId) {
this.staleRunCandidate = { sessionId, turnId: projectedTurnId, firstSeenAtMs: nowMs };
return false;
}
const inProcessTurnId = this.options.controller.snapshot().activeTurnId;
const observation = {
sessionId,
projectedTurnId,
runtimeState: activeRun.state,
queuedCount: queued.filter((item) => item.status === 'queued').length,
stalledForMs: Math.max(0, nowMs - candidate.firstSeenAtMs),
...(activeRun.turnId ? { turnId: activeRun.turnId } : {}),
...(inProcessTurnId ? { inProcessTurnId } : {}),
};
if (inProcessTurnId) {
if (this.reportedStaleInProcessTurnId !== inProcessTurnId) {
this.reportedStaleInProcessTurnId = inProcessTurnId;
this.recordRunLifecycle({ kind: 'stale-in-process-run', ...observation });
this.breadcrumb('cli.run.stale_in_process', {
runtimeState: activeRun.state,
stalledForMs: observation.stalledForMs,
});
}
return false;
}

this.staleRunCandidate = undefined;
const liveTurn = this.liveTurn;
if (liveTurn?.sessionId === sessionId && liveTurn.turnId === projectedTurnId) {
liveTurn.controller.abort();
void liveTurn.task.catch(() => undefined);
this.liveTurn = undefined;
this.options.controller.runtimeTurnSettlement.settleProjection(projectedTurnId, 'succeeded');
}
this.options.runProjection.reconcileRuntimeTurn(undefined);
this.recordRunLifecycle({ kind: 'stale-run-reconciled', ...observation });
this.breadcrumb('cli.run.stale_reconciled', {
runtimeState: activeRun.state,
stalledForMs: observation.stalledForMs,
});
this.options.onChanged();
await this.refreshRuntimeSessionProjection(sessionId, { includeDurableHistory: true }).catch(
() => false,
);
await this.options.activeRunFlow.refresh(true).catch(() => undefined);
this.options.onChanged();
return true;
}

private recordRunLifecycle(observation: TuiRunLifecycleObservation): void {
try {
this.options.observability?.recordRunLifecycle?.(observation);
} catch {
// Diagnostics must not affect Runtime event handling.
}
}

private recordTerminalNotOwned(
event: TuiSessionLifecycleEvent,
sessionId: string,
reason: 'in-process-turn' | 'projection-mismatch',
context: { readonly inProcessTurnId?: string; readonly projectedTurnId?: string },
): void {
// Only a terminal for a Turn other than the one the TUI holds can leave it stale.
const heldTurnId = context.inProcessTurnId ?? context.projectedTurnId;
if (!event.turnId || event.turnId === heldTurnId) return;
this.recordRunLifecycle({
kind: 'terminal-not-owned',
sessionId,
turnId: event.turnId,
reason,
...(context.inProcessTurnId ? { inProcessTurnId: context.inProcessTurnId } : {}),
...(context.projectedTurnId ? { projectedTurnId: context.projectedTurnId } : {}),
});
}

applyActiveSessionProjection(): void {
this.clearLlmRetry();
const view = selectActiveSessionView(this.options.stateStore.snapshot());
Expand Down Expand Up @@ -523,6 +708,14 @@ export class TuiRuntimeEventFlow {
) {
return;
}
this.recordRunLifecycle({
kind: 'runtime-turn-adopted',
sessionId,
turnId,
...(this.options.runProjection.snapshot().latestRuntimeTurnId
? { projectedTurnId: this.options.runProjection.snapshot().latestRuntimeTurnId }
: {}),
});
this.startLiveTurn(
sessionId,
turnId,
Expand Down Expand Up @@ -751,6 +944,14 @@ export class TuiRuntimeEventFlow {

private handleQueuedDrainStarted(event: TuiSessionLifecycleEvent): void {
this.options.runProjection.markQueueTurnStarted(event.turnId);
if (event.sessionId) {
this.recordRunLifecycle({
kind: 'queue-turn-started',
sessionId: event.sessionId,
...(event.turnId ? { turnId: event.turnId } : {}),
queuedCount: event.queueItemIds.length,
});
}
for (const itemId of event.queueItemIds) {
const cached = this.options.runProjection.findQueueItem(itemId);
if (!cached) continue;
Expand All @@ -767,7 +968,11 @@ export class TuiRuntimeEventFlow {
liveTurnDurationMs?: number,
interactionHandled = false,
): Promise<boolean> {
if (this.options.controller.snapshot().activeTurnId) return false;
const inProcessTurnId = this.options.controller.snapshot().activeTurnId;
if (inProcessTurnId) {
this.recordTerminalNotOwned(event, sessionId, 'in-process-turn', { inProcessTurnId });
return false;
}

const activeRun = this.options.activeRunFlow.currentSnapshot();
const activeRuntimeTurnId =
Expand All @@ -787,7 +992,10 @@ export class TuiRuntimeEventFlow {
!event.turnId ||
projectedTurnId === event.turnId ||
projectedTurnId.startsWith('session:');
if (!eventMatchesProjection) return false;
if (!eventMatchesProjection) {
this.recordTerminalNotOwned(event, sessionId, 'projection-mismatch', { projectedTurnId });
return false;
}

// Late or background terminal events cannot settle the foreground turn.
const ownsTerminal = Boolean(
Expand Down
15 changes: 10 additions & 5 deletions packages/tui/src/tui/features/transcript/panel.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import type { Component, Focusable } from '../../rendering/component.js';
import { truncateToWidth, visibleWidth } from '../../rendering/text.js';
import { renderTuiStructuredPreview } from '../../transcript/presentation/structured-preview.js';
import { renderDeliveredAssets } from '../../transcript/delivered-assets.js';
import { wrapLiteralUserText } from '../../transcript/presentation/literal-text.js';
import {
renderTuiActionHint,
tuiChalk as chalk,
Expand Down Expand Up @@ -506,11 +507,15 @@ export class TuiTranscriptPanel implements TuiFeatureScreen, Component, Focusabl
if (this.rawIds.has(cell.id)) {
return this.renderTextSection('raw', presentation.rawText || 'No raw content', contentWidth);
}
if (cell.kind === 'assistant' || cell.kind === 'assistant-preamble' || cell.kind === 'user') {
const content =
cell.kind === 'user'
? { text: cell.content, assets: [] }
: projectAssistantContentForTerminal(cell.content);
if (cell.kind === 'user') {
// User prompts render literally here too, matching the main rail; see
// `wrapLiteralUserText`.
return wrapLiteralUserText(cell.content, contentWidth).map((line) =>
this.detailLine(line, width),
);
}
if (cell.kind === 'assistant' || cell.kind === 'assistant-preamble') {
const content = projectAssistantContentForTerminal(cell.content);
const body = new Markdown(content.text, 0, 0, markdownTheme, {
color: (value) => chalk.hex(colors.text)(value),
}).render(contentWidth);
Expand Down
Loading
Loading