diff --git a/apps/cli/src/commands/session.test.ts b/apps/cli/src/commands/session.test.ts index 5d97136ca..f2157418c 100644 --- a/apps/cli/src/commands/session.test.ts +++ b/apps/cli/src/commands/session.test.ts @@ -46,6 +46,7 @@ import { resolveOpenedBySessionRelation, resolveSessionCreateOwnerUserId, selectDefaultAgentConfigForCreate, + resolveSessionRequester, resolveSessionCommandRequesterUserId, resolveChatArgs, resolveRenameArgs, @@ -347,6 +348,19 @@ describe('session command helpers', () => { ); }); + it('keeps delegated requester identity separate from the authenticated executor', () => { + const delegatedRequester = { userId: 'collaborator-b' }; + expect( + resolveSessionRequester({ userId: 'machine-owner-a' }, undefined, delegatedRequester) + ).toEqual({ userId: 'collaborator-b', isDelegated: true }); + expect( + resolveSessionRequester({ userId: 'machine-owner-c' }, undefined, delegatedRequester) + ).toEqual({ userId: 'collaborator-b', isDelegated: true }); + expect(() => + resolveSessionRequester({ userId: 'machine-owner-a' }, 'someone-else', delegatedRequester) + ).toThrow('Requester identity must match the delegated Session requester.'); + }); + it('keeps Session ownership separate from the authenticated requester', () => { expect(resolveSessionCreateOwnerUserId('machine-owner', 'session-owner')).toBe('session-owner'); expect(resolveSessionCreateOwnerUserId('machine-owner', ' ')).toBe('machine-owner'); diff --git a/apps/cli/src/commands/session.ts b/apps/cli/src/commands/session.ts index 198bf0d26..799a27ba3 100644 --- a/apps/cli/src/commands/session.ts +++ b/apps/cli/src/commands/session.ts @@ -100,6 +100,7 @@ import { LoroDocumentManager, type SessionDocument } from '@/lib/loro/doc'; import { renderTerminalTable } from '@/lib/terminal-table'; import { canRequestMachineForCliToken, + canUseMachineForCliToken, type WorkspaceBillingEntitlement, listWorkspaceGitHubRepositoriesForCliToken, listWorkspacesForToken, @@ -126,6 +127,15 @@ import { getCliHttpFetch } from '@/utils/http-transport'; type CommonOptions = CommonCommandOptions; +export type DelegatedSessionRequester = { + userId: string; +}; + +type ResolvedSessionRequester = { + userId: string; + isDelegated: boolean; +}; + export const DEFAULT_SESSION_LIST_LIMIT = 50; export const MAX_MCP_SESSION_LIST_LIMIT = 200; export const DEFAULT_SESSION_HISTORY_LIMIT = 50; @@ -145,10 +155,8 @@ export type CreateOptions = CommonOptions & currentSessionId?: SessionId; defaultMachineId?: MachineId; requesterUserId?: string; - /** - * Trusted Session attribution supplied by an internal caller. Access checks - * must continue to use requesterUserId, which is bound to CLI auth. - */ + /** Trusted human requester supplied by a delegated internal caller. */ + delegatedRequester?: DelegatedSessionRequester; sessionOwnerUserId?: string; parent?: string; useCurrentSessionAsParent?: boolean; @@ -1886,15 +1894,39 @@ export async function readSessionMachineAccess(args: { workspaceId: WorkspaceId; machineId: MachineId; requesterUserId?: string; + delegatedRequester?: DelegatedSessionRequester; + localProjectId?: string; +}): Promise { + const requester = resolveSessionRequester( + args.auth, + args.requesterUserId, + args.delegatedRequester + ); + return await readResolvedSessionMachineAccess({ + auth: args.auth, + workspaceId: args.workspaceId, + machineId: args.machineId, + requester, + ...(args.localProjectId ? { localProjectId: args.localProjectId } : {}), + }); +} + +async function readResolvedSessionMachineAccess(args: { + auth: AuthContext; + workspaceId: WorkspaceId; + machineId: MachineId; + requester: ResolvedSessionRequester; localProjectId?: string; }): Promise { - const requesterUserId = resolveSessionCommandRequesterUserId(args.auth, args.requesterUserId); try { - return await canRequestMachineForCliToken({ + const readAccess = args.requester.isDelegated + ? canUseMachineForCliToken + : canRequestMachineForCliToken; + return await readAccess({ token: args.auth.token, workspaceId: args.workspaceId, machineId: args.machineId, - requesterUserId, + requesterUserId: args.requester.userId, ...(args.localProjectId ? { localProjectId: args.localProjectId } : {}), }); } catch (error) { @@ -1908,10 +1940,10 @@ async function assertMachineAccess(args: { auth: AuthContext; workspaceId: WorkspaceId; machineId: MachineId; - requesterUserId?: string; + requester: ResolvedSessionRequester; localProjectId?: string; }): Promise { - const access = await readSessionMachineAccess(args); + const access = await readResolvedSessionMachineAccess(args); if (!access.allowed) { throw new Error(`Machine access denied for ${args.machineId}: ${access.reason}`); } @@ -2060,6 +2092,28 @@ export function resolveSessionCommandRequesterUserId( return auth.userId; } +export function resolveSessionRequester( + auth: Pick, + requesterUserId?: string, + delegatedRequester?: DelegatedSessionRequester +): ResolvedSessionRequester { + if (!delegatedRequester) { + return { + userId: resolveSessionCommandRequesterUserId(auth, requesterUserId), + isDelegated: false, + }; + } + const delegatedUserId = normalizeCliValue(delegatedRequester.userId); + if (!delegatedUserId) { + throw new Error('Delegated Session requester must identify a user.'); + } + const requested = normalizeCliValue(requesterUserId); + if (requested !== undefined && requested !== delegatedUserId) { + throw new Error('Requester identity must match the delegated Session requester.'); + } + return { userId: delegatedUserId, isDelegated: true }; +} + export function resolveSessionCreateOwnerUserId( requesterUserId: string, sessionOwnerUserId?: string @@ -2096,16 +2150,16 @@ async function listAuthorizedMachineMetasForCreate(args: { auth: AuthContext; workspaceId: WorkspaceId; machines: readonly MachineMeta[]; - requesterUserId?: string; + requester: ResolvedSessionRequester; }): Promise { const rows = await Promise.all( args.machines.map(async (machine) => ({ machine, - access: await readSessionMachineAccess({ + access: await readResolvedSessionMachineAccess({ auth: args.auth, workspaceId: args.workspaceId, machineId: machine.id, - requesterUserId: args.requesterUserId, + requester: args.requester, }), })) ); @@ -2122,16 +2176,16 @@ async function filterAuthorizedLocalProjectsForCreate< workspaceId: WorkspaceId; machineId: MachineId; localProjects: readonly T[]; - requesterUserId?: string; + requester: ResolvedSessionRequester; }): Promise { const rows = await Promise.all( args.localProjects.map(async (project) => ({ project, - access: await readSessionMachineAccess({ + access: await readResolvedSessionMachineAccess({ auth: args.auth, workspaceId: args.workspaceId, machineId: args.machineId, - requesterUserId: args.requesterUserId, + requester: args.requester, localProjectId: project.id, }), })) @@ -2148,7 +2202,7 @@ async function resolveTargetMachineForCreate(args: { auth: AuthContext; machineSelector?: string; defaultMachineId?: MachineId; - requesterUserId?: string; + requester: ResolvedSessionRequester; parentSessionId?: SessionId; }): Promise { const machines = await listMachineMetasForWorkspace(args.manager); @@ -2159,7 +2213,7 @@ async function resolveTargetMachineForCreate(args: { auth: args.auth, workspaceId: args.workspaceId, machines, - requesterUserId: args.requesterUserId, + requester: args.requester, }); if (authorizedMachines.length === 0) { throw new Error('No authorized machines are available in this workspace.'); @@ -2249,12 +2303,12 @@ async function assertGitHubRepoAccess(args: { auth: AuthContext; workspaceId: WorkspaceId; repoFullName: string; - requesterUserId?: string; + requesterUserId: string; }): Promise { const repos = await listWorkspaceGitHubRepositoriesForCliToken({ token: args.auth.token, workspaceId: args.workspaceId, - requesterUserId: resolveSessionCommandRequesterUserId(args.auth, args.requesterUserId), + requesterUserId: args.requesterUserId, enabledOnly: true, }); const normalized = args.repoFullName.toLowerCase(); @@ -2269,7 +2323,7 @@ export async function readLocalProjectGitStateOnMachine(args: { machineId: MachineId; localProjectId: string; localRootPath: string; - requesterUserId?: string; + requesterUserId: string; }): Promise< | { success: true; state: Awaited> } | { success: false; error: string; message?: string } @@ -2290,7 +2344,7 @@ export async function readLocalProjectGitStateOnMachine(args: { async (client) => await client.requestLocalProjectGitState({ localProjectId: args.localProjectId as LocalProjectId, - requestedByUserId: resolveSessionCommandRequesterUserId(args.auth, args.requesterUserId), + requestedByUserId: args.requesterUserId, timeoutMs: 30_000, }) ); @@ -2395,13 +2449,13 @@ export function resolveLocalProjectCreateGitContext(args: { async function listWorkspaceGitHubRepositoriesBestEffort(args: { auth: AuthContext; workspaceId: WorkspaceId; - requesterUserId?: string; + requesterUserId: string; }): Promise<{ fullName: string }[]> { try { return await listWorkspaceGitHubRepositoriesForCliToken({ token: args.auth.token, workspaceId: args.workspaceId, - requesterUserId: resolveSessionCommandRequesterUserId(args.auth, args.requesterUserId), + requesterUserId: args.requesterUserId, enabledOnly: true, }); } catch (error) { @@ -2420,7 +2474,7 @@ async function resolveLocalProjectCreateGitContextOnMachine(args: { machineId: MachineId; localProjectId: string; localRootPath: string; - requesterUserId?: string; + requesterUserId: string; requestedBranch?: string; useWorktree?: boolean; }): Promise { @@ -2453,7 +2507,7 @@ async function resolveLocalProjectRefOnMachineOrThrow( auth: AuthContext, machineId: MachineId, selector: string, - requesterUserId: string | undefined, + requester: ResolvedSessionRequester, requestedBranch?: string, useWorktree?: boolean ): Promise { @@ -2469,7 +2523,7 @@ async function resolveLocalProjectRefOnMachineOrThrow( workspaceId, machineId, localProjects, - requesterUserId, + requester, }); if (authorizedLocalProjects.length === 0) { throw new Error('No authorized local projects are available on the target machine.'); @@ -2498,7 +2552,7 @@ async function resolveLocalProjectRefOnMachineOrThrow( machineId, localProjectId: project.id, localRootPath: project.rootPath, - requesterUserId, + requesterUserId: requester.userId, requestedBranch, useWorktree, }); @@ -2570,14 +2624,12 @@ async function resolveCreateContext(args: { workspace: WorkspaceSummary; manager: LoroDocumentManager; options: CreateOptions; + requester: ResolvedSessionRequester; skipMachineAvailabilityCheck?: boolean; }): Promise { const workspaceId = args.workspace.id as WorkspaceId; const agentSelector = resolveCreateAgentSelector(args.options); - const requesterUserId = resolveSessionCommandRequesterUserId( - args.auth, - args.options.requesterUserId - ); + const requesterUserId = args.requester.userId; const parentSelector = normalizeCliValue(args.options.parent); const currentSessionId = resolveCreateCurrentSessionId(args.options); if (parentSelector && args.options.useCurrentSessionAsParent === true) { @@ -2622,14 +2674,14 @@ async function resolveCreateContext(args: { auth: args.auth, machineSelector: args.options.machine, defaultMachineId: args.options.defaultMachineId, - requesterUserId, + requester: args.requester, parentSessionId, }); await assertMachineAccess({ auth: args.auth, workspaceId, machineId: targetMachine.id, - requesterUserId, + requester: args.requester, }); if (args.skipMachineAvailabilityCheck !== true) { await ensureTargetMachineOnline({ @@ -2686,7 +2738,7 @@ async function resolveCreateContext(args: { args.auth, targetMachine.id, normalizedLocalProject, - requesterUserId, + args.requester, requestedBranch, args.options.worktree === true ); @@ -2696,7 +2748,7 @@ async function resolveCreateContext(args: { auth: args.auth, workspaceId, machineId: targetMachine.id, - requesterUserId, + requester: args.requester, localProjectId: project?.kind === 'local' ? project.localProjectId : undefined, }); @@ -2728,7 +2780,12 @@ export async function validateSessionCreateOptions(args: { */ dispatchConfig?: ResolvedTurnDispatchConfig; }): Promise { - const resolved = await resolveCreateContext(args); + const requester = resolveSessionRequester( + args.auth, + args.options.requesterUserId, + args.options.delegatedRequester + ); + const resolved = await resolveCreateContext({ ...args, requester }); return await resolveEffectiveSessionCreateDispatchConfig({ manager: args.manager, workspaceId: args.workspace.id as WorkspaceId, @@ -2915,12 +2972,17 @@ export async function createSessionResult( sessionId: options.sessionId, }); } - const requesterUserId = resolveSessionCommandRequesterUserId(auth, options.requesterUserId); + const requester = resolveSessionRequester( + auth, + options.requesterUserId, + options.delegatedRequester + ); + const requesterUserId = requester.userId; const sessionOwnerUserId = resolveSessionCreateOwnerUserId( requesterUserId, options.sessionOwnerUserId ); - const resolved = await resolveCreateContext({ auth, workspace, manager, options }); + const resolved = await resolveCreateContext({ auth, workspace, manager, options, requester }); const { targetMachine, agentConfig, @@ -3087,12 +3149,30 @@ export async function validateSessionChatTarget(args: { manager: LoroDocumentManager; sessionId: SessionId; requesterUserIdOverride?: string; + delegatedRequester?: DelegatedSessionRequester; }): Promise { - await syncWorkspaceMetaForRead(args.manager, `session.chat:${args.sessionId}:prewrite:meta`); - const requesterUserId = resolveSessionCommandRequesterUserId( + const requester = resolveSessionRequester( args.auth, - args.requesterUserIdOverride + args.requesterUserIdOverride, + args.delegatedRequester ); + return await validateSessionChatTargetForRequester({ + auth: args.auth, + workspace: args.workspace, + manager: args.manager, + sessionId: args.sessionId, + requester, + }); +} + +async function validateSessionChatTargetForRequester(args: { + auth: AuthContext; + workspace: WorkspaceSummary; + manager: LoroDocumentManager; + sessionId: SessionId; + requester: ResolvedSessionRequester; +}): Promise { + await syncWorkspaceMetaForRead(args.manager, `session.chat:${args.sessionId}:prewrite:meta`); const session = await resolveSessionMetaOrThrow(args.manager, args.sessionId); if (session.isArchived) { throw new Error(`Session ${args.sessionId} is archived. Restore it before chatting.`); @@ -3101,7 +3181,7 @@ export async function validateSessionChatTarget(args: { auth: args.auth, workspaceId: args.workspace.id as WorkspaceId, machineId: session.machineId, - requesterUserId, + requester: args.requester, localProjectId: session.project?.kind === 'local' ? session.project.localProjectId : undefined, }); await ensureTargetMachineOnline({ @@ -3129,7 +3209,8 @@ export async function sendSessionChatResult( userTurnId: string; chainDepth: number; bypassSessionQuota?: boolean; - } + }, + delegatedRequester?: DelegatedSessionRequester ): Promise<{ sessionId: SessionId; machineId: MachineId; @@ -3137,13 +3218,14 @@ export async function sendSessionChatResult( userTurnId: string; completionPromise?: Promise>>; }> { - const requesterUserId = resolveSessionCommandRequesterUserId(auth, requesterUserIdOverride); - const session = await validateSessionChatTarget({ + const requester = resolveSessionRequester(auth, requesterUserIdOverride, delegatedRequester); + const requesterUserId = requester.userId; + const session = await validateSessionChatTargetForRequester({ auth, workspace, manager, sessionId, - requesterUserIdOverride, + requester, }); if (dispatchConfig.modeId || dispatchConfig.modelId || dispatchConfig.configOptionValues) { const capability = await readAgentAcpCapability({ @@ -3970,6 +4052,7 @@ const sessionCancelCommand = new Command('cancel') auth, workspaceId, machineId: session.machineId, + requester: resolveSessionRequester(auth), localProjectId: session.project?.kind === 'local' ? session.project.localProjectId : undefined, }); diff --git a/apps/cli/src/lib/message-handler.ts b/apps/cli/src/lib/message-handler.ts index d7a65c8fa..03ecdadf6 100644 --- a/apps/cli/src/lib/message-handler.ts +++ b/apps/cli/src/lib/message-handler.ts @@ -2837,6 +2837,9 @@ export class MessageHandler { throw new Error(`Requester Session not found: ${operation.requesterSessionId}`); } const requester = requesterRecord.meta as SessionMeta; + const delegatedRequester = operation.frozenContinuationConfig.sourceTurnId + ? ({ userId: operation.requesterUserId } as const) + : undefined; if (operation.kind === 'session_create' || operation.kind === 'session_create_many') { const runConfig: AgentRunConfigSelection = { @@ -2861,8 +2864,12 @@ export class MessageHandler { workspace: this.workspaceId, currentSessionId: operation.requesterSessionId, workspaceMetaPrewriteSatisfied: true, - requesterUserId: operation.requesterUserId, - sessionOwnerUserId: requester.userId, + ...(delegatedRequester + ? { delegatedRequester } + : { + requesterUserId: operation.requesterUserId, + sessionOwnerUserId: requester.userId, + }), defaultMachineId: requester.machineId, sessionId: item.target.sessionId, userTurnId: item.target.userTurnId, @@ -2911,12 +2918,13 @@ export class MessageHandler { taskToolsEnabled: operation.frozenContinuationConfig.inputConfig.taskToolsEnabled === true, }, undefined, - operation.requesterUserId, + delegatedRequester ? undefined : operation.requesterUserId, { userTurnId: item.target.userTurnId, chainDepth: operation.initiatorChainDepth + 1, bypassSessionQuota: shouldBypassSessionQuota(operation.kind), - } + }, + delegatedRequester ); } @@ -6584,6 +6592,22 @@ export class MessageHandler { allowArbitraryPaths: true, sameMachine: true, }); + case 'session/get-active-invocation-context': { + const sessionId = request.params.sessionId as SessionId; + const invocation = this.executionService.getActiveInvocationContext(sessionId); + return invocation + ? { + type: 'session/active-invocation-context' as const, + sessionId, + active: true as const, + ...invocation, + } + : { + type: 'session/active-invocation-context' as const, + sessionId, + active: false as const, + }; + } case 'session/cancel': { const result = await this.executionService.cancelSession({ type: 'session/cancel', diff --git a/apps/cli/src/mcp/AGENTS.md b/apps/cli/src/mcp/AGENTS.md index 6f609459e..23b4276d9 100644 --- a/apps/cli/src/mcp/AGENTS.md +++ b/apps/cli/src/mcp/AGENTS.md @@ -28,6 +28,15 @@ Root and `apps/cli/AGENTS.md` instructions apply. the workspace catalog; no driving-Turn mention authorization is required. Resolve its target, Prompt prefix, revision, and concrete run config before Operation acceptance. Recovery uses the frozen canonical Prompt and target dispatch config and never rereads the mutable catalog. +- Session orchestration derives its human identity from the active execution runtime populated + by the dispatch payload, not from the daemon credential, Session owner, or observed history. + An absent active runtime fails closed; never reconstruct invocation identity from history. + Freeze the source Turn id and invoking user with every accepted Operation. Store the user + once as `requesterUserId` and the causal Turn as `sourceTurnId`. The + Operation's requester Session id already identifies the source Session, and a single-value + actor tag adds no information. Recovery uses the Operation's owner Machine plus current + authorization; it does not freeze the daemon account that originally accepted the Operation. + Every MCP Session path rejects a runtime invocation without userId. - Direct Role creation stays on the ordinary `lody_session_create` and `lody_session_create_many` tools. When `agentRoleId` is present, tolerate manual Machine, Agent, and run-config fields but remove them before resolution: the current Role row is authoritative diff --git a/apps/cli/src/mcp/lody-mcp-server-chat-sync.test.ts b/apps/cli/src/mcp/lody-mcp-server-chat-sync.test.ts index 22fb55d39..5989ae7fa 100644 --- a/apps/cli/src/mcp/lody-mcp-server-chat-sync.test.ts +++ b/apps/cli/src/mcp/lody-mcp-server-chat-sync.test.ts @@ -2,9 +2,9 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; const mocks = vi.hoisted(() => ({ accept: vi.fn(), + activeInvocation: vi.fn(), findMatchingRetry: vi.fn(), getDocMeta: vi.fn(), - getHistory: vi.fn(), validateSessionChatTarget: vi.fn(), })); @@ -24,7 +24,6 @@ vi.mock('@/lib/command-runtime', async (importOriginal) => { ) => await fn({ repo: { getDocMeta: mocks.getDocMeta }, - getOrCreateSessionDoc: vi.fn(async () => ({ getHistory: mocks.getHistory })), getOnlineMachineIds: vi.fn(async () => new Set(['machine-id'])), }) ), @@ -51,6 +50,22 @@ vi.mock('@/orchestration/operation-store', async (importOriginal) => { }; }); +vi.mock('@lody/shared/node/local-ipc', async (importOriginal) => { + const actual = await importOriginal(); + const { Effect } = await import('effect'); + return { + ...actual, + makeLocalControlClientAuto: vi.fn(() => ({ + machineRpc: vi.fn(() => + Effect.sync(() => ({ + ok: true as const, + result: mocks.activeInvocation(), + })) + ), + })), + }; +}); + import { WORKSPACE_SYNC_UNAVAILABLE_MESSAGE, WorkspaceSyncUnavailableError, @@ -105,8 +120,15 @@ describe('session chat prevalidation sync failures', () => { vi.stubEnv('LODY_MCP_MACHINE_ID', 'machine-id'); vi.stubEnv('LODY_MCP_WORKSPACE_ID', 'workspace-id'); vi.stubEnv('LODY_MCP_SESSION_ID', requesterSession.id); + mocks.activeInvocation.mockReturnValue({ + type: 'session/active-invocation-context' as const, + sessionId: requesterSession.id, + active: true as const, + requesterUserId: requesterSession.userId, + sourceTurnId: 'requester-turn-id', + inputConfig: {}, + }); mocks.findMatchingRetry.mockReturnValue(undefined); - mocks.getHistory.mockResolvedValue([]); mocks.getDocMeta .mockResolvedValueOnce({ meta: requesterSession }) .mockResolvedValueOnce({ meta: targetSession }); @@ -146,4 +168,32 @@ describe('session chat prevalidation sync failures', () => { expectRetryableSyncResult(result); expect(mocks.accept).not.toHaveBeenCalled(); }); + + it('fails closed when the active execution runtime has already been released', async () => { + mocks.activeInvocation.mockReturnValue({ + type: 'session/active-invocation-context' as const, + sessionId: requesterSession.id, + active: false as const, + }); + + const result = await callAndMapMcpError(() => + startSessionChatOperation({ + operationId: 'inactive-runtime-operation', + sessionId: targetSession.id, + prompt: 'continue', + }) + ); + const content = result.content[0]; + if (!content || content.type !== 'text') throw new Error('expected text result'); + expect(JSON.parse(content.text)).toEqual({ + ok: false, + error: { + code: 'INVOKING_TURN_NOT_FOUND', + message: 'The exact Turn driving this MCP invocation is no longer active.', + retryable: false, + }, + }); + expect(mocks.validateSessionChatTarget).not.toHaveBeenCalled(); + expect(mocks.accept).not.toHaveBeenCalled(); + }); }); diff --git a/apps/cli/src/mcp/lody-mcp-server.test.ts b/apps/cli/src/mcp/lody-mcp-server.test.ts index 38c3a41f4..4b1632a29 100644 --- a/apps/cli/src/mcp/lody-mcp-server.test.ts +++ b/apps/cli/src/mcp/lody-mcp-server.test.ts @@ -13,7 +13,6 @@ import { type AgentRole, type AgentRoleId, type MachineId, - type SessionHistoryInput, type SessionId, type SessionTurnInputConfig, type WorkspaceId, @@ -79,26 +78,13 @@ const { getSessionContext, resolveOperationStorePathForContext, resolveUploadPath, - resolveInvokingHistoryInput, + buildInvocationIdentity, summarizeProjectRefForMcp, resolveSessionExecutionSnapshot, makeMachineOnlineLookupForMcp, truncateUtf8HeadTail, } = __lodyMcpServerInternals; -const historyTurn = ( - id: string, - role: SessionHistoryInput['role'], - chainDepth?: number -): SessionHistoryInput => ({ - id, - role, - timestamp: '2026-07-20T00:00:00.000Z', - items: [], - fileDiff: [], - ...(chainDepth === undefined ? {} : { inputConfig: { chainDepth } }), -}); - const createMcpContext = (): ReturnType => ({ machineId: 'machine-id', workspaceId: 'workspace-id', @@ -287,14 +273,13 @@ describe('session MCP input schemas', () => { expect(result.isError).toBe(true); }); - it('anchors continuation depth before a later queued human input', () => { - const completion = historyTurn('operation-completion', 'system', 5); - const executingAssistant = historyTurn('assistant:continuation', 'assistant'); - const queuedHumanInput = historyTurn('queued-human', 'user', 0); - - expect(resolveInvokingHistoryInput([completion, executingAssistant, queuedHumanInput])).toBe( - completion - ); + it('derives delegated identity from the exact driving Turn', () => { + expect( + buildInvocationIdentity({ id: 'source-turn', userId: 'collaborator-b', inputConfig: {} }) + ).toEqual({ userId: 'collaborator-b', sourceTurnId: 'source-turn' }); + expect(() => + buildInvocationIdentity({ id: 'legacy-turn', userId: ' ', inputConfig: {} }) + ).toThrow('has no authenticated human identity'); }); it('uses stable ids and rejects legacy selector names', () => { @@ -477,8 +462,11 @@ describe('session MCP input schemas', () => { ); bindMcpCreateContext( options, - { userId: 'machine-owner' }, - { machineId: 'machine-id', userId: 'session-owner' } + { + userId: 'collaborator-b', + sourceTurnId: 'source-turn', + }, + { machineId: 'machine-id' } ); expect(options).toMatchObject({ @@ -487,8 +475,7 @@ describe('session MCP input schemas', () => { machine: 'machine-id', agentConfig: 'agent-config-id', useCurrentSessionAsParent: true, - requesterUserId: 'machine-owner', - sessionOwnerUserId: 'session-owner', + delegatedRequester: { userId: 'collaborator-b' }, defaultMachineId: 'machine-id', }); }); diff --git a/apps/cli/src/mcp/lody-mcp-server.ts b/apps/cli/src/mcp/lody-mcp-server.ts index 273a3d7c1..3b8433b30 100644 --- a/apps/cli/src/mcp/lody-mcp-server.ts +++ b/apps/cli/src/mcp/lody-mcp-server.ts @@ -48,9 +48,10 @@ import { type MachineId, type MachineMeta, type ProjectRef, - type SessionHistoryInput, type SessionId, type SessionMeta, + SessionActiveInvocationContextResultSchema, + type SessionActiveInvocationContextResult, type TaskId, type TaskIndexRow, type TaskPrProvider, @@ -76,6 +77,7 @@ import { REVIEW_VERDICT_VALUES, ReviewSubmissionSchema, hasPendingUserTurnActivation, + normalizeSessionTurnInputConfig, } from '@lody/shared'; import { makeLocalControlClientAuto } from '@lody/shared/node/local-ipc'; import { @@ -126,6 +128,7 @@ import { type CreateOptions, type ResolvedTurnDispatchConfig, type SessionLiveStatusBatchItem, + type DelegatedSessionRequester, } from '@/commands/session'; import type { SessionTurnOutputEvent, @@ -1045,6 +1048,45 @@ const postSessionControl = async ( return responses.map((response) => LocalSessionControlResponseSchema.parse(response)); }; +const readActiveInvocationContext = async ( + ctx: McpSessionContext +): Promise => { + const response = await Effect.runPromise( + makeLocalControlClientAuto({ socketPath: ctx.localControlSocketPath }) + .machineRpc( + { + method: 'session/get-active-invocation-context', + machineId: ctx.machineId, + workspaceId: ctx.workspaceId, + params: { sessionId: ctx.sessionId }, + }, + { timeoutMs: SESSION_CONTROL_TIMEOUT_MS } + ) + .pipe( + Effect.catchTag('IpcTimeoutError', (error) => + Effect.fail( + new Error(`local control timed out after ${SESSION_CONTROL_TIMEOUT_MS}ms`, { + cause: error, + }) + ) + ), + Effect.catchTag('IpcProtocolError', (error) => + Effect.fail(new Error(error.message, { cause: error })) + ) + ) + ); + if (!response.ok) { + throw new Error(response.error); + } + const invocation = SessionActiveInvocationContextResultSchema.parse(response.result); + if (invocation.sessionId !== ctx.sessionId) { + throw new Error( + `Active invocation context session mismatch: expected ${ctx.sessionId}, received ${invocation.sessionId}` + ); + } + return invocation; +}; + const pickResponse = ( responses: LocalSessionControlResponsePayload[], expectedType: TType, @@ -2087,14 +2129,14 @@ const canUseMachineForOptions = async (args: { auth: AuthContext; workspaceId: WorkspaceId; machineId: MachineId; - requesterUserId: string; + delegatedRequester: DelegatedSessionRequester; localProjectId?: string; }): Promise => { const access = await readSessionMachineAccess({ auth: args.auth, workspaceId: args.workspaceId, machineId: args.machineId, - requesterUserId: args.requesterUserId, + delegatedRequester: args.delegatedRequester, ...(args.localProjectId ? { localProjectId: args.localProjectId } : {}), }); return access.allowed; @@ -2104,7 +2146,7 @@ const filterAuthorizedMachinesForOptions = async ( auth: AuthContext, workspaceId: WorkspaceId, machines: readonly MachineMeta[], - requesterUserId: string + delegatedRequester: DelegatedSessionRequester ): Promise => { const rows = await Promise.all( machines.map(async (machine) => ({ @@ -2113,7 +2155,7 @@ const filterAuthorizedMachinesForOptions = async ( auth, workspaceId, machineId: machine.id, - requesterUserId, + delegatedRequester, }), })) ); @@ -2125,7 +2167,7 @@ const filterAuthorizedLocalProjectsForOptions = async ( workspaceId: WorkspaceId, machineId: MachineId, localProjects: readonly LocalProjectMeta[], - requesterUserId: string + delegatedRequester: DelegatedSessionRequester ): Promise => { const rows = await Promise.all( localProjects.map(async (project) => ({ @@ -2134,7 +2176,7 @@ const filterAuthorizedLocalProjectsForOptions = async ( auth, workspaceId, machineId, - requesterUserId, + delegatedRequester, localProjectId: project.id, }), })) @@ -2170,38 +2212,73 @@ const readCurrentSessionMeta = async ( const bindMcpCreateContext = ( options: CreateOptions, - auth: Pick, - requester: Pick + identity: InvocationIdentity, + requester: Pick ): void => { - options.requesterUserId = auth.userId; - options.sessionOwnerUserId = requester.userId; + options.delegatedRequester = toDelegatedSessionRequester(identity); options.defaultMachineId = requester.machineId; }; +type InvocationIdentity = { + userId: string; + sourceTurnId: string; +}; + +const toDelegatedSessionRequester = (identity: InvocationIdentity): DelegatedSessionRequester => ({ + userId: identity.userId, +}); + type InvokingTurnContext = { chainDepth: number; frozenInputConfig: SessionTurnInputConfig; + identity: InvocationIdentity; }; -const resolveInvokingHistoryInput = ( - history: SessionHistoryInput[] -): SessionHistoryInput | undefined => { - let assistantIndex = -1; - for (let index = history.length - 1; index >= 0; index -= 1) { - if (history[index]?.role === 'assistant') { - assistantIndex = index; - break; - } +const buildInvocationIdentity = (source: InvokingTurnSource): InvocationIdentity => { + const userId = source.userId?.trim(); + if (!userId) { + throw new LodyOperationStoreError( + 'INVOKING_USER_UNAVAILABLE', + `The driving Turn ${source.id} has no authenticated human identity.`, + false + ); } - const assistant = assistantIndex >= 0 ? history[assistantIndex] : undefined; - if (assistant?.userTurnId) { - const linked = history.find((entry) => entry.id === assistant.userTurnId); - if (linked) return linked; + return { + userId, + sourceTurnId: source.id, + }; +}; + +type InvokingTurnSource = { + id: string; + userId: string; + inputConfig: SessionTurnInputConfig; +}; + +const resolveInvokingTurnSource = async (): Promise => { + const active = await readActiveInvocationContext(getSessionContext()); + if (!active.active) { + throw new LodyOperationStoreError( + 'INVOKING_TURN_NOT_FOUND', + 'The exact Turn driving this MCP invocation is no longer active.', + false + ); } - const inputsBeforeExecution = assistantIndex >= 0 ? history.slice(0, assistantIndex) : history; - return [...inputsBeforeExecution] - .reverse() - .find((entry) => entry.role === 'user' || entry.role === 'system'); + const inputConfig = + normalizeSessionTurnInputConfig(active.inputConfig) ?? + (Object.keys(active.inputConfig).length === 0 ? {} : undefined); + if (!inputConfig) { + throw new LodyOperationStoreError( + 'INVOKING_TURN_NOT_FOUND', + `The active Turn ${active.sourceTurnId} has an invalid execution configuration.`, + false + ); + } + return { + id: active.sourceTurnId, + userId: active.requesterUserId, + inputConfig, + }; }; const assertInvokingTurnTaskToolsEnabled = async ( @@ -2216,9 +2293,8 @@ const assertInvokingTurnTaskToolsEnabled = async ( false ); } - const sessionDoc = await manager.getOrCreateSessionDoc(session.id); - const source = resolveInvokingHistoryInput(await sessionDoc.getHistory()); - if (source?.inputConfig?.taskToolsEnabled !== true) { + const source = await resolveInvokingTurnSource(); + if (source.inputConfig.taskToolsEnabled !== true) { throw new LodyOperationStoreError( 'TASK_TOOLS_DISABLED', 'Lody Task tools are disabled for the driving user turn.', @@ -2227,14 +2303,9 @@ const assertInvokingTurnTaskToolsEnabled = async ( } }; -const resolveInvokingTurnContext = async ( - manager: LoroDocumentManager, - session: SessionMeta -): Promise => { - const sessionDoc = await manager.getOrCreateSessionDoc(session.id); - const history = await sessionDoc.getHistory(); - const source = resolveInvokingHistoryInput(history); - const chainDepth = source?.inputConfig?.chainDepth ?? 0; +const resolveInvokingTurnContext = async (session: SessionMeta): Promise => { + const source = await resolveInvokingTurnSource(); + const chainDepth = source.inputConfig.chainDepth ?? 0; if (chainDepth >= LODY_MAX_CHAIN_DEPTH) { throw new LodyOperationStoreError( 'CHAIN_DEPTH_EXCEEDED', @@ -2244,10 +2315,11 @@ const resolveInvokingTurnContext = async ( } return { chainDepth, + identity: buildInvocationIdentity(source), frozenInputConfig: { - ...(source?.inputConfig ?? {}), - cliType: source?.inputConfig?.cliType ?? session.cliType, - agentType: source?.inputConfig?.agentType ?? session.agentType, + ...source.inputConfig, + cliType: source.inputConfig.cliType ?? session.cliType, + agentType: source.inputConfig.agentType ?? session.agentType, chainDepth, }, }; @@ -2474,7 +2546,9 @@ const buildSessionCreateOptions = async ( if (!currentSession) { throw new Error(`Session not found: ${ctx.sessionId}`); } - const requesterUserId = auth.userId; + const invoking = await resolveInvokingTurnContext(currentSession); + const requesterUserId = invoking.identity.userId; + const delegatedRequester = toDelegatedSessionRequester(invoking.identity); const machineEntries = await listAliveDocMetas(manager, isMachineDocRoomId); const onlineMachineIds = await manager.getOnlineMachineIds(); const isMachineOnline = (machineId: MachineId): boolean => @@ -2487,7 +2561,7 @@ const buildSessionCreateOptions = async ( auth, workspaceId, machineCandidates, - requesterUserId + delegatedRequester ); const selectedMachine = selectMachineForOptions( machines, @@ -2540,7 +2614,7 @@ const buildSessionCreateOptions = async ( workspaceId, selectedMachine.id, localProjectCandidates, - requesterUserId + delegatedRequester ) ).slice(0, MAX_MCP_CREATE_OPTION_MATCHES); const summarizedLocalProjects = await Promise.all( @@ -2608,7 +2682,7 @@ const startSessionCreateOperation = async (args: SessionCreateCommandInput): Pro false ); } - const invoking = await resolveInvokingTurnContext(manager, currentSession); + const invoking = await resolveInvokingTurnContext(currentSession); const roleCatalog = args.agentRoleId ? await loadWorkspaceAgentRoleCatalog(manager, workspace.id as WorkspaceId) : undefined; @@ -2624,14 +2698,18 @@ const startSessionCreateOperation = async (args: SessionCreateCommandInput): Pro ctx.sessionId as SessionId, args.operationId!, 'session_create', - canonicalCommand + canonicalCommand, + invoking.identity.userId, + invoking.identity.sourceTurnId ) ); - if (retry) return await withOperationStore((store) => store.snapshot(retry)); + if (retry) { + return await withOperationStore((store) => store.snapshot(retry)); + } const targetMachineId = (resolved.input.machineId ?? currentSession.machineId) as MachineId; await assertMachineOnlineForSingleCommand(manager, targetMachineId, ctx); const createOptions = buildMcpCreateOptions(resolved.input, ctx); - bindMcpCreateContext(createOptions, auth, currentSession); + bindMcpCreateContext(createOptions, invoking.identity, currentSession); bindAgentRoleCreateOptions(createOptions, resolved.role); createOptions.workspaceMetaPrewriteSatisfied = true; let effectiveDispatchConfig: ResolvedTurnDispatchConfig; @@ -2660,7 +2738,7 @@ const startSessionCreateOperation = async (args: SessionCreateCommandInput): Pro workspaceId: workspace.id as WorkspaceId, ownerMachineId: ctx.machineId as MachineId, requesterSessionId: ctx.sessionId as SessionId, - requesterUserId: auth.userId, + requesterUserId: invoking.identity.userId, operationId: args.operationId!, kind: 'session_create', canonicalCommand, @@ -2669,6 +2747,7 @@ const startSessionCreateOperation = async (args: SessionCreateCommandInput): Pro ? { agentConfigId: currentSession.agentConfigId } : {}), inputConfig: invoking.frozenInputConfig, + sourceTurnId: invoking.identity.sourceTurnId, targetDispatchConfigs: [effectiveDispatchConfig], }, initiatorChainDepth: invoking.chainDepth, @@ -2752,6 +2831,7 @@ const startSessionChatOperation = async (args: SessionChatToolInput): Promise store.snapshot(retry)); + if (retry) { + return await withOperationStore((store) => store.snapshot(retry)); + } const targetSession = await readCurrentSessionMeta(manager, args.sessionId as SessionId); if (!targetSession) { throw new LodyOperationStoreError( @@ -2782,7 +2866,7 @@ const startSessionChatOperation = async (args: SessionChatToolInput): Promise ({ ...(args.defaults ?? {}), ...item })); - const invoking = await resolveInvokingTurnContext(manager, requester); + const invoking = await resolveInvokingTurnContext(requester); const roleCatalog = expanded.some((item) => Boolean(item.agentRoleId)) ? await loadWorkspaceAgentRoleCatalog(manager, workspace.id as WorkspaceId) : undefined; @@ -3063,10 +3148,14 @@ const startSessionCreateManyOperation = async ( ctx.sessionId as SessionId, args.operationId, 'session_create_many', - canonicalCommand + canonicalCommand, + invoking.identity.userId, + invoking.identity.sourceTurnId ) ); - if (retry) return await withOperationStore((store) => store.snapshot(retry)); + if (retry) { + return await withOperationStore((store) => store.snapshot(retry)); + } const isMachineOnline = makeMachineOnlineLookupForMcp(manager, ctx); const validatedItems = await mapWithConcurrency( expanded, @@ -3137,7 +3226,7 @@ const startSessionCreateManyOperation = async ( }; } const options = buildMcpCreateOptions(resolved.input, ctx); - bindMcpCreateContext(options, auth, requester); + bindMcpCreateContext(options, invoking.identity, requester); bindAgentRoleCreateOptions(options, resolved.role); try { const effectiveDispatchConfig = await validateSessionCreateOptions({ @@ -3176,13 +3265,14 @@ const startSessionCreateManyOperation = async ( workspaceId: workspace.id as WorkspaceId, ownerMachineId: ctx.machineId as MachineId, requesterSessionId: ctx.sessionId as SessionId, - requesterUserId: auth.userId, + requesterUserId: invoking.identity.userId, operationId: args.operationId, kind: 'session_create_many', canonicalCommand, frozenContinuationConfig: { ...(requester.agentConfigId ? { agentConfigId: requester.agentConfigId } : {}), inputConfig: invoking.frozenInputConfig, + sourceTurnId: invoking.identity.sourceTurnId, targetDispatchConfigs, }, initiatorChainDepth: invoking.chainDepth, @@ -3219,7 +3309,7 @@ const startSessionCreateManyOperation = async ( } try { const options = buildMcpCreateOptions(resolved.input, ctx); - bindMcpCreateContext(options, auth, requester); + bindMcpCreateContext(options, invoking.identity, requester); bindAgentRoleCreateOptions(options, resolved.role); options.sessionId = storedItem.target.sessionId; options.userTurnId = storedItem.target.userTurnId; @@ -3288,6 +3378,7 @@ const startSessionChatManyOperation = async (args: SessionChatManyToolInput): Pr ); } const expanded = args.items.map((item) => ({ ...(args.defaults ?? {}), ...item })); + const invoking = await resolveInvokingTurnContext(requester); const canonicalCommand = { items: expanded, ...(args.deadlineSeconds !== undefined ? { deadlineSeconds: args.deadlineSeconds } : {}), @@ -3297,11 +3388,14 @@ const startSessionChatManyOperation = async (args: SessionChatManyToolInput): Pr ctx.sessionId as SessionId, args.operationId, 'session_chat_many', - canonicalCommand + canonicalCommand, + invoking.identity.userId, + invoking.identity.sourceTurnId ) ); - if (retry) return await withOperationStore((store) => store.snapshot(retry)); - const invoking = await resolveInvokingTurnContext(manager, requester); + if (retry) { + return await withOperationStore((store) => store.snapshot(retry)); + } const isMachineOnline = makeMachineOnlineLookupForMcp(manager, ctx); const initialItems = await mapWithConcurrency( expanded, @@ -3346,7 +3440,7 @@ const startSessionChatManyOperation = async (args: SessionChatManyToolInput): Pr workspace, manager, sessionId: target.id, - requesterUserIdOverride: auth.userId, + delegatedRequester: toDelegatedSessionRequester(invoking.identity), }); } catch (error) { if (error instanceof WorkspaceSyncUnavailableError) { @@ -3365,13 +3459,14 @@ const startSessionChatManyOperation = async (args: SessionChatManyToolInput): Pr workspaceId: workspace.id as WorkspaceId, ownerMachineId: ctx.machineId as MachineId, requesterSessionId: ctx.sessionId as SessionId, - requesterUserId: auth.userId, + requesterUserId: invoking.identity.userId, operationId: args.operationId, kind: 'session_chat_many', canonicalCommand, frozenContinuationConfig: { ...(requester.agentConfigId ? { agentConfigId: requester.agentConfigId } : {}), inputConfig: invoking.frozenInputConfig, + sourceTurnId: invoking.identity.sourceTurnId, }, initiatorChainDepth: invoking.chainDepth, ...timing, @@ -3416,12 +3511,13 @@ const startSessionChatManyOperation = async (args: SessionChatManyToolInput): Pr taskToolsEnabled: invoking.frozenInputConfig.taskToolsEnabled === true, }, undefined, - auth.userId, + undefined, { userTurnId: storedItem.target.userTurnId, chainDepth: invoking.chainDepth + 1, bypassSessionQuota: shouldBypassSessionQuota('session_chat_many'), - } + }, + toDelegatedSessionRequester(invoking.identity) ); await withOperationStore((store) => store.markItemInputDurable( @@ -3915,7 +4011,7 @@ export const __lodyMcpServerInternals = { resolveSessionRenameItems, applySessionRenameItems, persistSessionRenameItems, - resolveInvokingHistoryInput, + buildInvocationIdentity, buildOperationTargetCancelArgs, summarizeProjectRefForMcp, resolveSessionExecutionSnapshot, @@ -4240,11 +4336,7 @@ export function buildLodyMcpServer(config: { taskToolsEnabled?: boolean } = {}): if (!currentSession) { throw new Error(`Session not found: ${ctx.sessionId}`); } - // Only a Role needs the driving Turn here, for the task-tool flag. - // This legacy path does not persist a continuation config. - const invoking = args.agentRoleId - ? await resolveInvokingTurnContext(manager, currentSession) - : undefined; + const invoking = await resolveInvokingTurnContext(currentSession); const roleCatalog = args.agentRoleId ? await loadWorkspaceAgentRoleCatalog(manager, workspace.id as WorkspaceId) : undefined; @@ -4255,7 +4347,7 @@ export function buildLodyMcpServer(config: { taskToolsEnabled?: boolean } = {}): args.agentRoleId ? roleCatalog?.get(args.agentRoleId) : undefined ); const options = buildMcpCreateOptions(resolved.input, ctx); - bindMcpCreateContext(options, auth, currentSession); + bindMcpCreateContext(options, invoking.identity, currentSession); bindAgentRoleCreateOptions(options, resolved.role); options.workspaceMetaPrewriteSatisfied = true; const result = await createSessionResult( @@ -4337,6 +4429,7 @@ export function buildLodyMcpServer(config: { taskToolsEnabled?: boolean } = {}): throw new Error(`Session not found: ${sessionId}`); } assertDifferentMcpSession(currentSession, targetSession); + const invoking = await resolveInvokingTurnContext(currentSession); const result = await sendSessionChatResult( auth, workspace, @@ -4345,7 +4438,9 @@ export function buildLodyMcpServer(config: { taskToolsEnabled?: boolean } = {}): args.prompt, resolveTurnDispatchConfig({}), buildStructuredOutputOptions(args), - auth.userId + undefined, + undefined, + toDelegatedSessionRequester(invoking.identity) ); const response = { ok: true, diff --git a/apps/cli/src/orchestration/AGENTS.md b/apps/cli/src/orchestration/AGENTS.md index 786eda812..7850b8887 100644 --- a/apps/cli/src/orchestration/AGENTS.md +++ b/apps/cli/src/orchestration/AGENTS.md @@ -25,6 +25,13 @@ Root and `apps/cli/AGENTS.md` apply. Normative behavior lives in - Create Operations freeze each target's effective dispatch config at acceptance; recovery must not re-read mutable requester history defaults. Full content stays in the target Session history. +- Accepted Operations store the invoking user once as `requesterUserId` and bind the exact source + Turn as `sourceTurnId`. `requesterSessionId` already identifies the source Session. Recovery + routes attribution and member-scoped authorization through that frozen user while the current + owner Machine credential remains the executor credential. Completion system Turns retain the + same userId so a continuation cannot silently switch identities. The Operation store matches + requester user and source Turn together with kind and command fingerprint; a later Turn reusing + the id is `OPERATION_ID_REUSED`, not a retry. - `operation-coordinator.ts` is owned only by the local Host-lease Worker. MCP subprocesses may accept Operations but never schedule completion Turns. - Reconciliation is level-checked. Loro subscriptions and SQLite directory diff --git a/apps/cli/src/orchestration/operation-coordinator.test.ts b/apps/cli/src/orchestration/operation-coordinator.test.ts index 5811a626d..6e733657c 100644 --- a/apps/cli/src/orchestration/operation-coordinator.test.ts +++ b/apps/cli/src/orchestration/operation-coordinator.test.ts @@ -903,6 +903,7 @@ describe('LodyOperationCoordinator', () => { expect(harness.histories.get(harness.requesterSessionId)).toEqual([ expect.objectContaining({ role: 'system', + userId: 'user-1', items: [ expect.objectContaining({ type: 'operation_completion', diff --git a/apps/cli/src/orchestration/operation-coordinator.ts b/apps/cli/src/orchestration/operation-coordinator.ts index f069af2ee..200acc715 100644 --- a/apps/cli/src/orchestration/operation-coordinator.ts +++ b/apps/cli/src/orchestration/operation-coordinator.ts @@ -1006,6 +1006,7 @@ export class LodyOperationCoordinator { const turn: SessionHistoryInput = { id: delivery.systemTurnId, role: 'system', + userId: operation.requesterUserId, timestamp: new Date(this.now()).toISOString(), items: [item], fileDiff: [], diff --git a/apps/cli/src/orchestration/operation-store.test.ts b/apps/cli/src/orchestration/operation-store.test.ts index 2575012d2..4755933fb 100644 --- a/apps/cli/src/orchestration/operation-store.test.ts +++ b/apps/cli/src/orchestration/operation-store.test.ts @@ -39,6 +39,7 @@ const baseInput = () => ({ frozenContinuationConfig: { agentConfigId: 'agent-1', inputConfig: { cliType: 'builtin' as const, agentType: 'codex', chainDepth: 0 }, + sourceTurnId: 'source-turn-1', }, initiatorChainDepth: 0, createdAt: '2026-07-20T00:00:00.000Z', @@ -84,6 +85,7 @@ describe('LodyOperationStore', () => { expect(first.created).toBe(true); expect(retry.created).toBe(false); expect(retry.operation.operationId).toBe('review-round-1'); + expect(retry.operation.frozenContinuationConfig.sourceTurnId).toBe('source-turn-1'); expect(store.snapshot(retry.operation)).toMatchObject({ state: 'active' }); } finally { store.close(); @@ -140,6 +142,50 @@ describe('LodyOperationStore', () => { } }); + it('binds accept and retry lookup to requester and source Turn identity', async () => { + const store = await makeStore(); + try { + const input = baseInput(); + store.accept(input); + + expect( + store.findMatchingRetry( + input.requesterSessionId, + input.operationId, + input.kind, + input.canonicalCommand, + input.requesterUserId, + input.frozenContinuationConfig.sourceTurnId + ) + ).toMatchObject({ operationId: input.operationId }); + expect(() => + store.findMatchingRetry( + input.requesterSessionId, + input.operationId, + input.kind, + input.canonicalCommand, + 'user-2', + input.frozenContinuationConfig.sourceTurnId + ) + ).toThrowError( + expect.objectContaining({ code: 'OPERATION_ID_REUSED' }) + ); + expect(() => + store.accept({ + ...input, + frozenContinuationConfig: { + ...input.frozenContinuationConfig, + sourceTurnId: 'source-turn-2', + }, + }) + ).toThrowError( + expect.objectContaining({ code: 'OPERATION_ID_REUSED' }) + ); + } finally { + store.close(); + } + }); + it('scopes lookup to the requester Session', async () => { const store = await makeStore(); try { diff --git a/apps/cli/src/orchestration/operation-store.ts b/apps/cli/src/orchestration/operation-store.ts index 08a973965..e5c4651c0 100644 --- a/apps/cli/src/orchestration/operation-store.ts +++ b/apps/cli/src/orchestration/operation-store.ts @@ -112,6 +112,7 @@ const FrozenConfigSchema = z .object({ agentConfigId: z.string().optional(), inputConfig: z.record(z.string(), z.unknown()), + sourceTurnId: z.string().trim().min(1).optional(), targetDispatchConfigs: z .array( z @@ -438,7 +439,13 @@ export class LodyOperationStore { if (existing) return { created: false, - operation: this.assertMatching(existing, input.kind, fingerprint), + operation: this.assertMatching( + existing, + input.kind, + fingerprint, + input.requesterUserId, + frozenConfig.sourceTurnId + ), claimedItemIndexes: [], }; const inserted = this.db @@ -488,7 +495,13 @@ export class LodyOperationStore { } return { created: inserted.changes === 1, - operation: this.assertMatching(operation, input.kind, fingerprint), + operation: this.assertMatching( + operation, + input.kind, + fingerprint, + input.requesterUserId, + frozenConfig.sourceTurnId + ), claimedItemIndexes, }; } @@ -519,12 +532,14 @@ export class LodyOperationStore { requesterSessionId: SessionId, operationId: string, kind: LodyOperationKind, - canonicalCommand: unknown + canonicalCommand: unknown, + requesterUserId: string, + sourceTurnId?: string ): StoredLodyOperation | undefined { const existing = this.getStored(requesterSessionId, operationId); if (!existing) return undefined; const fingerprint = fingerprintLodyCommand(kind, canonicalizeLodyCommand(canonicalCommand)); - return this.assertMatching(existing, kind, fingerprint); + return this.assertMatching(existing, kind, fingerprint, requesterUserId, sourceTurnId); } listActive(workspaceId: WorkspaceId, ownerMachineId: MachineId): StoredLodyOperation[] { @@ -791,9 +806,16 @@ export class LodyOperationStore { private assertMatching( operation: StoredLodyOperation, kind: LodyOperationKind, - fingerprint: string + fingerprint: string, + requesterUserId: string, + sourceTurnId?: string ): StoredLodyOperation { - if (operation.kind !== kind || operation.fingerprint !== fingerprint) { + if ( + operation.kind !== kind || + operation.fingerprint !== fingerprint || + operation.requesterUserId !== requesterUserId || + operation.frozenContinuationConfig.sourceTurnId !== sourceTurnId + ) { throw new LodyOperationStoreError( 'OPERATION_ID_REUSED', `Operation id ${operation.operationId} is already bound to different input.`, diff --git a/apps/cli/src/session/AGENTS.md b/apps/cli/src/session/AGENTS.md index 370195ba4..2df397a72 100644 --- a/apps/cli/src/session/AGENTS.md +++ b/apps/cli/src/session/AGENTS.md @@ -9,13 +9,21 @@ message bus. The WS/DO path is DEPRECATED. Session CLI/MCP orchestration contract: specs/session-orchestration.md. Target-machine authorization is checked by the injected access capability with the source CLI token, -which derives the requester identity at the trusted boundary. Do not send a caller-supplied requester -through workspace Machine RPC: that transport does not authenticate member identity. +which derives the requester identity for ordinary CLI calls and verifies the frozen Turn +requester for MCP delegation. Session command boundaries receive an optional delegated requester; +its presence selects delegated access, while source-Turn provenance stays at the MCP/Operation +boundary. Do not send an untrusted requester through +workspace Machine RPC: that transport does not authenticate member identity. Live status is a target-daemon Machine RPC read, and durable session metadata is not a live-presence substitute. -Session orchestration MCP intentionally runs with the daemon owner's CLI credential, -including for teammate-started Sessions on a shared machine. Do not add requester -delegation proofs or a shared-machine gate without a new product and security decision. +Session orchestration MCP authenticates execution with the daemon owner's CLI credential, +but derives the human identity causally from the active dispatch/execution runtime. Persisted +history must never reconstruct a missing invocation; fail closed when no active runtime exists. Freeze that identity +and source Turn into durable Operations; the Operation already identifies the source +Session and stores the invoking user once as `requesterUserId`. Retries and recovery must not +reread mutable history. Machine and Provider credentials remain execution-host scoped, while +Session/Turn attribution, member authorization, GitHub access, and downstream Git identity use +the frozen identity. Never fall back to the Session owner when the driving Turn has no userId. - `session-dispatch-watcher.ts` — the current dispatch entry: watches `repo.watch('doc-metadata')` + per-session mirror subscribe; dispatches when @@ -111,7 +119,9 @@ delegation proofs or a shared-machine gate without a new product and security de quiescent. Use live execution/presence for current-work signals; goal activity may still protect history rewrites or an in-memory runtime that can resume autonomously. It is the per-session execution mutex: never mint a second visible turn while a - `TurnRuntimeState` is registered. User-dispatch turns derive assistant entry ids + `TurnRuntimeState` is registered. Its optional `invocation` atomically owns source Turn, + requester, and input config; steer replaces that object before tool execution can continue. + User-dispatch turns derive assistant entry ids from `userTurnId` (`assistant:`), so a retried/recovered dispatch reuses the same history entry. INVARIANT: a steer (guide) the agent never accepted must not stay parked in diff --git a/apps/cli/src/session/session-dispatch-watcher.ts b/apps/cli/src/session/session-dispatch-watcher.ts index 764f5ceda..52d3757f0 100644 --- a/apps/cli/src/session/session-dispatch-watcher.ts +++ b/apps/cli/src/session/session-dispatch-watcher.ts @@ -1574,6 +1574,11 @@ export class SessionDispatchWatcher { sessionId, sessionDoc, userTurnId: nextUserTurn.id, + invocation: { + sourceTurnId: nextUserTurn.id, + requesterUserId: nextUserTurn.userId, + inputConfig: nextUserTurn.inputConfig ?? {}, + }, dispatchSource, accessPromise: executionAccessPromise, requestPromise, diff --git a/apps/cli/src/session/session-execution-service.ts b/apps/cli/src/session/session-execution-service.ts index 4ca91a1d9..9c1f840c6 100644 --- a/apps/cli/src/session/session-execution-service.ts +++ b/apps/cli/src/session/session-execution-service.ts @@ -212,12 +212,19 @@ type PromptHandoffRun = { signalSuccessor: () => void; }; +type TurnInvocation = { + /** Causal input Turn for authorization and durable provenance. */ + sourceTurnId: string; + requesterUserId?: string; + inputConfig: SessionTurnInputConfig; +}; + type TurnRuntimeState = { sessionId: SessionId; /** Logical chain tail exposed to Web, cancel, and optimistic steer validation. */ turnId: string; userTurnId?: string; - requesterUserId?: string; + invocation?: TurnInvocation; session?: ISession; project?: ProjectRef; baseCommitHash?: string | null; @@ -309,6 +316,7 @@ type VisibleSessionTurnOptions = { sessionDoc: SessionDocument; session?: ISession; userTurnId?: string; + invocation?: TurnInvocation; /** * How the turn payload reached this machine. 'rpc' turns can start before the * user's history entry syncs locally, so their turn-scoped history writes go @@ -346,6 +354,7 @@ export type PreparedSessionDispatchOptions = { sessionId: SessionId; sessionDoc: SessionDocument; userTurnId: string; + invocation: TurnInvocation; dispatchSource: SessionDispatchSource; accessPromise: Promise; requestPromise: Promise; @@ -1273,6 +1282,14 @@ export class SessionExecutionService { return reject('stale-turn', 'Steer application arrived after ownership changed'); } + // The provider has accepted this steer and may execute tools before + // history/finalization catches up. Switch causal identity first. + runtime.invocation = { + sourceTurnId: options.userTurnId, + requesterUserId: options.userId, + inputConfig: options.inputConfig, + }; + try { await this.finalizeYieldedTurnOutput(runtime, options.sessionId, previousTurnId); } catch (error) { @@ -1325,7 +1342,6 @@ export class SessionExecutionService { runtime.activePromptRun = nextPromptRun; runtime.turnId = nextTurnId; runtime.userTurnId = options.userTurnId; - runtime.requesterUserId = options.userId; this.markCurrentTurn(options.sessionId, nextTurnId); ownedPromptRun.signalSuccessor(); return { @@ -1458,6 +1474,7 @@ export class SessionExecutionService { sessionId, sessionDoc: options.sessionDoc, userTurnId, + invocation: options.invocation, dispatchSource, unhandledErrorCode: 'session_chat_failed', describeUnhandledError: (error) => @@ -1529,16 +1546,17 @@ export class SessionExecutionService { } private createTurnRuntime( - sessionId: SessionId, - turnId: string, - userTurnId?: string, - session?: ISession + options: Pick< + VisibleSessionTurnOptions, + 'sessionId' | 'session' | 'userTurnId' | 'invocation' + > & { turnId: string } ): TurnRuntimeState { return { - sessionId, - turnId, - userTurnId, - session, + sessionId: options.sessionId, + turnId: options.turnId, + userTurnId: options.userTurnId, + invocation: options.invocation, + session: options.session, promptStarted: false, promptInFlight: false, autoPromptInFlight: false, @@ -2529,7 +2547,7 @@ export class SessionExecutionService { options: VisibleSessionTurnOptions, body: (ctx: VisibleSessionTurnContext) => Effect.Effect ): Promise { - const { sessionId, sessionDoc, session, userTurnId } = options; + const { sessionId, sessionDoc, userTurnId } = options; const span = startTraceSpan(this.deps.logger, 'execution.visible_turn', { sessionId, ...(userTurnId ? { userTurnId } : {}), @@ -2564,7 +2582,7 @@ export class SessionExecutionService { deferACPUpdateTarget: true, }); this.markCurrentTurn(sessionId, turnId); - runtime = this.createTurnRuntime(sessionId, turnId, userTurnId, session); + runtime = this.createTurnRuntime({ ...options, turnId }); this.registerTurnRuntime(runtime); } finally { releaseConflict(); @@ -2986,6 +3004,28 @@ export class SessionExecutionService { return this.turnRuntimeBySession.get(sessionId)?.userTurnId; } + getActiveInvocationContext(sessionId: SessionId): + | { + requesterUserId: string; + sourceTurnId: string; + inputConfig: SessionTurnInputConfig; + } + | undefined { + const runtime = this.turnRuntimeBySession.get(sessionId); + if (!runtime) { + return undefined; + } + const { invocation } = runtime; + if (!invocation?.requesterUserId) { + throw new Error(`Active invocation identity is unavailable for session ${sessionId}`); + } + return { + requesterUserId: invocation.requesterUserId, + sourceTurnId: invocation.sourceTurnId, + inputConfig: invocation.inputConfig, + }; + } + private async setDispatchProcessing( sessionId: SessionId, sessionDoc: SessionDocument, @@ -3605,7 +3645,6 @@ export class SessionExecutionService { ): Effect.Effect => Effect.gen(function* () { const { turnId, runtime, abortIfCancelled, openAssistantEntry, prompt } = ctx; - runtime.requesterUserId = message.userId; let activeSession = readySession; let staleAcpPromptRecoveryAttempted = false; let baseCommitHash: string | null = null; @@ -3873,7 +3912,7 @@ export class SessionExecutionService { const completedTurnId = runtime.turnId; const completedUserTurnId = runtime.userTurnId ?? executionUserTurnId; - const completedRequesterUserId = runtime.requesterUserId ?? userId; + const completedRequesterUserId = runtime.invocation?.requesterUserId ?? userId; // Read before finalization clears the turn's ACP update state. const producedOutput = self.turnProducedVisibleOutput(sessionId, completedTurnId); @@ -4002,6 +4041,11 @@ export class SessionExecutionService { sessionDoc, ...(session ? { session } : {}), userTurnId: executionUserTurnId, + invocation: { + sourceTurnId: userTurnId, + requesterUserId: userId, + inputConfig: acpSessionConfig, + }, ...(dispatchOptions?.dispatchSource ? { dispatchSource: dispatchOptions.dispatchSource } : {}), @@ -4322,6 +4366,15 @@ export class SessionExecutionService { sessionId, sessionDoc, userTurnId, + ...(userTurnId + ? { + invocation: { + sourceTurnId: userTurnId, + requesterUserId: message.userId, + inputConfig: acpSessionConfig, + }, + } + : {}), ...(dispatchOptions?.dispatchSource ? { dispatchSource: dispatchOptions.dispatchSource } : {}), @@ -4343,7 +4396,6 @@ export class SessionExecutionService { }) => Effect.gen(function* () { setUnhandledErrorContext(turnErrorContext); - runtime.requesterUserId = message.userId; const memoryPressureResult = yield* self.tryPromise(() => self.evictForTurnStart(sessionId) ); @@ -4545,7 +4597,8 @@ export class SessionExecutionService { const completedTurnId = runtime.turnId; const completedUserTurnId = runtime.userTurnId ?? userTurnId; - const completedRequesterUserId = runtime.requesterUserId ?? sessionConfig.requesterUserId; + const completedRequesterUserId = + runtime.invocation?.requesterUserId ?? sessionConfig.requesterUserId; // Read before finalization clears the turn's ACP update state. const producedOutput = self.turnProducedVisibleOutput(sessionId, completedTurnId); diff --git a/apps/cli/tests/session-execution-service.test.ts b/apps/cli/tests/session-execution-service.test.ts index c7dd6932f..20170acd4 100644 --- a/apps/cli/tests/session-execution-service.test.ts +++ b/apps/cli/tests/session-execution-service.test.ts @@ -270,7 +270,11 @@ describe('SessionExecutionService', () => { userTurnId: 'user-1', session: activeSession, promptInFlight: true, - requesterUserId: 'user-1', + invocation: { + sourceTurnId: 'user-1', + requesterUserId: 'user-1', + inputConfig: { prompt: 'initial prompt' }, + }, activePromptRun: initialPromptRun, yieldedFinalization: Promise.resolve(), }; @@ -306,6 +310,16 @@ describe('SessionExecutionService', () => { ); expect(runtime.turnId).toBe('assistant:user-2'); expect(runtime.userTurnId).toBe('user-2'); + expect(runtime.invocation).toEqual({ + requesterUserId: 'user-1', + sourceTurnId: 'user-2', + inputConfig: { prompt: 'change direction' }, + }); + expect(service.getActiveInvocationContext(sessionId)).toEqual({ + requesterUserId: 'user-1', + sourceTurnId: 'user-2', + inputConfig: { prompt: 'change direction' }, + }); expect(initialPromptRun.successor?.turnId).toBe('assistant:user-2'); expect(runtime.activePromptRun.turnId).toBe('assistant:user-2'); @@ -1343,7 +1357,7 @@ describe('SessionExecutionService', () => { ); }); - it('starts active presence before prepared dispatch awaits machine access', async () => { + it('exposes RPC invocation identity before prepared dispatch awaits machine access', async () => { let resolveAccess!: (value: { outcome: 'indeterminate'; cause: 'network'; @@ -1371,7 +1385,12 @@ describe('SessionExecutionService', () => { sessionId: 'session-prepared-presence' as SessionId, sessionDoc, userTurnId: 'turn-prepared-presence', - dispatchSource: 'crdt', + invocation: { + sourceTurnId: 'turn-prepared-presence', + requesterUserId: 'user-b', + inputConfig: { prompt: 'fast path prompt', taskToolsEnabled: true }, + }, + dispatchSource: 'rpc', accessPromise, requestPromise: new Promise(() => {}), onAccessAllowed, @@ -1386,12 +1405,17 @@ describe('SessionExecutionService', () => { expect(deps.beginConversationTurn).toHaveBeenCalledWith( 'session-prepared-presence', 'turn-prepared-presence', - { dispatchSource: 'crdt', sessionDoc, deferACPUpdateTarget: true } + { dispatchSource: 'rpc', sessionDoc, deferACPUpdateTarget: true } ); expect(service.getExecutionSnapshot('session-prepared-presence' as SessionId)).toMatchObject({ activeTurnId: 'assistant:turn-prepared-presence', hasActiveTurn: true, }); + expect(service.getActiveInvocationContext('session-prepared-presence' as SessionId)).toEqual({ + requesterUserId: 'user-b', + sourceTurnId: 'turn-prepared-presence', + inputConfig: { prompt: 'fast path prompt', taskToolsEnabled: true }, + }); expect(onAccessAllowed).not.toHaveBeenCalled(); resolveAccess({ outcome: 'indeterminate', cause: 'network', error: 'offline' }); @@ -1447,6 +1471,7 @@ describe('SessionExecutionService', () => { sessionId, sessionDoc: sessionDoc as never, userTurnId, + invocation: { sourceTurnId: userTurnId, inputConfig: {} }, dispatchSource: 'crdt', accessPromise: new Promise(() => {}), requestPromise: new Promise(() => {}), @@ -1459,6 +1484,9 @@ describe('SessionExecutionService', () => { activeTurnId: turnId, hasActiveTurn: true, }); + expect(() => service.getActiveInvocationContext(sessionId)).toThrow( + 'Active invocation identity is unavailable' + ); await expect( service.cancelSession({ @@ -1559,6 +1587,7 @@ describe('SessionExecutionService', () => { sessionId, sessionDoc: preparedSessionDoc as never, userTurnId, + invocation: { sourceTurnId: userTurnId, inputConfig: {} }, dispatchSource: 'rpc', accessPromise: Promise.resolve({ outcome: 'allowed' as const }), requestPromise: new Promise(() => {}), diff --git a/packages/shared/src/local-machine-rpc.ts b/packages/shared/src/local-machine-rpc.ts index 512028d7c..98264f1b8 100644 --- a/packages/shared/src/local-machine-rpc.ts +++ b/packages/shared/src/local-machine-rpc.ts @@ -50,7 +50,38 @@ const BaseLocalMachineRpcRequestSchema = z }) .strict(); +export const SessionActiveInvocationContextResultSchema = z.discriminatedUnion('active', [ + z + .object({ + type: z.literal('session/active-invocation-context'), + sessionId: SessionIdSchema, + active: z.literal(false), + }) + .strict(), + z + .object({ + type: z.literal('session/active-invocation-context'), + sessionId: SessionIdSchema, + active: z.literal(true), + requesterUserId: z.string().trim().min(1), + sourceTurnId: z.string().trim().min(1), + inputConfig: z.record(z.string(), z.unknown()), + }) + .strict(), +]); +export type SessionActiveInvocationContextResult = z.infer< + typeof SessionActiveInvocationContextResultSchema +>; + export const LocalMachineRpcRequestSchema = z.discriminatedUnion('method', [ + BaseLocalMachineRpcRequestSchema.extend({ + method: z.literal('session/get-active-invocation-context'), + params: z + .object({ + sessionId: SessionIdSchema, + }) + .strict(), + }).strict(), BaseLocalMachineRpcRequestSchema.extend({ method: z.literal('code-collab/get-file-index'), params: CodeCollabV2FileIndexRequestSchema, @@ -204,6 +235,7 @@ export type LocalMachineRpcRequest = z.infer { it.each([ + { + method: 'session/get-active-invocation-context', + params: { sessionId: 'session-1' }, + }, { method: 'session/fork', params: { @@ -230,4 +234,31 @@ describe('local Machine RPC', () => { }; expect(LocalMachineRpcResponseSchema.safeParse(endpointWithoutTarget).success).toBe(false); }); + + it('validates active invocation identity and its frozen input config', () => { + expect( + LocalMachineRpcResponseSchema.safeParse({ + ok: true, + result: { + type: 'session/active-invocation-context', + sessionId: 'session-1', + active: true, + requesterUserId: 'user-b', + sourceTurnId: 'turn-b', + inputConfig: { chainDepth: 1, taskToolsEnabled: true }, + }, + }).success + ).toBe(true); + expect( + LocalMachineRpcResponseSchema.safeParse({ + ok: true, + result: { + type: 'session/active-invocation-context', + sessionId: 'session-1', + active: true, + requesterUserId: 'user-b', + }, + }).success + ).toBe(false); + }); });