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
Original file line number Diff line number Diff line change
Expand Up @@ -209,7 +209,7 @@ test('reseeds an evicted replica on the same live subscription', async () => {
await owner.close();
});

for (const reason of ['slow_consumer', 'transcript_changed'] as const) test(`a reseed superseded by subscription recovery does not displace the new replica (${reason})`, async () => {
test('a reseed superseded by subscription recovery does not displace the new replica', async () => {
const firstEvents = new AsyncFrameQueue();
const secondEvents = new AsyncFrameQueue();
const reseedFetch = deferred<void>();
Expand Down Expand Up @@ -270,7 +270,7 @@ for (const reason of ['slow_consumer', 'transcript_changed'] as const) test(`a r
hostEpoch: 'host-1',
subscriptionId: 'subscription-1',
sequence: 1,
reason,
reason: 'slow_consumer',
});
await pollFor(() => opens === 2);
reseedFetch.resolve(undefined);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -443,7 +443,6 @@ function isRecoverableSubscriptionFailure(error: unknown): boolean {
if (!(error instanceof RuntimeHostSubscriptionError)) return false;
return (
error.reason === 'slow_consumer' ||
error.reason === 'transcript_changed' ||
error.reason === 'sequence_gap' ||
error.reason === 'projection_revision_invalid'
);
Expand Down
3 changes: 1 addition & 2 deletions docs/astryx-surface-file-inventory.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ Generated against `@astryxdesign/core@0.6.2` (195 component exports).

Wiki bar: Design Conventions · API Use-the-System · Theming · Container Padding.

**Totals:** 299 files — blocker 0, reimplementation 0, polish 4, aligned 295.
**Totals:** 298 files — blocker 0, reimplementation 0, polish 4, aligned 294.

## Exclusions (explicit)

Expand Down Expand Up @@ -274,7 +274,6 @@ Wiki bar: Design Conventions · API Use-the-System · Theming · Container Paddi
| `packages/ui/src/daily-review-panel.tsx` | module-hub | Banner, Button, Divider, EmptyState, HStack, Heading, List, ListItem, SegmentedControl, SegmentedControlItem, Skeleton, StackItem, Text, Toolbar, VStack | aligned — uses Astryx (Banner, Button, Divider, EmptyState, HStack, Heading, List, ListItem) | aligned |
| `packages/ui/src/directory-reference-chip.tsx` | ui-composition | Token, Tooltip | aligned — uses Astryx (Token, Tooltip) | aligned |
| `packages/ui/src/executor-model-picker.tsx` | ui-composition | Button, Popover | aligned — uses Astryx (Button, Popover) | aligned |
| `packages/ui/src/form-interaction-history.tsx` | ui-composition | none | aligned — no raw controls; no Astryx JSX usage | aligned |
| `packages/ui/src/form-interaction-prompt.tsx` | ui-composition | Button, CheckboxInput, RadioList, RadioListItem, Selector, Text, TextInput | aligned — uses Astryx (Button, CheckboxInput, RadioList, RadioListItem, Selector, Text, TextInput) | aligned |
| `packages/ui/src/icons.tsx` | ui-composition | none | aligned — no raw controls; no Astryx JSX usage | aligned |
| `packages/ui/src/inline-reference.tsx` | ui-composition | ChatTokenizedText | aligned — uses Astryx (ChatTokenizedText) | aligned |
Expand Down
1 change: 0 additions & 1 deletion docs/astryx-surface-file-inventory.paths
Original file line number Diff line number Diff line change
Expand Up @@ -244,7 +244,6 @@ packages/ui/src/composer.tsx
packages/ui/src/daily-review-panel.tsx
packages/ui/src/directory-reference-chip.tsx
packages/ui/src/executor-model-picker.tsx
packages/ui/src/form-interaction-history.tsx
packages/ui/src/form-interaction-prompt.tsx
packages/ui/src/icons.tsx
packages/ui/src/inline-reference.tsx
Expand Down
37 changes: 0 additions & 37 deletions packages/cli/src/__tests__/runtime-host-prompt-transcript.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@ import { setImmediate } from 'node:timers/promises';
import { deferred } from '@maka/core/test-only/async-primitives';
import type { StoredMessage } from '@maka/core/session';
import {
RuntimeHostOperationError,
RuntimeHostSubscriptionError,
type DecodedSessionTranscriptPage,
type RuntimeHostSessionSubscription,
Expand Down Expand Up @@ -318,42 +317,6 @@ test('recovery across root turns rereads the old prompt below the new bootstrap
await channel.close();
});

test('prompt transcript recovers when a retired subscription page fails before its close frame', async () => {
const first = new TranscriptSubscription('first', 7);
const second = new TranscriptSubscription('second', 31);
let opens = 0;
const channel = await openChannel(async () => (++opens === 1 ? first : second));
const transcript = channel.trackPromptTranscript('turn');
try {
first.advance(31);
await setImmediate();
// The Host has removed the subscription but its close frame has not arrived.
first.readPage = async () => {
throw new RuntimeHostOperationError(
'session.transcript.page',
'not_found',
'Session subscription was not found',
);
};
const expected = [result('tool', 'Form completed'), terminal()];
second.readPage = async (input) =>
second.page(input, [
{ identity: 16, message: expected[0]! },
{ identity: 24, message: expected[1]! },
]);
const observed: StoredMessage[] = [];
await transcript.reconcile(async (messages) => {
observed.push(...messages);
});
assert.equal(opens, 2);
assert.deepEqual(observed, expected);
assert.equal(second.pages[0]?.anchorSequence, 7);
} finally {
transcript.dispose();
await channel.close();
}
});

test('prompt transcript rejects nonadvancing cursors and propagates consumer failures', async () => {
for (const failure of ['cursor', 'consumer'] as const) {
const subscription = new TranscriptSubscription('first', 7);
Expand Down
62 changes: 29 additions & 33 deletions packages/cli/src/__tests__/runtime-host-session-driver.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2461,41 +2461,37 @@ describe('Runtime Host Maka Session driver', () => {
assert.equal(statuses.at(-1), undefined);
});

for (const reason of ['slow_consumer', 'transcript_changed'] as const)
test(`reopens a failed Session channel before starting the next turn (${reason})`, async () => {
const first = new FakeSubscription(
continuitySnapshot({ rootTurn: null }),
Promise.resolve([]),
);
const second = new FakeSubscription(
continuitySnapshot({ rootTurn: null }),
Promise.resolve([]),
'subscription-2',
);
const connection = new FakeConnection([first, second]);
const driver = createRuntimeHostMakaSessionDriver({
connection: connection.value,
cwd: '/tmp',
llmConnectionId: 'connection-1',
llmConnectionSlug: 'openai-main',
model: 'gpt-5',
newId: sequenceIds('turn-2'),
});
await driver.switchSession('session-1');

first.push({
kind: 'subscription.closed',
hostEpoch: 'host-1',
subscriptionId: 'subscription-1',
sequence: 1,
reason,
});
await new Promise((resolve) => setImmediate(resolve));
test('reopens a failed Session channel before starting the next turn', async () => {
const first = new FakeSubscription(continuitySnapshot({ rootTurn: null }), Promise.resolve([]));
const second = new FakeSubscription(
continuitySnapshot({ rootTurn: null }),
Promise.resolve([]),
'subscription-2',
);
const connection = new FakeConnection([first, second]);
const driver = createRuntimeHostMakaSessionDriver({
connection: connection.value,
cwd: '/tmp',
llmConnectionId: 'connection-1',
llmConnectionSlug: 'openai-main',
model: 'gpt-5',
newId: sequenceIds('turn-2'),
});
await driver.switchSession('session-1');

const turn = await driver.preparePrompt('Continue');
second.push(deltaFrame(1, 'turn-2', 0, 'Recovered', 'subscription-2', 'run-2'));
assert.equal((await nextEvent(turn.events)).text, 'Recovered');
first.push({
kind: 'subscription.closed',
hostEpoch: 'host-1',
subscriptionId: 'subscription-1',
sequence: 1,
reason: 'slow_consumer',
});
await new Promise((resolve) => setImmediate(resolve));

const turn = await driver.preparePrompt('Continue');
second.push(deltaFrame(1, 'turn-2', 0, 'Recovered', 'subscription-2', 'run-2'));
assert.equal((await nextEvent(turn.events)).text, 'Recovered');
});

test('starts explicit Skills through the Host command and preserves its typed feedback', async () => {
const subscription = new FakeSubscription(
Expand Down
15 changes: 12 additions & 3 deletions packages/cli/src/pi-transcript-tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,11 @@
*/

import type { ToolOutputStream, ToolResultContent } from '@maka/core/events';
import { formatQuietJsonValue, formatToolInvocationLine } from '@maka/core/tool-quiet-preview';
import {
formatQuietJsonValue,
formatToolInvocationLine,
formatUserQuestionResult,
} from '@maka/core/tool-quiet-preview';
import { redactSecrets } from '@maka/core/display-redaction';
import {
isActiveShellRunStatus,
Expand Down Expand Up @@ -600,6 +604,11 @@ function plainResultText(entry: MakaPiToolEntry): string {
if (result?.kind === 'text') return typeof result.text === 'string' ? result.text : '';
if (result?.kind === 'json') {
const value = result.value;
const answers =
entry.toolName === 'AskUserQuestion'
? formatUserQuestionResult(entry.input, value, 'en')
: undefined;
if (answers) return answers;
if (value !== null && typeof value === 'object') {
const content = (value as { content?: unknown }).content;
if (typeof content === 'string') return content;
Expand All @@ -610,8 +619,8 @@ function plainResultText(entry: MakaPiToolEntry): string {
// Generic json fallback: use the shared quiet-value formatter instead of
// dumping a single-line JSON blob. It extracts headline + body from
// known shapes (lists, text payloads, Write/Edit results, key-value) and
// never produces escaped JSON braces (#1065). AskUserQuestion, GoalSet,
// ScheduledTask, and any future tool without a custom case render
// never produces escaped JSON braces (#1065). GoalSet, ScheduledTask,
// and any future tool without a custom case render
// human-readable text here.
const preview = formatQuietJsonValue(value, 'en');
return preview.headline ? `${preview.headline}\n${preview.body}` : preview.body;
Expand Down
13 changes: 3 additions & 10 deletions packages/cli/src/runtime-host-session-channel.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@ import {
} from '@maka/runtime-host/adapter';
import {
isRuntimeHostReconnectingConnection,
RuntimeHostOperationError,
RuntimeHostRequestInterruptedError,
RuntimeHostSubscriptionError,
type RuntimeHostConnection,
Expand Down Expand Up @@ -852,11 +851,6 @@ export class RuntimeHostSessionChannel {

#canRecover(error: unknown): boolean {
if (!isRuntimeHostReconnectingConnection(this.#connection)) return false;
// The Host can retire a transcript subscription before its close frame is
// consumed. Recover the invalidated read through the same bounded policy.
if (error instanceof RuntimeHostOperationError) {
return error.operation === 'session.transcript.page' && error.code === 'not_found';
}
if (error instanceof RuntimeHostRequestInterruptedError) {
return error.reason === 'connection_lost';
}
Expand All @@ -865,7 +859,6 @@ export class RuntimeHostSessionChannel {
(error.reason === 'connection_closed' ||
error.reason === 'sequence_gap' ||
error.reason === 'projection_revision_invalid' ||
error.reason === 'transcript_changed' ||
error.reason === 'slow_consumer')
);
}
Expand All @@ -886,10 +879,10 @@ export class RuntimeHostSessionChannel {
return;
}
if (frame.kind === 'subscription.closed') {
Comment thread
Astro-Han marked this conversation as resolved.
if (frame.reason === 'slow_consumer' || frame.reason === 'transcript_changed') {
if (frame.reason === 'slow_consumer') {
throw new RuntimeHostSubscriptionError(
frame.reason,
'Runtime Host Session transcript requires a fresh subscription',
'slow_consumer',
'Runtime Host Session subscription consumer fell behind',
);
}
this.#fail(new Error(`Runtime Host Session subscription closed: ${frame.reason}`));
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/__tests__/tool-quiet-preview.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import {
formatAsKeyValueLines,
formatQuietJsonValue,
formatToolInvocationLine,
formatUserQuestionResult,
projectToolArgsPreview,
} from '../tool-quiet-preview.js';
import { projectToolActivityArgs } from '../tool-activity-args.js';
Expand All @@ -36,6 +37,36 @@ describe('tool quiet preview', () => {
assert.doesNotMatch(key, /secret/);
assert.match(key, /redacted/i);
});

it('lists AskUserQuestion answers against the offered options', () => {
const args = {
questions: [
{ question: 'Which client?', options: [{ label: 'claude' }, { label: 'maka' }] },
{ question: 'Which scope?', options: [{ label: 'user' }, { label: 'project' }] },
{ question: 'Anything else?', options: [{ label: 'no' }, { label: 'yes' }] },
],
};
const value = {
answers: [
{ question: 'Which client?', answer: 'maka' },
{ question: 'Which scope?', answer: null },
{ question: 'Anything else?', answer: 'typed reply' },
],
};
assert.equal(
formatUserQuestionResult(projectToolActivityArgs('AskUserQuestion', args), value, 'zh-CN'),
'Which client?\n claude\n✓ maka\n\nWhich scope?\n user\n project\n 未回答\n\nAnything else?\n no\n yes\n✓ typed reply',
);
// Live rows carry only the question-text args preview.
assert.equal(
formatUserQuestionResult(projectToolArgsPreview('AskUserQuestion', args), value, 'en'),
'Which client?\n✓ maka\n\nWhich scope?\n Not answered\n\nAnything else?\n✓ typed reply',
);
assert.equal(
formatUserQuestionResult(args, { answers: [{ question: 'Which client?' }] }, 'en'),
undefined,
);
});
});

describe('formatToolInvocationLine', () => {
Expand Down
43 changes: 0 additions & 43 deletions packages/core/src/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,14 +20,6 @@
import type { ExecutorConfiguration } from './executor-catalog.js';

import { isWorkHubActionReceipt, type WorkHubActionReceipt } from './workhub-action-result.js';
import {
decodeInteractionRequest,
decodeInteractionCanonicalOutcome,
isInteractionCanonicalOutcomeValidForRequest,
type InteractionFormRequest,
type InteractionQuestionRequest,
type InteractionCanonicalOutcome,
} from './interaction.js';
import { isExecutorId } from './executor-id.js';
import { isThinkingLevel, type ThinkingLevel } from './model-thinking.js';

Expand Down Expand Up @@ -782,7 +774,6 @@ export type StoredMessage =
| AssistantMessage
| ToolCallMessage
| ToolResultMessage
| FormInteractionMessage
| PermissionDecisionMessage
| TokenUsageMessage
| TurnStateMessage
Expand Down Expand Up @@ -927,19 +918,6 @@ export interface ToolResultMessage {
parentOperationId?: string;
}

/** Read projection of a canonical answered or closed form or question, never model-authored text. */
export interface FormInteractionMessage {
type: 'form_interaction';
id: string;
turnId: string;
ts: number;
request: InteractionFormRequest | InteractionQuestionRequest;
outcome: Extract<
InteractionCanonicalOutcome,
{ kind: 'form_answer' | 'question_answer' | 'closure' }
>;
}

export interface PermissionDecisionMessage {
type: 'permission_decision';
/** Equals PermissionRequestEvent.requestId for audit correlation. */
Expand Down Expand Up @@ -1369,10 +1347,6 @@ const TOOL_RESULT_MESSAGE_SHAPE = defineObjectShape<ToolResultMessage>()(
'parentOperationId',
],
);
const FORM_INTERACTION_MESSAGE_SHAPE = defineObjectShape<FormInteractionMessage>()(
['type', 'id', 'turnId', 'ts', 'request', 'outcome'],
[],
);
const PERMISSION_DECISION_MESSAGE_SHAPE = defineObjectShape<PermissionDecisionMessage>()(
['type', 'id', 'turnId', 'ts', 'toolUseId', 'toolName', 'decision'],
['rememberForTurn', 'reviewer', 'rationale', 'riskLevel', 'hint'],
Expand Down Expand Up @@ -1673,23 +1647,6 @@ function decodeMessage(
)
return message as unknown as ToolResultMessage;
break;
case 'form_interaction':
if (
hasMessageEnvelope(message, true) &&
hasExactShape(message, FORM_INTERACTION_MESSAGE_SHAPE)
) {
const request = decodeInteractionRequest(message.request);
const outcome = decodeInteractionCanonicalOutcome(message.outcome);
if (
(request.kind === 'form' || request.kind === 'question') &&
(outcome.kind === 'form_answer' ||
outcome.kind === 'question_answer' ||
outcome.kind === 'closure') &&
isInteractionCanonicalOutcomeValidForRequest(request, outcome)
)
return { ...message, request, outcome } as unknown as FormInteractionMessage;
}
break;
case 'permission_decision':
if (
hasExactShape(message, PERMISSION_DECISION_MESSAGE_SHAPE) &&
Expand Down
Loading