Skip to content
Closed
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
609 changes: 609 additions & 0 deletions apps/server/src/mcp/ConversationConfigurationMcpService.test.ts

Large diffs are not rendered by default.

737 changes: 737 additions & 0 deletions apps/server/src/mcp/ConversationConfigurationMcpService.ts

Large diffs are not rendered by default.

2 changes: 2 additions & 0 deletions apps/server/src/mcp/McpHttpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import { HttpRouter, HttpServerRequest, HttpServerResponse } from "effect/unstab
import packageJson from "../../package.json" with { type: "json" };
import * as McpInvocationContext from "./McpInvocationContext.ts";
import * as OrchestratorMcpService from "./OrchestratorMcpService.ts";
import * as ConversationConfigurationMcpService from "./ConversationConfigurationMcpService.ts";
import * as ThreadMetadataMcpService from "./ThreadMetadataMcpService.ts";
import * as McpSessionRegistry from "./McpSessionRegistry.ts";
import * as PreviewAutomationBroker from "./PreviewAutomationBroker.ts";
Expand Down Expand Up @@ -224,6 +225,7 @@ export const PreviewToolkitRegistrationLive = Layer.mergeAll(

export const OrchestratorToolkitRegistrationLive = McpServer.toolkit(OrchestratorToolkit).pipe(
Layer.provide(OrchestratorToolkitHandlersLive),
Layer.provide(ConversationConfigurationMcpService.layer),
Layer.provide(OrchestratorMcpService.layer),
Layer.provide(ThreadMetadataMcpService.layer),
);
Expand Down
645 changes: 640 additions & 5 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts

Large diffs are not rendered by default.

13 changes: 13 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import * as Effect from "effect/Effect";

import { McpInvocationContext } from "../../McpInvocationContext.ts";
import { OrchestratorMcpService } from "../../OrchestratorMcpService.ts";
import { ConversationConfigurationMcpService } from "../../ConversationConfigurationMcpService.ts";
import { ThreadMetadataMcpService } from "../../ThreadMetadataMcpService.ts";

const handlers = {
Expand Down Expand Up @@ -116,6 +117,18 @@ const handlers = {
const service = yield* OrchestratorMcpService;
return yield* service.interruptThread(scope, input);
}),
t3_thread_configuration: (input) =>
Effect.gen(function* () {
const scope = yield* McpInvocationContext;
const service = yield* ConversationConfigurationMcpService;
return yield* service.read(scope, input);
}),
t3_thread_configure: (input) =>
Effect.gen(function* () {
const scope = yield* McpInvocationContext;
const service = yield* ConversationConfigurationMcpService;
return yield* service.configure(scope, input);
}),
} satisfies Parameters<typeof OrchestratorToolkit.toLayer>[0];

export const OrchestratorToolkitHandlersLive = OrchestratorToolkit.toLayer(handlers);
26 changes: 26 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/tools.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ import {
CreateThreadsTool,
DelegateTaskTool,
ScheduleTaskTool,
ThreadConfigurationTool,
ThreadConfigureTool,
ThreadUpdateTool,
} from "./tools.ts";

Expand Down Expand Up @@ -77,4 +79,28 @@ describe("orchestrator MCP tool guidance", () => {
]);
assert.include(ThreadUpdateTool.description ?? "", "Workspace and branch changes");
});

