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 @@ -591,7 +591,7 @@ function createMessages(
stores: ExecutionStoresWriter<'interactive'>,
): HostMessageCoordinator {
const root: HostMessageRootPort = {
readSessionHeader: async () => ({ isArchived: false }),
readSessionAvailability: async () => ({ isArchived: false }),
readRootState: () => ({ kind: 'active', sessionId, turnId: 'turn-1', runId: 'run-1' }),
claimStopFence: async () => ({
ready: Promise.resolve(),
Expand Down
38 changes: 38 additions & 0 deletions packages/runtime-host/src/__tests__/execution-composition.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -470,6 +470,44 @@ test('production composition commits automatic titles through Host-owned Session
});
});

test('production composition lazily adopts an unlocked legacy Session before turn start', async () => {
await withCompositionRoot(async ({ root, owner }) => {
const connectionId = await configureFakeDefaultTarget(owner);
const stores = await openInteractiveExecutionStoresForWrite(owner.lease);
const session = await stores.sessionStore.create({
cwd: root,
llmConnectionSlug: 'fake',
model: 'fake-model',
permissionMode: 'ask',
});
const { composition } = await createCapturedExecutionComposition(owner);
const context = {
hostEpoch: 'execution-composition-test',
connectionId: 'legacy-adoption-client',
principal: 'local_os_user' as const,
acquireResidency: () => ({ release() {} }),
};
try {
const started = await composition.handlers['turn.start'](
{
sessionId: session.id,
turnId: 'legacy-adoption-turn',
content: { text: 'adopt this legacy Session' },
},
context,
);
assert.equal(started.ok, true, JSON.stringify(started));
await waitFor(
async () =>
(await stores.sessionStore.readHeaderSnapshot(session.id)).llmConnectionId ===
connectionId,
);
} finally {
await composition.close();
}
});
});

test('WorkHub creates new work through the production assignment composition', async () => {
await withCompositionRoot(async ({ root, owner }) => {
const connectionId = await configureFakeDefaultTarget(owner);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -557,7 +557,8 @@ async function createFixture(options: { recoverAdmissions?: boolean } = {}): Pro
let requestedDrain = false;
const goalChangeListeners = new Set<() => void>();
const rootPort: HostMessageRootPort = {
readSessionHeader: (sessionId) => requireCoordinator(coordinator).readSessionHeader(sessionId),
readSessionAvailability: (sessionId, admission) =>
requireCoordinator(coordinator).readSessionAvailability(sessionId, admission),
readRootState: (sessionId) => requireCoordinator(coordinator).readRootState(sessionId),
claimStopFence: (input, commitQueueFence, lease) =>
requireCoordinator(coordinator).claimStopFence(input, commitQueueFence, lease),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3629,7 +3629,7 @@ function createFixture(
const terminal = deferred<TurnSnapshot>();
let coordinator: HostMessageCoordinator;
const root: HostMessageRootPort = {
readSessionHeader: async () => {
readSessionAvailability: async () => {
return { isArchived: false };
},
readRootState: async () => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1704,7 +1704,7 @@ test('worktree child Sessions reject roots outside managed child execution', asy
},
);
assert.equal(
(await fixture.coordinator.readSessionHeader(child.id))?.unavailableReason,
(await fixture.coordinator.readSessionAvailability(child.id))?.unavailableReason,
unavailableMessage,
);
await assert.rejects(
Expand Down Expand Up @@ -2582,8 +2582,8 @@ test('hosted linked child roots share admission, message, terminal, and stop aut
let drainRequested = false;
let stopClosureSignal: ReturnType<typeof deferred<void>> | undefined;
const rootPort: HostMessageRootPort = {
readSessionHeader: (sessionId) =>
requireCoordinator(coordinator).readSessionHeader(sessionId),
readSessionAvailability: (sessionId, admission) =>
requireCoordinator(coordinator).readSessionAvailability(sessionId, admission),
readRootState: (sessionId) => requireCoordinator(coordinator).readRootState(sessionId),
claimStopFence: (input, commitQueueFence, admission) =>
requireCoordinator(coordinator).claimStopFence(input, commitQueueFence, admission),
Expand Down Expand Up @@ -5078,7 +5078,8 @@ async function createFailureFixture(options: {
let messages!: HostMessageCoordinator;
let interactions: HostInteractionCoordinator | undefined;
const rootPort: HostMessageRootPort = {
readSessionHeader: (sessionId) => requireCoordinator(coordinator).readSessionHeader(sessionId),
readSessionAvailability: (sessionId, admission) =>
requireCoordinator(coordinator).readSessionAvailability(sessionId, admission),
readRootState: (sessionId) => requireCoordinator(coordinator).readRootState(sessionId),
claimStopFence: (input, commitQueueFence, admission) =>
requireCoordinator(coordinator).claimStopFence(input, commitQueueFence, admission),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1075,6 +1075,91 @@ test('explicit recovery persists the selected Connection entity identity', async
assert.equal(fixture.header().model, 'model-1');
});

test('unlocked legacy Sessions lazily adopt the current slug owner with CAS revalidation', async () => {
const observed: unknown[] = [];
const fixture = createFixture({
legacyConnectionIdentity: true,
header: { connectionLocked: false },
connection: {
onResolve: (ref) => observed.push(ref),
},
});

const adopted = await fixture.admission.run(fixture.sessionId, (lease) =>
fixture.coordinator.adoptLegacySessionConnectionIdentity(fixture.sessionId, lease),
);

assert.ok(adopted);
assert.equal(fixture.header().llmConnectionId, 'connection-1');
assert.deepEqual(observed, [
{
kind: 'catalog_slug',
connectionSlug: 'test',
},
]);
});

test('locked legacy Sessions never guess an identity, even when the slug exists today', async () => {
const observed: unknown[] = [];
const fixture = createFixture({
legacyConnectionIdentity: true,
connection: {
onResolve: (ref) => observed.push(ref),
},
});

const adopted = await fixture.admission.run(fixture.sessionId, (lease) =>
fixture.coordinator.adoptLegacySessionConnectionIdentity(fixture.sessionId, lease),
);

assert.equal(adopted, undefined);
assert.equal(fixture.header().llmConnectionId, undefined);
assert.deepEqual(observed, []);
});

test('legacy adoption leaves a Session unbound when its metadata CAS loses a race', async () => {
const fixture = createFixture({
legacyConnectionIdentity: true,
header: { connectionLocked: false },
stores: {
updateHeaderVersioned: async (_sessionId, _patch, expectedRevision) => {
throw new SessionMetadataVersionConflictError(
'session-1',
expectedRevision,
expectedRevision + 1,
);
},
},
});

const adopted = await fixture.admission.run(fixture.sessionId, (lease) =>
fixture.coordinator.adoptLegacySessionConnectionIdentity(fixture.sessionId, lease),
);

assert.equal(adopted, undefined);
assert.equal(fixture.header().llmConnectionId, undefined);
});

for (const executionResolution of [
{ kind: 'not_found' as const },
{ kind: 'identity_mismatch' as const },
]) {
test(`legacy adoption stays unbound when slug preflight is ${executionResolution.kind}`, async () => {
const fixture = createFixture({
legacyConnectionIdentity: true,
header: { connectionLocked: false },
connection: { executionResolution },
});

const adopted = await fixture.admission.run(fixture.sessionId, (lease) =>
fixture.coordinator.adoptLegacySessionConnectionIdentity(fixture.sessionId, lease),
);

assert.equal(adopted, undefined);
assert.equal(fixture.header().llmConnectionId, undefined);
});
}

test('creation persists a canonical cwd while fingerprints retain exact target intent', async () => {
const root = await mkdtemp(join(tmpdir(), 'maka-session-create-cwd-'));
const target = join(root, 'target');
Expand Down Expand Up @@ -1498,6 +1583,7 @@ function createFixture(
},
...options.manager,
};
const admission = new SessionAdmissionGate();
const continuity: SessionContinuity = {
refreshCanonical: async () => undefined,
...options.continuity,
Expand All @@ -1506,7 +1592,7 @@ function createFixture(
stores,
runtimePolicy,
manager,
admission: new SessionAdmissionGate(),
admission,
continuity,
workspaceResolver: new HostWorkspaceResolver(
options.projectCatalog ?? ({ list: async () => [] } as never),
Expand All @@ -1519,6 +1605,7 @@ function createFixture(
});
return {
coordinator,
admission,
sessionId,
revision: () => revision,
header: () => header,
Expand Down
35 changes: 21 additions & 14 deletions packages/runtime-host/src/server/execution-composition.ts
Original file line number Diff line number Diff line change
Expand Up @@ -557,8 +557,8 @@ export async function createExecutionRuntimeHostComposition(
let deepResearch: HostDeepResearchCoordinator | undefined;
let dailyReview: HostDailyReviewCoordinator | undefined;
const rootPort: HostMessageRootPort = {
readSessionHeader: (sessionId) =>
requireRootCoordinator(rootCoordinator).readSessionHeader(sessionId),
readSessionAvailability: (sessionId, admission) =>
requireRootCoordinator(rootCoordinator).readSessionAvailability(sessionId, admission),
readRootState: (sessionId) =>
requireRootCoordinator(rootCoordinator).readRootState(sessionId),
claimStopFence: (input, commitQueueFence, admission) =>
Expand Down Expand Up @@ -1163,6 +1163,18 @@ export async function createExecutionRuntimeHostComposition(
Date.now,
context.sessionAccessAuthority,
);
const sessionCatalog = new HostSessionCatalogCoordinator({
stores: stores.sessionStore,
runtimePolicy: runtimePolicyStores,
manager,
admission: sessionAdmission,
continuity: continuityCoordinator,
workspaceResolver,
requestDrain: context.requestDrain,
...(context.sessionAccessAuthority
? { sessionAccessAuthority: context.sessionAccessAuthority }
: {}),
});
rootCoordinator = new RootTurnCoordinator(
manager,
stores,
Expand Down Expand Up @@ -1209,6 +1221,13 @@ export async function createExecutionRuntimeHostComposition(
},
(input) => sessionEffectCoordinator.nameSessionFromRootMessage(input),
context.owner.capability.rootId,
async (sessionId, header, admissionLease) => {
const adopted = await sessionCatalog.adoptLegacySessionConnectionIdentity(
sessionId,
admissionLease,
);
return adopted?.header ?? header;
},
);
const coordinator = rootCoordinator;
const contextOperations = new HostContextCoordinator({
Expand Down Expand Up @@ -1327,18 +1346,6 @@ export async function createExecutionRuntimeHostComposition(
oauthCredentials,
onCommittedMutation: registerConfigurationMutation,
});
const sessionCatalog = new HostSessionCatalogCoordinator({
stores: stores.sessionStore,
runtimePolicy: runtimePolicyStores,
manager,
admission: sessionAdmission,
continuity: continuityCoordinator,
workspaceResolver,
requestDrain: context.requestDrain,
...(context.sessionAccessAuthority
? { sessionAccessAuthority: context.sessionAccessAuthority }
: {}),
});
const workHubCoordination = new HostWorkHubCoordinationCoordinator({
stateRoot: context.owner.capability.canonicalPath,
stores: stores.sessionStore,
Expand Down
19 changes: 11 additions & 8 deletions packages/runtime-host/src/server/message-coordinator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -190,7 +190,10 @@ export type HostMessageExecutionDisposition =

/** Root execution operations that must share the message coordinator's Session gate. */
export interface HostMessageRootPort {
readSessionHeader(sessionId: string): Promise<HostMessageSessionHeader | null>;
readSessionAvailability(
sessionId: string,
admission?: SessionAdmissionLease,
): Promise<HostMessageSessionHeader | null>;
readRootState(sessionId: string): Promise<HostMessageRootState> | HostMessageRootState;
claimStopFence(
input: Omit<TurnInterruptInput, 'originHostEpoch' | 'interruptId'>,
Expand Down Expand Up @@ -1094,7 +1097,7 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority {
}
| undefined;
for (let attempt = 0; ; attempt++) {
const header = await this.#root.readSessionHeader(input.sessionId);
const header = await this.#root.readSessionAvailability(input.sessionId, admission);
if (this.#failStopped) {
return failure('host_draining', 'Runtime Host message authority has failed');
}
Expand Down Expand Up @@ -1412,7 +1415,7 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority {
}

async #retractAdmitted(input: QueueRetractInput): Promise<MessageOutcome<QueueRetractResult>> {
const header = await this.#root.readSessionHeader(input.sessionId);
const header = await this.#root.readSessionAvailability(input.sessionId);
if (this.#failStopped) {
return failure('host_draining', 'Runtime Host message authority has failed');
}
Expand Down Expand Up @@ -1585,7 +1588,7 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority {
async #retractQueuedEntryAdmitted(
input: QueueEntryRetractInput,
): Promise<MessageOutcome<QueueMutationResult>> {
const header = await this.#root.readSessionHeader(input.sessionId);
const header = await this.#root.readSessionAvailability(input.sessionId);
if (this.#failStopped) {
return failure('host_draining', 'Runtime Host message authority has failed');
}
Expand Down Expand Up @@ -1621,7 +1624,7 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority {
async #promoteQueuedEntryAdmitted(
input: QueueEntryPromoteInput,
): Promise<MessageOutcome<QueueMutationResult>> {
const header = await this.#root.readSessionHeader(input.sessionId);
const header = await this.#root.readSessionAvailability(input.sessionId);
if (this.#failStopped) {
return failure('host_draining', 'Runtime Host message authority has failed');
}
Expand Down Expand Up @@ -1703,7 +1706,7 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority {
async #updateQueuedEntryAdmitted(
input: QueueEntryUpdateInput,
): Promise<MessageOutcome<QueueMutationResult>> {
const header = await this.#root.readSessionHeader(input.sessionId);
const header = await this.#root.readSessionAvailability(input.sessionId);
if (this.#failStopped) {
return failure('host_draining', 'Runtime Host message authority has failed');
}
Expand Down Expand Up @@ -1821,7 +1824,7 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority {
async #reorderQueuedEntriesAdmitted(
input: QueueEntriesReorderInput,
): Promise<MessageOutcome<QueueMutationResult>> {
const header = await this.#root.readSessionHeader(input.sessionId);
const header = await this.#root.readSessionAvailability(input.sessionId);
if (this.#failStopped) {
return failure('host_draining', 'Runtime Host message authority has failed');
}
Expand Down Expand Up @@ -1887,7 +1890,7 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority {
};
}

const header = await this.#root.readSessionHeader(input.sessionId);
const header = await this.#root.readSessionAvailability(input.sessionId);
if (this.#failStopped) {
return {
kind: 'conflict' as const,
Expand Down
Loading