blob: c9c8f5e7ca26bab13563c529d6aa7e939aa46dbe [file]
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;
}