Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -1034,6 +1034,68 @@ test('keeps a restored observation retryable until replacement seeding succeeds'
assert.equal(pendingAt >= 0 && readyAt > pendingAt, true);
await replacementCandidate.close();
}); // A failed replacement leaves the registry available to the next candidate.
test('keeps replacement Host ready when one transcript observation cannot seed', async () => {
const observations = new RuntimeHostSessionObservationRegistry();
const firstIpc = ipcHarness();
const firstHost = connectionHarness('transcript-source', {
sessionId: 'session-1',
});
const firstCandidate = await createDesktopRuntimeHostCandidate(
firstHost.connection,
deps(firstIpc),
observations,
);
await firstIpc.invoke('sessions:observe', 'session-1', 'observer-1');
await firstIpc.invoke('sessions:observe', 'session-2', 'observer-2');
await firstCandidate.close();

const events: Array<{ channel: string; payload: unknown }> = [];
const failingHost = connectionHarness('transcript-seed-failure', {
sessionId: 'session-1',
subscribeFailure: new RuntimeHostOperationError(
'subscription.open',
'transcript_unavailable',
'Session transcript is unavailable',
),
subscribeFailureSessionId: 'session-1',
subscriptionSnapshots: {
'session-2': continuitySnapshot({
session: { ...continuitySnapshot().session, sessionId: 'session-2' },
rootTurn: {
sessionId: 'session-2',
turnId: 'turn-2',
runId: 'run-2',
status: 'running',
},
}),
},
});
const candidate = await createDesktopRuntimeHostCandidate(
failingHost.connection,
{
...deps(ipcHarness()),
renderer: {
send(channel, _host, payload) {
events.push({ channel, payload });
},
},
},
observations,
);
assert.deepEqual(observations.observationSessionIds(), ['session-1', 'session-2']);
await waitFor(() => events.some(({ channel, payload }) =>
channel === 'sessions:event:session-1' &&
(payload as { type?: string }).type === 'host_observation_error'));
assert.ok(events.some(({ channel, payload }) =>
channel === 'sessions:active-interactions-changed' &&
(payload as { sessionId?: string }).sessionId === 'session-2'));
assert.ok(!events.some(({ channel, payload }) =>
channel === 'sessions:event:session-2' &&
(payload as { type?: string }).type === 'host_observation_error'));
await candidate.close();
await observations.close();
});

