| import { randomBytes, randomUUID } from 'node:crypto'; |
| import { chmod, mkdir, readFile, writeFile } from 'node:fs/promises'; |
| import { dirname, join } from 'node:path'; |
| import { |
| AiSdkBackend, |
| AgentGraphCoordinator, |
| AgentGraphSupervisorWakeCoordinator, |
| AutomationManager, |
| AutomationScheduler, |
| BackendRegistry, |
| GoalManager, |
| RuntimeReadModel, |
| SessionManager, |
| ShellRunProcessManager, |
| buildAutomationTool, |
| buildAskUserQuestionTool, |
| buildRequestSandboxBoundaryTool, |
| buildBuiltinTools, |
| buildSessionRecapMessages, |
| buildChildAgentTools, |
| createBuiltinSandboxManager, |
| isBuiltinFilesystemWorkerSandboxAvailable, |
| createSandboxDiagnosticsProvider, |
| createProviderRequestCaptureRecorder, |
| createFilesystemWorkerLaunchSpecProvider, |
| createLocalContinuationSafetyInspector, |
| createConfiguredSubagentCatalog, |
| drainGoalTurn, |
| FilesystemWorkerClient, |
| buildDefaultContextBudgetPolicy, |
| buildSkillAgentTool, |
| buildSkillSearchAgentTool, |
| SkillShadowSelectionTracker, |
| buildGoalTools, |
| buildParentAgentTools, |
| assertProductBindingCatalogClean, |
| AGENT_TOOL_GROUP_ID, |
| buildLlmHistorySummarizer, |
| buildNativeWebSearchTool, |
| cleanupLegacyHistoryCompactArtifacts, |
| buildProviderOptions, |
| buildSubscriptionModelFetch, |
| evaluateAutomationCanFire, |
| getAIModel, |
| generateSessionTitle as generateRuntimeSessionTitle, |
| loadHistoryCompactBlocksFromArtifacts, |
| listRunnableBuiltinAgentDefinitions, |
| replayPlanItemsToModelMessages, |
| recoverAgentGraphSupervisorContextOverflow, |
| renderAgentSwarmSupervisorWake, |
| resolveSkillDiscoveryPaths, |
| resolveSelectedModelContextWindow, |
| projectEffectiveProductToolSurface, |
| routeWebSearchTools, |
| shouldWakeAgentSwarmSupervisor, |
| type AutomationDefinition, |
| type EffectiveProductToolSurface, |
| type HostCapabilitiesResolver, |
| type MakaTool, |
| type InvocationResult, |
| type InvocationSource, |
| type ShellRunUpdate, |
| type ModelMessage, |
| } from '@maka/runtime'; |
| import { |
| createSqliteAgentRunStore, |
| createAttachmentByteReader, |
| createSqliteArtifactStore, |
| createAutomationStore, |
| createConnectionStore, |
| createFileCredentialStore, |
| openRuntimeEventPersistence, |
| createForeignSessionStore, |
| createGitWorktreeChildExecutor, |
| createReadImageSnapshotter, |
| createSessionStore, |
| isSessionNotFoundError, |
| createSettingsStore, |
| createSqliteShellRunStore, |
| assertSessionBundleRootLayout, |
| type ForeignSessionStore, |
| persistProviderRequestCaptureArtifact, |
| } from '@maka/storage'; |
| import { createAgentGraphControlStore } from '@maka/storage/agent-graph-control-store'; |
| import { resolveStorageRoot } from '@maka/storage/root-authority'; |
| import { resolveWorkspaceIdentity } from '@maka/storage/workspace-identity'; |
| import { fetchProviderModels } from '@maka/runtime'; |
| import { createApiKeyOnboardingSurface } from './onboarding.js'; |
| import { relayModelProfile, isActiveShellRunStatus, resolveModelVisionSupport } from '@maka/core'; |
| import type { ReadySessionTarget } from './connection-target.js'; |
| import { |
| listReadyModelChoices, |
| resolveDefaultSessionTarget, |
| resolveSessionTargetForSlug, |
| } from './connection-target.js'; |
| import { buildCliSystemPrompt, buildCliTurnTailPrompt } from './cli-system-prompt.js'; |
| import { CliGoalContinuation } from './cli-goal-continuation.js'; |
| import type { |
| MakaCliSkillSurface, |
| MakaOnboardingSurface, |
| ModelChoice, |
| SessionRecapGenerator, |
| } from './pi-tui-contracts.js'; |
| import { cleanRecapText } from './session-recap.js'; |
| |
| export interface MakaCliRuntimeContext { |
| /** Legacy shared-root input retained for callers that do not split roots. */ |
| workspaceRoot: string; |
| /** Durable session-owned state root used by the runtime stores. */ |
| stateRoot: string; |
| /** Host-injected configuration root used by connections, credentials, and settings. */ |
| configRoot: string; |
| cwd: string; |
| runtime: SessionManager; |
| target: ReadySessionTarget; |
| /** Selectable models across every ready connection, for the `/model` picker. */ |
| modelChoices: ModelChoice[]; |
| /** Tools passed to the backend, including TUI-only interactive and subagent tools. */ |
| tools: MakaTool[]; |
| /** |
| * Explicit skill invocation surface (issue #1148): the discovery source + |
| * host gate shared with the Skill tool and the system-prompt catalog, so |
| * `/skill:<name>` autocomplete, highlight, and submit-time injection all |
| * resolve against exactly what the host can load. |
| */ |
| skills: MakaCliSkillSurface; |
| automationManager: AutomationManager; |
| automationScheduler: AutomationScheduler; |
| subscribeShellRunUpdates(listener: (update: ShellRunUpdate) => void): () => void; |
| listShellRunUpdates(sessionId: string): Promise<ShellRunUpdate[]>; |
| goalManager: GoalManager; |
| goalContinuation: CliGoalContinuation; |
| /** One-sentence session recap generator (issue #1055), shared by `/recap` and idle-return auto-recap. */ |
| recap: SessionRecapGenerator; |
| /** Read-only scanner for other agents' sessions (Claude Code, Codex), for the resume picker (#1057). */ |
| foreignSessions: ForeignSessionStore; |
| /** Host-owned Graph runtime used by the TUI and `maka run --graph`. */ |
| agentGraph?: { |
| reserveActivity(sessionId: string): { release(): void }; |
| waitForCompletion(sessionId: string): Promise<void>; |
| }; |
| close(): Promise<void>; |
| /** API-key onboarding surface for the /setup wizard (#1098). */ |
| onboarding: MakaOnboardingSurface; |
| } |
| |
| export function resolveCliStreamConnectTimeoutMs( |
| env: NodeJS.ProcessEnv = process.env, |
| ): number | undefined { |
| const raw = env.MAKA_STREAM_CONNECT_TIMEOUT_MS; |
| if (raw === undefined || raw.trim() === '') return undefined; |
| if (!/^[1-9]\d*$/.test(raw)) { |
| throw new Error('MAKA_STREAM_CONNECT_TIMEOUT_MS must be a positive integer'); |
| } |
| const timeoutMs = Number(raw); |
| if (!Number.isSafeInteger(timeoutMs)) { |
| throw new Error('MAKA_STREAM_CONNECT_TIMEOUT_MS must be a positive safe integer'); |
| } |
| return timeoutMs; |
| } |
| |
| /** |
| * Generates a one-sentence recap of a session so far, using a tool-free model |
| * call whose exchange is never written to the session's own history. Never |
| * throws — failures resolve to `{ ok: false }` so callers can surface them |
| * without a try/catch. |
| */ |
| export interface CreateMakaCliRuntimeContextInput { |
| surface: 'tui' | 'run' | 'activation'; |
| /** Legacy root; both new roots default to this path. */ |
| workspaceRoot: string; |
| /** Optional portable session-owned state root. */ |
| stateRoot?: string; |
| /** Optional host-injected configuration root. */ |
| configRoot?: string; |
| cwd: string; |
| requestedConnectionSlug?: string; |
| requestedModel?: string; |
| maxSteps?: number; |
| /** Compose the durable Graph control plane and Git-worktree child executor. */ |
| enableAgentGraph?: boolean; |
| /** Canonical cwd used for one resumed session without rewriting its stored header. */ |
| sessionCwdOverride?: { sessionId: string; cwd: string }; |
| runtimeInvocationObserver?: (result: InvocationResult) => void | Promise<void>; |
| /** Invocation provenance used by hosted activation callers. */ |
| runtimeSource?: InvocationSource; |
| /** Enables authoritative safe-boundary continuation for hosted activation callers. */ |
| safeBoundaryResumeEnabled?: boolean; |
| onSessionTitleChanged?: (sessionId: string) => void; |
| /** |
| * Optional cron executor. When provided, the Automation tool advertises the |
| * cron kind and cron fires spawn a fresh session + run via this callback |
| * (reviewer G1: a host derives cron support from the executor it passes in). |
| * Omitted by the default CLI (no multi-session surface) — heartbeat only. |
| */ |
| automationCreateFreshRun?: ( |
| prompt: string, |
| automationId: string, |
| ) => Promise<import('@maka/runtime').AutomationFireResult>; |
| } |
| |
| export interface GetOrCreateCliClaudeDeviceIdDeps { |
| newId?: () => string; |
| } |
| |
| export function isMakaClaudeSubscriptionCloakEnabled( |
| env: { MAKA_CLAUDE_SUBSCRIPTION_CLOAK?: string } = process.env, |
| ): boolean { |
| return env.MAKA_CLAUDE_SUBSCRIPTION_CLOAK !== '0'; |
| } |
| |
| export async function createMakaCliRuntimeContext( |
| input: CreateMakaCliRuntimeContextInput, |
| ): Promise<MakaCliRuntimeContext> { |
| const stateRoot = input.stateRoot ?? input.workspaceRoot; |
| const configRoot = input.configRoot ?? input.workspaceRoot; |
| const agentGraphEnabled = input.surface === 'tui' || input.enableAgentGraph === true; |
| if (input.stateRoot !== undefined || input.configRoot !== undefined) { |
| await assertSessionBundleRootLayout({ |
| stateRoot, |
| configRoot, |
| allowShared: false, |
| }); |
| } |
| await resolveStorageRoot({ path: stateRoot, kind: 'interactive' }); |
| const store = createSessionStore(stateRoot); |
| const runStore = createSqliteAgentRunStore(stateRoot); |
| const runtimePersistence = await openRuntimeEventPersistence({ |
| workspaceRoot: stateRoot, |
| }); |
| const runtimeEventStore = runtimePersistence.runtimeEventStore; |
| const shellRunStore = createSqliteShellRunStore(stateRoot); |
| await Promise.all([runStore.ready?.(), shellRunStore.ready()]).catch(async (error) => { |
| await store.close?.().catch(() => {}); |
| runtimePersistence.close(); |
| runStore.close?.(); |
| shellRunStore.close(); |
| throw error; |
| }); |
| const artifactStore = createSqliteArtifactStore(stateRoot); |
| const agentGraphControlStore = agentGraphEnabled |
| ? createAgentGraphControlStore(stateRoot) |
| : undefined; |
| const worktreeChildExecutor = agentGraphEnabled |
| ? createGitWorktreeChildExecutor({ storageRoot: stateRoot }) |
| : undefined; |
| const agentGraphErrors = new Map<string, unknown>(); |
| let agentGraphCoordinator: AgentGraphCoordinator | undefined; |
| let agentGraphSupervisorWakeCoordinator: AgentGraphSupervisorWakeCoordinator | undefined; |
| const connectionStore = createConnectionStore(configRoot); |
| const credentialStore = createFileCredentialStore(configRoot); |
| const settingsStore = createSettingsStore(configRoot); |
| const subagentCatalog = createConfiguredSubagentCatalog({ |
| getSettings: () => settingsStore.get(), |
| getConnection: (slug) => connectionStore.get(slug), |
| }); |
| // Read-only scanner over other agents' local session stores (~/.claude, |
| // ~/.codex). Independent of the Maka workspace — takes no workspaceRoot. |
| const foreignSessions = createForeignSessionStore(); |
| // Authoritative RuntimeEvent read model (issue #1055's session-recap |
| // generator projects through this instead of re-deriving its own lossy |
| // StoredMessage-based projection). Built once and shared — mirrors |
| // SessionManager's own construction in session-manager.ts's readModel(). |
| const runtimeReadModel = new RuntimeReadModel({ |
| runStore, |
| runtimeEventStore, |
| projectionCache: store, |
| }); |
| const targetInput = { |
| connectionStore, |
| credentialStore, |
| requestedModel: input.requestedModel, |
| }; |
| const target = input.requestedConnectionSlug |
| ? await resolveSessionTargetForSlug(input.requestedConnectionSlug, targetInput) |
| : await resolveDefaultSessionTarget(targetInput); |
| const modelChoices = await listReadyModelChoices({ connectionStore, credentialStore }); |
| const backends = new BackendRegistry(); |
| const shellRunListeners = new Set<(update: ShellRunUpdate) => void>(); |
| const shellRuns = new ShellRunProcessManager({ |
| store: shellRunStore, |
| newId: randomUUID, |
| now: Date.now, |
| onShellRunUpdate: (update) => { |
| for (const listener of shellRunListeners) { |
| try { |
| listener(update); |
| } catch { |
| // One UI observer must not suppress updates for the rest. |
| } |
| } |
| }, |
| }); |
| const sandboxManager = createBuiltinSandboxManager(); |
| const filesystemWorkerLaunchSpecProvider = |
| sandboxManager && isBuiltinFilesystemWorkerSandboxAvailable() |
| ? createFilesystemWorkerLaunchSpecProvider({ |
| runtime: 'node', |
| platform: process.platform, |
| resourceLocation: { kind: 'runtime' }, |
| }) |
| : undefined; |
| const filesystemWorker = |
| sandboxManager && filesystemWorkerLaunchSpecProvider |
| ? new FilesystemWorkerClient({ |
| sandboxManager, |
| getLaunchSpec: filesystemWorkerLaunchSpecProvider, |
| }) |
| : undefined; |
| const sandboxDiagnosticsProvider = createSandboxDiagnosticsProvider({ |
| ...(sandboxManager ? { sandboxManager } : {}), |
| ...(filesystemWorkerLaunchSpecProvider |
| ? { getFilesystemWorkerLaunchSpec: filesystemWorkerLaunchSpecProvider } |
| : {}), |
| }); |
| const tools = buildBuiltinTools({ |
| shellRuns, |
| runtimeResources: shellRuns, |
| backgroundTasks: shellRuns, |
| ptyControls: shellRuns, |
| snapshotImage: createReadImageSnapshotter(artifactStore), |
| ...(sandboxManager ? { sandboxManager } : {}), |
| ...(filesystemWorker |
| ? { |
| filesystemWorker, |
| } |
| : {}), |
| }); |
| // Child sessions get fresh catalog tools. Their Read tool cannot inspect |
| // parent runtime resources, and the SessionManager narrows this union to the |
| // selected profile after checking host capabilities (for example, a |
| // worktree executor for implementation children). Agent tools are excluded |
| // so children cannot recursively spawn from this surface. |
| const childAgentTools = agentGraphEnabled |
| ? buildChildAgentTools([ |
| ...buildBuiltinTools({ |
| snapshotImage: createReadImageSnapshotter(artifactStore), |
| ...(sandboxManager ? { sandboxManager } : {}), |
| ...(filesystemWorker |
| ? { |
| filesystemWorker, |
| } |
| : {}), |
| }), |
| buildNativeWebSearchTool(), |
| ]) |
| : []; |
| const automationManager = new AutomationManager({ |
| generateId: () => randomUUID(), |
| now: () => Date.now(), |
| }); |
| // A heartbeat-only CLI owns no durable Automations and must not reconcile the |
| // shared authority. A cron-enabled host persists through the same operational |
| // SQLite authority as Desktop. |
| const cronEnabled = input.automationCreateFreshRun !== undefined; |
| const automationStore = createAutomationStore(stateRoot); |
| // If the authority cannot be read, do not attempt a later replacement write. |
| let durableStoreReadable = true; |
| const syncAutomations = cronEnabled |
| ? (): void => { |
| if (!durableStoreReadable) return; |
| const durable = automationManager |
| .listAll() |
| .filter((a) => a.durable && (a.status === 'active' || a.status === 'paused')); |
| automationStore.sync(durable).catch((err) => { |
| console.warn('[runtime-bootstrap] failed to persist durable automations:', err); |
| }); |
| } |
| : (): void => { |
| /* heartbeat-only host owns no durable automations; never overwrite the shared store */ |
| }; |
| const automationTool = buildAutomationTool({ |
| automationManager, |
| onAutomationChange: syncAutomations, |
| cronEnabled, |
| }); |
| |
| // Load durable automations only on a host that can run them — a cron-disabled |
| // host must not adopt/reconcile crons it doesn't own (see above). |
| if (cronEnabled) { |
| try { |
| const saved = await automationStore.loadAll(); |
| automationManager.registerAll(saved); |
| } catch (err) { |
| durableStoreReadable = false; |
| console.error( |
| '[runtime-bootstrap] durable automation store unreadable; persistence disabled to avoid data loss:', |
| err, |
| ); |
| } |
| } |
| |
| const goalManager = new GoalManager({ generateId: () => randomUUID(), now: () => Date.now() }); |
| const goalTokenCache = new Map<string, number>(); |
| let runtime!: SessionManager; |
| // Construct the lifecycle authority before exposing Goal tools. The runtime |
| // reference is only read after context creation, when a real turn settles. |
| const goalContinuation = new CliGoalContinuation({ |
| goalManager, |
| evaluator: { |
| async evaluate(prompt: string, sessionId: string): Promise<string> { |
| const ai = (await import('ai')) as unknown as { |
| generateText(opts: Record<string, unknown>): Promise<{ text: string }>; |
| }; |
| const header = await store.readHeader(sessionId); |
| const ready = await resolveSessionTargetForSlug(header.llmConnectionSlug, { |
| connectionStore, |
| credentialStore, |
| requestedModel: header.model, |
| }); |
| const modelFetch = buildSubscriptionModelFetch({ |
| connection: ready.connection, |
| sessionId: 'goal-evaluator', |
| modelId: ready.model, |
| ...(ready.connection.providerType === 'claude-subscription' |
| ? { |
| claude: { |
| cloakEnabled: isMakaClaudeSubscriptionCloakEnabled(), |
| deviceId: await getOrCreateCliClaudeDeviceId(configRoot), |
| accountUuid: ready.oauthTokens?.account_uuid ?? '', |
| }, |
| } |
| : {}), |
| }); |
| const result = await ai.generateText({ |
| model: getAIModel({ |
| connection: ready.connection, |
| apiKey: ready.apiKey ?? '', |
| modelId: ready.model, |
| fetch: modelFetch, |
| }), |
| prompt, |
| providerOptions: buildProviderOptions( |
| ready.connection, |
| ready.model, |
| header.thinkingLevel, |
| ), |
| maxOutputTokens: 1024, |
| }); |
| return result.text; |
| }, |
| }, |
| async getRecentContext(sessionId: string): Promise<string> { |
| const messages = await runtime.getMessages(sessionId); |
| let total = 0; |
| for (const message of messages) { |
| if (message.type === 'token_usage') |
| total += message.total ?? message.input + message.output; |
| } |
| goalTokenCache.set(sessionId, total); |
| return messages |
| .slice(-10) |
| .filter((message) => message.type === 'user' || message.type === 'assistant') |
| .slice(-6) |
| .map( |
| (message) => |
| `[${message.type}]: ${(message.type === 'user' || message.type === 'assistant' ? message.text : '').slice(0, 500)}`, |
| ) |
| .join('\n'); |
| }, |
| getTokenCount: (sessionId: string) => goalTokenCache.get(sessionId) ?? 0, |
| }); |
| |
| // One-sentence session recap (issue #1055): a tool-free model call over the |
| // session's own history, never written back to it. Mirrors the goal |
| // evaluator's connection resolution + call shape above. |
| const recap: SessionRecapGenerator = { |
| async generate(sessionId, reason) { |
| let modelId = ''; |
| // The actual bounded request sent to the model (projection + budget trim |
| // + trailing instruction), persisted verbatim to the artifact below. |
| let requestMessages: ModelMessage[] = []; |
| let rawText = ''; |
| let cleaned = ''; |
| let errorMessage: string | undefined; |
| try { |
| // Authoritative RuntimeEvent projection (issue #1182 review): reuses |
| // the same read model, budget policy, and replay-plan projection the |
| // backend uses for its own history, instead of re-deriving a lossy |
| // StoredMessage-based one. |
| const view = await runtimeReadModel.getSessionView(sessionId); |
| const header = await store.readHeader(sessionId); |
| const ready = await resolveSessionTargetForSlug(header.llmConnectionSlug, { |
| connectionStore, |
| credentialStore, |
| requestedModel: header.model, |
| }); |
| modelId = ready.model; |
| const modelFetch = buildSubscriptionModelFetch({ |
| connection: ready.connection, |
| sessionId: 'session-recap', |
| modelId: ready.model, |
| ...(ready.connection.providerType === 'claude-subscription' |
| ? { |
| claude: { |
| cloakEnabled: isMakaClaudeSubscriptionCloakEnabled(), |
| deviceId: await getOrCreateCliClaudeDeviceId(configRoot), |
| accountUuid: ready.oauthTokens?.account_uuid ?? '', |
| }, |
| } |
| : {}), |
| }); |
| requestMessages = buildSessionRecapMessages({ |
| events: view.events, |
| connection: ready.connection, |
| modelId: ready.model, |
| }); |
| const ai = (await import('ai')) as unknown as { |
| generateText(opts: Record<string, unknown>): Promise<{ text: string }>; |
| }; |
| const result = await ai.generateText({ |
| model: getAIModel({ |
| connection: ready.connection, |
| apiKey: ready.apiKey ?? '', |
| modelId: ready.model, |
| fetch: modelFetch, |
| }), |
| messages: requestMessages, |
| providerOptions: buildProviderOptions( |
| ready.connection, |
| ready.model, |
| header.thinkingLevel, |
| ), |
| maxOutputTokens: 1024, |
| }); |
| rawText = result.text; |
| cleaned = cleanRecapText(rawText); |
| return { ok: true as const, text: cleaned, raw: rawText }; |
| } catch (error) { |
| errorMessage = error instanceof Error ? error.message : String(error); |
| return { ok: false as const, error: errorMessage }; |
| } finally { |
| try { |
| await artifactStore.create({ |
| sessionId, |
| turnId: randomUUID(), |
| name: 'recap-request.json', |
| kind: 'file', |
| content: JSON.stringify( |
| { |
| reason, |
| model: modelId, |
| messageCount: requestMessages.length, |
| messages: requestMessages, |
| raw: rawText, |
| cleaned, |
| ...(errorMessage ? { error: errorMessage } : {}), |
| }, |
| null, |
| 2, |
| ), |
| ...(cleaned ? { summary: cleaned.slice(0, 100) } : {}), |
| }); |
| } catch { |
| // Best-effort persistence; recap must work even if the artifact store fails. |
| } |
| } |
| }, |
| }; |
| |
| const goalTools = |
| input.surface === 'tui' |
| ? buildGoalTools({ |
| goalManager, |
| goalContinuation, |
| getTokenCount: (sessionId: string) => goalTokenCache.get(sessionId) ?? 0, |
| }) |
| : []; |
| const subagentTools = agentGraphEnabled |
| ? buildParentAgentTools({ |
| definitions: listRunnableBuiltinAgentDefinitions({ |
| tools: childAgentTools.filter((tool) => tool.name !== 'WebSearch'), |
| worktreeChildExecutorAvailable: worktreeChildExecutor !== undefined, |
| }), |
| }) |
| : []; |
| const subagentToolNames = new Set(subagentTools.map((tool) => tool.name)); |
| const surfaceTools = |
| input.surface === 'tui' ? [buildAskUserQuestionTool(), buildRequestSandboxBoundaryTool()] : []; |
| let cliProductToolSurface: EffectiveProductToolSurface; |
| const resolveCliSkillHost: HostCapabilitiesResolver = () => |
| cliProductToolSurface.hostCapabilities; |
| const skillShadowTracker = new SkillShadowSelectionTracker(); |
| const skillTool = buildSkillAgentTool( |
| ({ cwd }) => resolveSkillDiscoveryPaths(cwd, configRoot), |
| resolveCliSkillHost, |
| { shadowTracker: skillShadowTracker }, |
| ); |
| const skillSearchTool = buildSkillSearchAgentTool( |
| ({ cwd }) => resolveSkillDiscoveryPaths(cwd, configRoot), |
| resolveCliSkillHost, |
| { shadowTracker: skillShadowTracker }, |
| ); |
| const boundTools = [ |
| ...tools, |
| automationTool, |
| ...goalTools, |
| skillTool, |
| skillSearchTool, |
| ...subagentTools, |
| ...surfaceTools, |
| ]; |
| assertProductBindingCatalogClean( |
| 'cli', |
| boundTools.map((tool) => tool.name), |
| ); |
| cliProductToolSurface = projectEffectiveProductToolSurface({ |
| host: 'cli', |
| tools: boundTools, |
| policy: { |
| economy: input.surface === 'tui' && !process.env.MAKA_DISABLE_DEFERRED_TOOLS, |
| }, |
| }); |
| const allTools = [...cliProductToolSurface.tools]; |
| |
| backends.register('ai-sdk', async (ctx) => { |
| const header = |
| input.sessionCwdOverride?.sessionId === ctx.sessionId |
| ? { ...ctx.header, cwd: input.sessionCwdOverride.cwd } |
| : ctx.header; |
| // Resolve the session's own connection — not the global default — so a |
| // /model switch that rebinds the session to another provider actually runs |
| // on that provider (the desktop app resolves the backend the same way). |
| const ready = await resolveSessionTargetForSlug(header.llmConnectionSlug, { |
| connectionStore, |
| credentialStore, |
| requestedModel: header.model, |
| }); |
| const modelFetch = buildSubscriptionModelFetch({ |
| connection: ready.connection, |
| sessionId: ctx.sessionId, |
| modelId: ready.model, |
| ...(ready.connection.providerType === 'claude-subscription' |
| ? { |
| claude: { |
| cloakEnabled: isMakaClaudeSubscriptionCloakEnabled(), |
| deviceId: await getOrCreateCliClaudeDeviceId(configRoot), |
| accountUuid: ready.oauthTokens?.account_uuid ?? '', |
| }, |
| } |
| : {}), |
| }); |
| const sandboxDiagnosticsSnapshot = await sandboxDiagnosticsProvider.resolve({ |
| mode: header.permissionMode, |
| cwd: header.cwd, |
| }); |
| const streamConnectTimeoutMs = resolveCliStreamConnectTimeoutMs(); |
| const agentGraphSupervisorTools = |
| !ctx.tools && agentGraphEnabled |
| ? await agentGraphCoordinator!.toolsForSession(ctx.sessionId) |
| : []; |
| const settings = await settingsStore.get(); |
| const routedChildTools = routeWebSearchTools({ |
| tools: |
| settings.webSearch.defaultProvider === 'model' |
| ? childAgentTools |
| : childAgentTools.filter((tool) => tool.name !== 'WebSearch'), |
| settings: settings.webSearch, |
| connection: ready.connection, |
| model: ready.model, |
| privacy: settings.privacy, |
| }); |
| const targetSubagentTools = |
| !ctx.tools && agentGraphEnabled |
| ? buildParentAgentTools({ |
| definitions: listRunnableBuiltinAgentDefinitions({ |
| tools: routedChildTools, |
| worktreeChildExecutorAvailable: worktreeChildExecutor !== undefined, |
| }), |
| }) |
| : []; |
| const routedTools = routeWebSearchTools({ |
| tools: ctx.tools |
| ? ctx.tools |
| : [ |
| ...allTools.filter((tool) => !subagentToolNames.has(tool.name)), |
| ...targetSubagentTools, |
| ...agentGraphSupervisorTools, |
| ], |
| settings: settings.webSearch, |
| connection: ready.connection, |
| model: ready.model, |
| privacy: settings.privacy, |
| allowAddNative: ctx.tools === undefined, |
| }); |
| const productToolSurface = projectEffectiveProductToolSurface({ |
| host: 'cli', |
| tools: routedTools, |
| policy: cliProductToolSurface.identity.policy, |
| }); |
| const backendTools = [...productToolSurface.tools]; |
| const admitsAgentChildren = productToolSurface.boundSurfaceIds.includes(AGENT_TOOL_GROUP_ID); |
| return new AiSdkBackend({ |
| sessionId: ctx.sessionId, |
| header: { ...header, model: ready.model }, |
| appendMessage: |
| ctx.appendMessage ?? ((message) => ctx.store.appendMessage(ctx.sessionId, message)), |
| readExecutionBoundary: () => ctx.store.readExecutionBoundary!(ctx.sessionId), |
| ...(input.surface === 'tui' |
| ? { |
| createSandboxBoundaryRequest: (request) => |
| ctx.store.createSandboxBoundaryRequest!(request), |
| settleSandboxBoundaryRequest: (request) => |
| ctx.store.settleSandboxBoundaryRequest!(request), |
| } |
| : {}), |
| connection: ready.connection, |
| apiKey: ready.apiKey, |
| modelId: ready.model, |
| modelFactory: (modelInput) => getAIModel({ ...modelInput, fetch: modelFetch }), |
| tools: backendTools, |
| sandboxDiagnosticsSnapshot, |
| toolAvailability: productToolSurface.toolAvailability, |
| ...(admitsAgentChildren |
| ? { |
| spawnChildAgent: (childInput) => runtime.spawnChildAgent(ctx.sessionId, childInput), |
| spawnChildSession: (childInput) => |
| runtime.spawnChildSession(ctx.sessionId, { |
| spawnedBy: { |
| parentRunId: childInput.parentRunId, |
| parentTurnId: childInput.parentTurnId, |
| toolCallId: childInput.toolCallId, |
| }, |
| agentProfile: childInput.agentProfile, |
| ...(childInput.subagentId ? { subagentId: childInput.subagentId } : {}), |
| prompt: childInput.prompt, |
| ...(childInput.swarm ? { swarm: childInput.swarm } : {}), |
| abortSignal: childInput.abortSignal, |
| ...(childInput.onReady ? { onReady: childInput.onReady } : {}), |
| ...(childInput.onEvent ? { onEvent: childInput.onEvent } : {}), |
| }), |
| prepareChildAgentResume: (sourceRunId) => |
| runtime.prepareChildAgentResume(ctx.sessionId, sourceRunId), |
| resumeChildAgent: (childInput) => runtime.resumeChildAgent(ctx.sessionId, childInput), |
| retryChildAgent: (childInput) => runtime.retryChildAgent(ctx.sessionId, childInput), |
| listChildAgents: () => runtime.listChildAgents(ctx.sessionId), |
| readChildAgentOutput: (childInput) => |
| runtime.readChildAgentOutput(ctx.sessionId, childInput), |
| } |
| : {}), |
| providerOptions: buildProviderOptions(ready.connection, ready.model, header.thinkingLevel), |
| ...(streamConnectTimeoutMs !== undefined ? { streamConnectTimeoutMs } : {}), |
| contextBudget: buildDefaultContextBudgetPolicy(ready.connection, { |
| name: 'cli-default-history-budget', |
| modelId: ready.model, |
| }), |
| supportsVision: resolveModelVisionSupport( |
| ready.connection.providerType, |
| ready.connection.models, |
| ready.model, |
| relayModelProfile(ready.connection, ready.model)?.vision, |
| ), |
| readAttachmentBytes: createAttachmentByteReader({ artifactStore, sessionId: ctx.sessionId }), |
| loadHistoryCompact: (event) => loadHistoryCompactBlocksFromArtifacts(artifactStore, event), |
| loadHistoryCompactCheckpoint: ctx.loadHistoryCompactCheckpoint, |
| summarizeHistoryCompact: buildLlmHistorySummarizer({ |
| // Reuse the same connection/model the session already drives, so the |
| // summary stays consistent with the model that will consume it. |
| resolveModel: () => |
| getAIModel({ |
| connection: ready.connection, |
| apiKey: ready.apiKey, |
| modelId: ready.model, |
| fetch: modelFetch, |
| }), |
| providerOptions: buildProviderOptions(ready.connection, ready.model, header.thinkingLevel), |
| }), |
| // The canonical metering sink (#1679). Without it this composition root |
| // produces diagnostics and no accounting at all — for `/compact` and for |
| // ordinary sends alike. |
| ...(ctx.recordModelCallAttempt ? { recordModelCallAttempt: ctx.recordModelCallAttempt } : {}), |
| recordHistoryCompactCheckpoint: ctx.recordHistoryCompactCheckpoint, |
| loadTurnRuntimeEvents: ctx.loadTurnRuntimeEvents, |
| allowMidTurnHistoryCompaction: ctx.allowMidTurnHistoryCompaction, |
| systemPrompt: |
| ctx.systemPrompt ?? |
| (async ({ cwd, emitSkillCatalogTrace }) => { |
| const settings = await settingsStore.get(); |
| return buildCliSystemPrompt({ |
| settings, |
| cwd, |
| workspaceRoot: configRoot, |
| host: productToolSurface.hostCapabilities, |
| modelContextWindow: resolveSelectedModelContextWindow(ready.connection, ready.model), |
| onSkillSelection: (report) => |
| emitSkillCatalogTrace?.('Skill catalog selection completed', { |
| policyVersion: report.policyVersion, |
| budgetChars: report.budgetChars, |
| usedChars: report.usedChars, |
| totalCount: report.totalCount, |
| eligibleCount: report.eligibleCount, |
| advertisedCount: report.advertisedCount, |
| omittedCount: report.omittedCount, |
| }), |
| }); |
| }), |
| turnTailPrompt: ({ cwd }) => |
| buildCliTurnTailPrompt({ cwd, sessionId: ctx.sessionId, automationManager, goalManager }), |
| shellRunContextSummary: ctx.shellRunContextSummary, |
| recordRunTrace: ctx.recordRunTrace, |
| ...(ctx.recordProviderRequestCapture |
| ? { |
| recordProviderRequestCapture: createProviderRequestCaptureRecorder({ |
| persistArtifact: async (capture) => { |
| const artifact = await persistProviderRequestCaptureArtifact(artifactStore, { |
| sessionId: ctx.sessionId, |
| turnId: capture.turnId, |
| captureId: capture.captureId, |
| step: capture.step, |
| serializedRequest: capture.serializedRequest, |
| now: Date.now(), |
| }); |
| return { artifactId: artifact.id }; |
| }, |
| recordLedger: ctx.recordProviderRequestCapture, |
| }), |
| recordProviderRequestAttempt: ctx.recordProviderRequestAttempt, |
| } |
| : {}), |
| newId: randomUUID, |
| now: Date.now, |
| ...(input.maxSteps !== undefined ? { maxSteps: input.maxSteps } : {}), |
| ...(runtimePersistence.runtimeCommitStore |
| ? { runtimeCommitSink: runtimePersistence.runtimeCommitStore } |
| : {}), |
| }); |
| }); |
| |
| const resolveChildTools = async (sessionId: string): Promise<readonly MakaTool[]> => { |
| const header = await store.readHeader(sessionId); |
| const ready = await resolveSessionTargetForSlug(header.llmConnectionSlug, { |
| connectionStore, |
| credentialStore, |
| requestedModel: header.model, |
| }); |
| const settings = await settingsStore.get(); |
| return routeWebSearchTools({ |
| tools: |
| settings.webSearch.defaultProvider === 'model' |
| ? childAgentTools |
| : childAgentTools.filter((tool) => tool.name !== 'WebSearch'), |
| settings: settings.webSearch, |
| connection: ready.connection, |
| model: ready.model, |
| privacy: settings.privacy, |
| }); |
| }; |
| |
| runtime = new SessionManager({ |
| store, |
| runStore, |
| runtimeEventStore, |
| ...(runtimePersistence.runtimeCommitStore |
| ? { |
| runtimeCommitSink: runtimePersistence.runtimeCommitStore, |
| toolBoundaryProtocol: runtimePersistence.runtimeCommitStore.toolBoundaryProtocol, |
| } |
| : {}), |
| shellRuns, |
| backends, |
| subagentCatalog, |
| runtimeSource: input.runtimeSource ?? (input.surface === 'activation' ? 'gateway' : undefined), |
| safeBoundaryResumeEnabled: |
| input.safeBoundaryResumeEnabled ?? |
| (input.surface === 'activation' || process.env.MAKA_RUNTIME_SAFE_BOUNDARY_RESUME === '1'), |
| onContinuationLifecycleEvent: (event) => { |
| const writeDiagnostic = input.surface === 'activation' ? console.error : console.info; |
| writeDiagnostic('[runtime-resume]', JSON.stringify(event)); |
| }, |
| inspectContinuationSafety: createLocalContinuationSafetyInspector({ |
| readSessionCwd: async (sessionId) => (await store.readHeader(sessionId)).cwd, |
| resolveWorkspaceIdentity: async (cwd) => resolveWorkspaceIdentity({ path: cwd }), |
| listAvailableToolNames: async () => allTools.map((tool) => tool.name), |
| hasPendingBackgroundOperations: async (sessionId) => { |
| const [shellUpdates, runs] = await Promise.all([ |
| shellRuns.listSessionUpdates(sessionId), |
| runStore.listSessionRuns(sessionId), |
| ]); |
| return ( |
| shellUpdates.some((update) => isActiveShellRunStatus(update.result.status)) || |
| runs.some( |
| (run) => |
| run.parentRunId !== undefined && |
| ['created', 'running', 'waiting_for_user'].includes(run.status), |
| ) |
| ); |
| }, |
| }), |
| ...(agentGraphEnabled |
| ? { childTools: childAgentTools, resolveChildTools, worktreeChildExecutor } |
| : {}), |
| runtimeInvocationObserver: input.runtimeInvocationObserver, |
| onSessionTitleChanged: input.onSessionTitleChanged, |
| ...(input.surface === 'tui' |
| ? { |
| generateSessionTitle: async ({ sessionId, header, sourceText }) => { |
| const ready = await resolveSessionTargetForSlug(header.llmConnectionSlug, { |
| connectionStore, |
| credentialStore, |
| requestedModel: header.model, |
| }); |
| const modelFetch = buildSubscriptionModelFetch({ |
| connection: ready.connection, |
| sessionId, |
| modelId: ready.model, |
| ...(ready.connection.providerType === 'claude-subscription' |
| ? { |
| claude: { |
| cloakEnabled: isMakaClaudeSubscriptionCloakEnabled(), |
| deviceId: await getOrCreateCliClaudeDeviceId(configRoot), |
| accountUuid: ready.oauthTokens?.account_uuid ?? '', |
| }, |
| } |
| : {}), |
| }); |
| return generateRuntimeSessionTitle({ |
| model: getAIModel({ |
| connection: ready.connection, |
| apiKey: ready.apiKey ?? '', |
| modelId: ready.model, |
| fetch: modelFetch, |
| }), |
| providerOptions: buildProviderOptions(ready.connection, ready.model), |
| sourceText, |
| }); |
| }, |
| } |
| : {}), |
| cleanupHistoryCompactArtifacts: async (cleanupInput) => { |
| await cleanupLegacyHistoryCompactArtifacts({ |
| ...cleanupInput, |
| artifactStore, |
| onDiagnostic: (diagnostic) => console.warn('[history-compact-cleanup]', diagnostic), |
| }); |
| }, |
| newId: randomUUID, |
| now: Date.now, |
| }); |
| if (agentGraphControlStore) { |
| agentGraphSupervisorWakeCoordinator = new AgentGraphSupervisorWakeCoordinator({ |
| activityRegistry: goalContinuation.activities, |
| wakeStore: agentGraphControlStore, |
| readSnapshot: (rootSessionId) => agentGraphCoordinator!.getSnapshot(rootSessionId), |
| startTurn: async (sessionId, message, activity, abortSignal, isCurrent) => { |
| if (!(await isCurrent())) { |
| return { |
| kind: 'superseded', |
| turnId: message.turnId, |
| reason: 'Agent graph supervisor checkpoint was superseded before execution.', |
| }; |
| } |
| let stopPromise: Promise<void> | undefined; |
| const stop = (): void => { |
| stopPromise ??= runtime.stopSession(sessionId, { source: 'graph_supervisor' }); |
| }; |
| abortSignal.addEventListener('abort', stop, { once: true }); |
| if (abortSignal.aborted) stop(); |
| try { |
| return await drainGoalTurn({ |
| events: runtime.sendMessage(sessionId, message), |
| turnId: message.turnId, |
| activity, |
| }); |
| } finally { |
| abortSignal.removeEventListener('abort', stop); |
| await stopPromise; |
| } |
| }, |
| isSessionDeliverable: async (sessionId) => { |
| try { |
| const header = await store.readHeader(sessionId); |
| return !header.isArchived && header.status !== 'archived'; |
| } catch (error) { |
| if (isSessionNotFoundError(error)) return false; |
| throw error; |
| } |
| }, |
| inspectAttempt: async (rootSessionId, attemptId, turnId) => { |
| const runs = (await runStore.listSessionRuns(rootSessionId)).filter( |
| (run) => run.agentGraphWakeAttemptId === attemptId && run.turnId === turnId, |
| ); |
| if (runs.length > 1) { |
| throw new Error( |
| `Agent graph supervisor wake attempt ${attemptId} has multiple AgentRuns`, |
| ); |
| } |
| return runs[0]?.status ?? 'missing'; |
| }, |
| recoverContextOverflow: (rootSessionId, { abortSignal }) => |
| recoverAgentGraphSupervisorContextOverflow({ |
| rootSessionId, |
| compactTurnId: randomUUID(), |
| abortSignal, |
| compactSession: (sessionId, input) => runtime.compactSession(sessionId, input), |
| }), |
| shouldWake: shouldWakeAgentSwarmSupervisor, |
| renderWake: renderAgentSwarmSupervisorWake, |
| newId: randomUUID, |
| onDiagnostic: (diagnostic) => { |
| console.warn('[agent-graph-supervisor-wake]', JSON.stringify(diagnostic)); |
| }, |
| onError: (rootSessionId, error) => { |
| agentGraphErrors.set(rootSessionId, error); |
| }, |
| }); |
| agentGraphCoordinator = new AgentGraphCoordinator({ |
| sessionStore: store, |
| runStore, |
| runtimeEventStore, |
| controlStore: agentGraphControlStore, |
| runtime, |
| newId: randomUUID, |
| onReconciliation: (rootSessionId, result) => { |
| agentGraphSupervisorWakeCoordinator!.notify(rootSessionId, result); |
| }, |
| onCheckpoint: (rootSessionId) => { |
| agentGraphSupervisorWakeCoordinator!.notify(rootSessionId); |
| }, |
| onError: (rootSessionId, error) => { |
| agentGraphErrors.set(rootSessionId, error); |
| }, |
| }); |
| } |
| await runtime.recoverInterruptedSessions(); |
| |
| const automationScheduler = new AutomationScheduler({ |
| automationManager, |
| canFire: async (automation) => { |
| if ( |
| automation.kind === 'heartbeat' && |
| goalContinuation.activities.whenIdle(automation.sessionId) |
| ) { |
| return false; |
| } |
| return evaluateAutomationCanFire(automation, { |
| // The CLI has no incognito UI, but the setting is shared — honour it if set. |
| isIncognitoActive: async () => |
| (await settingsStore.get()).privacy?.incognitoActive === true, |
| readSessionHeader: (sessionId) => store.readHeader(sessionId), |
| // Default idle set {active, done, waiting_for_user} — a session parked |
| // waiting for the user IS the wakeup's home scenario (#639): the |
| // heartbeat starts a turn in place of the user. It still never fires |
| // into a 'running' (mid-turn) session. |
| // Cron is disabled here (createFreshRun omitted); the scheduler ignores it. |
| }); |
| }, |
| // Heartbeat: inject into the automation's session; resolve after the drain. |
| // The CLI has no multi-session UI, so cron (fresh-session) is disabled — |
| // createFreshRun is omitted, so the tool advertises heartbeat only. |
| injectTurn: async (sessionId, prompt, automationId) => { |
| const turnId = randomUUID(); |
| const outcome = await goalContinuation.runAutomationTurn({ |
| sessionId, |
| turnId, |
| start: () => |
| runtime.sendMessage(sessionId, { |
| turnId, |
| text: prompt, |
| origin: { kind: 'automation', automationId }, |
| }), |
| }); |
| const error = |
| outcome.kind === 'errored' || outcome.kind === 'suspended' |
| ? outcome.reason |
| : outcome.kind === 'aborted' |
| ? 'Automation turn was aborted.' |
| : undefined; |
| return { |
| runId: turnId, |
| ok: outcome.kind === 'completed', |
| ...(error ? { error } : {}), |
| }; |
| }, |
| createFreshRun: input.automationCreateFreshRun, |
| // unref() the tick timer: a background poll must never hold the CLI |
| // process open. Without this, any bootstrap consumer that exits without |
| // close() (a finished one-shot run, a test) hangs on the 5s tick forever. |
| setTimeout: (fn, ms) => { |
| const timer = setTimeout(fn, ms); |
| timer.unref?.(); |
| return timer; |
| }, |
| clearTimeout: (timer) => clearTimeout(timer as ReturnType<typeof setTimeout>), |
| onStateChange: syncAutomations, |
| }); |
| |
| automationScheduler.start(); |
| |
| return { |
| workspaceRoot: input.workspaceRoot, |
| stateRoot, |
| configRoot, |
| cwd: input.cwd, |
| runtime, |
| target, |
| modelChoices, |
| tools: allTools, |
| skills: { |
| source: (cwd) => resolveSkillDiscoveryPaths(cwd, configRoot), |
| host: cliProductToolSurface.hostCapabilities, |
| }, |
| automationManager, |
| automationScheduler, |
| subscribeShellRunUpdates: (listener) => { |
| shellRunListeners.add(listener); |
| return () => shellRunListeners.delete(listener); |
| }, |
| listShellRunUpdates: (sessionId) => runtime.listShellRunUpdates(sessionId), |
| goalManager, |
| goalContinuation, |
| onboarding: createApiKeyOnboardingSurface({ |
| connectionStore, |
| credentialStore, |
| fetchModels: fetchProviderModels, |
| }), |
| recap, |
| foreignSessions, |
| ...(agentGraphCoordinator && agentGraphSupervisorWakeCoordinator |
| ? { |
| agentGraph: { |
| reserveActivity: (sessionId: string) => goalContinuation.activities.reserve(sessionId), |
| waitForCompletion: async (sessionId: string) => { |
| for (;;) { |
| await agentGraphCoordinator!.waitForIdle(sessionId); |
| await agentGraphSupervisorWakeCoordinator!.waitForIdle(); |
| await agentGraphCoordinator!.waitForIdle(sessionId); |
| const error = agentGraphErrors.get(sessionId); |
| if (error !== undefined) throw error; |
| const snapshot = await agentGraphCoordinator!.getSnapshot(sessionId); |
| if (snapshot.closed || snapshot.scheduleRevision === 0) return; |
| await agentGraphSupervisorWakeCoordinator!.waitForIdle(); |
| await agentGraphCoordinator!.waitForIdle(sessionId); |
| const settled = await agentGraphCoordinator!.getSnapshot(sessionId); |
| if (settled.closed) return; |
| throw new Error( |
| `Agent graph became idle without finish (status=${settled.status}, revision=${settled.scheduleRevision}, snapshot=${settled.snapshotVersion})`, |
| ); |
| } |
| }, |
| }, |
| } |
| : {}), |
| close: async () => { |
| // Stop the automation scheduler's timer (else it keeps the process alive |
| // and ticks into a stopped session), then terminate background shell runs. |
| automationScheduler.dispose(); |
| goalContinuation.dispose(); |
| goalManager.dispose(); |
| await agentGraphSupervisorWakeCoordinator?.close(); |
| await agentGraphCoordinator?.close(); |
| agentGraphControlStore?.close(); |
| await shellRuns.terminateAll(); |
| shellRunListeners.clear(); |
| await store.close?.(); |
| runtimePersistence.close(); |
| runStore.close?.(); |
| shellRunStore.close(); |
| artifactStore.close?.(); |
| }, |
| }; |
| } |
| |
| export async function getOrCreateCliClaudeDeviceId( |
| workspaceRoot: string, |
| deps: GetOrCreateCliClaudeDeviceIdDeps = {}, |
| ): Promise<string> { |
| const deviceIdFilePath = join(workspaceRoot, '.maka_cli_claude_device_id'); |
| try { |
| const existing = (await readFile(deviceIdFilePath, 'utf8')).trim(); |
| if (/^[a-f0-9]{64}$/i.test(existing)) return existing.toLowerCase(); |
| } catch { |
| // fall through to create; device id persistence is best-effort metadata. |
| } |
| |
| const next = (deps.newId ?? (() => randomBytes(32).toString('hex')))().toLowerCase(); |
| try { |
| await mkdir(dirname(deviceIdFilePath), { recursive: true }); |
| await writeFile(deviceIdFilePath, next, { mode: 0o600 }); |
| await chmod(deviceIdFilePath, 0o600); |
| } catch { |
| // best-effort persistence; use the generated id for this process if disk fails. |
| } |
| return next; |
| } |