mirror of
https://github.com/VibedByKaKi/t3-code-android-nightly.git
synced 2026-10-09 03:41:17 +02:00
3807 lines
150 KiB
TypeScript
3807 lines
150 KiB
TypeScript
import type {
|
|
Event as OpenCodeEvent,
|
|
Message as OpenCodeMessage,
|
|
Part as OpenCodePart,
|
|
PermissionRequest,
|
|
PermissionRuleset,
|
|
QuestionRequest,
|
|
Session as OpenCodeSession,
|
|
Todo as OpenCodeTodo,
|
|
ToolPart,
|
|
} from "@opencode-ai/sdk/v2";
|
|
import { HostProcessEnvironment } from "@t3tools/shared/hostProcess";
|
|
import { getModelSelectionStringOptionValue } from "@t3tools/shared/model";
|
|
import { causeErrorTag } from "@t3tools/shared/observability";
|
|
import {
|
|
defaultInstanceIdForDriver,
|
|
type ModelSelection,
|
|
type OpenCodeSettings,
|
|
type OrchestrationV2AppThread,
|
|
type OrchestrationV2ConversationMessage,
|
|
type OrchestrationV2ExecutionNode,
|
|
type OrchestrationV2PlanStep,
|
|
type OrchestrationV2ProviderCapabilities,
|
|
type OrchestrationV2ProviderFailure,
|
|
type OrchestrationV2ProviderRef,
|
|
type OrchestrationV2ProviderSession,
|
|
type OrchestrationV2ProviderThread,
|
|
type OrchestrationV2ProviderTurn,
|
|
type OrchestrationV2RuntimeRequest,
|
|
type OrchestrationV2Subagent,
|
|
type OrchestrationV2TurnItem,
|
|
OpenCodeSettings as OpenCodeSettingsSchema,
|
|
type PlanId,
|
|
ProviderDriverKind,
|
|
type ProviderInstanceId,
|
|
type ProviderRequestKind,
|
|
type ProviderSessionId,
|
|
type RuntimeRequestId,
|
|
type ThreadId,
|
|
} from "@t3tools/contracts";
|
|
import * as Cause from "effect/Cause";
|
|
import * as DateTime from "effect/DateTime";
|
|
import * as Deferred from "effect/Deferred";
|
|
import * as Duration from "effect/Duration";
|
|
import * as Effect from "effect/Effect";
|
|
import * as Exit from "effect/Exit";
|
|
import * as Fiber from "effect/Fiber";
|
|
import * as Layer from "effect/Layer";
|
|
import * as Option from "effect/Option";
|
|
import * as Queue from "effect/Queue";
|
|
import * as Random from "effect/Random";
|
|
import * as Schema from "effect/Schema";
|
|
import * as Semaphore from "effect/Semaphore";
|
|
import * as Scope from "effect/Scope";
|
|
import * as Stream from "effect/Stream";
|
|
|
|
import { resolveAttachmentPath } from "../../attachmentStore.ts";
|
|
import * as ServerConfig from "../../config.ts";
|
|
import { mcpToolPresentation } from "@t3tools/provider-core/server/mcpToolPresentation";
|
|
import * as McpProviderSession from "@t3tools/provider-core/server/mcpSession";
|
|
import type { EventNdjsonLogger } from "../../provider/EventNdjsonLogger.ts";
|
|
import * as ProviderEventLoggers from "../../provider/ProviderEventLoggers.ts";
|
|
import {
|
|
structuralProtocolMethod,
|
|
summarizeNativeProtocolPayload,
|
|
} from "../../provider/NativeProtocolLogging.ts";
|
|
import { mergeProviderInstanceEnvironment } from "@t3tools/provider-core/server/instanceEnvironment";
|
|
import { t3OrchestrationSystemPrompt } from "@t3tools/provider-core/server/orchestrationInstructions";
|
|
import { buildRuntimeInstructions } from "@t3tools/provider-core/server/runtimeInstructions";
|
|
import * as OpenCodeRuntime from "../../provider/opencodeRuntime.ts";
|
|
import * as IdAllocator from "@t3tools/provider-core/server/IdAllocator";
|
|
import { makeProviderFailure } from "@t3tools/provider-core/server/failure";
|
|
import { turnScopedSelectionTransition } from "@t3tools/provider-core/server/selectionTransition";
|
|
import { providerMessageTextWithAttachmentPaths } from "@t3tools/provider-core/server/attachmentPrompt";
|
|
import * as ProviderAdapter from "@t3tools/provider-core/server/ProviderAdapter";
|
|
import {
|
|
ProviderAdapterDriverCreateError,
|
|
type ProviderAdapterDriver,
|
|
type ProviderAdapterDriverCreateInput,
|
|
} from "@t3tools/provider-core/server/adapterDriver";
|
|
import { makeSubagentChildThread, subagentThreadTitle } from "../SubagentProjection.ts";
|
|
import { openCodeToolTurnItem } from "./OpenCodeToolItems.ts";
|
|
|
|
export { openCodeToolProjectionKind } from "./OpenCodeToolItems.ts";
|
|
|
|
export const OPENCODE_PROVIDER = ProviderDriverKind.make("opencode");
|
|
export const OPENCODE_DEFAULT_INSTANCE_ID = defaultInstanceIdForDriver(OPENCODE_PROVIDER);
|
|
export const OPENCODE_SDK_PROTOCOL = "opencode-sdk.sse" as const;
|
|
const DEFAULT_OPENCODE_SETTINGS = Schema.decodeSync(OpenCodeSettingsSchema)({});
|
|
|
|
let openCodeMessageIdEpochMillis = -1;
|
|
let openCodeMessageIdCounter = 0;
|
|
|
|
const makeOpenCodeMessageId = Effect.fnUntraced(function* () {
|
|
const epochMillis = DateTime.toEpochMillis(yield* DateTime.now);
|
|
if (epochMillis !== openCodeMessageIdEpochMillis) {
|
|
openCodeMessageIdEpochMillis = epochMillis;
|
|
openCodeMessageIdCounter = 0;
|
|
}
|
|
openCodeMessageIdCounter += 1;
|
|
const encodedTime = BigInt.asUintN(
|
|
48,
|
|
BigInt(epochMillis) * 0x1000n + BigInt(openCodeMessageIdCounter),
|
|
)
|
|
.toString(16)
|
|
.padStart(12, "0");
|
|
const alphabet = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz";
|
|
const random = (yield* Effect.forEach(Array.from({ length: 28 }), () =>
|
|
Random.nextIntBetween(0, alphabet.length - 1),
|
|
))
|
|
.map((index) => alphabet[index])
|
|
.join("");
|
|
return `msg_${encodedTime}${random}`;
|
|
});
|
|
|
|
/**
|
|
* OpenCode's session, message, part, and interaction-request identifiers are
|
|
* durable. It does not expose a first-class turn object: the initiating user
|
|
* message is the best native turn correlation point, and session idle is the
|
|
* authoritative terminal signal.
|
|
*/
|
|
const OpenCodeProviderCapabilitiesV2 = {
|
|
sessions: {
|
|
// The current adapter owns one directory-bound client/server per session.
|
|
// Keep it isolated until its runtime is made safe for cross-thread pooling.
|
|
supportsMultipleProviderThreadsPerSession: false,
|
|
supportsModelSwitchInSession: true,
|
|
supportsProviderSwitchingViaHandoff: true,
|
|
supportsRuntimeModeSwitchInSession: false,
|
|
pendingRequestsSurviveRestart: false,
|
|
},
|
|
threads: {
|
|
canCreateEmptyThread: true,
|
|
canReadThreadSnapshot: true,
|
|
canRollbackThread: true,
|
|
canForkThread: true,
|
|
canForkFromTurn: true,
|
|
canForkFromSubagentThread: true,
|
|
exposesNativeThreadId: true,
|
|
},
|
|
turns: {
|
|
exposesNativeTurnId: false,
|
|
emitsTurnStarted: true,
|
|
emitsTurnCompleted: true,
|
|
supportsInterrupt: true,
|
|
supportsActiveSteering: true,
|
|
supportsSteeringByInterruptRestart: true,
|
|
supportsQueuedMessages: true,
|
|
terminalStatusQuality: "strong",
|
|
},
|
|
streaming: {
|
|
streamsAssistantText: true,
|
|
streamsReasoning: true,
|
|
streamsToolOutput: true,
|
|
streamsPlanText: false,
|
|
emitsMessageCompleted: true,
|
|
},
|
|
tools: {
|
|
exposesToolItemIds: true,
|
|
emitsToolStarted: true,
|
|
emitsToolCompleted: true,
|
|
emitsToolOutput: true,
|
|
supportsMcpTools: true,
|
|
supportsDynamicToolCallbacks: false,
|
|
},
|
|
approvals: {
|
|
supportsCommandApproval: true,
|
|
supportsFileReadApproval: true,
|
|
supportsFileChangeApproval: true,
|
|
supportsApplyPatchApproval: true,
|
|
approvalsHaveNativeRequestIds: true,
|
|
approvalCallbacksAreLiveOnly: true,
|
|
approvalsCanOriginateFromSubagents: true,
|
|
},
|
|
planning: {
|
|
emitsPlanUpdated: true,
|
|
emitsTodoList: true,
|
|
emitsProposedPlan: false,
|
|
supportsStructuredQuestions: true,
|
|
planDeltasHaveItemIds: false,
|
|
},
|
|
subagents: {
|
|
supportsSubagents: true,
|
|
exposesSubagentThreadIds: true,
|
|
emitsSubagentLifecycle: true,
|
|
canWaitForSubagents: true,
|
|
canCloseSubagents: false,
|
|
canForkSubagentThread: true,
|
|
},
|
|
context: {
|
|
acceptsSystemContext: false,
|
|
acceptsDeveloperContext: false,
|
|
acceptsSyntheticUserContext: true,
|
|
canGenerateSummaries: true,
|
|
canConsumeHandoffSummaries: true,
|
|
supportsDeltaHandoff: true,
|
|
supportsFullThreadHandoff: true,
|
|
maxRecommendedHandoffChars: null,
|
|
},
|
|
checkpointing: {
|
|
appCanCheckpointFilesystem: true,
|
|
supportsNestedCheckpointScopes: true,
|
|
providerCanRollbackConversation: true,
|
|
providerRollbackReturnsSnapshot: true,
|
|
providerCanReadConversationSnapshot: true,
|
|
},
|
|
identity: {
|
|
nativeThreadIds: "strong",
|
|
nativeTurnIds: "weak",
|
|
nativeItemIds: "strong",
|
|
nativeRequestIds: "strong",
|
|
},
|
|
runtimePolicy: {
|
|
enforcement: "native",
|
|
},
|
|
} satisfies OrchestrationV2ProviderCapabilities;
|
|
|
|
type TerminalTurnStatus = Extract<
|
|
OrchestrationV2ProviderTurn["status"],
|
|
"completed" | "interrupted" | "failed" | "cancelled"
|
|
>;
|
|
|
|
type OpenCodeStepUsage = Pick<
|
|
Extract<OpenCodePart, { readonly type: "step-finish" }>,
|
|
"id" | "tokens"
|
|
>;
|
|
|
|
interface OpenCodeTurnTokenUsageAccumulator {
|
|
readonly partIds: Set<string>;
|
|
readonly promptMessageIds: Set<string>;
|
|
readonly assistantOwnershipByMessageId: Map<string, "owned" | "other" | "unknown">;
|
|
// Native removal does not undo usage. Keep unresolved counts until this turn settles.
|
|
readonly unresolvedStepsByMessageId: Map<string, Map<string, OpenCodeStepUsage>>;
|
|
inputTokens: number;
|
|
cachedInputTokens: number;
|
|
cacheCreationTokens: number;
|
|
outputTokens: number;
|
|
reasoningTokens: number;
|
|
hasSubagents: boolean;
|
|
complete: boolean;
|
|
}
|
|
|
|
function makeOpenCodeTurnTokenUsageAccumulator(): OpenCodeTurnTokenUsageAccumulator {
|
|
return {
|
|
partIds: new Set(),
|
|
promptMessageIds: new Set(),
|
|
assistantOwnershipByMessageId: new Map(),
|
|
unresolvedStepsByMessageId: new Map(),
|
|
inputTokens: 0,
|
|
cachedInputTokens: 0,
|
|
cacheCreationTokens: 0,
|
|
outputTokens: 0,
|
|
reasoningTokens: 0,
|
|
hasSubagents: false,
|
|
complete: true,
|
|
};
|
|
}
|
|
|
|
function accumulateOpenCodeStepUsage(
|
|
accumulator: OpenCodeTurnTokenUsageAccumulator,
|
|
part: OpenCodeStepUsage,
|
|
): void {
|
|
if (accumulator.partIds.has(part.id)) return;
|
|
accumulator.partIds.add(part.id);
|
|
accumulator.inputTokens += part.tokens.input + part.tokens.cache.read + part.tokens.cache.write;
|
|
accumulator.cachedInputTokens += part.tokens.cache.read;
|
|
accumulator.cacheCreationTokens += part.tokens.cache.write;
|
|
accumulator.outputTokens += part.tokens.output + part.tokens.reasoning;
|
|
accumulator.reasoningTokens += part.tokens.reasoning;
|
|
}
|
|
|
|
interface ActiveOpenCodeTurn {
|
|
readonly usage: OpenCodeTurnTokenUsageAccumulator;
|
|
readonly isRoot: boolean;
|
|
readonly threadId: ThreadId;
|
|
readonly runId: ProviderAdapter.ProviderAdapterV2TurnInput["runId"] | null;
|
|
readonly rootNodeId: ProviderAdapter.ProviderAdapterV2TurnInput["rootNodeId"];
|
|
readonly appThread: OrchestrationV2AppThread;
|
|
readonly modelSelection: ModelSelection;
|
|
readonly runtimePolicy: ProviderAdapter.ProviderAdapterV2RuntimePolicy;
|
|
readonly providerTurnId: OrchestrationV2ProviderTurn["id"];
|
|
/** The provider thread the turn started on, which its terminal names. */
|
|
readonly providerThreadId: OrchestrationV2ProviderTurn["providerThreadId"];
|
|
readonly providerTurnOrdinal: number;
|
|
readonly runOrdinal: number;
|
|
readonly runAttemptId: OrchestrationV2ProviderTurn["runAttemptId"];
|
|
readonly startedAt: DateTime.Utc;
|
|
readonly itemOrdinals: Map<string, number>;
|
|
readonly parts: Map<string, Exclude<OpenCodePart, ToolPart>>;
|
|
readonly partIdsByMessage: Map<string, Set<string>>;
|
|
readonly toolNamesByCallId: Map<string, string>;
|
|
mcpServerNames?: ReadonlyArray<string>;
|
|
readonly providerTurn: OrchestrationV2ProviderTurn;
|
|
nextItemOrdinal: number;
|
|
nativeUserMessageId: string | null;
|
|
admissionMessageId: string | null;
|
|
interrupted: boolean;
|
|
finalized: boolean;
|
|
planId: PlanId | null;
|
|
admissionGeneration: number;
|
|
admissionReconciliationGeneration: number | null;
|
|
admissionPending: boolean;
|
|
admissionAccepted: boolean;
|
|
admissionMessageObserved: boolean;
|
|
idleDuringAdmission: boolean;
|
|
admissionSettled: Deferred.Deferred<void>;
|
|
admissionAbortController: AbortController | null;
|
|
}
|
|
|
|
type OpenCodeAdmissionSignal =
|
|
| "accepted"
|
|
| "assistant-completed"
|
|
| "busy"
|
|
| "idle"
|
|
| "user-message";
|
|
type OpenCodeAdmissionAction = "hold" | "reconcile-idle" | "release";
|
|
|
|
export function advanceOpenCodePromptAdmission(
|
|
admission: Pick<
|
|
ActiveOpenCodeTurn,
|
|
"admissionAccepted" | "admissionMessageObserved" | "admissionPending" | "idleDuringAdmission"
|
|
>,
|
|
signal: OpenCodeAdmissionSignal,
|
|
): OpenCodeAdmissionAction {
|
|
if (!admission.admissionPending) return "release";
|
|
if (signal === "assistant-completed") {
|
|
admission.admissionAccepted = true;
|
|
admission.admissionMessageObserved = true;
|
|
admission.admissionPending = false;
|
|
return "release";
|
|
}
|
|
if (signal === "idle") {
|
|
admission.idleDuringAdmission = true;
|
|
return "hold";
|
|
}
|
|
if (signal === "accepted") admission.admissionAccepted = true;
|
|
if (signal === "busy" || signal === "user-message") admission.admissionMessageObserved = true;
|
|
if (!admission.admissionAccepted || !admission.admissionMessageObserved) return "hold";
|
|
if (admission.idleDuringAdmission) return "reconcile-idle";
|
|
admission.admissionPending = false;
|
|
return "release";
|
|
}
|
|
|
|
export function cancelOpenCodePromptAdmission(
|
|
admission: Pick<ActiveOpenCodeTurn, "admissionGeneration" | "admissionPending">,
|
|
nextGeneration: number,
|
|
): void {
|
|
admission.admissionGeneration = nextGeneration;
|
|
admission.admissionPending = false;
|
|
}
|
|
|
|
export const reconcileOpenCodePromptAdmissionStatus = Effect.fn(
|
|
"reconcileOpenCodePromptAdmissionStatus",
|
|
)(function* (
|
|
admission: Pick<ActiveOpenCodeTurn, "admissionGeneration" | "admissionPending">,
|
|
generation: number,
|
|
readStatus: Effect.Effect<"busy" | "idle" | "unknown">,
|
|
) {
|
|
if (admission.admissionGeneration !== generation || !admission.admissionPending) {
|
|
return "stale" as const;
|
|
}
|
|
const status = yield* readStatus;
|
|
if (admission.admissionGeneration !== generation || !admission.admissionPending) {
|
|
return "stale" as const;
|
|
}
|
|
if (status !== "unknown") admission.admissionPending = false;
|
|
return status;
|
|
});
|
|
|
|
interface OpenCodeSubagentContext {
|
|
readonly nativeItemId: string;
|
|
readonly nodeId: OrchestrationV2Subagent["id"];
|
|
readonly parentState: OpenCodeThreadState;
|
|
readonly parentTurn: ActiveOpenCodeTurn;
|
|
readonly prompt: string;
|
|
readonly title: string | null;
|
|
readonly startedAt: DateTime.Utc;
|
|
childSessionId: string | null;
|
|
childThreadId: ThreadId | null;
|
|
childProviderThreadId: OrchestrationV2ProviderThread["id"] | null;
|
|
model: string | null;
|
|
result: string | null;
|
|
}
|
|
|
|
interface OpenCodeThreadState {
|
|
readonly nativeSessionId: string;
|
|
providerThread: OrchestrationV2ProviderThread;
|
|
appThread: OrchestrationV2AppThread | null;
|
|
activeTurn: ActiveOpenCodeTurn | null;
|
|
readonly providerTurns: Map<string, OrchestrationV2ProviderTurn>;
|
|
readonly messages: Map<string, OrchestrationV2ConversationMessage>;
|
|
readonly runtimeRequests: Map<string, OrchestrationV2RuntimeRequest>;
|
|
readonly messageRoles: Map<string, "user" | "assistant">;
|
|
readonly userMessageIds: Array<string>;
|
|
parentSubagent: OpenCodeSubagentContext | null;
|
|
nextChildTurnOrdinal: number;
|
|
nextAdmissionGeneration: number;
|
|
}
|
|
|
|
interface OpenCodeRequestOwner {
|
|
readonly state: OpenCodeThreadState;
|
|
readonly turn: ActiveOpenCodeTurn;
|
|
/** The subagent on the owner's thread that leads to the asking session. */
|
|
readonly subagent: OpenCodeSubagentContext | null;
|
|
}
|
|
|
|
/**
|
|
* The top-level thread and active turn that own a session's requests, walking
|
|
* up through (possibly nested) subagents.
|
|
*/
|
|
function topLevelRequestOwner(state: OpenCodeThreadState): OpenCodeRequestOwner | undefined {
|
|
let owner = state;
|
|
let subagent: OpenCodeSubagentContext | null = null;
|
|
while (owner.parentSubagent !== null) {
|
|
subagent = owner.parentSubagent;
|
|
owner = subagent.parentState;
|
|
}
|
|
return owner.activeTurn === null ? undefined : { state: owner, turn: owner.activeTurn, subagent };
|
|
}
|
|
|
|
interface PendingOpenCodeRequest {
|
|
readonly requestId: RuntimeRequestId;
|
|
readonly nativeRequestId: string;
|
|
/** The session that asked, which may be a subagent's. */
|
|
readonly nativeSessionId: string;
|
|
readonly turn: ActiveOpenCodeTurn;
|
|
readonly state: OpenCodeThreadState;
|
|
readonly nodeId: OrchestrationV2ExecutionNode["id"];
|
|
readonly parentNodeId: OrchestrationV2ExecutionNode["id"];
|
|
readonly turnItemId: OrchestrationV2TurnItem["id"];
|
|
readonly requestKind: OpenCodePermissionRequestKind | "user_input";
|
|
readonly createdAt: DateTime.Utc;
|
|
readonly permission?: PermissionRequest;
|
|
readonly question?: QuestionRequest;
|
|
}
|
|
|
|
export interface OpenCodeAdapterV2Options {
|
|
readonly instanceId: ProviderInstanceId;
|
|
readonly settings: OpenCodeSettings;
|
|
readonly environment: NodeJS.ProcessEnv;
|
|
readonly runtime: OpenCodeRuntime.OpenCodeRuntimeShape;
|
|
readonly idAllocator: IdAllocator.IdAllocatorV2Shape;
|
|
readonly serverConfig: ServerConfig.ServerConfig["Service"];
|
|
readonly nativeEventLogger?: EventNdjsonLogger;
|
|
}
|
|
|
|
export interface OpenCodeProtocolLogEvent {
|
|
readonly direction: "incoming" | "outgoing";
|
|
readonly messageKind: "request" | "response" | "notification" | "error";
|
|
readonly method: string;
|
|
readonly payload: unknown;
|
|
}
|
|
|
|
function formatOpenCodeProtocolLogPayload(event: OpenCodeProtocolLogEvent) {
|
|
return {
|
|
direction: event.direction,
|
|
messageKind: event.messageKind,
|
|
method: structuralProtocolMethod(event.method),
|
|
payload: summarizeNativeProtocolPayload(event.payload),
|
|
};
|
|
}
|
|
|
|
export function makeOpenCodeProtocolLogger(input: {
|
|
readonly nativeEventLogger: EventNdjsonLogger | undefined;
|
|
readonly idAllocator: IdAllocator.IdAllocatorV2Shape;
|
|
readonly providerInstanceId: ProviderInstanceId;
|
|
readonly providerSessionId: ProviderSessionId;
|
|
readonly threadId: ThreadId;
|
|
}): (event: OpenCodeProtocolLogEvent) => Effect.Effect<void, never> {
|
|
return (event) =>
|
|
Effect.gen(function* () {
|
|
if (!input.nativeEventLogger) return;
|
|
const observedAt = DateTime.formatIso(yield* DateTime.now);
|
|
const method = structuralProtocolMethod(event.method);
|
|
yield* input.nativeEventLogger.write(
|
|
{
|
|
observedAt,
|
|
event: {
|
|
id: yield* input.idAllocator.allocate.rawEvent({
|
|
providerSessionId: input.providerSessionId,
|
|
method,
|
|
}),
|
|
kind: "protocol",
|
|
protocol: OPENCODE_SDK_PROTOCOL,
|
|
provider: OPENCODE_PROVIDER,
|
|
providerInstanceId: input.providerInstanceId,
|
|
providerSessionId: input.providerSessionId,
|
|
createdAt: observedAt,
|
|
threadId: input.threadId,
|
|
payload: formatOpenCodeProtocolLogPayload(event),
|
|
},
|
|
},
|
|
input.threadId,
|
|
);
|
|
}).pipe(
|
|
Effect.catchCause((cause) =>
|
|
Cause.hasInterrupts(cause)
|
|
? Effect.interrupt
|
|
: Effect.logWarning("Failed to write native OpenCode event log.", {
|
|
errorTag: causeErrorTag(cause),
|
|
reasonCount: cause.reasons.length,
|
|
provider: OPENCODE_PROVIDER,
|
|
threadId: input.threadId,
|
|
}),
|
|
),
|
|
);
|
|
}
|
|
|
|
function protocolError(
|
|
detail: string,
|
|
payload?: unknown,
|
|
): ProviderAdapter.ProviderAdapterProtocolError {
|
|
return new ProviderAdapter.ProviderAdapterProtocolError({
|
|
driver: OPENCODE_PROVIDER,
|
|
detail,
|
|
...(payload === undefined ? {} : { payload }),
|
|
});
|
|
}
|
|
|
|
function nativeThreadId(providerThread: OrchestrationV2ProviderThread): string {
|
|
const nativeId = providerThread.nativeThreadRef?.nativeId;
|
|
if (nativeId === null || nativeId === undefined) {
|
|
throw protocolError(`Provider thread ${providerThread.id} has no OpenCode session id`);
|
|
}
|
|
return nativeId;
|
|
}
|
|
|
|
function dateTimeFromEpoch(value: number | undefined, fallback: DateTime.Utc): DateTime.Utc {
|
|
if (value === undefined) return fallback;
|
|
return Option.getOrElse(DateTime.make(value), () => fallback);
|
|
}
|
|
|
|
function nonEmptyString(value: unknown): string | undefined {
|
|
return typeof value === "string" && value.trim().length > 0 ? value : undefined;
|
|
}
|
|
|
|
function recordValue(input: unknown, key: string): unknown {
|
|
return typeof input === "object" && input !== null && key in input
|
|
? (input as Record<string, unknown>)[key]
|
|
: undefined;
|
|
}
|
|
|
|
function recordString(input: unknown, ...keys: ReadonlyArray<string>): string | undefined {
|
|
for (const key of keys) {
|
|
const value = nonEmptyString(recordValue(input, key));
|
|
if (value !== undefined) return value;
|
|
}
|
|
return undefined;
|
|
}
|
|
|
|
function sdkResponseForRawLog(value: unknown): unknown {
|
|
if (typeof value !== "object" || value === null) return value;
|
|
if ("data" in value) return { data: (value as { readonly data?: unknown }).data ?? null };
|
|
if ("stream" in value) return { subscribed: true };
|
|
return value;
|
|
}
|
|
|
|
type OpenCodePermissionRequestKind = Extract<
|
|
ProviderRequestKind,
|
|
"command" | "file-read" | "file-change"
|
|
>;
|
|
|
|
export function openCodePermissionRequestKind(
|
|
permission: string,
|
|
toolName?: string,
|
|
): OpenCodePermissionRequestKind {
|
|
const normalized = permission.toLowerCase();
|
|
const normalizedTool = toolName?.toLowerCase() ?? "";
|
|
if (
|
|
normalized === "edit" ||
|
|
normalized === "write" ||
|
|
normalized === "patch" ||
|
|
normalized === "apply_patch" ||
|
|
normalizedTool === "edit" ||
|
|
normalizedTool === "write" ||
|
|
normalizedTool === "patch" ||
|
|
normalizedTool === "apply_patch"
|
|
) {
|
|
return "file-change";
|
|
}
|
|
if (
|
|
normalized === "read" ||
|
|
normalized === "glob" ||
|
|
normalized === "grep" ||
|
|
normalized === "lsp" ||
|
|
normalized === "external_directory" ||
|
|
normalizedTool === "read" ||
|
|
normalizedTool.includes("glob") ||
|
|
normalizedTool.includes("grep") ||
|
|
normalizedTool.includes("search")
|
|
) {
|
|
return "file-read";
|
|
}
|
|
return "command";
|
|
}
|
|
|
|
const OPENCODE_ALWAYS_ALLOWED_PERMISSIONS = [
|
|
"question",
|
|
"read",
|
|
"glob",
|
|
"grep",
|
|
"lsp",
|
|
"todowrite",
|
|
"task",
|
|
"skill",
|
|
] as const;
|
|
|
|
const OPENCODE_RESTRICTED_PERMISSIONS = [
|
|
"bash",
|
|
"edit",
|
|
"webfetch",
|
|
"websearch",
|
|
"codesearch",
|
|
"external_directory",
|
|
"doom_loop",
|
|
] as const;
|
|
|
|
/**
|
|
* OpenCode does not provide an OS sandbox, so permission rules are also the
|
|
* enforcement boundary for non-interactive policies. Read/planning tools are
|
|
* safe by default; edits are auto-approved only for workspace-write modes,
|
|
* while shell/network/external access remains gated unless policy explicitly
|
|
* allows it.
|
|
*/
|
|
export function openCodePermissionRules(
|
|
runtimePolicy: ProviderAdapter.ProviderAdapterV2RuntimePolicy,
|
|
): PermissionRuleset {
|
|
const sandboxPolicy = recordValue(runtimePolicy, "sandboxPolicy");
|
|
const sandboxType = recordString(sandboxPolicy, "type");
|
|
const rawApprovalPolicy = runtimePolicy.approvalPolicy;
|
|
const approvalPolicy = nonEmptyString(rawApprovalPolicy);
|
|
const requiresApproval =
|
|
approvalPolicy === undefined
|
|
? (typeof rawApprovalPolicy === "object" && rawApprovalPolicy !== null) ||
|
|
runtimePolicy.runtimeMode !== "full-access"
|
|
: approvalPolicy !== "never";
|
|
const externallySandboxed = sandboxType === "externalSandbox";
|
|
const dangerFullAccess = sandboxType === "dangerFullAccess";
|
|
const implicitFullAccess =
|
|
sandboxType === undefined && runtimePolicy.runtimeMode === "full-access";
|
|
|
|
if (!requiresApproval && (externallySandboxed || dangerFullAccess || implicitFullAccess)) {
|
|
return [{ permission: "*", pattern: "*", action: "allow" }];
|
|
}
|
|
|
|
// Task sessions initially inherit only parent deny rules. Seed explicit
|
|
// denies before the effective ask/allow overrides so a child is safe during
|
|
// the short interval before emitSubagent installs its complete policy.
|
|
const rules: PermissionRuleset = [
|
|
{ permission: "*", pattern: "*", action: "deny" },
|
|
...OPENCODE_RESTRICTED_PERMISSIONS.map((permission) => ({
|
|
permission,
|
|
pattern: "*",
|
|
action: "deny" as const,
|
|
})),
|
|
];
|
|
|
|
if (requiresApproval) {
|
|
rules.push({ permission: "*", pattern: "*", action: "ask" });
|
|
for (const permission of OPENCODE_RESTRICTED_PERMISSIONS) {
|
|
rules.push({ permission, pattern: "*", action: "ask" });
|
|
}
|
|
}
|
|
|
|
rules.push(
|
|
...OPENCODE_ALWAYS_ALLOWED_PERMISSIONS.map((permission) => ({
|
|
permission,
|
|
pattern: "*",
|
|
action: "allow" as const,
|
|
})),
|
|
);
|
|
|
|
rules.push(
|
|
{ permission: "read", pattern: "*.env", action: requiresApproval ? "ask" : "deny" },
|
|
{ permission: "read", pattern: "*.env.*", action: requiresApproval ? "ask" : "deny" },
|
|
{ permission: "read", pattern: "*.env.example", action: "allow" },
|
|
);
|
|
|
|
if (
|
|
runtimePolicy.runtimeMode === "auto-accept-edits" ||
|
|
(!requiresApproval && sandboxType === "workspaceWrite")
|
|
) {
|
|
rules.push({ permission: "edit", pattern: "*", action: "allow" });
|
|
}
|
|
|
|
if (!requiresApproval && recordValue(sandboxPolicy, "networkAccess") === true) {
|
|
for (const permission of ["webfetch", "websearch", "codesearch"] as const) {
|
|
rules.push({ permission, pattern: "*", action: "allow" });
|
|
}
|
|
}
|
|
|
|
if (!requiresApproval && sandboxType === "readOnly") {
|
|
const access = recordValue(sandboxPolicy, "access");
|
|
if (recordString(access, "type") === "fullAccess") {
|
|
rules.push({ permission: "external_directory", pattern: "*", action: "allow" });
|
|
}
|
|
}
|
|
|
|
if (!requiresApproval && sandboxType === "workspaceWrite") {
|
|
const writableRoots = recordValue(sandboxPolicy, "writableRoots");
|
|
if (Array.isArray(writableRoots)) {
|
|
for (const root of writableRoots) {
|
|
if (typeof root === "string" && root.trim().length > 0) {
|
|
rules.push({
|
|
permission: "external_directory",
|
|
pattern: `${root.replace(/\/$/, "")}/*`,
|
|
action: "allow",
|
|
});
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
return rules;
|
|
}
|
|
|
|
function permissionRuleEquals(
|
|
left: PermissionRuleset[number],
|
|
right: PermissionRuleset[number],
|
|
): boolean {
|
|
return (
|
|
left.permission === right.permission &&
|
|
left.pattern === right.pattern &&
|
|
left.action === right.action
|
|
);
|
|
}
|
|
|
|
/**
|
|
* OpenCode task sessions inherit only the parent's deny/external-directory
|
|
* rules and then add agent-specific restrictions such as disabling nested
|
|
* tasks. Install the complete parent policy while retaining only rules that
|
|
* were added specifically for the selected child agent.
|
|
*/
|
|
export function openCodeChildPermissionRules(
|
|
runtimePolicy: ProviderAdapter.ProviderAdapterV2RuntimePolicy,
|
|
nativeChildRules: PermissionRuleset,
|
|
): PermissionRuleset {
|
|
const parentRules = openCodePermissionRules(runtimePolicy);
|
|
const inheritedRules = parentRules.filter(
|
|
(rule) => rule.permission === "external_directory" || rule.action === "deny",
|
|
);
|
|
const childSpecificRules = nativeChildRules.filter(
|
|
(childRule) =>
|
|
!inheritedRules.some((inheritedRule) => permissionRuleEquals(childRule, inheritedRule)),
|
|
);
|
|
return [...parentRules, ...childSpecificRules];
|
|
}
|
|
|
|
/**
|
|
* OpenCode's fork/revert boundary is exclusive. To retain the selected app
|
|
* turn, address the next native user message; omitting a boundary retains the
|
|
* current head when the selected turn is already last.
|
|
*/
|
|
export function openCodeBoundaryAfterProviderTurn(
|
|
providerTurns: ReadonlyArray<OrchestrationV2ProviderTurn>,
|
|
selectedProviderTurnId: OrchestrationV2ProviderTurn["id"],
|
|
): string | undefined {
|
|
const selected = providerTurns.find((turn) => turn.id === selectedProviderTurnId);
|
|
if (selected === undefined) return undefined;
|
|
return providerTurns
|
|
.filter((turn) => turn.ordinal > selected.ordinal)
|
|
.toSorted((left, right) => left.ordinal - right.ordinal)
|
|
.map((turn) => turn.nativeTurnRef?.nativeId)
|
|
.find((nativeId): nativeId is string => nativeId !== null && nativeId !== undefined);
|
|
}
|
|
|
|
function toolStatus(part: ToolPart): {
|
|
readonly node: OrchestrationV2ExecutionNode["status"];
|
|
readonly item: OrchestrationV2TurnItem["status"];
|
|
} {
|
|
switch (part.state.status) {
|
|
case "pending":
|
|
return { node: "pending", item: "pending" };
|
|
case "running":
|
|
return { node: "running", item: "running" };
|
|
case "completed":
|
|
return { node: "completed", item: "completed" };
|
|
case "error":
|
|
return { node: "failed", item: "failed" };
|
|
}
|
|
}
|
|
|
|
function toolInput(part: ToolPart): Record<string, unknown> {
|
|
return part.state.input;
|
|
}
|
|
|
|
function toolOutput(part: ToolPart): string | undefined {
|
|
if (part.state.status === "completed") return part.state.output;
|
|
if (part.state.status === "error") return part.state.error;
|
|
return undefined;
|
|
}
|
|
|
|
function toolStartedAt(part: ToolPart, now: DateTime.Utc): DateTime.Utc {
|
|
return dateTimeFromEpoch(
|
|
part.state.status === "pending" ? undefined : part.state.time.start,
|
|
now,
|
|
);
|
|
}
|
|
|
|
function toolCompletedAt(part: ToolPart, now: DateTime.Utc): DateTime.Utc | null {
|
|
return part.state.status === "completed" || part.state.status === "error"
|
|
? dateTimeFromEpoch(part.state.time.end, now)
|
|
: null;
|
|
}
|
|
|
|
function toolTitle(part: ToolPart): string | null {
|
|
return part.state.status === "running" || part.state.status === "completed"
|
|
? (part.state.title ?? null)
|
|
: null;
|
|
}
|
|
|
|
function toolModel(part: ToolPart): string | null {
|
|
const metadata =
|
|
part.state.status === "running" ||
|
|
part.state.status === "completed" ||
|
|
part.state.status === "error"
|
|
? part.state.metadata
|
|
: undefined;
|
|
const model = recordValue(metadata, "model");
|
|
const providerId = recordString(model, "providerID", "providerId");
|
|
const modelId = recordString(model, "modelID", "modelId", "id");
|
|
return providerId !== undefined && modelId !== undefined ? `${providerId}/${modelId}` : null;
|
|
}
|
|
|
|
function taskSessionId(part: ToolPart): string | null {
|
|
const metadata =
|
|
part.state.status === "running" ||
|
|
part.state.status === "completed" ||
|
|
part.state.status === "error"
|
|
? part.state.metadata
|
|
: undefined;
|
|
return recordString(metadata, "sessionId", "sessionID") ?? null;
|
|
}
|
|
|
|
function makeProviderThread(input: {
|
|
readonly idAllocator: IdAllocator.IdAllocatorV2Shape;
|
|
readonly providerInstanceId: ProviderInstanceId;
|
|
readonly providerSessionId: OrchestrationV2ProviderThread["providerSessionId"];
|
|
readonly appThreadId: OrchestrationV2ProviderThread["appThreadId"];
|
|
readonly ownerNodeId?: OrchestrationV2ProviderThread["ownerNodeId"];
|
|
readonly nativeSession: OpenCodeSession;
|
|
readonly forkedFrom?: OrchestrationV2ProviderThread["forkedFrom"];
|
|
readonly now: DateTime.Utc;
|
|
}): OrchestrationV2ProviderThread {
|
|
const createdAt = dateTimeFromEpoch(input.nativeSession.time.created, input.now);
|
|
return {
|
|
id: input.idAllocator.derive.providerThread({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeThreadId: input.nativeSession.id,
|
|
}),
|
|
driver: OPENCODE_PROVIDER,
|
|
providerInstanceId: input.providerInstanceId,
|
|
providerSessionId: input.providerSessionId,
|
|
appThreadId: input.appThreadId,
|
|
ownerNodeId: input.ownerNodeId ?? null,
|
|
nativeThreadRef: {
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeId: input.nativeSession.id,
|
|
strength: "strong",
|
|
},
|
|
nativeConversationHeadRef: null,
|
|
status: "idle",
|
|
firstRunOrdinal: null,
|
|
lastRunOrdinal: null,
|
|
handoffIds: [],
|
|
forkedFrom: input.forkedFrom ?? null,
|
|
createdAt,
|
|
updatedAt: dateTimeFromEpoch(input.nativeSession.time.updated, input.now),
|
|
};
|
|
}
|
|
|
|
function providerRef(nativeId: string, strength: "strong" | "weak" = "strong") {
|
|
return {
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeId,
|
|
strength,
|
|
} satisfies OrchestrationV2ProviderRef;
|
|
}
|
|
|
|
function openCodeErrorMessage(event: Extract<OpenCodeEvent, { type: "session.error" }>): string {
|
|
const error = event.properties.error;
|
|
if (error === undefined) return "OpenCode session failed without an error payload.";
|
|
return recordString(error.data, "message") ?? error.name;
|
|
}
|
|
|
|
function terminalStatusForError(
|
|
event: Extract<OpenCodeEvent, { type: "session.error" }>,
|
|
turn: ActiveOpenCodeTurn,
|
|
): TerminalTurnStatus {
|
|
return turn.interrupted || isMessageAbortedError(event) ? "interrupted" : "failed";
|
|
}
|
|
|
|
function isMessageAbortedError(event: Extract<OpenCodeEvent, { type: "session.error" }>): boolean {
|
|
return event.properties.error?.name === "MessageAbortedError";
|
|
}
|
|
|
|
function isOpenCodeNotFound(cause: unknown): boolean {
|
|
const seen = new Set<unknown>();
|
|
const queue: Array<unknown> = [cause];
|
|
for (let steps = 0; queue.length > 0 && steps < 32; steps += 1) {
|
|
const node = queue.shift();
|
|
if (node === null || typeof node !== "object" || seen.has(node)) {
|
|
continue;
|
|
}
|
|
seen.add(node);
|
|
const record = node as Record<string, unknown>;
|
|
|
|
const response = record.response;
|
|
const statuses = [
|
|
record.status,
|
|
record.statusCode,
|
|
response !== null && typeof response === "object"
|
|
? (response as { readonly status?: unknown }).status
|
|
: undefined,
|
|
].filter((status): status is number => typeof status === "number");
|
|
if (statuses.includes(404)) {
|
|
return true;
|
|
}
|
|
if (statuses.length > 0) {
|
|
continue;
|
|
}
|
|
|
|
const name = record.name;
|
|
if (typeof name === "string" && name.toLowerCase() === "notfounderror") {
|
|
return true;
|
|
}
|
|
|
|
for (const key of ["cause", "body", "error", "data"] as const) {
|
|
if (record[key] !== undefined) {
|
|
queue.push(record[key]);
|
|
}
|
|
}
|
|
}
|
|
return false;
|
|
}
|
|
|
|
function unwrapData<A>(operation: string, result: { readonly data?: A }): NonNullable<A> {
|
|
if (result.data === undefined) {
|
|
throw new OpenCodeRuntime.OpenCodeRuntimeError({
|
|
operation,
|
|
detail: `OpenCode ${operation} returned no response payload.`,
|
|
});
|
|
}
|
|
return result.data as NonNullable<A>;
|
|
}
|
|
|
|
export function makeOpenCodeAdapterV2(
|
|
options: OpenCodeAdapterV2Options,
|
|
): ProviderAdapter.ProviderAdapterV2Shape {
|
|
const { idAllocator, runtime, serverConfig } = options;
|
|
|
|
return ProviderAdapter.ProviderAdapterV2.of({
|
|
instanceId: options.instanceId,
|
|
driver: OPENCODE_PROVIDER,
|
|
getCapabilities: () => Effect.succeed(OpenCodeProviderCapabilitiesV2),
|
|
planSelectionTransition: () => Effect.succeed(turnScopedSelectionTransition()),
|
|
openSession: Effect.fn("OpenCodeAdapterV2.openSession")(
|
|
function* (input: ProviderAdapter.ProviderAdapterV2OpenSessionInput) {
|
|
const scope = yield* Effect.scope;
|
|
const cwd = input.runtimePolicy.cwd ?? serverConfig.cwd;
|
|
const connection = yield* runtime.connectToOpenCodeServer({
|
|
binaryPath: options.settings.binaryPath,
|
|
directory: cwd,
|
|
serverUrl: options.settings.serverUrl,
|
|
environment: options.environment,
|
|
});
|
|
const client = runtime.createOpenCodeSdkClient({
|
|
baseUrl: connection.url,
|
|
directory: cwd,
|
|
...(connection.external && options.settings.serverPassword
|
|
? { serverPassword: options.settings.serverPassword }
|
|
: {}),
|
|
});
|
|
|
|
const mcpSession = McpProviderSession.readMcpProviderSession(input.threadId);
|
|
const hasT3Mcp = mcpSession !== undefined && !connection.external;
|
|
const orchestrationSystemPrompt = t3OrchestrationSystemPrompt(hasT3Mcp);
|
|
if (hasT3Mcp) {
|
|
yield* OpenCodeRuntime.runOpenCodeSdk("mcp.add", () =>
|
|
client.mcp.add({
|
|
name: "t3-code",
|
|
config: {
|
|
type: "remote",
|
|
url: mcpSession.endpoint,
|
|
headers: { Authorization: mcpSession.authorizationHeader },
|
|
oauth: false,
|
|
},
|
|
}),
|
|
);
|
|
}
|
|
|
|
const now = yield* DateTime.now;
|
|
let sessionEntity: OrchestrationV2ProviderSession = {
|
|
id: input.providerSessionId,
|
|
driver: OPENCODE_PROVIDER,
|
|
providerInstanceId: options.instanceId,
|
|
status: "ready",
|
|
cwd,
|
|
model: input.modelSelection.model,
|
|
capabilities: OpenCodeProviderCapabilitiesV2,
|
|
createdAt: now,
|
|
updatedAt: now,
|
|
lastError: null,
|
|
};
|
|
const events = yield* Queue.unbounded<ProviderAdapter.ProviderAdapterV2Event, Cause.Done>();
|
|
let nativeStreamFailure: OrchestrationV2ProviderFailure | null = null;
|
|
const threads = new Map<string, OpenCodeThreadState>();
|
|
const commandReceipts = new Map<string, Deferred.Deferred<void>>();
|
|
const commandControllers = new Map<string, Set<AbortController>>();
|
|
const pendingRequests = new Map<string, PendingOpenCodeRequest>();
|
|
const pendingRequestsByNativeId = new Map<string, PendingOpenCodeRequest>();
|
|
const subagentsByNativeItemId = new Map<string, OpenCodeSubagentContext>();
|
|
const subagentsByChildSessionId = new Map<string, OpenCodeSubagentContext>();
|
|
// Permission and question requests can originate from child sessions
|
|
// (task subagents and their descendants). Related sessions map back
|
|
// to the root thread state whose active turn owns the request; asks
|
|
// arriving before the relation is known resolve it via session.get.
|
|
const relatedSessionOwners = new Map<string, OpenCodeThreadState>();
|
|
// Sessions of this runtime that OpenCode reports busy. A background
|
|
// task child keeps running after its parent turn settles, so idle
|
|
// release must not close the server under it.
|
|
const busySessionIds = new Set<string>();
|
|
// Requests settled before their session relation resolved: a late
|
|
// routing attempt must never resurrect them. Bounded because entries
|
|
// only matter for the seconds a routing retry can still be running.
|
|
const settledNativeRequestIds = new Set<string>();
|
|
const pendingChildRequestRoutes = new Set<string>();
|
|
const rememberSettledRequest = (nativeRequestId: string) => {
|
|
if (settledNativeRequestIds.size >= 2_048) settledNativeRequestIds.clear();
|
|
settledNativeRequestIds.add(nativeRequestId);
|
|
};
|
|
const abortController = new AbortController();
|
|
let closing = false;
|
|
let hasConnected = false;
|
|
|
|
const emitProviderEvent = (event: ProviderAdapter.ProviderAdapterV2Event) =>
|
|
Queue.offer(events, event).pipe(Effect.asVoid);
|
|
|
|
const logProtocolEvent = makeOpenCodeProtocolLogger({
|
|
nativeEventLogger: options.nativeEventLogger,
|
|
idAllocator,
|
|
providerInstanceId: options.instanceId,
|
|
providerSessionId: input.providerSessionId,
|
|
threadId: input.threadId,
|
|
});
|
|
|
|
const sdkCall = <A>(
|
|
method: string,
|
|
payload: unknown,
|
|
call: (signal: AbortSignal) => Promise<A>,
|
|
): Effect.Effect<A, OpenCodeRuntime.OpenCodeRuntimeError> =>
|
|
logProtocolEvent({
|
|
direction: "outgoing",
|
|
messageKind: "request",
|
|
method,
|
|
payload,
|
|
}).pipe(
|
|
Effect.andThen(OpenCodeRuntime.runOpenCodeSdk(method, call)),
|
|
Effect.tap((response) =>
|
|
logProtocolEvent({
|
|
direction: "incoming",
|
|
messageKind: "response",
|
|
method,
|
|
payload: sdkResponseForRawLog(response),
|
|
}),
|
|
),
|
|
);
|
|
|
|
const abortDescendants = (rootId: string) =>
|
|
Effect.gen(function* () {
|
|
const visited = new Set([rootId]);
|
|
const semaphore = Semaphore.makeUnsafe(8);
|
|
const visit = (
|
|
sessionId: string,
|
|
abort: boolean,
|
|
): Effect.Effect<OpenCodeRuntime.OpenCodeRuntimeError | undefined> =>
|
|
Effect.gen(function* () {
|
|
const abortResult = abort
|
|
? yield* sdkCall("session.abort", { sessionID: sessionId }, (signal) =>
|
|
client.session.abort({ sessionID: sessionId }, { signal }),
|
|
).pipe(
|
|
semaphore.withPermit,
|
|
Effect.catchIf(isOpenCodeNotFound, () => Effect.void),
|
|
Effect.result,
|
|
)
|
|
: undefined;
|
|
const childrenResult = yield* sdkCall(
|
|
"session.children",
|
|
{ sessionID: sessionId },
|
|
(signal) => client.session.children({ sessionID: sessionId }, { signal }),
|
|
).pipe(
|
|
semaphore.withPermit,
|
|
Effect.catchIf(isOpenCodeNotFound, () => Effect.void),
|
|
Effect.result,
|
|
);
|
|
const firstFailure =
|
|
abortResult?._tag === "Failure" ? abortResult.failure : undefined;
|
|
if (childrenResult._tag === "Failure")
|
|
return firstFailure ?? childrenResult.failure;
|
|
const fresh = (childrenResult.success?.data ?? []).filter((child) => {
|
|
if (visited.has(child.id)) return false;
|
|
visited.add(child.id);
|
|
return true;
|
|
});
|
|
const failures = yield* Effect.forEach(fresh, (child) => visit(child.id, true), {
|
|
concurrency: 8,
|
|
});
|
|
return firstFailure ?? failures.find((failure) => failure !== undefined);
|
|
});
|
|
const failure = yield* visit(rootId, false);
|
|
if (failure) return yield* Effect.fail(failure);
|
|
}).pipe(Effect.timeout("15 seconds"));
|
|
|
|
const updateProviderSession = (
|
|
status: OrchestrationV2ProviderSession["status"],
|
|
lastError: string | null = sessionEntity.lastError,
|
|
) =>
|
|
Effect.gen(function* () {
|
|
const updatedAt = yield* DateTime.now;
|
|
sessionEntity = { ...sessionEntity, status, lastError, updatedAt };
|
|
yield* emitProviderEvent({
|
|
type: "provider_session.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
providerSession: sessionEntity,
|
|
});
|
|
});
|
|
|
|
const updateProviderThread = (
|
|
state: OpenCodeThreadState,
|
|
patch: Partial<OrchestrationV2ProviderThread>,
|
|
) =>
|
|
Effect.gen(function* () {
|
|
const updatedAt = yield* DateTime.now;
|
|
state.providerThread = { ...state.providerThread, ...patch, updatedAt };
|
|
yield* emitProviderEvent({
|
|
type: "provider_thread.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThread: state.providerThread,
|
|
});
|
|
});
|
|
|
|
const itemOrdinal = (turn: ActiveOpenCodeTurn, nativeItemId: string): number => {
|
|
const existing = turn.itemOrdinals.get(nativeItemId);
|
|
if (existing !== undefined) return existing;
|
|
const ordinal = turn.nextItemOrdinal++;
|
|
turn.itemOrdinals.set(nativeItemId, ordinal);
|
|
return ordinal;
|
|
};
|
|
|
|
const emitProviderTurn = (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
status: OrchestrationV2ProviderTurn["status"],
|
|
completedAt: DateTime.Utc | null,
|
|
) => {
|
|
const providerTurn: OrchestrationV2ProviderTurn = {
|
|
...turn.providerTurn,
|
|
nativeTurnRef:
|
|
turn.nativeUserMessageId === null
|
|
? turn.providerTurn.nativeTurnRef
|
|
: providerRef(turn.nativeUserMessageId, "weak"),
|
|
status,
|
|
completedAt,
|
|
...(completedAt === null
|
|
? {}
|
|
: {
|
|
turnTokenUsage:
|
|
turn.usage.partIds.size === 0
|
|
? {
|
|
usageScope: "main_agent" as const,
|
|
usageStatus: "unavailable" as const,
|
|
hasSubagents: turn.usage.hasSubagents,
|
|
}
|
|
: {
|
|
usageScope: "main_agent" as const,
|
|
usageStatus:
|
|
status === "completed" &&
|
|
turn.usage.complete &&
|
|
turn.usage.unresolvedStepsByMessageId.size === 0
|
|
? ("complete" as const)
|
|
: ("partial" as const),
|
|
inputTokens: turn.usage.inputTokens,
|
|
cachedInputTokens: turn.usage.cachedInputTokens,
|
|
cacheCreationTokens: turn.usage.cacheCreationTokens,
|
|
outputTokens: turn.usage.outputTokens,
|
|
reasoningTokens: turn.usage.reasoningTokens,
|
|
hasSubagents: turn.usage.hasSubagents,
|
|
},
|
|
}),
|
|
};
|
|
Object.assign(turn.providerTurn, providerTurn);
|
|
state.providerTurns.set(String(providerTurn.id), providerTurn);
|
|
return emitProviderEvent({
|
|
type: "provider_turn.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
threadId: turn.threadId,
|
|
providerTurn,
|
|
});
|
|
};
|
|
|
|
const emitTextPart = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
part: Extract<OpenCodePart, { type: "text" | "reasoning" }>,
|
|
forceCompleted = false,
|
|
) {
|
|
if (part.type === "text" && (part.ignored === true || part.synthetic === true)) return;
|
|
if (part.text.length === 0) return;
|
|
const emittedAt = yield* DateTime.now;
|
|
const isCompleted = forceCompleted || part.time?.end !== undefined;
|
|
const startedAt = dateTimeFromEpoch(part.time?.start, emittedAt);
|
|
const completedAt = isCompleted ? dateTimeFromEpoch(part.time?.end, emittedAt) : null;
|
|
const nativeItemRef = providerRef(part.id);
|
|
const nodeId = idAllocator.derive.nodeFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.id,
|
|
});
|
|
const turnItemId = idAllocator.derive.turnItemFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.id,
|
|
});
|
|
const ordinal = itemOrdinal(turn, part.id);
|
|
yield* emitProviderEvent({
|
|
type: "node.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
node: {
|
|
id: nodeId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
parentNodeId: turn.rootNodeId,
|
|
rootNodeId: turn.rootNodeId,
|
|
kind: part.type === "text" ? "assistant_message" : "reasoning",
|
|
status: isCompleted ? "completed" : "running",
|
|
countsForRun: false,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef,
|
|
runtimeRequestId: null,
|
|
checkpointScopeId: null,
|
|
startedAt,
|
|
completedAt,
|
|
},
|
|
});
|
|
if (part.type === "text") {
|
|
const messageId = idAllocator.derive.messageFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.id,
|
|
});
|
|
const message: OrchestrationV2ConversationMessage = {
|
|
createdBy: "agent",
|
|
creationSource: "provider",
|
|
id: messageId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
nodeId,
|
|
role: "assistant",
|
|
text: part.text,
|
|
attachments: [],
|
|
streaming: !isCompleted,
|
|
createdAt: startedAt,
|
|
updatedAt: emittedAt,
|
|
};
|
|
state.messages.set(String(message.id), message);
|
|
yield* emitProviderEvent({
|
|
type: "message.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
message,
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem: {
|
|
id: turnItemId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
nodeId,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef,
|
|
parentItemId: null,
|
|
ordinal,
|
|
status: isCompleted ? "completed" : "running",
|
|
title: null,
|
|
startedAt,
|
|
completedAt,
|
|
updatedAt: emittedAt,
|
|
type: "assistant_message",
|
|
messageId,
|
|
text: part.text,
|
|
streaming: !isCompleted,
|
|
},
|
|
});
|
|
return;
|
|
}
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem: {
|
|
id: turnItemId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
nodeId,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef,
|
|
parentItemId: null,
|
|
ordinal,
|
|
status: isCompleted ? "completed" : "running",
|
|
title: null,
|
|
startedAt,
|
|
completedAt,
|
|
updatedAt: emittedAt,
|
|
type: "reasoning",
|
|
text: part.text,
|
|
streaming: !isCompleted,
|
|
},
|
|
});
|
|
});
|
|
|
|
const emitSubagent = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
part: ToolPart,
|
|
) {
|
|
const now = yield* DateTime.now;
|
|
const nativeItemRef = providerRef(part.id);
|
|
const nodeId = idAllocator.derive.nodeFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.id,
|
|
});
|
|
const turnItemId = idAllocator.derive.turnItemFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.id,
|
|
});
|
|
turn.usage.hasSubagents = true;
|
|
const input = toolInput(part);
|
|
const prompt = recordString(input, "prompt") ?? "";
|
|
const title = toolTitle(part) ?? recordString(input, "description") ?? null;
|
|
let context = subagentsByNativeItemId.get(part.id);
|
|
if (context === undefined) {
|
|
context = {
|
|
nativeItemId: part.id,
|
|
nodeId,
|
|
parentState: state,
|
|
parentTurn: turn,
|
|
prompt,
|
|
title,
|
|
startedAt: toolStartedAt(part, now),
|
|
childSessionId: null,
|
|
childThreadId: null,
|
|
childProviderThreadId: null,
|
|
model: null,
|
|
result: null,
|
|
};
|
|
subagentsByNativeItemId.set(part.id, context);
|
|
}
|
|
context.model = toolModel(part) ?? context.model;
|
|
const childSessionId = taskSessionId(part);
|
|
if (childSessionId !== null && context.childSessionId === null) {
|
|
context.childSessionId = childSessionId;
|
|
context.childThreadId = idAllocator.derive.threadFromProviderThread({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeThreadId: childSessionId,
|
|
});
|
|
context.childProviderThreadId = idAllocator.derive.providerThread({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeThreadId: childSessionId,
|
|
});
|
|
subagentsByChildSessionId.set(childSessionId, context);
|
|
relatedSessionOwners.set(childSessionId, state);
|
|
const childModelSelection: ModelSelection = {
|
|
instanceId: options.instanceId,
|
|
model: context.model ?? turn.modelSelection.model,
|
|
};
|
|
const childThread = makeSubagentChildThread({
|
|
parentThread: turn.appThread,
|
|
childThreadId: context.childThreadId,
|
|
parentNodeId: nodeId,
|
|
activeProviderThreadId: context.childProviderThreadId,
|
|
providerInstanceId: options.instanceId,
|
|
modelSelection: childModelSelection,
|
|
title: subagentThreadTitle({
|
|
parentTitle: turn.appThread.title,
|
|
title,
|
|
prompt,
|
|
ordinal: itemOrdinal(turn, part.id),
|
|
}),
|
|
now,
|
|
createdBy: "agent",
|
|
creationSource: "provider",
|
|
});
|
|
const childProviderThread: OrchestrationV2ProviderThread = {
|
|
id: context.childProviderThreadId,
|
|
driver: OPENCODE_PROVIDER,
|
|
providerInstanceId: options.instanceId,
|
|
providerSessionId: inputProviderSessionId,
|
|
appThreadId: context.childThreadId,
|
|
ownerNodeId: nodeId,
|
|
nativeThreadRef: providerRef(childSessionId),
|
|
nativeConversationHeadRef: null,
|
|
status: "active",
|
|
firstRunOrdinal: null,
|
|
lastRunOrdinal: null,
|
|
handoffIds: [],
|
|
forkedFrom: null,
|
|
createdAt: now,
|
|
updatedAt: now,
|
|
};
|
|
const childSessionResponse = yield* sdkCall(
|
|
"session.get",
|
|
{ sessionID: childSessionId },
|
|
() => client.session.get({ sessionID: childSessionId }),
|
|
);
|
|
const nativeChildSession = unwrapData("session.get", childSessionResponse);
|
|
const childPermission = openCodeChildPermissionRules(
|
|
turn.runtimePolicy,
|
|
nativeChildSession.permission ?? [],
|
|
);
|
|
yield* sdkCall(
|
|
"session.update",
|
|
{ sessionID: childSessionId, permission: childPermission },
|
|
() =>
|
|
client.session.update({
|
|
sessionID: childSessionId,
|
|
permission: childPermission,
|
|
}),
|
|
);
|
|
threads.set(childSessionId, {
|
|
nativeSessionId: childSessionId,
|
|
providerThread: childProviderThread,
|
|
appThread: childThread,
|
|
activeTurn: null,
|
|
providerTurns: new Map(),
|
|
messages: new Map(),
|
|
runtimeRequests: new Map(),
|
|
messageRoles: new Map(),
|
|
userMessageIds: [],
|
|
parentSubagent: context,
|
|
nextChildTurnOrdinal: 1,
|
|
nextAdmissionGeneration: 1,
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "app_thread.created",
|
|
driver: OPENCODE_PROVIDER,
|
|
appThread: childThread,
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "provider_thread.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThread: childProviderThread,
|
|
});
|
|
}
|
|
const output = toolOutput(part);
|
|
if (part.state.status === "completed" && output !== undefined) context.result = output;
|
|
const status = toolStatus(part);
|
|
const completedAt = toolCompletedAt(part, now);
|
|
const subagentStatus: OrchestrationV2Subagent["status"] =
|
|
status.item === "failed"
|
|
? "failed"
|
|
: status.item === "completed"
|
|
? "completed"
|
|
: status.item === "pending"
|
|
? "pending"
|
|
: "running";
|
|
const subagent: OrchestrationV2Subagent = {
|
|
id: nodeId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
parentNodeId: turn.rootNodeId,
|
|
origin: "provider_native",
|
|
createdBy: "agent",
|
|
driver: OPENCODE_PROVIDER,
|
|
providerInstanceId: options.instanceId,
|
|
providerThreadId: context.childProviderThreadId,
|
|
childThreadId: context.childThreadId,
|
|
nativeTaskRef: nativeItemRef,
|
|
prompt,
|
|
title,
|
|
model: context.model,
|
|
status: subagentStatus,
|
|
result: context.result,
|
|
startedAt: context.startedAt,
|
|
completedAt,
|
|
updatedAt: now,
|
|
};
|
|
yield* emitProviderEvent({
|
|
type: "node.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
node: {
|
|
id: nodeId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
parentNodeId: turn.rootNodeId,
|
|
rootNodeId: turn.rootNodeId,
|
|
kind: "subagent",
|
|
status: status.node,
|
|
countsForRun: false,
|
|
providerThreadId: context.childProviderThreadId ?? state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef,
|
|
runtimeRequestId: null,
|
|
checkpointScopeId: null,
|
|
startedAt: context.startedAt,
|
|
completedAt,
|
|
},
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "subagent.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
subagent,
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem: {
|
|
id: turnItemId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
nodeId,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef,
|
|
parentItemId: null,
|
|
ordinal: itemOrdinal(turn, part.id),
|
|
status: status.item,
|
|
title,
|
|
startedAt: context.startedAt,
|
|
completedAt,
|
|
updatedAt: now,
|
|
type: "subagent",
|
|
subagentId: nodeId,
|
|
origin: "provider_native",
|
|
driver: OPENCODE_PROVIDER,
|
|
providerInstanceId: options.instanceId,
|
|
childThreadId: context.childThreadId,
|
|
prompt,
|
|
result: context.result,
|
|
},
|
|
});
|
|
});
|
|
|
|
const emitToolPart = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
part: ToolPart,
|
|
) {
|
|
const normalizedTool = part.tool.toLowerCase();
|
|
if (normalizedTool === "task") {
|
|
yield* emitSubagent(state, turn, part);
|
|
return;
|
|
}
|
|
// question.asked carries the respondable semantic item. Projecting
|
|
// the implementation tool as well would duplicate it in the UI.
|
|
if (normalizedTool === "question") return;
|
|
const now = yield* DateTime.now;
|
|
const status = toolStatus(part);
|
|
const startedAt = toolStartedAt(part, now);
|
|
const completedAt = toolCompletedAt(part, now);
|
|
const nativeItemRef = providerRef(part.id);
|
|
const nodeId = idAllocator.derive.nodeFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.id,
|
|
});
|
|
const turnItemId = idAllocator.derive.turnItemFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.id,
|
|
});
|
|
const base = {
|
|
id: turnItemId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
nodeId,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef,
|
|
parentItemId: null,
|
|
ordinal: itemOrdinal(turn, part.id),
|
|
status: status.item,
|
|
title: toolTitle(part),
|
|
startedAt,
|
|
completedAt,
|
|
updatedAt: now,
|
|
} satisfies Pick<
|
|
OrchestrationV2TurnItem,
|
|
| "id"
|
|
| "threadId"
|
|
| "runId"
|
|
| "nodeId"
|
|
| "providerThreadId"
|
|
| "providerTurnId"
|
|
| "nativeItemRef"
|
|
| "parentItemId"
|
|
| "ordinal"
|
|
| "status"
|
|
| "title"
|
|
| "startedAt"
|
|
| "completedAt"
|
|
| "updatedAt"
|
|
>;
|
|
const input = toolInput(part);
|
|
const output = toolOutput(part);
|
|
const isNativeTool = part.tool === "code_search" || part.tool === "apply_patch";
|
|
if (!isNativeTool && part.tool.includes("_") && turn.mcpServerNames === undefined) {
|
|
const serverNames = yield* OpenCodeRuntime.runOpenCodeSdk("mcp.status", (signal) =>
|
|
client.mcp.status(undefined, { signal, throwOnError: true }),
|
|
).pipe(
|
|
Effect.timeout("1 second"),
|
|
Effect.map((response) => Object.keys(response.data ?? {})),
|
|
Effect.catch(() => Effect.succeed(undefined)),
|
|
);
|
|
if (serverNames !== undefined) turn.mcpServerNames = serverNames;
|
|
}
|
|
const matchingServers = turn.mcpServerNames?.filter(
|
|
(name) =>
|
|
!isNativeTool && part.tool.startsWith(`${name.replace(/[^a-zA-Z0-9_-]/g, "_")}_`),
|
|
);
|
|
const serverName = matchingServers?.length === 1 ? matchingServers[0] : undefined;
|
|
const presentation =
|
|
serverName === undefined
|
|
? {}
|
|
: mcpToolPresentation({
|
|
serverName,
|
|
toolName: part.tool.slice(serverName.replace(/[^a-zA-Z0-9_-]/g, "_").length + 1),
|
|
title: toolTitle(part) === part.tool ? undefined : toolTitle(part),
|
|
});
|
|
const turnItem: OrchestrationV2TurnItem = matchingServers?.length
|
|
? {
|
|
...base,
|
|
type: "dynamic_tool",
|
|
...presentation,
|
|
toolName: part.tool,
|
|
input,
|
|
...(output === undefined ? {} : { output }),
|
|
}
|
|
: openCodeToolTurnItem(base, {
|
|
name: part.tool,
|
|
input,
|
|
output,
|
|
completedMetadata:
|
|
part.state.status === "completed" ? part.state.metadata : undefined,
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "node.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
node: {
|
|
id: nodeId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
parentNodeId: turn.rootNodeId,
|
|
rootNodeId: turn.rootNodeId,
|
|
kind: "tool_call",
|
|
status: status.node,
|
|
countsForRun: false,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef,
|
|
runtimeRequestId: null,
|
|
checkpointScopeId: null,
|
|
startedAt,
|
|
completedAt,
|
|
},
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem,
|
|
});
|
|
});
|
|
|
|
const emitTodo = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
todos: ReadonlyArray<OpenCodeTodo>,
|
|
) {
|
|
const now = yield* DateTime.now;
|
|
if (turn.planId === null) {
|
|
turn.planId = yield* idAllocator.allocate.plan({
|
|
threadId: turn.threadId,
|
|
...(turn.runId === null ? {} : { runId: turn.runId }),
|
|
driver: OPENCODE_PROVIDER,
|
|
});
|
|
}
|
|
const planId = turn.planId;
|
|
const nativeItemId = `${state.nativeSessionId}:todo:${turn.providerTurnId}`;
|
|
const nodeId = idAllocator.derive.nodeFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId,
|
|
});
|
|
const turnItemId = idAllocator.derive.turnItemFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId,
|
|
});
|
|
const steps: Array<OrchestrationV2PlanStep> = todos.map((todo, index) => ({
|
|
id: `${nativeItemId}:${index + 1}`,
|
|
text: todo.content.trim() || `Todo ${index + 1}`,
|
|
status:
|
|
todo.status === "completed"
|
|
? "completed"
|
|
: todo.status === "in_progress"
|
|
? "running"
|
|
: "pending",
|
|
}));
|
|
const completed = steps.length > 0 && steps.every((step) => step.status === "completed");
|
|
yield* emitProviderEvent({
|
|
type: "node.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
node: {
|
|
id: nodeId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
parentNodeId: turn.rootNodeId,
|
|
rootNodeId: turn.rootNodeId,
|
|
kind: "todo_list",
|
|
status: completed ? "completed" : "running",
|
|
countsForRun: false,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef: providerRef(nativeItemId, "weak"),
|
|
runtimeRequestId: null,
|
|
checkpointScopeId: null,
|
|
startedAt: turn.startedAt,
|
|
completedAt: completed ? now : null,
|
|
},
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "plan.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
plan: {
|
|
id: planId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
nodeId,
|
|
status: completed ? "completed" : "active",
|
|
kind: "todo_list",
|
|
steps,
|
|
},
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem: {
|
|
id: turnItemId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
nodeId,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef: providerRef(nativeItemId, "weak"),
|
|
parentItemId: null,
|
|
ordinal: itemOrdinal(turn, nativeItemId),
|
|
status: completed ? "completed" : "running",
|
|
title: "Todo list",
|
|
startedAt: turn.startedAt,
|
|
completedAt: completed ? now : null,
|
|
updatedAt: now,
|
|
type: "todo_list",
|
|
planId,
|
|
steps,
|
|
},
|
|
});
|
|
});
|
|
|
|
const requestQuestions = (request: QuestionRequest) =>
|
|
request.questions.map((question, index) => ({
|
|
id: OpenCodeRuntime.openCodeQuestionId(index, question),
|
|
header: question.header.trim() || `Question ${index + 1}`,
|
|
question: question.question.trim() || question.header.trim() || `Question ${index + 1}`,
|
|
options: question.options.map((option) => ({
|
|
label: option.label.trim() || "Option",
|
|
description: option.description.trim() || option.label.trim() || "Option",
|
|
})),
|
|
}));
|
|
|
|
const runtimeRequestTurnItem = (
|
|
pending: PendingOpenCodeRequest,
|
|
status: OrchestrationV2TurnItem["status"],
|
|
completedAt: DateTime.Utc | null,
|
|
updatedAt: DateTime.Utc,
|
|
): OrchestrationV2TurnItem => {
|
|
const base = {
|
|
id: pending.turnItemId,
|
|
threadId: pending.turn.threadId,
|
|
runId: pending.turn.runId,
|
|
nodeId: pending.nodeId,
|
|
providerThreadId: pending.state.providerThread.id,
|
|
providerTurnId: pending.turn.providerTurnId,
|
|
nativeItemRef: providerRef(pending.nativeRequestId),
|
|
parentItemId: null,
|
|
ordinal: itemOrdinal(pending.turn, pending.nativeRequestId),
|
|
status,
|
|
startedAt: pending.createdAt,
|
|
completedAt,
|
|
updatedAt,
|
|
};
|
|
if (pending.question !== undefined) {
|
|
return {
|
|
...base,
|
|
title: "User input",
|
|
type: "user_input_request",
|
|
requestId: pending.requestId,
|
|
questions: requestQuestions(pending.question),
|
|
};
|
|
}
|
|
const permission = pending.permission;
|
|
if (permission === undefined) {
|
|
throw protocolError(`OpenCode request ${pending.requestId} has no native payload`);
|
|
}
|
|
return {
|
|
...base,
|
|
title: permission.permission,
|
|
type: "approval_request",
|
|
requestId: pending.requestId,
|
|
requestKind: pending.requestKind === "user_input" ? "command" : pending.requestKind,
|
|
prompt:
|
|
permission.patterns.length === 0
|
|
? permission.permission
|
|
: permission.patterns.join("\n"),
|
|
};
|
|
};
|
|
|
|
const emitRuntimeRequest = Effect.fnUntraced(function* (
|
|
owner: OpenCodeRequestOwner,
|
|
nativeRequestId: string,
|
|
request:
|
|
| { readonly type: "permission"; readonly value: PermissionRequest }
|
|
| { readonly type: "question"; readonly value: QuestionRequest },
|
|
) {
|
|
if (pendingRequestsByNativeId.has(nativeRequestId)) return;
|
|
const { state, turn, subagent } = owner;
|
|
const now = yield* DateTime.now;
|
|
const requestId = yield* idAllocator.allocate.runtimeRequest({
|
|
driver: OPENCODE_PROVIDER,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeRequestId,
|
|
});
|
|
const nodeId = idAllocator.derive.approvalNode({ requestId });
|
|
const turnItemId = idAllocator.derive.approvalTurnItem({ requestId });
|
|
// The tool call belongs to the asking session's turn, not the owner's.
|
|
const permissionToolName =
|
|
request.type === "permission" && request.value.tool !== undefined
|
|
? threads
|
|
.get(request.value.sessionID)
|
|
?.activeTurn?.toolNamesByCallId.get(request.value.tool.callID)
|
|
: undefined;
|
|
const permissionRequestKind =
|
|
request.type === "permission"
|
|
? openCodePermissionRequestKind(request.value.permission, permissionToolName)
|
|
: undefined;
|
|
const requestKind: OrchestrationV2RuntimeRequest["kind"] =
|
|
permissionRequestKind ?? "user_input";
|
|
const pending: PendingOpenCodeRequest = {
|
|
requestId,
|
|
nativeRequestId,
|
|
nativeSessionId: request.value.sessionID,
|
|
turn,
|
|
state,
|
|
nodeId,
|
|
parentNodeId: subagent?.nodeId ?? turn.rootNodeId,
|
|
turnItemId,
|
|
requestKind,
|
|
createdAt: now,
|
|
...(request.type === "permission"
|
|
? { permission: request.value }
|
|
: { question: request.value }),
|
|
};
|
|
pendingRequests.set(String(requestId), pending);
|
|
pendingRequestsByNativeId.set(nativeRequestId, pending);
|
|
const runtimeRequest: OrchestrationV2RuntimeRequest = {
|
|
id: requestId,
|
|
nodeId,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeRequestRef: providerRef(nativeRequestId),
|
|
kind: requestKind,
|
|
status: "pending",
|
|
responseCapability: {
|
|
type: "live",
|
|
providerSessionId: inputProviderSessionId,
|
|
},
|
|
createdAt: now,
|
|
resolvedAt: null,
|
|
};
|
|
state.runtimeRequests.set(String(requestId), runtimeRequest);
|
|
yield* emitProviderEvent({
|
|
type: "node.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
node: {
|
|
id: nodeId,
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
parentNodeId: pending.parentNodeId,
|
|
rootNodeId: turn.rootNodeId,
|
|
kind: request.type === "question" ? "user_input_request" : "approval_request",
|
|
status: "waiting",
|
|
countsForRun: false,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef: providerRef(nativeRequestId),
|
|
runtimeRequestId: requestId,
|
|
checkpointScopeId: null,
|
|
startedAt: now,
|
|
completedAt: null,
|
|
},
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "runtime_request.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
threadId: turn.threadId,
|
|
runtimeRequest,
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem: runtimeRequestTurnItem(pending, "waiting", null, now),
|
|
});
|
|
yield* updateProviderSession("waiting", null);
|
|
});
|
|
|
|
const resolveRuntimeRequest = Effect.fnUntraced(function* (
|
|
nativeRequestId: string,
|
|
status: "resolved" | "cancelled",
|
|
) {
|
|
const pending = pendingRequestsByNativeId.get(nativeRequestId);
|
|
if (pending === undefined) return;
|
|
const now = yield* DateTime.now;
|
|
const current = pending.state.runtimeRequests.get(String(pending.requestId));
|
|
if (current !== undefined) {
|
|
const resolved: OrchestrationV2RuntimeRequest = {
|
|
...current,
|
|
status,
|
|
resolvedAt: now,
|
|
};
|
|
pending.state.runtimeRequests.set(String(pending.requestId), resolved);
|
|
yield* emitProviderEvent({
|
|
type: "runtime_request.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
threadId: pending.turn.threadId,
|
|
runtimeRequest: resolved,
|
|
});
|
|
}
|
|
yield* emitProviderEvent({
|
|
type: "node.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
node: {
|
|
id: pending.nodeId,
|
|
threadId: pending.turn.threadId,
|
|
runId: pending.turn.runId,
|
|
parentNodeId: pending.parentNodeId,
|
|
rootNodeId: pending.turn.rootNodeId,
|
|
kind: pending.question === undefined ? "approval_request" : "user_input_request",
|
|
status: status === "resolved" ? "completed" : "cancelled",
|
|
countsForRun: false,
|
|
providerThreadId: pending.state.providerThread.id,
|
|
providerTurnId: pending.turn.providerTurnId,
|
|
nativeItemRef: providerRef(nativeRequestId),
|
|
runtimeRequestId: pending.requestId,
|
|
checkpointScopeId: null,
|
|
startedAt: pending.createdAt,
|
|
completedAt: now,
|
|
},
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem: runtimeRequestTurnItem(
|
|
pending,
|
|
status === "resolved" ? "completed" : "cancelled",
|
|
now,
|
|
now,
|
|
),
|
|
});
|
|
pendingRequests.delete(String(pending.requestId));
|
|
pendingRequestsByNativeId.delete(nativeRequestId);
|
|
const hasOtherPending = Array.from(pendingRequests.values()).some(
|
|
(candidate) => candidate.turn.isRoot,
|
|
);
|
|
if (!hasOtherPending) yield* updateProviderSession("running", null);
|
|
});
|
|
|
|
/** Resolve the thread state a session belongs to: its own, or for a
|
|
* child not registered yet, the nearest known ancestor found by walking
|
|
* the native parent chain. Registers every hop so later requests from
|
|
* the same child resolve without another lookup. */
|
|
const resolveSessionOwner = Effect.fnUntraced(function* (sessionId: string) {
|
|
const known = threads.get(sessionId) ?? relatedSessionOwners.get(sessionId);
|
|
if (known !== undefined) return known;
|
|
let cursor = sessionId;
|
|
const hops: string[] = [];
|
|
for (let depth = 0; depth < 5; depth += 1) {
|
|
const response = yield* sdkCall("session.get", { sessionID: cursor }, () =>
|
|
client.session.get({ sessionID: cursor }),
|
|
).pipe(Effect.option);
|
|
const info = Option.getOrUndefined(response)?.data;
|
|
const parentId = info?.parentID;
|
|
if (parentId === undefined) return undefined;
|
|
hops.push(cursor);
|
|
const owner = threads.get(parentId) ?? relatedSessionOwners.get(parentId);
|
|
if (owner !== undefined) {
|
|
for (const hop of hops) relatedSessionOwners.set(hop, owner);
|
|
return owner;
|
|
}
|
|
cursor = parentId;
|
|
}
|
|
return undefined;
|
|
});
|
|
|
|
/** Every request is asked on the top-level thread and its active turn,
|
|
* under the subagent that leads to the asking session, because native
|
|
* subagent threads are hidden from the sidebar. A child's request can
|
|
* arrive before the task part or session.created event that reveals
|
|
* its relation to a thread. The first resolution attempt runs inline
|
|
* (the replayable path); if the relation or the owning turn is not
|
|
* established yet, a short forked backoff keeps trying instead of
|
|
* dropping the request. */
|
|
const routeRuntimeRequest = Effect.fnUntraced(function* (
|
|
nativeRequestId: string,
|
|
sessionId: string,
|
|
request:
|
|
| { readonly type: "permission"; readonly value: PermissionRequest }
|
|
| { readonly type: "question"; readonly value: QuestionRequest },
|
|
) {
|
|
if (pendingChildRequestRoutes.has(nativeRequestId)) return;
|
|
const attempt = Effect.gen(function* () {
|
|
if (
|
|
settledNativeRequestIds.has(nativeRequestId) ||
|
|
pendingRequestsByNativeId.has(nativeRequestId)
|
|
) {
|
|
return true;
|
|
}
|
|
const state = yield* resolveSessionOwner(sessionId);
|
|
const owner = state === undefined ? undefined : topLevelRequestOwner(state);
|
|
if (owner === undefined) return false;
|
|
yield* emitRuntimeRequest(owner, nativeRequestId, request);
|
|
return true;
|
|
});
|
|
if (yield* attempt) return;
|
|
pendingChildRequestRoutes.add(nativeRequestId);
|
|
yield* Effect.gen(function* () {
|
|
for (let retry = 0; retry < 5; retry += 1) {
|
|
yield* Effect.sleep(Duration.millis(Math.min(200 * 2 ** retry, 2_000)));
|
|
if (yield* attempt) return;
|
|
}
|
|
}).pipe(
|
|
Effect.ensuring(Effect.sync(() => pendingChildRequestRoutes.delete(nativeRequestId))),
|
|
Effect.forkIn(scope),
|
|
);
|
|
});
|
|
|
|
const finalizeTurn = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
status: TerminalTurnStatus,
|
|
terminal?: {
|
|
readonly failure?: OrchestrationV2ProviderFailure;
|
|
readonly threadDisposition?: "reusable" | "broken";
|
|
},
|
|
) {
|
|
if (turn.finalized) return;
|
|
if (nativeStreamFailure !== null) {
|
|
status = "failed";
|
|
terminal = { failure: nativeStreamFailure, threadDisposition: "broken" };
|
|
}
|
|
turn.finalized = true;
|
|
const completedAt = yield* DateTime.now;
|
|
for (const part of turn.parts.values()) {
|
|
if (part.type === "text" || part.type === "reasoning") {
|
|
yield* emitTextPart(state, turn, part, true);
|
|
}
|
|
}
|
|
for (const pending of Array.from(pendingRequests.values())) {
|
|
if (
|
|
pending.turn.providerTurnId === turn.providerTurnId ||
|
|
pending.nativeSessionId === state.nativeSessionId
|
|
) {
|
|
yield* resolveRuntimeRequest(pending.nativeRequestId, "cancelled");
|
|
}
|
|
}
|
|
yield* emitProviderTurn(state, turn, status, completedAt);
|
|
const threadDisposition = terminal?.threadDisposition ?? "reusable";
|
|
yield* updateProviderThread(state, {
|
|
status: turn.isRoot ? "active" : threadDisposition === "broken" ? "error" : "idle",
|
|
nativeConversationHeadRef:
|
|
turn.nativeUserMessageId === null
|
|
? state.providerThread.nativeConversationHeadRef
|
|
: providerRef(turn.nativeUserMessageId, "weak"),
|
|
});
|
|
state.activeTurn = null;
|
|
if (!turn.isRoot) {
|
|
yield* emitProviderEvent({
|
|
type: "node.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
node: {
|
|
id: turn.rootNodeId,
|
|
threadId: turn.threadId,
|
|
runId: null,
|
|
parentNodeId: null,
|
|
rootNodeId: turn.rootNodeId,
|
|
kind: "root_turn",
|
|
status,
|
|
countsForRun: false,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef: providerRef(state.nativeSessionId),
|
|
runtimeRequestId: null,
|
|
checkpointScopeId: null,
|
|
startedAt: turn.startedAt,
|
|
completedAt,
|
|
},
|
|
});
|
|
return;
|
|
}
|
|
const anotherTurnIsActive = Array.from(threads.values()).some(
|
|
(candidate) => candidate.activeTurn?.isRoot === true,
|
|
);
|
|
yield* updateProviderSession(
|
|
anotherTurnIsActive ? "running" : status === "failed" ? "error" : "ready",
|
|
status === "failed" ? sessionEntity.lastError : null,
|
|
);
|
|
yield* emitProviderEvent(
|
|
status === "failed"
|
|
? {
|
|
type: "turn.terminal",
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThreadId: turn.providerThreadId,
|
|
providerTurnId: turn.providerTurnId,
|
|
runOrdinal: turn.runOrdinal,
|
|
failureItemOrdinal: itemOrdinal(turn, `terminal-failure:${turn.providerTurnId}`),
|
|
status,
|
|
failure:
|
|
terminal?.failure ??
|
|
makeProviderFailure({
|
|
message: sessionEntity.lastError ?? undefined,
|
|
class: "provider_error",
|
|
}),
|
|
threadDisposition,
|
|
}
|
|
: {
|
|
type: "turn.terminal",
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThreadId: turn.providerThreadId,
|
|
providerTurnId: turn.providerTurnId,
|
|
runOrdinal: turn.runOrdinal,
|
|
status,
|
|
failure: null,
|
|
threadDisposition,
|
|
},
|
|
);
|
|
});
|
|
|
|
const promptAdmissionIsCurrent = (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
generation: number,
|
|
) =>
|
|
state.activeTurn === turn &&
|
|
!turn.finalized &&
|
|
!turn.interrupted &&
|
|
turn.admissionPending &&
|
|
turn.admissionGeneration === generation;
|
|
|
|
const runPromptAdmissionReconciliation = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
generation: number,
|
|
) {
|
|
while (promptAdmissionIsCurrent(state, turn, generation)) {
|
|
const status = yield* reconcileOpenCodePromptAdmissionStatus(
|
|
turn,
|
|
generation,
|
|
sdkCall("session.status", { sessionID: state.nativeSessionId, generation }, () =>
|
|
client.session.status(),
|
|
).pipe(
|
|
Effect.match({
|
|
onFailure: () => "unknown" as const,
|
|
onSuccess: (response) => {
|
|
const statuses = unwrapData("session.status", response);
|
|
const sessionStatus = statuses[state.nativeSessionId];
|
|
return sessionStatus === undefined || sessionStatus.type === "idle"
|
|
? ("idle" as const)
|
|
: ("busy" as const);
|
|
},
|
|
}),
|
|
),
|
|
);
|
|
if (
|
|
state.activeTurn !== turn ||
|
|
status === "stale" ||
|
|
turn.finalized ||
|
|
turn.interrupted ||
|
|
turn.admissionGeneration !== generation
|
|
) {
|
|
return;
|
|
}
|
|
if (status === "idle") {
|
|
yield* finalizeTurn(state, turn, "completed");
|
|
return;
|
|
}
|
|
if (status === "busy") return;
|
|
yield* Effect.sleep("250 millis");
|
|
}
|
|
});
|
|
|
|
const reconcilePromptAdmission = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
) {
|
|
const generation = turn.admissionGeneration;
|
|
if (
|
|
!promptAdmissionIsCurrent(state, turn, generation) ||
|
|
turn.admissionReconciliationGeneration === generation
|
|
) {
|
|
return;
|
|
}
|
|
turn.admissionReconciliationGeneration = generation;
|
|
yield* runPromptAdmissionReconciliation(state, turn, generation).pipe(
|
|
Effect.ensuring(
|
|
Effect.sync(() => {
|
|
if (turn.admissionReconciliationGeneration === generation) {
|
|
turn.admissionReconciliationGeneration = null;
|
|
}
|
|
}),
|
|
),
|
|
Effect.forkIn(scope),
|
|
);
|
|
});
|
|
|
|
const createChildTurn = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
message: Extract<OpenCodeMessage, { role: "user" }>,
|
|
) {
|
|
if (state.appThread === null || state.parentSubagent === null) return null;
|
|
const now = yield* DateTime.now;
|
|
const startedAt = dateTimeFromEpoch(message.time.created, now);
|
|
const rootNodeId = idAllocator.derive.nodeFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: `${state.nativeSessionId}:root:${message.id}`,
|
|
});
|
|
const providerTurnId = idAllocator.derive.providerTurn({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeTurnId: message.id,
|
|
});
|
|
const providerTurn: OrchestrationV2ProviderTurn = {
|
|
id: providerTurnId,
|
|
providerThreadId: state.providerThread.id,
|
|
nodeId: rootNodeId,
|
|
runAttemptId: null,
|
|
nativeTurnRef: providerRef(message.id, "weak"),
|
|
ordinal: state.nextChildTurnOrdinal++,
|
|
status: "running",
|
|
startedAt,
|
|
completedAt: null,
|
|
};
|
|
const turn: ActiveOpenCodeTurn = {
|
|
isRoot: false,
|
|
threadId: state.appThread.id,
|
|
runId: null,
|
|
rootNodeId,
|
|
appThread: state.appThread,
|
|
modelSelection: state.appThread.modelSelection,
|
|
runtimePolicy: state.parentSubagent.parentTurn.runtimePolicy,
|
|
providerTurnId,
|
|
providerThreadId: providerTurn.providerThreadId,
|
|
providerTurnOrdinal: providerTurn.ordinal,
|
|
runOrdinal: state.parentSubagent.parentTurn.runOrdinal,
|
|
runAttemptId: null,
|
|
startedAt,
|
|
itemOrdinals: new Map(),
|
|
usage: makeOpenCodeTurnTokenUsageAccumulator(),
|
|
parts: new Map(),
|
|
partIdsByMessage: new Map(),
|
|
toolNamesByCallId: new Map(),
|
|
providerTurn,
|
|
nextItemOrdinal: 1,
|
|
nativeUserMessageId: message.id,
|
|
admissionMessageId: null,
|
|
interrupted: false,
|
|
finalized: false,
|
|
planId: null,
|
|
admissionGeneration: 0,
|
|
admissionReconciliationGeneration: null,
|
|
admissionPending: false,
|
|
admissionAccepted: true,
|
|
admissionMessageObserved: true,
|
|
idleDuringAdmission: false,
|
|
admissionSettled: Deferred.makeUnsafe<void>(),
|
|
admissionAbortController: null,
|
|
};
|
|
state.activeTurn = turn;
|
|
state.providerTurns.set(String(providerTurnId), providerTurn);
|
|
yield* emitProviderEvent({
|
|
type: "node.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
node: {
|
|
id: rootNodeId,
|
|
threadId: turn.threadId,
|
|
runId: null,
|
|
parentNodeId: null,
|
|
rootNodeId,
|
|
kind: "root_turn",
|
|
status: "running",
|
|
countsForRun: false,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId,
|
|
nativeItemRef: providerRef(message.id, "weak"),
|
|
runtimeRequestId: null,
|
|
checkpointScopeId: null,
|
|
startedAt,
|
|
completedAt: null,
|
|
},
|
|
});
|
|
yield* emitProviderTurn(state, turn, "running", null);
|
|
return turn;
|
|
});
|
|
|
|
const projectChildUserPart = Effect.fnUntraced(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
part: Extract<OpenCodePart, { type: "text" }>,
|
|
) {
|
|
const now = yield* DateTime.now;
|
|
const messageId = idAllocator.derive.messageFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.messageID,
|
|
});
|
|
const turnItemId = idAllocator.derive.turnItemFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: part.messageID,
|
|
});
|
|
const projected: OrchestrationV2ConversationMessage = {
|
|
createdBy: "agent",
|
|
creationSource: "provider",
|
|
id: messageId,
|
|
threadId: turn.threadId,
|
|
runId: null,
|
|
nodeId: turn.rootNodeId,
|
|
role: "user",
|
|
text: part.text,
|
|
attachments: [],
|
|
streaming: false,
|
|
createdAt: turn.startedAt,
|
|
updatedAt: now,
|
|
};
|
|
state.messages.set(String(messageId), projected);
|
|
yield* emitProviderEvent({
|
|
type: "message.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
message: projected,
|
|
});
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem: {
|
|
createdBy: "agent",
|
|
creationSource: "provider",
|
|
id: turnItemId,
|
|
threadId: turn.threadId,
|
|
runId: null,
|
|
nodeId: turn.rootNodeId,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef: providerRef(part.messageID),
|
|
parentItemId: null,
|
|
ordinal: itemOrdinal(turn, part.messageID),
|
|
status: "completed",
|
|
title: null,
|
|
startedAt: projected.createdAt,
|
|
completedAt: now,
|
|
updatedAt: now,
|
|
type: "user_message",
|
|
messageId,
|
|
inputIntent: "turn_start",
|
|
text: part.text,
|
|
attachments: [],
|
|
},
|
|
});
|
|
});
|
|
|
|
const emitCompactionItem = Effect.fn("OpenCodeAdapterV2.emitCompactionItem")(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
) {
|
|
const nativeItemId = `${turn.providerTurnId}:compaction`;
|
|
if (turn.itemOrdinals.has(nativeItemId)) return;
|
|
const now = yield* DateTime.now;
|
|
yield* emitProviderEvent({
|
|
type: "turn_item.updated",
|
|
driver: OPENCODE_PROVIDER,
|
|
turnItem: {
|
|
id: idAllocator.derive.turnItemFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId,
|
|
}),
|
|
threadId: turn.threadId,
|
|
runId: turn.runId,
|
|
nodeId: turn.rootNodeId,
|
|
providerThreadId: state.providerThread.id,
|
|
providerTurnId: turn.providerTurnId,
|
|
nativeItemRef: providerRef(nativeItemId, "weak"),
|
|
parentItemId: null,
|
|
ordinal: itemOrdinal(turn, nativeItemId),
|
|
type: "compaction",
|
|
driver: OPENCODE_PROVIDER,
|
|
status: "completed",
|
|
title: "Context compacted",
|
|
startedAt: turn.startedAt,
|
|
completedAt: now,
|
|
updatedAt: now,
|
|
},
|
|
});
|
|
});
|
|
|
|
const handleMessageUpdated = Effect.fnUntraced(function* (
|
|
event: Extract<OpenCodeEvent, { type: "message.updated" }>,
|
|
) {
|
|
const state = threads.get(event.properties.sessionID);
|
|
if (state === undefined) return;
|
|
const message = event.properties.info;
|
|
state.messageRoles.set(message.id, message.role);
|
|
if (message.role === "assistant") {
|
|
const usage = state.activeTurn?.usage;
|
|
if (usage === undefined) return;
|
|
const prior = usage.assistantOwnershipByMessageId.get(message.id);
|
|
const ownership =
|
|
prior !== undefined && prior !== "unknown"
|
|
? prior
|
|
: !message.parentID
|
|
? "unknown"
|
|
: usage.promptMessageIds.has(message.parentID)
|
|
? "owned"
|
|
: "other";
|
|
usage.assistantOwnershipByMessageId.set(message.id, ownership);
|
|
if (ownership !== "unknown") {
|
|
if (ownership === "owned") {
|
|
for (const step of usage.unresolvedStepsByMessageId.get(message.id)?.values() ??
|
|
[]) {
|
|
accumulateOpenCodeStepUsage(usage, step);
|
|
}
|
|
}
|
|
usage.unresolvedStepsByMessageId.delete(message.id);
|
|
}
|
|
return;
|
|
}
|
|
const commandReceipt = commandReceipts.get(message.id);
|
|
if (commandReceipt) yield* Deferred.succeed(commandReceipt, undefined);
|
|
const isNewUserMessage = !state.userMessageIds.includes(message.id);
|
|
if (isNewUserMessage) state.userMessageIds.push(message.id);
|
|
let turn = state.activeTurn;
|
|
if (turn === null && state.parentSubagent !== null && isNewUserMessage) {
|
|
turn = yield* createChildTurn(state, message);
|
|
}
|
|
const matchesAdmission =
|
|
turn !== null &&
|
|
(turn.admissionMessageId === null || turn.admissionMessageId === message.id);
|
|
if (turn !== null && matchesAdmission) turn.usage.promptMessageIds.add(message.id);
|
|
if (turn !== null && matchesAdmission && turn.nativeUserMessageId === null) {
|
|
turn.nativeUserMessageId = message.id;
|
|
yield* emitProviderTurn(state, turn, "running", null);
|
|
}
|
|
if (turn !== null && matchesAdmission && turn.admissionPending) {
|
|
if (advanceOpenCodePromptAdmission(turn, "user-message") === "reconcile-idle") {
|
|
yield* reconcilePromptAdmission(state, turn);
|
|
}
|
|
}
|
|
});
|
|
|
|
const handlePartUpdated = Effect.fnUntraced(function* (
|
|
event: Extract<OpenCodeEvent, { type: "message.part.updated" }>,
|
|
) {
|
|
const part = event.properties.part;
|
|
const state = threads.get(part.sessionID);
|
|
const turn = state?.activeTurn;
|
|
if (state === undefined || turn === null || turn === undefined || turn.finalized) return;
|
|
if (part.type === "text" && state.messageRoles.get(part.messageID) === "user") {
|
|
if (!turn.isRoot) yield* projectChildUserPart(state, turn, part);
|
|
return;
|
|
}
|
|
if (part.type === "step-finish") {
|
|
const usage = turn.usage;
|
|
const ownership = usage.assistantOwnershipByMessageId.get(part.messageID);
|
|
if (ownership === "owned") accumulateOpenCodeStepUsage(usage, part);
|
|
else if (
|
|
ownership === "unknown" ||
|
|
(ownership === undefined && !state.messageRoles.has(part.messageID))
|
|
) {
|
|
const steps =
|
|
usage.unresolvedStepsByMessageId.get(part.messageID) ??
|
|
new Map<string, OpenCodeStepUsage>();
|
|
steps.set(part.id, { id: part.id, tokens: part.tokens });
|
|
usage.unresolvedStepsByMessageId.set(part.messageID, steps);
|
|
}
|
|
return;
|
|
}
|
|
if (part.type === "tool") {
|
|
// Approval routing needs the tool name, without retaining its input and output.
|
|
turn.toolNamesByCallId.set(part.callID, part.tool);
|
|
} else {
|
|
turn.parts.set(part.id, part);
|
|
const ids = turn.partIdsByMessage.get(part.messageID) ?? new Set<string>();
|
|
ids.add(part.id);
|
|
turn.partIdsByMessage.set(part.messageID, ids);
|
|
}
|
|
switch (part.type) {
|
|
case "text":
|
|
case "reasoning":
|
|
yield* emitTextPart(state, turn, part);
|
|
return;
|
|
case "tool":
|
|
yield* emitToolPart(state, turn, part);
|
|
return;
|
|
default:
|
|
return;
|
|
}
|
|
});
|
|
|
|
const handlePartDelta = Effect.fnUntraced(function* (
|
|
event: Extract<OpenCodeEvent, { type: "message.part.delta" }>,
|
|
) {
|
|
if (event.properties.field !== "text") return;
|
|
const state = threads.get(event.properties.sessionID);
|
|
const turn = state?.activeTurn;
|
|
const current = turn?.parts.get(event.properties.partID);
|
|
if (
|
|
state === undefined ||
|
|
turn === null ||
|
|
turn === undefined ||
|
|
current === undefined ||
|
|
(current.type !== "text" && current.type !== "reasoning")
|
|
) {
|
|
return;
|
|
}
|
|
const updated = { ...current, text: current.text + event.properties.delta };
|
|
turn.parts.set(updated.id, updated);
|
|
yield* emitTextPart(state, turn, updated);
|
|
});
|
|
|
|
const handleAssistantCompleted = Effect.fnUntraced(function* (
|
|
event: Extract<OpenCodeEvent, { type: "message.updated" }>,
|
|
) {
|
|
const message = event.properties.info;
|
|
if (message.role !== "assistant" || message.time.completed === undefined) return;
|
|
const state = threads.get(message.sessionID);
|
|
const turn = state?.activeTurn;
|
|
if (state === undefined || turn === null || turn === undefined) return;
|
|
// Some OpenCode versions ignore the client-provided message ID. A
|
|
// completed assistant message is definitive admission evidence, so
|
|
// let the following idle event settle the turn without weakening
|
|
// the stale-user-message guard.
|
|
advanceOpenCodePromptAdmission(turn, "assistant-completed");
|
|
for (const partId of turn.partIdsByMessage.get(message.id) ?? []) {
|
|
const part = turn.parts.get(partId);
|
|
if (part?.type === "text" || part?.type === "reasoning") {
|
|
yield* emitTextPart(state, turn, part, true);
|
|
}
|
|
}
|
|
});
|
|
|
|
const handleEvent = Effect.fnUntraced(function* (event: OpenCodeEvent) {
|
|
yield* logProtocolEvent({
|
|
direction: "incoming",
|
|
messageKind: "notification",
|
|
method: event.type,
|
|
payload: event,
|
|
});
|
|
switch (event.type) {
|
|
case "server.connected":
|
|
if (hasConnected) {
|
|
for (const state of threads.values()) {
|
|
if (state.activeTurn !== null) state.activeTurn.usage.complete = false;
|
|
}
|
|
}
|
|
hasConnected = true;
|
|
return;
|
|
case "message.updated":
|
|
yield* handleMessageUpdated(event);
|
|
yield* handleAssistantCompleted(event);
|
|
return;
|
|
case "session.compacted": {
|
|
const state = threads.get(event.properties.sessionID);
|
|
if (state?.activeTurn !== undefined && state.activeTurn !== null) {
|
|
yield* emitCompactionItem(state, state.activeTurn);
|
|
}
|
|
return;
|
|
}
|
|
case "message.part.updated":
|
|
yield* handlePartUpdated(event);
|
|
return;
|
|
case "message.part.delta":
|
|
yield* handlePartDelta(event);
|
|
return;
|
|
case "todo.updated": {
|
|
const state = threads.get(event.properties.sessionID);
|
|
if (state?.activeTurn !== null && state?.activeTurn !== undefined) {
|
|
yield* emitTodo(state, state.activeTurn, event.properties.todos);
|
|
}
|
|
return;
|
|
}
|
|
case "permission.asked":
|
|
yield* routeRuntimeRequest(event.properties.id, event.properties.sessionID, {
|
|
type: "permission",
|
|
value: event.properties,
|
|
});
|
|
return;
|
|
case "question.asked":
|
|
yield* routeRuntimeRequest(event.properties.id, event.properties.sessionID, {
|
|
type: "question",
|
|
value: event.properties,
|
|
});
|
|
return;
|
|
case "permission.replied":
|
|
case "question.replied":
|
|
rememberSettledRequest(event.properties.requestID);
|
|
yield* resolveRuntimeRequest(event.properties.requestID, "resolved");
|
|
return;
|
|
case "question.rejected":
|
|
rememberSettledRequest(event.properties.requestID);
|
|
yield* resolveRuntimeRequest(event.properties.requestID, "cancelled");
|
|
return;
|
|
case "session.created":
|
|
case "session.updated": {
|
|
const info = event.properties.info;
|
|
if (info.parentID !== undefined && !threads.has(info.id)) {
|
|
const owner = threads.get(info.parentID) ?? relatedSessionOwners.get(info.parentID);
|
|
if (owner !== undefined) {
|
|
relatedSessionOwners.set(info.id, owner);
|
|
}
|
|
}
|
|
return;
|
|
}
|
|
case "session.deleted":
|
|
relatedSessionOwners.delete(event.properties.info.id);
|
|
busySessionIds.delete(event.properties.info.id);
|
|
return;
|
|
case "session.status": {
|
|
const sessionId = event.properties.sessionID;
|
|
switch (event.properties.status.type) {
|
|
case "busy":
|
|
case "retry":
|
|
if (threads.has(sessionId) || relatedSessionOwners.has(sessionId)) {
|
|
busySessionIds.add(sessionId);
|
|
}
|
|
break;
|
|
case "idle":
|
|
busySessionIds.delete(sessionId);
|
|
break;
|
|
}
|
|
const state = threads.get(sessionId);
|
|
if (state === undefined) return;
|
|
if (event.properties.status.type === "busy") {
|
|
yield* updateProviderThread(state, { status: "active" });
|
|
return;
|
|
}
|
|
if (event.properties.status.type === "idle" && state.activeTurn !== null) {
|
|
if (state.activeTurn.admissionPending) {
|
|
advanceOpenCodePromptAdmission(state.activeTurn, "idle");
|
|
return;
|
|
}
|
|
yield* finalizeTurn(
|
|
state,
|
|
state.activeTurn,
|
|
state.activeTurn.interrupted ? "interrupted" : "completed",
|
|
);
|
|
}
|
|
return;
|
|
}
|
|
case "session.idle": {
|
|
busySessionIds.delete(event.properties.sessionID);
|
|
const state = threads.get(event.properties.sessionID);
|
|
if (state?.activeTurn !== null && state?.activeTurn !== undefined) {
|
|
if (state.activeTurn.admissionPending) {
|
|
advanceOpenCodePromptAdmission(state.activeTurn, "idle");
|
|
return;
|
|
}
|
|
yield* finalizeTurn(
|
|
state,
|
|
state.activeTurn,
|
|
state.activeTurn.interrupted ? "interrupted" : "completed",
|
|
);
|
|
}
|
|
return;
|
|
}
|
|
case "session.error": {
|
|
const states =
|
|
event.properties.sessionID === undefined
|
|
? Array.from(threads.values()).filter((state) => state.activeTurn !== null)
|
|
: [threads.get(event.properties.sessionID)].filter(
|
|
(state): state is OpenCodeThreadState => state !== undefined,
|
|
);
|
|
const message = openCodeErrorMessage(event);
|
|
if (
|
|
!isMessageAbortedError(event) &&
|
|
(event.properties.sessionID === undefined ||
|
|
states.some((state) => state.parentSubagent === null))
|
|
) {
|
|
yield* updateProviderSession("error", message);
|
|
}
|
|
for (const state of states) {
|
|
if (state.activeTurn !== null) {
|
|
yield* finalizeTurn(
|
|
state,
|
|
state.activeTurn,
|
|
terminalStatusForError(event, state.activeTurn),
|
|
{
|
|
failure: makeProviderFailure({
|
|
message,
|
|
code: event.properties.error?.name ?? null,
|
|
class: "provider_error",
|
|
}),
|
|
threadDisposition:
|
|
event.properties.sessionID === undefined ? "broken" : "reusable",
|
|
},
|
|
);
|
|
}
|
|
}
|
|
return;
|
|
}
|
|
default:
|
|
return;
|
|
}
|
|
});
|
|
|
|
const inputProviderSessionId = input.providerSessionId;
|
|
|
|
const subscription = yield* sdkCall("event.subscribe", {}, () =>
|
|
client.event.subscribe(undefined, { signal: abortController.signal }),
|
|
);
|
|
yield* Scope.addFinalizer(
|
|
scope,
|
|
Effect.sync(() => abortController.abort()),
|
|
);
|
|
yield* Stream.fromAsyncIterable(
|
|
subscription.stream,
|
|
(cause) =>
|
|
new OpenCodeRuntime.OpenCodeRuntimeError({
|
|
operation: "event.subscribe",
|
|
detail: OpenCodeRuntime.openCodeRuntimeErrorDetail(cause),
|
|
cause,
|
|
}),
|
|
).pipe(
|
|
Stream.runForEach(handleEvent),
|
|
Effect.exit,
|
|
Effect.flatMap((exit) =>
|
|
Effect.gen(function* () {
|
|
// No further status can clear a busy session once the stream ends.
|
|
busySessionIds.clear();
|
|
if (closing || abortController.signal.aborted) return;
|
|
const detail = Exit.isSuccess(exit)
|
|
? "OpenCode event stream ended unexpectedly."
|
|
: OpenCodeRuntime.openCodeRuntimeErrorDetail(Cause.squash(exit.cause));
|
|
nativeStreamFailure = makeProviderFailure({
|
|
message: detail,
|
|
class: "transport_error",
|
|
});
|
|
yield* updateProviderSession("error", detail);
|
|
for (const state of threads.values()) {
|
|
if (state.activeTurn !== null)
|
|
yield* finalizeTurn(state, state.activeTurn, "failed", {
|
|
failure: makeProviderFailure({ message: detail, class: "transport_error" }),
|
|
threadDisposition: "broken",
|
|
});
|
|
}
|
|
yield* Queue.end(events);
|
|
}),
|
|
),
|
|
Effect.forkIn(scope),
|
|
);
|
|
|
|
if (!connection.external && connection.exitCode !== null) {
|
|
yield* connection.exitCode.pipe(
|
|
Effect.flatMap((code) =>
|
|
abortController.signal.aborted
|
|
? Effect.void
|
|
: Effect.gen(function* () {
|
|
const detail = `OpenCode server exited unexpectedly (${code}).`;
|
|
yield* updateProviderSession("error", detail);
|
|
for (const state of threads.values()) {
|
|
if (state.activeTurn !== null) {
|
|
yield* finalizeTurn(state, state.activeTurn, "failed", {
|
|
failure: makeProviderFailure({
|
|
message: detail,
|
|
class: "transport_error",
|
|
}),
|
|
threadDisposition: "broken",
|
|
});
|
|
}
|
|
}
|
|
}),
|
|
),
|
|
Effect.forkIn(scope),
|
|
);
|
|
}
|
|
|
|
if (connection.external) {
|
|
yield* Scope.addFinalizer(
|
|
scope,
|
|
Effect.suspend(() =>
|
|
Effect.forEach(
|
|
Array.from(threads.values()).filter((state) => state.parentSubagent === null),
|
|
(state) =>
|
|
sdkCall("session.abort", { sessionID: state.nativeSessionId }, (signal) =>
|
|
client.session.abort({ sessionID: state.nativeSessionId }, { signal }),
|
|
).pipe(
|
|
Effect.timeout("1 second"),
|
|
Effect.ignore({ log: true }),
|
|
Effect.andThen(
|
|
abortDescendants(state.nativeSessionId).pipe(
|
|
Effect.timeout("1 second"),
|
|
Effect.ignore({ log: true }),
|
|
),
|
|
),
|
|
),
|
|
{ concurrency: 8, discard: true },
|
|
),
|
|
),
|
|
);
|
|
}
|
|
|
|
yield* Scope.addFinalizer(
|
|
scope,
|
|
Effect.sync(() => {
|
|
closing = true;
|
|
}),
|
|
);
|
|
|
|
const registerThread = (
|
|
nativeSession: OpenCodeSession,
|
|
providerThread: OrchestrationV2ProviderThread,
|
|
appThread: OrchestrationV2AppThread | null,
|
|
): OpenCodeThreadState => {
|
|
const existing = threads.get(nativeSession.id);
|
|
if (existing !== undefined) {
|
|
existing.providerThread = providerThread;
|
|
if (appThread !== null) existing.appThread = appThread;
|
|
return existing;
|
|
}
|
|
const state: OpenCodeThreadState = {
|
|
nativeSessionId: nativeSession.id,
|
|
providerThread,
|
|
appThread,
|
|
activeTurn: null,
|
|
providerTurns: new Map(),
|
|
messages: new Map(),
|
|
runtimeRequests: new Map(),
|
|
messageRoles: new Map(),
|
|
userMessageIds: [],
|
|
parentSubagent: subagentsByChildSessionId.get(nativeSession.id) ?? null,
|
|
nextChildTurnOrdinal: 1,
|
|
nextAdmissionGeneration: 1,
|
|
};
|
|
threads.set(nativeSession.id, state);
|
|
return state;
|
|
};
|
|
|
|
const resolvePromptParts = (turnInput: ProviderAdapter.ProviderAdapterV2TurnInput) => {
|
|
const text = providerMessageTextWithAttachmentPaths({
|
|
text: turnInput.message.text,
|
|
attachments: turnInput.message.attachments,
|
|
resolveAttachmentPath: (attachment) =>
|
|
resolveAttachmentPath({ attachmentsDir: serverConfig.attachmentsDir, attachment }),
|
|
}).trim();
|
|
const files = OpenCodeRuntime.toOpenCodeFileParts({
|
|
attachments: turnInput.message.attachments,
|
|
resolveAttachmentPath: (attachment) =>
|
|
resolveAttachmentPath({ attachmentsDir: serverConfig.attachmentsDir, attachment }),
|
|
});
|
|
if (text.length === 0 && files.length === 0) {
|
|
throw protocolError("OpenCode turns require text or at least one valid attachment");
|
|
}
|
|
return [...(text.length === 0 ? [] : [{ type: "text" as const, text }]), ...files];
|
|
};
|
|
|
|
const readSnapshot = Effect.fnUntraced(function* (
|
|
providerThread: OrchestrationV2ProviderThread,
|
|
) {
|
|
const sessionId = nativeThreadId(providerThread);
|
|
const sessionResponse = yield* sdkCall("session.get", { sessionID: sessionId }, () =>
|
|
client.session.get({ sessionID: sessionId }),
|
|
);
|
|
const nativeSession = unwrapData("session.get", sessionResponse);
|
|
const response = yield* sdkCall("session.messages", { sessionID: sessionId }, () =>
|
|
client.session.messages({ sessionID: sessionId }),
|
|
);
|
|
const nativeMessages = [];
|
|
for (const entry of unwrapData("session.messages", response)) {
|
|
// OpenCode retains reverted records in session.messages. Its
|
|
// revert marker identifies the first removed message, so stop
|
|
// before it or a restart resurrects rolled-back assistant work.
|
|
if (entry.info.id === nativeSession.revert?.messageID) break;
|
|
nativeMessages.push(entry);
|
|
}
|
|
const state = threads.get(sessionId);
|
|
const snapshotNow = yield* DateTime.now;
|
|
const messages: Array<OrchestrationV2ConversationMessage> = nativeMessages.flatMap(
|
|
({ info, parts }) => {
|
|
const text = parts
|
|
.filter(
|
|
(part): part is Extract<OpenCodePart, { type: "text" }> => part.type === "text",
|
|
)
|
|
.filter((part) => part.ignored !== true && part.synthetic !== true)
|
|
.map((part) => part.text)
|
|
.join("\n");
|
|
if (text.length === 0) return [];
|
|
const createdAt = dateTimeFromEpoch(info.time.created, snapshotNow);
|
|
return [
|
|
{
|
|
createdBy: info.role === "user" ? "user" : "agent",
|
|
creationSource: "provider",
|
|
id: idAllocator.derive.messageFromProviderItem({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeItemId: info.id,
|
|
}),
|
|
threadId: providerThread.appThreadId ?? input.threadId,
|
|
runId: null,
|
|
nodeId: null,
|
|
role: info.role,
|
|
text,
|
|
attachments: [],
|
|
streaming: false,
|
|
createdAt,
|
|
updatedAt:
|
|
info.role === "assistant"
|
|
? dateTimeFromEpoch(info.time.completed, createdAt)
|
|
: createdAt,
|
|
},
|
|
];
|
|
},
|
|
);
|
|
const lastUser = nativeMessages.findLast(({ info }) => info.role === "user")?.info.id;
|
|
return {
|
|
providerThread: {
|
|
...providerThread,
|
|
providerSessionId: input.providerSessionId,
|
|
nativeConversationHeadRef:
|
|
lastUser === undefined ? null : providerRef(lastUser, "weak"),
|
|
status: "idle" as const,
|
|
updatedAt: snapshotNow,
|
|
},
|
|
providerTurns: state === undefined ? [] : [...state.providerTurns.values()],
|
|
messages,
|
|
runtimeRequests: state === undefined ? [] : [...state.runtimeRequests.values()],
|
|
providerPayload: nativeMessages,
|
|
};
|
|
});
|
|
|
|
// Native commands wait for generation, unlike promptAsync. A user-message
|
|
// receipt admits the run while the scoped request continues in the background.
|
|
const submitPrompt = Effect.fn("OpenCode.submitPrompt")(function* (
|
|
state: OpenCodeThreadState,
|
|
turn: ActiveOpenCodeTurn,
|
|
payload: Parameters<typeof client.session.promptAsync>[0] & {
|
|
messageID: string;
|
|
sessionID: string;
|
|
},
|
|
abortController: AbortController,
|
|
) {
|
|
const text = payload.parts
|
|
?.filter((part) => part.type === "text")
|
|
.map((part) => part.text)
|
|
.join("\n")
|
|
.trim();
|
|
const match = text?.match(/^\/([^\s/]+)(?:\s+([\s\S]*))?$/);
|
|
const command = match
|
|
? (yield* OpenCodeRuntime.loadOpenCodeCommands(client).pipe(
|
|
Effect.timeout("10 seconds"),
|
|
Effect.orElseSucceed(() => []),
|
|
)).find((entry) => entry.name === match[1])
|
|
: undefined;
|
|
if (!command) {
|
|
return yield* sdkCall("session.promptAsync", payload, (signal) =>
|
|
client.session.promptAsync(payload, {
|
|
signal: AbortSignal.any([signal, abortController.signal]),
|
|
}),
|
|
).pipe(Effect.asVoid);
|
|
}
|
|
const receipt = Deferred.makeUnsafe<void>();
|
|
commandReceipts.set(payload.messageID, receipt);
|
|
const controllers =
|
|
commandControllers.get(payload.sessionID) ?? new Set<AbortController>();
|
|
controllers.add(abortController);
|
|
commandControllers.set(payload.sessionID, controllers);
|
|
const generation = turn.admissionGeneration;
|
|
const commandPayload = {
|
|
sessionID: payload.sessionID,
|
|
messageID: payload.messageID,
|
|
command: command.name,
|
|
arguments: match?.[2] ?? "",
|
|
...(payload.model
|
|
? { model: `${payload.model.providerID}/${payload.model.modelID}` }
|
|
: {}),
|
|
...(payload.agent ? { agent: payload.agent } : {}),
|
|
...(payload.variant ? { variant: payload.variant } : {}),
|
|
parts: payload.parts?.filter((part) => part.type === "file") ?? [],
|
|
};
|
|
const request = yield* sdkCall("session.command", commandPayload, (signal) =>
|
|
client.session.command(commandPayload, {
|
|
signal: AbortSignal.any([signal, abortController.signal]),
|
|
}),
|
|
).pipe(
|
|
Effect.asVoid,
|
|
Effect.tapError((cause) =>
|
|
abortController.signal.aborted ||
|
|
turn.finalized ||
|
|
turn.admissionGeneration !== generation
|
|
? Effect.void
|
|
: finalizeTurn(state, turn, "failed", {
|
|
failure: makeProviderFailure({ cause, class: "provider_error" }),
|
|
}),
|
|
),
|
|
Effect.ensuring(
|
|
Effect.sync(() => {
|
|
commandReceipts.delete(payload.messageID);
|
|
controllers.delete(abortController);
|
|
if (controllers.size === 0) commandControllers.delete(payload.sessionID);
|
|
}),
|
|
),
|
|
Effect.forkIn(scope),
|
|
);
|
|
yield* Effect.raceFirst(Fiber.join(request), Deferred.await(receipt)).pipe(
|
|
Effect.timeout("10 seconds"),
|
|
Effect.catchTags({
|
|
TimeoutError: (cause) => {
|
|
const error = new OpenCodeRuntime.OpenCodeRuntimeError({
|
|
operation: "session.command",
|
|
detail: "OpenCode command admission did not complete within 10 seconds.",
|
|
cause,
|
|
});
|
|
abortController.abort();
|
|
return finalizeTurn(state, turn, "failed", {
|
|
failure: makeProviderFailure({ cause: error, class: "provider_error" }),
|
|
}).pipe(Effect.andThen(Effect.fail(error)));
|
|
},
|
|
}),
|
|
);
|
|
});
|
|
|
|
const runtimeSession: ProviderAdapter.ProviderAdapterV2SessionRuntime = {
|
|
instanceId: options.instanceId,
|
|
driver: OPENCODE_PROVIDER,
|
|
providerSessionId: input.providerSessionId,
|
|
providerSession: sessionEntity,
|
|
events: Stream.fromEffectRepeat(Queue.take(events)),
|
|
hasPendingBackgroundWork: Effect.sync(() => busySessionIds.size > 0),
|
|
ensureThread: (threadInput) =>
|
|
Effect.gen(function* () {
|
|
// Only a row that already carries a native session can be
|
|
// resumed; a placeholder without one still needs session.create.
|
|
if (threadInput.existingProviderThread?.nativeThreadRef != null) {
|
|
return yield* runtimeSession.resumeThread({
|
|
providerThread: threadInput.existingProviderThread,
|
|
});
|
|
}
|
|
// No title: OpenCode generates one from the first prompt only when
|
|
// session.create leaves it unset (SessionPrompt.ensureTitle).
|
|
const response = yield* sdkCall(
|
|
"session.create",
|
|
{ permission: openCodePermissionRules(threadInput.runtimePolicy) },
|
|
() =>
|
|
client.session.create({
|
|
permission: openCodePermissionRules(threadInput.runtimePolicy),
|
|
}),
|
|
);
|
|
const nativeSession = unwrapData("session.create", response);
|
|
const createdAt = yield* DateTime.now;
|
|
const created = makeProviderThread({
|
|
idAllocator,
|
|
providerInstanceId: options.instanceId,
|
|
providerSessionId: input.providerSessionId,
|
|
appThreadId: threadInput.threadId,
|
|
nativeSession,
|
|
now: createdAt,
|
|
});
|
|
const existing = threadInput.existingProviderThread;
|
|
// Bind the new native session to the caller's row when one was
|
|
// handed over: a second live row per app thread would make
|
|
// `activeProviderThreadId` flap between the two on every update.
|
|
const providerThread =
|
|
existing === undefined
|
|
? created
|
|
: {
|
|
...existing,
|
|
providerSessionId: input.providerSessionId,
|
|
nativeThreadRef: created.nativeThreadRef,
|
|
nativeConversationHeadRef: created.nativeConversationHeadRef,
|
|
status: created.status,
|
|
updatedAt: created.updatedAt,
|
|
};
|
|
registerThread(nativeSession, providerThread, null);
|
|
return providerThread;
|
|
}).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterEnsureThreadError({
|
|
driver: OPENCODE_PROVIDER,
|
|
threadId: threadInput.threadId,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
resumeThread: (threadInput) =>
|
|
Effect.gen(function* () {
|
|
const sessionId = nativeThreadId(threadInput.providerThread);
|
|
const response = yield* sdkCall("session.get", { sessionID: sessionId }, () =>
|
|
client.session.get({ sessionID: sessionId }),
|
|
);
|
|
const nativeSession = unwrapData("session.get", response);
|
|
const resumedAt = yield* DateTime.now;
|
|
const providerThread = {
|
|
...threadInput.providerThread,
|
|
providerSessionId: input.providerSessionId,
|
|
status: "idle" as const,
|
|
updatedAt: dateTimeFromEpoch(nativeSession.time.updated, resumedAt),
|
|
};
|
|
registerThread(nativeSession, providerThread, null);
|
|
return providerThread;
|
|
}).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterResumeThreadError({
|
|
driver: OPENCODE_PROVIDER,
|
|
providerSessionId: input.providerSessionId,
|
|
providerThreadId: threadInput.providerThread.id,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
compactThread: (turnInput) =>
|
|
runtimeSession.startTurn({
|
|
...turnInput,
|
|
message: { ...turnInput.message, text: "/compact" },
|
|
}),
|
|
startTurn: (turnInput) =>
|
|
Effect.gen(function* () {
|
|
if (nativeStreamFailure !== null) {
|
|
return yield* protocolError(
|
|
"OpenCode event stream has ended; reconnect the provider session before starting another turn.",
|
|
);
|
|
}
|
|
const sessionId = nativeThreadId(turnInput.providerThread);
|
|
const state = threads.get(sessionId);
|
|
if (state === undefined) {
|
|
return yield* protocolError(`OpenCode session ${sessionId} is not registered`);
|
|
}
|
|
if (state.activeTurn !== null) {
|
|
return yield* protocolError(
|
|
`OpenCode provider thread ${turnInput.providerThread.id} already has an active turn`,
|
|
);
|
|
}
|
|
const parsedModel = OpenCodeRuntime.parseOpenCodeModelSlug(
|
|
turnInput.modelSelection.model,
|
|
);
|
|
if (parsedModel === null) {
|
|
return yield* protocolError(
|
|
`OpenCode model '${turnInput.modelSelection.model}' must use provider/model format`,
|
|
);
|
|
}
|
|
const isCompaction =
|
|
turnInput.message.text.trim() === "/compact" &&
|
|
turnInput.message.attachments.length === 0;
|
|
const parts = isCompaction ? [] : resolvePromptParts(turnInput);
|
|
const startedAt = yield* DateTime.now;
|
|
const syntheticNativeTurnId = `${sessionId}:attempt:${turnInput.attemptId}`;
|
|
const providerTurnId = idAllocator.derive.providerTurn({
|
|
driver: OPENCODE_PROVIDER,
|
|
nativeTurnId: syntheticNativeTurnId,
|
|
});
|
|
const providerTurn: OrchestrationV2ProviderTurn = {
|
|
id: providerTurnId,
|
|
providerThreadId: turnInput.providerThread.id,
|
|
nodeId: turnInput.rootNodeId,
|
|
runAttemptId: turnInput.attemptId,
|
|
nativeTurnRef: providerRef(syntheticNativeTurnId, "weak"),
|
|
ordinal: turnInput.providerTurnOrdinal,
|
|
status: "running",
|
|
startedAt,
|
|
completedAt: null,
|
|
};
|
|
const admissionMessageId = yield* makeOpenCodeMessageId();
|
|
// No Effect may be yielded between this check and installing the
|
|
// turn. If the event stream ended while IDs were being prepared,
|
|
// registering afterward would leave a running turn that the EOF
|
|
// handler had already finished scanning.
|
|
if (nativeStreamFailure !== null) {
|
|
return yield* protocolError(
|
|
"OpenCode event stream has ended; reconnect the provider session before starting another turn.",
|
|
);
|
|
}
|
|
if (state.activeTurn !== null) {
|
|
return yield* protocolError(
|
|
`OpenCode provider thread ${turnInput.providerThread.id} already has an active turn`,
|
|
);
|
|
}
|
|
const turn: ActiveOpenCodeTurn = {
|
|
isRoot: true,
|
|
threadId: turnInput.threadId,
|
|
runId: turnInput.runId,
|
|
rootNodeId: turnInput.rootNodeId,
|
|
appThread: turnInput.appThread,
|
|
modelSelection: turnInput.modelSelection,
|
|
runtimePolicy: turnInput.runtimePolicy,
|
|
providerTurnId,
|
|
providerThreadId: turnInput.providerThread.id,
|
|
providerTurnOrdinal: turnInput.providerTurnOrdinal,
|
|
runOrdinal: turnInput.runOrdinal,
|
|
runAttemptId: turnInput.attemptId,
|
|
startedAt,
|
|
itemOrdinals: new Map(),
|
|
usage: makeOpenCodeTurnTokenUsageAccumulator(),
|
|
parts: new Map(),
|
|
partIdsByMessage: new Map(),
|
|
toolNamesByCallId: new Map(),
|
|
providerTurn,
|
|
nextItemOrdinal: turnInput.providerTurnOrdinal * 100 + 1,
|
|
nativeUserMessageId: null,
|
|
admissionMessageId,
|
|
interrupted: false,
|
|
finalized: false,
|
|
planId: null,
|
|
admissionGeneration: state.nextAdmissionGeneration++,
|
|
admissionReconciliationGeneration: null,
|
|
admissionPending: !isCompaction,
|
|
admissionAccepted: isCompaction,
|
|
admissionMessageObserved: false,
|
|
idleDuringAdmission: false,
|
|
admissionSettled: Deferred.makeUnsafe<void>(),
|
|
admissionAbortController: new AbortController(),
|
|
};
|
|
if (turn.admissionMessageId !== null)
|
|
turn.usage.promptMessageIds.add(turn.admissionMessageId);
|
|
const admissionSettled = turn.admissionSettled;
|
|
const admissionAbortController = turn.admissionAbortController;
|
|
state.appThread = turnInput.appThread;
|
|
state.activeTurn = turn;
|
|
state.providerTurns.set(String(providerTurnId), providerTurn);
|
|
yield* emitProviderTurn(state, turn, "running", null);
|
|
yield* updateProviderThread(state, {
|
|
status: "active",
|
|
firstRunOrdinal: state.providerThread.firstRunOrdinal ?? turnInput.runOrdinal,
|
|
lastRunOrdinal: turnInput.runOrdinal,
|
|
});
|
|
yield* updateProviderSession("running", null);
|
|
if (isCompaction) {
|
|
yield* sdkCall(
|
|
"session.summarize",
|
|
{ sessionID: sessionId, ...parsedModel, auto: false },
|
|
(signal) =>
|
|
client.session.summarize(
|
|
{ sessionID: sessionId, ...parsedModel, auto: false },
|
|
{ signal: AbortSignal.any([signal, admissionAbortController!.signal]) },
|
|
),
|
|
).pipe(
|
|
Effect.tap(() =>
|
|
turn.interrupted || turn.finalized || nativeStreamFailure !== null
|
|
? Effect.void
|
|
: emitCompactionItem(state, turn),
|
|
),
|
|
Effect.tap(() =>
|
|
finalizeTurn(state, turn, turn.interrupted ? "interrupted" : "completed"),
|
|
),
|
|
Effect.tapError((cause) =>
|
|
finalizeTurn(state, turn, turn.interrupted ? "interrupted" : "failed", {
|
|
failure: makeProviderFailure({ cause, class: "provider_error" }),
|
|
}),
|
|
),
|
|
Effect.ensuring(Deferred.succeed(admissionSettled, undefined)),
|
|
);
|
|
return;
|
|
}
|
|
const systemPrompt = [
|
|
orchestrationSystemPrompt,
|
|
buildRuntimeInstructions({
|
|
harness: "OpenCode",
|
|
model: turnInput.modelSelection.model,
|
|
}),
|
|
]
|
|
.filter(Boolean)
|
|
.join("\n\n");
|
|
const agent =
|
|
getModelSelectionStringOptionValue(turnInput.modelSelection, "agent") ??
|
|
(turnInput.runtimePolicy.interactionMode === "plan" ? "plan" : undefined);
|
|
const variant = getModelSelectionStringOptionValue(
|
|
turnInput.modelSelection,
|
|
"variant",
|
|
);
|
|
yield* submitPrompt(
|
|
state,
|
|
turn,
|
|
{
|
|
sessionID: sessionId,
|
|
messageID: turn.admissionMessageId!,
|
|
model: parsedModel,
|
|
...(agent === undefined ? {} : { agent }),
|
|
...(variant === undefined ? {} : { variant }),
|
|
system: systemPrompt,
|
|
parts,
|
|
},
|
|
admissionAbortController!,
|
|
).pipe(
|
|
Effect.tapError((cause) =>
|
|
admissionAbortController!.signal.aborted
|
|
? Effect.void
|
|
: finalizeTurn(state, turn, "failed", {
|
|
failure: makeProviderFailure({ cause, class: "provider_error" }),
|
|
}),
|
|
),
|
|
Effect.catch((cause) =>
|
|
admissionAbortController!.signal.aborted ? Effect.void : Effect.fail(cause),
|
|
),
|
|
Effect.ensuring(
|
|
Effect.all([
|
|
Deferred.succeed(admissionSettled, undefined).pipe(Effect.ignore),
|
|
Effect.sync(() => {
|
|
if (turn.admissionAbortController === admissionAbortController) {
|
|
turn.admissionAbortController = null;
|
|
}
|
|
}),
|
|
]).pipe(Effect.asVoid),
|
|
),
|
|
);
|
|
if (state.activeTurn === turn && !turn.finalized && !turn.interrupted) {
|
|
const admissionAction = advanceOpenCodePromptAdmission(turn, "accepted");
|
|
if (admissionAction === "reconcile-idle") {
|
|
yield* reconcilePromptAdmission(state, turn);
|
|
}
|
|
}
|
|
}).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterTurnStartError({
|
|
driver: OPENCODE_PROVIDER,
|
|
threadId: turnInput.threadId,
|
|
providerThreadId: turnInput.providerThread.id,
|
|
runId: turnInput.runId,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
steerTurn: (steerInput) =>
|
|
Effect.gen(function* () {
|
|
const sessionId = nativeThreadId(steerInput.providerThread);
|
|
const state = threads.get(sessionId);
|
|
const turn = state?.activeTurn;
|
|
if (
|
|
state === undefined ||
|
|
turn === undefined ||
|
|
turn === null ||
|
|
turn.providerTurnId !== steerInput.providerTurnId ||
|
|
turn.interrupted
|
|
) {
|
|
return yield* protocolError(
|
|
`OpenCode turn ${steerInput.providerTurnId} is not active`,
|
|
);
|
|
}
|
|
const parsedModel = OpenCodeRuntime.parseOpenCodeModelSlug(turn.modelSelection.model);
|
|
if (parsedModel === null) {
|
|
return yield* protocolError(
|
|
`OpenCode model '${turn.modelSelection.model}' must use provider/model format`,
|
|
);
|
|
}
|
|
const text = providerMessageTextWithAttachmentPaths({
|
|
text: steerInput.message.text,
|
|
attachments: steerInput.message.attachments,
|
|
resolveAttachmentPath: (attachment) =>
|
|
resolveAttachmentPath({
|
|
attachmentsDir: serverConfig.attachmentsDir,
|
|
attachment,
|
|
}),
|
|
}).trim();
|
|
const files = OpenCodeRuntime.toOpenCodeFileParts({
|
|
attachments: steerInput.message.attachments,
|
|
resolveAttachmentPath: (attachment) =>
|
|
resolveAttachmentPath({
|
|
attachmentsDir: serverConfig.attachmentsDir,
|
|
attachment,
|
|
}),
|
|
});
|
|
if (text.length === 0 && files.length === 0) {
|
|
return yield* protocolError("OpenCode steering requires text or an attachment");
|
|
}
|
|
const parts = [
|
|
...(text.length === 0 ? [] : [{ type: "text" as const, text }]),
|
|
...files,
|
|
];
|
|
turn.admissionGeneration = state.nextAdmissionGeneration++;
|
|
turn.admissionMessageId = yield* makeOpenCodeMessageId();
|
|
turn.admissionPending = true;
|
|
turn.admissionAccepted = false;
|
|
turn.admissionMessageObserved = false;
|
|
turn.idleDuringAdmission = false;
|
|
turn.admissionSettled = Deferred.makeUnsafe<void>();
|
|
turn.admissionAbortController = new AbortController();
|
|
if (turn.admissionMessageId !== null)
|
|
turn.usage.promptMessageIds.add(turn.admissionMessageId);
|
|
const admissionSettled = turn.admissionSettled;
|
|
const admissionAbortController = turn.admissionAbortController;
|
|
yield* submitPrompt(
|
|
state,
|
|
turn,
|
|
{
|
|
sessionID: sessionId,
|
|
messageID: turn.admissionMessageId,
|
|
model: parsedModel,
|
|
parts,
|
|
},
|
|
admissionAbortController!,
|
|
).pipe(
|
|
Effect.ensuring(
|
|
Effect.all([
|
|
Deferred.succeed(admissionSettled, undefined).pipe(Effect.ignore),
|
|
Effect.sync(() => {
|
|
if (turn.admissionAbortController === admissionAbortController) {
|
|
turn.admissionAbortController = null;
|
|
}
|
|
}),
|
|
]).pipe(Effect.asVoid),
|
|
),
|
|
);
|
|
if (state.activeTurn === turn && !turn.finalized) {
|
|
const admissionAction = advanceOpenCodePromptAdmission(turn, "accepted");
|
|
if (admissionAction === "reconcile-idle") {
|
|
yield* reconcilePromptAdmission(state, turn);
|
|
}
|
|
}
|
|
}).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterSteerRunError({
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThreadId: steerInput.providerThread.id,
|
|
providerTurnId: steerInput.providerTurnId,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
interruptTurn: (interruptInput) =>
|
|
Effect.gen(function* () {
|
|
const sessionId = nativeThreadId(interruptInput.providerThread);
|
|
const state = threads.get(sessionId);
|
|
const turn = state?.activeTurn;
|
|
// Stop on a settled turn: a background task child can outlive
|
|
// it (busySessionIds), and aborting the descendants ends it.
|
|
if (
|
|
(turn === undefined || turn === null) &&
|
|
interruptInput.requestRuntimeRestart === true
|
|
) {
|
|
for (const controller of commandControllers.get(sessionId) ?? [])
|
|
controller.abort();
|
|
yield* abortDescendants(sessionId);
|
|
return;
|
|
}
|
|
if (
|
|
turn === undefined ||
|
|
turn === null ||
|
|
turn.providerTurnId !== interruptInput.providerTurnId
|
|
) {
|
|
return yield* protocolError(
|
|
`OpenCode turn ${interruptInput.providerTurnId} is not active`,
|
|
);
|
|
}
|
|
turn.interrupted = true;
|
|
for (const controller of commandControllers.get(sessionId) ?? []) controller.abort();
|
|
const admissionWasPending = turn.admissionPending;
|
|
cancelOpenCodePromptAdmission(turn, state!.nextAdmissionGeneration++);
|
|
turn.admissionAbortController?.abort();
|
|
if (admissionWasPending) {
|
|
// Settle the cancelled request before sending the session abort.
|
|
// A timed-out local request must not prevent the definitive
|
|
// server-side abort below from running.
|
|
yield* Deferred.await(turn.admissionSettled).pipe(
|
|
Effect.timeout("10 seconds"),
|
|
Effect.ignore,
|
|
);
|
|
}
|
|
yield* sdkCall("session.abort", { sessionID: sessionId }, (signal) =>
|
|
client.session.abort({ sessionID: sessionId }, { signal }),
|
|
).pipe(
|
|
Effect.timeout("10 seconds"),
|
|
// The turn can settle while the abort is in flight, and
|
|
// aborting an already-idle session fails. That stop still
|
|
// succeeded; only surface failures for a turn that is
|
|
// genuinely still running.
|
|
Effect.catch((cause) =>
|
|
turn.finalized
|
|
? Effect.void
|
|
: Effect.suspend(() => {
|
|
turn.interrupted = false;
|
|
return Effect.fail(cause);
|
|
}),
|
|
),
|
|
);
|
|
// Root abort does not stop child sessions. Report incomplete cleanup
|
|
// even if the root has already emitted its terminal event.
|
|
yield* abortDescendants(sessionId);
|
|
}).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterInterruptError({
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThreadId: interruptInput.providerThread.id,
|
|
providerTurnId: interruptInput.providerTurnId,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
respondToRuntimeRequest: (requestInput) =>
|
|
Effect.gen(function* () {
|
|
const pending = pendingRequests.get(String(requestInput.requestId));
|
|
if (pending === undefined) {
|
|
return yield* protocolError(
|
|
`No pending OpenCode request ${requestInput.requestId}`,
|
|
);
|
|
}
|
|
if (pending.question !== undefined) {
|
|
if (requestInput.answers === undefined) {
|
|
return yield* protocolError(
|
|
`OpenCode question request ${requestInput.requestId} requires answers`,
|
|
);
|
|
}
|
|
const answers = OpenCodeRuntime.toOpenCodeQuestionAnswers(
|
|
pending.question,
|
|
requestInput.answers,
|
|
);
|
|
yield* sdkCall(
|
|
"question.reply",
|
|
{ requestID: pending.nativeRequestId, answers },
|
|
(signal) =>
|
|
client.question.reply(
|
|
{
|
|
requestID: pending.nativeRequestId,
|
|
answers,
|
|
},
|
|
{ signal },
|
|
),
|
|
).pipe(Effect.timeout("10 seconds"));
|
|
return;
|
|
}
|
|
if (requestInput.decision === undefined) {
|
|
return yield* protocolError(
|
|
`OpenCode approval request ${requestInput.requestId} requires a decision`,
|
|
);
|
|
}
|
|
const reply = OpenCodeRuntime.toOpenCodePermissionReply(requestInput.decision);
|
|
yield* sdkCall(
|
|
"permission.reply",
|
|
{ requestID: pending.nativeRequestId, reply },
|
|
(signal) =>
|
|
client.permission.reply(
|
|
{
|
|
requestID: pending.nativeRequestId,
|
|
reply,
|
|
},
|
|
{ signal },
|
|
),
|
|
).pipe(Effect.timeout("10 seconds"));
|
|
}).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterRuntimeRequestResponseError({
|
|
driver: OPENCODE_PROVIDER,
|
|
requestId: requestInput.requestId,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
readThreadSnapshot: (snapshotInput) =>
|
|
readSnapshot(snapshotInput.providerThread).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterReadThreadSnapshotError({
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThreadId: snapshotInput.providerThread.id,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
rollbackThread: (rollbackInput) =>
|
|
Effect.gen(function* () {
|
|
const sessionId = nativeThreadId(rollbackInput.providerThread);
|
|
const state = threads.get(sessionId);
|
|
if (state?.activeTurn !== null && state?.activeTurn !== undefined) {
|
|
return yield* protocolError(
|
|
`Cannot roll back OpenCode thread ${rollbackInput.providerThread.id} while a turn is active`,
|
|
);
|
|
}
|
|
const response = yield* sdkCall("session.messages", { sessionID: sessionId }, () =>
|
|
client.session.messages({ sessionID: sessionId }),
|
|
);
|
|
const messages = unwrapData("session.messages", response);
|
|
let boundaryMessageId: string | undefined;
|
|
if (rollbackInput.target.type === "thread_start") {
|
|
boundaryMessageId = messages.find(({ info }) => info.role === "user")?.info.id;
|
|
} else {
|
|
boundaryMessageId = openCodeBoundaryAfterProviderTurn(
|
|
rollbackInput.providerThreadTurns,
|
|
rollbackInput.target.providerTurn.id,
|
|
);
|
|
}
|
|
let retainedThread = rollbackInput.providerThread;
|
|
if (boundaryMessageId !== undefined) {
|
|
const boundaryIndex = messages.findIndex(
|
|
({ info }) => info.id === boundaryMessageId,
|
|
);
|
|
if (boundaryIndex < 0)
|
|
return yield* protocolError(
|
|
"The OpenCode rewind boundary is no longer available.",
|
|
);
|
|
const fork = unwrapData(
|
|
"session.fork",
|
|
yield* sdkCall(
|
|
"session.fork",
|
|
{ sessionID: sessionId, messageID: boundaryMessageId },
|
|
() =>
|
|
client.session.fork({ sessionID: sessionId, messageID: boundaryMessageId }),
|
|
),
|
|
);
|
|
const retained = unwrapData(
|
|
"session.messages",
|
|
yield* sdkCall("session.messages", { sessionID: fork.id }, () =>
|
|
client.session.messages({ sessionID: fork.id }),
|
|
),
|
|
);
|
|
if (retained.length !== boundaryIndex)
|
|
return yield* protocolError(
|
|
"OpenCode did not preserve the requested rewind boundary.",
|
|
);
|
|
yield* sdkCall("session.update", { sessionID: fork.id }, () =>
|
|
client.session.update({
|
|
sessionID: fork.id,
|
|
permission: openCodePermissionRules(input.runtimePolicy),
|
|
}),
|
|
);
|
|
retainedThread = {
|
|
...rollbackInput.providerThread,
|
|
nativeThreadRef: providerRef(fork.id),
|
|
};
|
|
registerThread(fork, retainedThread, state?.appThread ?? null);
|
|
}
|
|
const snapshot = yield* readSnapshot(retainedThread);
|
|
return {
|
|
...snapshot,
|
|
providerThread: {
|
|
...snapshot.providerThread,
|
|
nativeConversationHeadRef:
|
|
rollbackInput.target.type === "provider_turn"
|
|
? rollbackInput.target.providerTurn.nativeTurnRef
|
|
: null,
|
|
},
|
|
};
|
|
}).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterRollbackThreadError({
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThreadId: rollbackInput.providerThread.id,
|
|
checkpointId: rollbackInput.target.checkpointId,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
forkThread: (forkInput) =>
|
|
Effect.gen(function* () {
|
|
const sourceSessionId = nativeThreadId(forkInput.sourceProviderThread);
|
|
const sourceState = threads.get(sourceSessionId);
|
|
if (sourceState?.activeTurn !== null && sourceState?.activeTurn !== undefined) {
|
|
return yield* protocolError(
|
|
`Cannot fork OpenCode thread ${forkInput.sourceProviderThread.id} while a turn is active`,
|
|
);
|
|
}
|
|
let boundaryMessageId: string | undefined;
|
|
if (forkInput.providerTurnId !== undefined) {
|
|
const sourceTurns = forkInput.sourceProviderTurns ?? [];
|
|
const selected = sourceTurns.find((turn) => turn.id === forkInput.providerTurnId);
|
|
if (selected === undefined) {
|
|
return yield* protocolError(
|
|
`OpenCode fork boundary turn ${forkInput.providerTurnId} was not found`,
|
|
);
|
|
}
|
|
boundaryMessageId = openCodeBoundaryAfterProviderTurn(sourceTurns, selected.id);
|
|
}
|
|
const response = yield* sdkCall(
|
|
"session.fork",
|
|
{
|
|
sessionID: sourceSessionId,
|
|
...(boundaryMessageId === undefined ? {} : { messageID: boundaryMessageId }),
|
|
},
|
|
() =>
|
|
client.session.fork({
|
|
sessionID: sourceSessionId,
|
|
...(boundaryMessageId === undefined ? {} : { messageID: boundaryMessageId }),
|
|
}),
|
|
);
|
|
const nativeSession = unwrapData("session.fork", response);
|
|
const forkedAt = yield* DateTime.now;
|
|
const providerThread = makeProviderThread({
|
|
idAllocator,
|
|
providerInstanceId: options.instanceId,
|
|
providerSessionId: input.providerSessionId,
|
|
appThreadId: forkInput.targetThreadId,
|
|
...(forkInput.ownerNodeId === undefined
|
|
? {}
|
|
: { ownerNodeId: forkInput.ownerNodeId }),
|
|
nativeSession,
|
|
forkedFrom: {
|
|
providerThreadId: forkInput.sourceProviderThread.id,
|
|
...(forkInput.providerTurnId === undefined
|
|
? {}
|
|
: { providerTurnId: forkInput.providerTurnId }),
|
|
},
|
|
now: forkedAt,
|
|
});
|
|
registerThread(nativeSession, providerThread, null);
|
|
return providerThread;
|
|
}).pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterForkThreadError({
|
|
driver: OPENCODE_PROVIDER,
|
|
providerThreadId: forkInput.sourceProviderThread.id,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
};
|
|
|
|
return runtimeSession;
|
|
},
|
|
(effect, input) =>
|
|
effect.pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapter.ProviderAdapterOpenSessionError({
|
|
driver: OPENCODE_PROVIDER,
|
|
providerSessionId: input.providerSessionId,
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
),
|
|
});
|
|
}
|
|
|
|
export type OpenCodeAdapterV2DriverEnv =
|
|
| OpenCodeRuntime.OpenCodeRuntime
|
|
| IdAllocator.IdAllocatorV2
|
|
| ProviderEventLoggers.ProviderEventLoggers
|
|
| ServerConfig.ServerConfig;
|
|
|
|
export const OpenCodeAdapterV2Driver: ProviderAdapterDriver<
|
|
OpenCodeSettings,
|
|
OpenCodeAdapterV2DriverEnv
|
|
> = {
|
|
driverKind: OPENCODE_PROVIDER,
|
|
configSchema: OpenCodeSettingsSchema,
|
|
defaultConfig: (): OpenCodeSettings => DEFAULT_OPENCODE_SETTINGS,
|
|
create: Effect.fn("OpenCodeAdapterV2Driver.create")(
|
|
function* (input: ProviderAdapterDriverCreateInput<OpenCodeSettings>) {
|
|
const hostEnvironment = yield* HostProcessEnvironment;
|
|
const openCodeRuntime = yield* OpenCodeRuntime.OpenCodeRuntime;
|
|
const idAllocator = yield* IdAllocator.IdAllocatorV2;
|
|
const providerEventLoggers = yield* ProviderEventLoggers.ProviderEventLoggers;
|
|
const serverConfig = yield* ServerConfig.ServerConfig;
|
|
return makeOpenCodeAdapterV2({
|
|
instanceId: input.instanceId,
|
|
settings: { ...input.config, enabled: input.enabled },
|
|
environment: mergeProviderInstanceEnvironment(input.environment, hostEnvironment),
|
|
runtime: openCodeRuntime,
|
|
idAllocator,
|
|
serverConfig,
|
|
...(providerEventLoggers.native === undefined
|
|
? {}
|
|
: { nativeEventLogger: providerEventLoggers.native }),
|
|
});
|
|
},
|
|
(effect, input) =>
|
|
effect.pipe(
|
|
Effect.mapError(
|
|
(cause) =>
|
|
new ProviderAdapterDriverCreateError({
|
|
driver: OPENCODE_PROVIDER,
|
|
instanceId: input.instanceId,
|
|
detail: "Failed to create OpenCode v2 adapter.",
|
|
cause,
|
|
}),
|
|
),
|
|
),
|
|
),
|
|
};
|
|
|
|
const layer: Layer.Layer<ProviderAdapter.ProviderAdapterV2, never, OpenCodeAdapterV2DriverEnv> =
|
|
Layer.effect(
|
|
ProviderAdapter.ProviderAdapterV2,
|
|
Effect.gen(function* () {
|
|
const hostEnvironment = yield* HostProcessEnvironment;
|
|
const openCodeRuntime = yield* OpenCodeRuntime.OpenCodeRuntime;
|
|
const idAllocator = yield* IdAllocator.IdAllocatorV2;
|
|
const providerEventLoggers = yield* ProviderEventLoggers.ProviderEventLoggers;
|
|
const serverConfig = yield* ServerConfig.ServerConfig;
|
|
return makeOpenCodeAdapterV2({
|
|
instanceId: OPENCODE_DEFAULT_INSTANCE_ID,
|
|
settings: DEFAULT_OPENCODE_SETTINGS,
|
|
environment: hostEnvironment,
|
|
runtime: openCodeRuntime,
|
|
idAllocator,
|
|
serverConfig,
|
|
...(providerEventLoggers.native === undefined
|
|
? {}
|
|
: { nativeEventLogger: providerEventLoggers.native }),
|
|
});
|
|
}),
|
|
);
|