test('drops a stale shared Session observation when Guest access is gone', async () => {
const observations = new RuntimeHostSessionObservationRegistry();
const firstIpc = ipcHarness();
Expand Down Expand Up @@ -1259,8 +1321,10 @@ function connectionHarness(
sessionId?: string;
revisionAbandon?: 'abandoned' | 'retained';
subscriptionSnapshot?: SessionContinuitySnapshot;
subscriptionSnapshots?: Record<string, SessionContinuitySnapshot>;
assistantStreams?: readonly SessionAssistantStreamIdentity[];
subscribeFailure?: Error;
subscribeFailureSessionId?: string;
runtimeResourcePty?: ReturnType<typeof ptySnapshot>;
runtimeResourceUpdate?: ShellRunUpdate;
sharedSessionAvailable?: boolean;
Expand Down Expand Up @@ -1402,7 +1466,9 @@ function connectionHarness(
throw new Error(`Unexpected operation: ${operation}`);
},
openSessionSubscription: async ({ sessionId }: { sessionId: string }) => {
if (options.subscribeFailure) throw options.subscribeFailure;
if (options.subscribeFailure &&
(!options.subscribeFailureSessionId || options.subscribeFailureSessionId === sessionId))
throw options.subscribeFailure;
const subscriptionFrames = new AsyncFrameQueue();
activeSubscriptionFrames = subscriptionFrames;
// The Host holds a subscription's frames until the subscriber calls
Expand Down Expand Up @@ -1435,7 +1501,7 @@ function connectionHarness(
ptyListeners.add(listener);
return () => ptyListeners.delete(listener);
},
snapshot: options.subscriptionSnapshot ?? {
snapshot: options.subscriptionSnapshots?.[sessionId] ?? options.subscriptionSnapshot ?? {
projectionRevision: 1,
session: { sessionId },
},
Expand Down
12 changes: 11 additions & 1 deletion apps/desktop/src/main/runtime-host-desktop-candidate.ts
Original file line number Diff line number Diff line change
Expand Up @@ -420,6 +420,7 @@ async function restoreSessionObservations(input: {
sessionIds(): string[];
announcePending(sessionId: string): void;
attach(): Promise<string[]>;
isRecoverable?(sessionId: string): boolean;
}): Promise<string[]> {
const requested = input.sessionIds();
for (const sessionId of requested) input.announcePending(sessionId);
Expand All @@ -428,7 +429,7 @@ async function restoreSessionObservations(input: {
const restoredSet = new Set(restored);
const registeredSet = new Set(input.sessionIds());
const failed = requested.filter(
(sessionId) => registeredSet.has(sessionId) && !restoredSet.has(sessionId),
(sessionId) => registeredSet.has(sessionId) && !restoredSet.has(sessionId) && !input.isRecoverable?.(sessionId),
);
if (failed.length > 0) {
throw new Error(`Failed to restore Session observations: ${failed.join(', ')}`);
Expand Down Expand Up @@ -768,10 +769,12 @@ export async function createDesktopRuntimeHostCandidate(
}
}
observationsAttached = Boolean(sessionObserver);
const transcriptSeedFailures = new Set<string>();
const restoredSessionIds = await restoreSessionObservations({
sessionIds: () => sessionObservations.observationSessionIds(),
announcePending: (sessionId) =>
sendToRenderer(`sessions:event:${sessionId}`, { type: 'host_observation_pending' }),
isRecoverable: (sessionId) => transcriptSeedFailures.has(sessionId),
attach: () => sessionObservations.attach(
sessionObserver,
(target) => ({
Expand All @@ -786,6 +789,13 @@ export async function createDesktopRuntimeHostCandidate(
off: target.off.bind(target),
}),
(missingSessionId) => emitSessionsChanged("deleted", missingSessionId),
(failedSessionId) => {
transcriptSeedFailures.add(failedSessionId);
sendToRenderer(`sessions:event:${failedSessionId}`, {
type: 'host_observation_error',
message: 'Session transcript is unavailable',
});
},
),
});
for (const sessionId of restoredSessionIds) {
Expand Down
12 changes: 12 additions & 0 deletions apps/desktop/src/main/runtime-host-session-observation-registry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,14 @@ function isMissingRuntimeHostSessionError(error: unknown): boolean {
return error.code === "not_found";
}

function isTranscriptSeedFailure(error: unknown): boolean {
return (
error instanceof RuntimeHostOperationError &&
error.operation === 'subscription.open' &&
error.code === 'transcript_unavailable'
);
}
Comment on lines +102 to +108

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This tolerance is keyed on the Host's exact sentence, which makes it a contract that only the two source files and a test can see. The Host builds that message in several places (packages/runtime-host/src/server/session-continuity-coordinator.ts:978, :1011, :1168), and the test here hard-codes the same string — so a reword on either side, or any wrapper that prefixes it along the way, silently turns the tolerance off and restores the hard startup failure this change exists to avoid. Nothing in the type system or in CI connects the two.

operation and code narrow it usefully, but persistence_failed is too broad to be the discriminator on its own. A dedicated code or reason on the error — something like transcript_unavailable — would let both sides agree without depending on wording, and the Desktop could keep matching on the operation plus that code.


interface SessionObservationRegistration {
readonly sessionId: string;
readonly messageAdmissions: boolean;
Expand Down Expand Up @@ -175,6 +183,7 @@ export class RuntimeHostSessionObservationRegistry {
source: SessionObservationSource,
bindTarget: ObservationTargetBinding = (target) => target,
onSessionMissing?: (sessionId: string) => void,
onTranscriptSeedFailure?: (sessionId: string) => void,
): Promise<string[]> {
this.#assertOpen();
if (this.#source && this.#source !== source) {
Expand Down Expand Up @@ -224,6 +233,9 @@ export class RuntimeHostSessionObservationRegistry {
if (registration.lifecycle === "pending") {
this.#deleteRegistration(observerId, registration);
}
if (isTranscriptSeedFailure(error)) {
onTranscriptSeedFailure?.(registration.sessionId);
}
this.#onError(error);
}
return undefined;
Expand Down
6 changes: 4 additions & 2 deletions packages/cli/src/pi-transcript.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ import type {
} from '@maka/core/events';
import {
deriveTurnRecords,
isRuntimeSystemNoteKind,
isUserVisibleSessionSystemNote,
STEP_LIMIT_NOTICE_TEXT,
type StoredMessage,
type SystemNoteMessage,
Expand Down Expand Up @@ -1348,7 +1348,7 @@ function tokenDelta(before: number | undefined, after: number | undefined): numb
function systemNoteText(message: SystemNoteMessage): string | undefined {
// Retired kinds are still decoded off legacy transcript rows, and none of
// them ever had a line here worth reading.
if (!isRuntimeSystemNoteKind(message.kind)) return undefined;
if (!isUserVisibleSessionSystemNote(message.kind)) return undefined;
switch (message.kind) {
case 'context_compacted':
return 'Context compacted to keep this task within the model window.';
Expand Down Expand Up @@ -1395,6 +1395,8 @@ function systemNoteText(message: SystemNoteMessage): string | undefined {
}
case 'step_limit':
return STEP_LIMIT_NOTICE_TEXT;
case 'transcript_omitted':
return 'Part of this session transcript was omitted because it exceeded the display limit.';
}
}

Expand Down
13 changes: 12 additions & 1 deletion packages/core/src/__tests__/runtime-event.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,12 @@ import {
type RuntimeEvent,
type RuntimeEventActions,
} from '../runtime-event.js';
import { decodeCanonicalMessage, isRuntimeSystemNoteKind } from '../session.js';
import {
decodeCanonicalMessage,
isRuntimeSystemNoteKind,
isSystemNoteKind,
isUserVisibleSessionSystemNote,
} from '../session.js';
import { decodeTurnOrigin } from '../turn-origin.js';

/** Minimal valid RuntimeEvent; callers spread overrides on top. */
Expand Down Expand Up @@ -146,6 +151,12 @@ test('decodes a released provider dropping note that nothing writes any more', (
assert.equal(isRuntimeSystemNoteKind('context_provider_dropping'), false);
});

test('keeps reader omission notices outside the runtime-written note kinds', () => {
assert.equal(isSystemNoteKind('transcript_omitted'), true);
assert.equal(isUserVisibleSessionSystemNote('transcript_omitted'), true);
assert.equal(isRuntimeSystemNoteKind('transcript_omitted'), false);
});

test('shares one decoder across all TurnOrigin variants', () => {
const origins = [
{ kind: 'scheduled_task', scheduledTaskId: 'task-1' },
Expand Down
16 changes: 12 additions & 4 deletions packages/core/src/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -832,10 +832,10 @@ export function userFacingText(message: Pick<UserMessage, 'text' | 'displayText'

/**
* Closed policy for system notes that are part of the user-visible transcript:
* exactly the notes the runtime writes.
* runtime notes plus the reader's explicit omitted-transcript placeholder.
*/
export function isUserVisibleSessionSystemNote(kind: string): boolean {
return isRuntimeSystemNoteKind(kind);
return isRuntimeSystemNoteKind(kind) || kind === 'transcript_omitted';
}

/**
Expand Down Expand Up @@ -1289,6 +1289,9 @@ export const RUNTIME_SYSTEM_NOTE_KINDS = [
'step_limit',
] as const;

/** Reader-generated notices are never written back into a RuntimeEvent ledger. */
export const SYNTHETIC_SYSTEM_NOTE_KINDS = ['transcript_omitted'] as const;

/**
* Notes nothing writes any more, still decoded so old transcripts and run
* ledgers stay readable, and never shown. The session-level ones are owned by
Expand All @@ -1305,15 +1308,20 @@ export const RETIRED_SYSTEM_NOTE_KINDS = [
] as const;

export type RuntimeSystemNoteKind = (typeof RUNTIME_SYSTEM_NOTE_KINDS)[number];
export type SystemNoteKind = RuntimeSystemNoteKind | (typeof RETIRED_SYSTEM_NOTE_KINDS)[number];
export type SystemNoteKind =
| RuntimeSystemNoteKind
| (typeof SYNTHETIC_SYSTEM_NOTE_KINDS)[number]
| (typeof RETIRED_SYSTEM_NOTE_KINDS)[number];

export function isRuntimeSystemNoteKind(kind: string): kind is RuntimeSystemNoteKind {
return (RUNTIME_SYSTEM_NOTE_KINDS as readonly string[]).includes(kind);
}

export function isSystemNoteKind(kind: string): kind is SystemNoteKind {
return (
isRuntimeSystemNoteKind(kind) || (RETIRED_SYSTEM_NOTE_KINDS as readonly string[]).includes(kind)
isRuntimeSystemNoteKind(kind) ||
(SYNTHETIC_SYSTEM_NOTE_KINDS as readonly string[]).includes(kind) ||
(RETIRED_SYSTEM_NOTE_KINDS as readonly string[]).includes(kind)
);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,10 @@ import {
} from '@maka/core/runtime-logical-execution';
import { runtimeInvocationFailureClass } from '@maka/runtime/runtime-event-read-model';
import { randomUUID } from 'node:crypto';
import { createSessionTranscriptReader } from '../server/session-transcript-reader.js';
import {
createSessionTranscriptReader,
OmittedTurnResultError,
} from '../server/session-transcript-reader.js';
import { mkdtemp, rm } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
Expand Down Expand Up @@ -1383,6 +1386,36 @@ test('a failed WorkHub final target check rejects continuation without draining
}
});

test('an omitted delegated result remains retryable without draining the Host', async () => {
const fixture = await createFailureFixture({
registerBackend: (backends) =>
backends.register('ai-sdk', (context) => new FakeBackend(context)),
});
try {
await assert.rejects(
fixture.coordinator.startWorkHubResult(
{
kind: 'workhub_result',
eventId: 'omitted-result',
actionId: 'action',
delegationId: 'delegation',
targetSessionId: fixture.sessionId,
targetTurnId: 'delegated-turn',
},
async () => {
throw new OmittedTurnResultError('Delegated turn has an omitted transcript result');
},
),
OmittedTurnResultError,
);
assert.equal(fixture.drainRequested(), false);
} finally {
await fixture.coordinator.close();
await fixture.messages.close();
await fixture.dispose();
}
});

test('resume query preserves Session-before-activation lock ordering', async () => {
const activation = new RuntimePolicyActivationGate();
const capabilities = new HostClientCapabilityCoordinator({
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1184,7 +1184,7 @@ test('open returns a bounded immutable durable tail', async () => {
coordinator.close();
});

test('logs a durable bootstrap failure before reporting unavailable persistence', async (t) => {
test('isolates a durable bootstrap failure without draining the Host', async (t) => {
const logs: string[] = [];
t.mock.method(console, 'error', (...args: unknown[]) => {
logs.push(args.map(String).join(' '));
Expand All @@ -1193,19 +1193,21 @@ test('logs a durable bootstrap failure before reporting unavailable persistence'
`injected durable bootstrap failure: api_key=sk-secretvalue123\n${'细'.repeat(4096)}`,
);
const publicationFailures: unknown[] = [];
const baseReader = transcriptReader([]);
let bootstrapReads = 0;
const reader: SessionTranscriptReader = {
...transcriptReader([]),
...baseReader,
readDurableHighWater: async () => 0,
readDurablePage: async () => {
throw failure;
readDurablePage: async (sessionId, request, project) => {
if (bootstrapReads++ === 0) throw failure;
return baseReader.readDurablePage(sessionId, request, project);
},
};
const coordinator = new SessionContinuityCoordinator(
HOST_EPOCH,
async () => canonical(),
new SessionAdmissionGate(),
(error) => {
assert.equal(logs.length, 1, 'log the cause before the publication-failure hook');
publicationFailures.push(error);
},
reader,
Expand All @@ -1221,9 +1223,17 @@ test('logs a durable bootstrap failure before reporting unavailable persistence'
);
assert.deepEqual(outcome, {
ok: false,
error: { code: 'persistence_failed', message: 'Session transcript is unavailable' },
error: { code: 'transcript_unavailable', message: 'Session transcript is unavailable' },
});
assert.deepEqual(publicationFailures, [failure]);
assert.deepEqual(publicationFailures, []);
const retry = await coordinator.handlers['subscription.open'](
{
sessionId: SESSION_ID,
transcript: { kind: 'tail', maxBytes: SESSION_TRANSCRIPT_BOOTSTRAP_MAX_BYTES },
},
connectionContext('connection-failed-bootstrap'),
);
assert.equal(retry.ok, true, 'the same Host can still serve a subsequent open');
assert.equal(logs.length, 1);
const prefix = '[runtime-host] subscription.open transcript bootstrap failed: ';
const diagnostic = logs[0] ?? '';
Expand Down
Loading
Loading