Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -1,15 +1,24 @@
import type {
DriverCommandUpdateInput,
DriverExternalToolEffectClaimInput,
DriverExternalToolEffectClaimOutput,
DriverExternalToolEffectCompleteInput,
DriverExternalToolEffectUnknownInput,
DriverNextCommandInput,
DriverNextCommandOutput,
} from "@mosoo/agent-driver/orpc";
import { RuntimeCommandResult } from "@mosoo/contracts/runtime-command";
import { McpExecuteCommandResult, RuntimeCommandResult } from "@mosoo/contracts/runtime-command";
import type { RuntimeCommand } from "@mosoo/contracts/runtime-command";
import { parseSchemaValue } from "@mosoo/contracts/validation";
import { parsePlatformId } from "@mosoo/id";
import type { DriverCommandId } from "@mosoo/id";

import { createErrorLogContext, logError } from "../../../../platform/cloudflare/logger";
import {
claimExternalToolEffect,
completeExternalToolEffect,
markExternalToolEffectUnknown,
} from "../session-runs/external-tool-effect-store.repository";
import {
claimNextQueuedRuntimeCommandRecord,
createRuntimeCommandRecord,
Expand Down Expand Up @@ -110,6 +119,63 @@ export class DriverInstanceRpcCommandController {
return { ok: true };
}

async handleClaimExternalToolEffect(
input: DriverExternalToolEffectClaimInput,
context: DriverInstanceRpcOperationContext,
): Promise<DriverExternalToolEffectClaimOutput> {
const { env, state } = this.#dependencies;

if (input.driverInstanceId !== state.requireDriverInstanceId()) {
throw new Error("Driver instance id mismatch.");
}
context.assertActiveConnection();

return claimExternalToolEffect(env.DB, {
commandId: parsePlatformId<DriverCommandId>(input.commandId, "driver command id"),
driverInstanceId: state.requireDriverInstanceId(),
});
}

async handleCompleteExternalToolEffect(
input: DriverExternalToolEffectCompleteInput,
context: DriverInstanceRpcOperationContext,
): Promise<{ ok: true }> {
const { env, state } = this.#dependencies;

if (input.driverInstanceId !== state.requireDriverInstanceId()) {
throw new Error("Driver instance id mismatch.");
}
context.assertActiveConnection();
await completeExternalToolEffect(env.DB, {
commandId: parsePlatformId<DriverCommandId>(input.commandId, "driver command id"),
driverInstanceId: state.requireDriverInstanceId(),
...(input.providerReceiptJson === undefined
? {}
: { providerReceiptJson: input.providerReceiptJson }),
result: parseSchemaValue(McpExecuteCommandResult, input.result),
});
context.assertActiveConnection();
return { ok: true };
}

async handleMarkExternalToolEffectUnknown(
input: DriverExternalToolEffectUnknownInput,
context: DriverInstanceRpcOperationContext,
): Promise<{ ok: true }> {
const { env, state } = this.#dependencies;

if (input.driverInstanceId !== state.requireDriverInstanceId()) {
throw new Error("Driver instance id mismatch.");
}
context.assertActiveConnection();
await markExternalToolEffectUnknown(env.DB, {
commandId: parsePlatformId<DriverCommandId>(input.commandId, "driver command id"),
driverInstanceId: state.requireDriverInstanceId(),
});
context.assertActiveConnection();
return { ok: true };
}

async handleNextCommand(
input: DriverNextCommandInput,
context: DriverInstanceRpcOperationContext,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,10 @@ import type {
DriverCompletionInput,
DriverEventBatchInput,
DriverEventBatchOutput,
DriverExternalToolEffectClaimInput,
DriverExternalToolEffectClaimOutput,
DriverExternalToolEffectCompleteInput,
DriverExternalToolEffectUnknownInput,
DriverFailureInput,
DriverHeartbeatInput,
DriverHelloInput,
Expand Down Expand Up @@ -46,6 +50,20 @@ export class DriverInstanceRpcController implements DriverInstanceRpcHandler {
return this.#commands.handleCommandUpdate(input, context);
}

async handleClaimExternalToolEffect(
input: DriverExternalToolEffectClaimInput,
context: DriverInstanceRpcOperationContext,
): Promise<DriverExternalToolEffectClaimOutput> {
return this.#commands.handleClaimExternalToolEffect(input, context);
}

async handleCompleteExternalToolEffect(
input: DriverExternalToolEffectCompleteInput,
context: DriverInstanceRpcOperationContext,
): Promise<{ ok: true }> {
return this.#commands.handleCompleteExternalToolEffect(input, context);
}

async handleCompleteRun(
input: DriverCompletionInput,
context: DriverInstanceRpcOperationContext,
Expand Down Expand Up @@ -101,6 +119,13 @@ export class DriverInstanceRpcController implements DriverInstanceRpcHandler {
return this.#events.handlePushLogs(input, context);
}

async handleMarkExternalToolEffectUnknown(
input: DriverExternalToolEffectUnknownInput,
context: DriverInstanceRpcOperationContext,
): Promise<{ ok: true }> {
return this.#commands.handleMarkExternalToolEffectUnknown(input, context);
}

async handleReady(
input: DriverReadyInput,
context: DriverInstanceRpcOperationContext,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,10 @@ import type {
DriverCompletionInput,
DriverEventBatchInput,
DriverEventBatchOutput,
DriverExternalToolEffectClaimInput,
DriverExternalToolEffectClaimOutput,
DriverExternalToolEffectCompleteInput,
DriverExternalToolEffectUnknownInput,
DriverFailureInput,
DriverHeartbeatInput,
DriverHeartbeatOutput,
Expand All @@ -16,8 +20,10 @@ import type {
DriverReadyInput,
} from "@mosoo/agent-driver/orpc";
import { DriverCapability } from "@mosoo/contracts/driver-instance";
import { ExternalToolEffectClaim } from "@mosoo/contracts/external-tool-effect";
import {
RuntimeCommand,
McpExecuteCommandResult,
RuntimeCommandResult,
RuntimeCommandStatus,
} from "@mosoo/contracts/runtime-command";
Expand All @@ -30,7 +36,7 @@ const DriverHelloInputWire = type({
capabilities: DriverCapability.array(),
driverVersion: NonEmptyString,
pid: "number",
protocolVersion: "1",
protocolVersion: "2",
runtime: '"openai-runtime" | "claude-agent-sdk" | "acp-fallback"',
startedAt: "string",
});
Expand Down Expand Up @@ -131,6 +137,23 @@ const DriverCommandUpdateInputWire = type({
status: RuntimeCommandStatus,
});

const DriverExternalToolEffectClaimInputWire = type({
commandId: NonEmptyString,
driverInstanceId: NonEmptyString,
});

const DriverExternalToolEffectCompleteInputWire = type({
commandId: NonEmptyString,
driverInstanceId: NonEmptyString,
"providerReceiptJson?": "string | null | undefined",
result: McpExecuteCommandResult,
});

const DriverExternalToolEffectUnknownInputWire = type({
commandId: NonEmptyString,
driverInstanceId: NonEmptyString,
});

const DriverNextCommandInputWire = type({
driverInstanceId: NonEmptyString,
});
Expand All @@ -154,13 +177,20 @@ type DriverNextCommandOutputWireValue = typeof DriverNextCommandOutputWire.infer

export interface RuntimeOrpcContext {
onCommandUpdate(input: DriverCommandUpdateInput): Promise<{ ok: true }>;
onClaimExternalToolEffect(
input: DriverExternalToolEffectClaimInput,
): Promise<DriverExternalToolEffectClaimOutput>;
onCompleteExternalToolEffect(input: DriverExternalToolEffectCompleteInput): Promise<{ ok: true }>;
onCompleteRun(input: DriverCompletionInput): Promise<{ ok: true }>;
onFailRun(input: DriverFailureInput): Promise<{ ok: true }>;
onHeartbeat(input: DriverHeartbeatInput): Promise<DriverHeartbeatOutput>;
onHello(input: DriverHelloInput): Promise<DriverHelloOutput>;
onNextCommand(input: DriverNextCommandInput): Promise<DriverNextCommandOutput>;
onPushEvents(input: DriverEventBatchInput): Promise<DriverEventBatchOutput>;
onPushLogs(input: DriverLogBatchInput): Promise<DriverLogBatchOutput>;
onMarkExternalToolEffectUnknown(
input: DriverExternalToolEffectUnknownInput,
): Promise<{ ok: true }>;
onReady(input: DriverReadyInput): Promise<{ ok: true }>;
onWatchCommands(): AsyncIteratorObject<RuntimeCommand>;
}
Expand All @@ -169,6 +199,24 @@ function parseDriverCommandUpdateInput(input: unknown): DriverCommandUpdateInput
return parseSchemaValue(DriverCommandUpdateInputWire, input);
}

function parseDriverExternalToolEffectClaimInput(
input: unknown,
): DriverExternalToolEffectClaimInput {
return parseSchemaValue(DriverExternalToolEffectClaimInputWire, input);
}

function parseDriverExternalToolEffectCompleteInput(
input: unknown,
): DriverExternalToolEffectCompleteInput {
return parseSchemaValue(DriverExternalToolEffectCompleteInputWire, input);
}

function parseDriverExternalToolEffectUnknownInput(
input: unknown,
): DriverExternalToolEffectUnknownInput {
return parseSchemaValue(DriverExternalToolEffectUnknownInputWire, input);
}

function parseDriverCompletionInput(input: unknown): DriverCompletionInput {
return parseSchemaValue(DriverCompletionInputWire, input);
}
Expand Down Expand Up @@ -218,12 +266,24 @@ const base = os.$context<RuntimeOrpcContext>();

export const runtimeOrpcRouter = {
driver: {
claimExternalToolEffect: base
.input(DriverExternalToolEffectClaimInputWire)
.output(ExternalToolEffectClaim)
.handler(async ({ context, input }) =>
context.onClaimExternalToolEffect(parseDriverExternalToolEffectClaimInput(input)),
),
commandUpdate: base
.input(DriverCommandUpdateInputWire)
.output(type({ ok: "true" }))
.handler(async ({ context, input }) =>
context.onCommandUpdate(parseDriverCommandUpdateInput(input)),
),
completeExternalToolEffect: base
.input(DriverExternalToolEffectCompleteInputWire)
.output(type({ ok: "true" }))
.handler(async ({ context, input }) =>
context.onCompleteExternalToolEffect(parseDriverExternalToolEffectCompleteInput(input)),
),
completeRun: base
.input(DriverCompletionInputWire)
.output(type({ ok: "true" }))
Expand Down Expand Up @@ -256,6 +316,12 @@ export const runtimeOrpcRouter = {
.input(DriverReadyInputWire)
.output(type({ ok: "true" }))
.handler(async ({ context, input }) => context.onReady(parseDriverReadyInput(input))),
markExternalToolEffectUnknown: base
.input(DriverExternalToolEffectUnknownInputWire)
.output(type({ ok: "true" }))
.handler(async ({ context, input }) =>
context.onMarkExternalToolEffectUnknown(parseDriverExternalToolEffectUnknownInput(input)),
),
},
driverInstance: {
nextCommand: base
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,10 @@ import type {
DriverCompletionInput,
DriverEventBatchInput,
DriverEventBatchOutput,
DriverExternalToolEffectClaimInput,
DriverExternalToolEffectClaimOutput,
DriverExternalToolEffectCompleteInput,
DriverExternalToolEffectUnknownInput,
DriverFailureInput,
DriverHeartbeatInput,
DriverHelloInput,
Expand All @@ -29,6 +33,14 @@ export interface DriverInstanceRpcHandler {
input: DriverCommandUpdateInput,
context: DriverInstanceRpcOperationContext,
): Promise<{ ok: true }>;
handleClaimExternalToolEffect(
input: DriverExternalToolEffectClaimInput,
context: DriverInstanceRpcOperationContext,
): Promise<DriverExternalToolEffectClaimOutput>;
handleCompleteExternalToolEffect(
input: DriverExternalToolEffectCompleteInput,
context: DriverInstanceRpcOperationContext,
): Promise<{ ok: true }>;
handleCompleteRun(
input: DriverCompletionInput,
context: DriverInstanceRpcOperationContext,
Expand Down Expand Up @@ -57,6 +69,10 @@ export interface DriverInstanceRpcHandler {
input: DriverLogBatchInput,
context: DriverInstanceRpcOperationContext,
): Promise<DriverLogBatchOutput>;
handleMarkExternalToolEffectUnknown(
input: DriverExternalToolEffectUnknownInput,
context: DriverInstanceRpcOperationContext,
): Promise<{ ok: true }>;
handleReady(
input: DriverReadyInput,
context: DriverInstanceRpcOperationContext,
Expand All @@ -69,14 +85,20 @@ export function createDriverInstanceRpcContext(
context: DriverInstanceRpcOperationContext,
): DriverInstanceRpcContext {
return {
onClaimExternalToolEffect: async (input) =>
handler.handleClaimExternalToolEffect(input, context),
onCommandUpdate: async (input) => handler.handleCommandUpdate(input, context),
onCompleteExternalToolEffect: async (input) =>
handler.handleCompleteExternalToolEffect(input, context),
onCompleteRun: async (input) => handler.handleCompleteRun(input, context),
onFailRun: async (input) => handler.handleFailRun(input, context),
onHeartbeat: async (input) => handler.handleHeartbeat(input, context),
onHello: async (input) => handler.handleHello(input, context),
onNextCommand: async (input) => handler.handleNextCommand(input, context),
onPushEvents: async (input) => handler.handlePushEvents(input, context),
onPushLogs: async (input) => handler.handlePushLogs(input, context),
onMarkExternalToolEffectUnknown: async (input) =>
handler.handleMarkExternalToolEffectUnknown(input, context),
onReady: async (input) => handler.handleReady(input, context),
onWatchCommands: () =>
handler.watchCommands(context)[Symbol.asyncIterator]() as ReturnType<
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import { classifyReclaim, decideReclaimRecovery } from "../../domain/session-run
import { isTerminalSessionRunStatus } from "../../domain/session-run-status";
import { createSessionRunTerminalFailureSourceId } from "../../domain/session-run-terminal-event-id";
import { recordRuntimeRunLeaseReleasedOutcome } from "../runtime-subject-lifecycle/runtime-run-lease-store";
import { markExecutingExternalToolEffectsUnknownForDriver } from "../session-runs/external-tool-effect-store.repository";
import { failAcceptedRuntimeCommandsForTerminalDriver } from "../session-runs/runtime-command-store.repository";
import { setSessionRunStatus } from "../session-runs/session-run-store.repository";
import type { SessionRunTransitionOutcome } from "../session-runs/session-run-store.repository";
Expand Down Expand Up @@ -110,6 +111,7 @@ export async function repairFinalizedTerminalDriverRunState(
status: "failed" | "stopped";
},
): Promise<TerminalDriverInstanceSessionRunReleaseResult> {
await markExecutingExternalToolEffectsUnknownForDriver(bindings.DB, input.driverInstanceId);
await failAcceptedRuntimeCommandsForTerminalDriver(bindings.DB, {
driverInstanceId: input.driverInstanceId,
});
Expand Down
Loading
Loading