it("publishes discoverable root objects for conversation configuration", () => {
const readSchema = Tool.getJsonSchema(ThreadConfigurationTool) as {
readonly type?: unknown;
readonly properties?: Readonly<Record<string, unknown>>;
};
const configureSchema = Tool.getJsonSchema(ThreadConfigureTool) as {
readonly type?: unknown;
readonly properties?: Readonly<Record<string, unknown>>;
};

assert.equal(readSchema.type, "object");
assert.hasAllKeys(readSchema.properties ?? {}, ["threadId"]);
assert.equal(configureSchema.type, "object");
assert.hasAllKeys(configureSchema.properties ?? {}, [
"threadId",
"providerInstanceId",
"model",
"options",
"runtimeMode",
"interactionMode",
"clientRequestId",
]);
});
});
37 changes: 37 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/tools.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,8 @@
import {
ConversationConfigurationInput,
ConversationConfigurationResult,
ConversationConfigureInput,
ConversationConfigureResult,
OrchestratorMcpCapabilitiesResult,
OrchestratorMcpCreatedThread,
OrchestratorMcpCreateThreadsInput,
Expand Down Expand Up @@ -32,6 +36,7 @@ import {
import { Tool, Toolkit } from "effect/unstable/ai";

import * as McpInvocationContext from "../../McpInvocationContext.ts";
import { ConversationConfigurationMcpService } from "../../ConversationConfigurationMcpService.ts";
import { OrchestratorMcpService } from "../../OrchestratorMcpService.ts";
import { ThreadMetadataMcpService } from "../../ThreadMetadataMcpService.ts";

Expand All @@ -40,6 +45,10 @@ const threadMetadataDependencies = [
McpInvocationContext.McpInvocationContext,
ThreadMetadataMcpService,
];
const configurationDependencies = [
McpInvocationContext.McpInvocationContext,
ConversationConfigurationMcpService,
];

export const OrchestratorCapabilitiesTool = Tool.make("orchestrator_capabilities", {
description:
Expand Down Expand Up @@ -249,6 +258,32 @@ export const ThreadInterruptTool = Tool.make("t3_thread_interrupt", {
.annotate(Tool.Title, "Interrupt a T3 thread")
.annotate(Tool.Destructive, true);

export const ThreadConfigurationTool = Tool.make("t3_thread_configuration", {
description:
"Read the current provider instance, model, model options, runtime mode, and interaction mode for a T3 thread in the calling project. Omit threadId for this thread. The result lists provider models and option descriptors plus the runtime and interaction modes allowed by this caller's permission ceiling.",
parameters: ConversationConfigurationInput,
success: ConversationConfigurationResult,
failure: OrchestratorMcpFailure,
failureMode: "return",
dependencies: configurationDependencies,
})
.annotate(Tool.Title, "Get thread configuration")
.annotate(Tool.Readonly, true)
.annotate(Tool.Destructive, false)
.annotate(Tool.Idempotent, true);

export const ThreadConfigureTool = Tool.make("t3_thread_configure", {
description:
"Change an existing T3 thread's provider, model, full model-option selection, runtime mode, or interaction mode through V2 orchestration. Omit threadId for this thread. Read t3_thread_configuration first to choose advertised values. The calling thread's runtime and interaction modes are hard ceilings. The result separates committed settings from requested provider-session detaches and next-turn context-handoff requirements; a detach can interrupt active provider work. It also reports current active and queued runs, durable command receipts, and any partial failure. Reusing clientRequestId replays accepted or rejected decisions without reapplying accepted legs. After resolving the cause of a rejected leg, use a new clientRequestId for a new attempt.",
parameters: ConversationConfigureInput,
success: ConversationConfigureResult,
failure: OrchestratorMcpFailure,
failureMode: "return",
dependencies: configurationDependencies,
})
.annotate(Tool.Title, "Configure a T3 thread")
.annotate(Tool.Destructive, true);

export const OrchestratorToolkit = Toolkit.make(
OrchestratorCapabilitiesTool,
DelegateTaskTool,
Expand All @@ -266,4 +301,6 @@ export const OrchestratorToolkit = Toolkit.make(
ThreadSendTool,
ThreadWaitTool,
ThreadInterruptTool,
ThreadConfigurationTool,
ThreadConfigureTool,
);
7 changes: 7 additions & 0 deletions apps/server/src/mcp/toolkits/worktree/registration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,15 @@ import * as NodeServices from "@effect/platform-node/NodeServices";
import { ProviderInstanceId, ThreadId } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Schema from "effect/Schema";
import { HttpBody, HttpClient, HttpRouter } from "effect/unstable/http";

import * as ServerEnvironment from "../../../environment/ServerEnvironment.ts";
import * as GitWorkflowService from "../../../git/GitWorkflowService.ts";
import { ThreadManagementService } from "../../../orchestration-v2/ThreadManagementService.ts";
import { CommandReceiptStoreV2 } from "../../../orchestration-v2/CommandReceiptStore.ts";
import { ProviderSwitchServiceV2 } from "../../../orchestration-v2/ProviderSwitchService.ts";
import * as ProjectService from "../../../project/ProjectService.ts";
import * as ProjectSetupScriptRunner from "../../../project/ProjectSetupScriptRunner.ts";
import { ProviderRegistry } from "../../../provider/Services/ProviderRegistry.ts";
Expand All @@ -22,6 +25,10 @@ import * as PreviewAutomationBroker from "../../PreviewAutomationBroker.ts";

const StubServicesLive = Layer.mergeAll(
Layer.mock(ThreadManagementService)({}),
Layer.mock(CommandReceiptStoreV2)({
getByCommandId: () => Effect.succeed(Option.none()),
}),
Layer.mock(ProviderSwitchServiceV2)({}),
Layer.mock(ProviderRegistry)({}),
Layer.mock(ScheduledTaskService)({}),
Layer.mock(ProjectService.ProjectService)({}),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -808,6 +808,7 @@ export const CLAUDE_READ_ONLY_T3_MCP_ALLOWED_TOOLS: ReadonlyArray<string> = [
"mcp__t3-code__list_scheduled_tasks",
"mcp__t3-code__t3_thread_list",
"mcp__t3-code__t3_thread_wait",
"mcp__t3-code__t3_thread_configuration",
];

// The SDK's `allowedTools` only pre-approves tool calls; availability is the
Expand Down
9 changes: 3 additions & 6 deletions apps/server/src/orchestration-v2/EventSink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -483,12 +483,9 @@ const baseLayer: Layer.Layer<
commandId: input.commandId,
events: normalized,
});
const sequence = storedEvents.at(-1)?.sequence;
if (sequence === undefined) {
return yield* Effect.die(
new Error(`Command ${input.commandId} produced no orchestration events.`),
);
}
const sequence =
storedEvents.at(-1)?.sequence ??
(yield* eventStore.latestSequence({ threadId: input.threadId }));
yield* applyStoredEvents(storedEvents);
yield* effectOutbox.enqueue(input.effects);
const receipt: CommandReceiptV2 = {
Expand Down
130 changes: 127 additions & 3 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,10 @@ import {
type OrchestrationV2ThreadProjection,
type OrchestrationV2TurnItem,
ProviderInstanceId,
type ProviderInteractionMode,
type ProviderSessionId,
RunId,
type RuntimeMode,
ThreadId,
} from "@t3tools/contracts";
import { modelSelectionsEqual } from "@t3tools/shared/model";
Expand Down Expand Up @@ -288,6 +290,42 @@ function commandThreadId(command: OrchestrationV2Command): ThreadId {
}
}

function runtimeModeRank(mode: RuntimeMode): number {
switch (mode) {
case "approval-required":
return 0;
case "auto-accept-edits":
return 1;
case "auto":
return 2;
case "full-access":
return 3;
}
}

function interactionModeRank(mode: ProviderInteractionMode): number {
return mode === "plan" ? 0 : 1;
}

function commandPolicyCeiling(command: OrchestrationV2Command) {
switch (command.type) {
case "thread.runtime-mode.set":
case "thread.interaction-mode.set":
case "thread.model-selection.set":
case "provider.switch":
return command.policyCeiling;
default:
return undefined;
}
}

function dispatchLockKeys(command: OrchestrationV2Command): ReadonlyArray<ThreadId> {
const policyCeiling = commandPolicyCeiling(command);
return [...new Set([commandThreadId(command), policyCeiling?.callerThreadId])]
.filter((threadId): threadId is ThreadId => threadId !== undefined)
.toSorted((left, right) => (left < right ? -1 : left > right ? 1 : 0));
}

function pendingThreadTitleGenerationEffect(
commandId: CommandId,
threadId: ThreadId,
Expand Down Expand Up @@ -1476,6 +1514,65 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
cause: `Thread ${command.threadId} is deleted.`,
});
}
const policyCeiling = commandPolicyCeiling(command);
if (policyCeiling !== undefined) {
const callerProjection =
policyCeiling.callerThreadId === command.threadId
? projection
: yield* projectionStore.getThreadProjection(policyCeiling.callerThreadId).pipe(
Effect.mapError(
(cause) =>
new OrchestratorProjectionError({
threadId: policyCeiling.callerThreadId,
cause,
}),
),
);
if (callerProjection.thread.deletedAt !== null) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: `Calling thread ${policyCeiling.callerThreadId} is deleted.`,
});
}
const prospectiveRuntimeMode =
command.type === "thread.runtime-mode.set" ? command.runtimeMode : thread.runtimeMode;
const prospectiveInteractionMode =
command.type === "thread.interaction-mode.set"
? command.interactionMode
: thread.interactionMode;
const isStrictlyRestrictiveModeChange =
(command.type === "thread.runtime-mode.set" &&
runtimeModeRank(command.runtimeMode) < runtimeModeRank(thread.runtimeMode)) ||
(command.type === "thread.interaction-mode.set" &&
interactionModeRank(command.interactionMode) <
interactionModeRank(thread.interactionMode));
if (
!isStrictlyRestrictiveModeChange &&
(runtimeModeRank(prospectiveRuntimeMode) > runtimeModeRank(policyCeiling.runtimeMode) ||
runtimeModeRank(prospectiveRuntimeMode) >
runtimeModeRank(callerProjection.thread.runtimeMode))
) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: `Target runtime mode ${prospectiveRuntimeMode} exceeds the calling thread ceiling.`,
});
}
if (
!isStrictlyRestrictiveModeChange &&
(interactionModeRank(prospectiveInteractionMode) >
interactionModeRank(policyCeiling.interactionMode) ||
interactionModeRank(prospectiveInteractionMode) >
interactionModeRank(callerProjection.thread.interactionMode))
) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: `Target interaction mode ${prospectiveInteractionMode} exceeds the calling thread ceiling.`,
});
}
}
if (
command.type === "thread.metadata.update" &&
command.expectedWorktreePath !== undefined &&
Expand All @@ -1487,6 +1584,26 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
cause: `Thread ${command.threadId} worktree changed before the metadata update could be applied.`,
});
}
if (
(command.type === "thread.model-selection.set" || command.type === "provider.switch") &&
command.expectedModelSelection !== undefined &&
!modelSelectionsEqual(command.expectedModelSelection, thread.modelSelection)
) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: `Thread ${command.threadId} model selection changed before the partial selection update could be applied.`,
});
}
if (
(command.type === "thread.runtime-mode.set" && command.runtimeMode === thread.runtimeMode) ||
(command.type === "thread.interaction-mode.set" &&
command.interactionMode === thread.interactionMode) ||
((command.type === "thread.model-selection.set" || command.type === "provider.switch") &&
modelSelectionsEqual(command.modelSelection, thread.modelSelection))
) {
return;
}
if (command.type === "thread.archive" && thread.archivedAt !== null) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
Expand Down Expand Up @@ -1873,7 +1990,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
command.worktreePath !== undefined &&
command.worktreePath !== thread.worktreePath
? projection.providerSessions.map((session) => session.id)
: command.type === "thread.runtime-mode.set"
: command.type === "thread.runtime-mode.set" && command.runtimeMode !== thread.runtimeMode
? projection.providerSessions
.filter(
(session) => !session.capabilities.sessions.supportsRuntimeModeSwitchInSession,
Expand Down Expand Up @@ -7106,7 +7223,11 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio

const plan = yield* dispatchOnce(command).pipe(
Effect.flatMap((planned) =>
planned.events.length > 0
planned.events.length > 0 ||
command.type === "thread.runtime-mode.set" ||
command.type === "thread.interaction-mode.set" ||
command.type === "thread.model-selection.set" ||
command.type === "provider.switch"
? Effect.succeed(planned)
: Effect.fail(
new OrchestratorDispatchError({
Expand Down Expand Up @@ -7184,7 +7305,10 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
});

const dispatchWithReceipt = (command: OrchestrationV2Command) =>
threadDispatch.withLock(commandThreadId(command), dispatchWithReceiptEffect(command));
dispatchLockKeys(command).reduceRight(
(effect, threadId) => threadDispatch.withLock(threadId, effect),
dispatchWithReceiptEffect(command),
);

const handleTerminalRun = (stored: OrchestrationV2StoredEvent) =>
Effect.gen(function* () {
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration-v2/runtimeLayer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -262,7 +262,9 @@ const providerRuntimeRecoveryProvided = providerRuntimeRecoveryLayer.pipe(

export const OrchestrationV2LayerLive = Layer.mergeAll(
orchestratorProvided,
commandReceiptStoreProvided,
threadManagementProvided,
providerSwitchServiceProvided,
effectWorkerProvided,
providerSessionManagerProvided,
providerAuthServiceProvided,
Expand Down
Loading
Loading