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
150 changes: 82 additions & 68 deletions apps/server/src/mcp/toolkits/environment/handlers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,75 +26,89 @@ import { EnvironmentToolkit } from "./tools.ts";
const environmentId = EnvironmentId.make("environment:preferences");
const threadId = ThreadId.make("thread:preferences");

it.effect("refuses a preferences update when the caller's turn ends while it waits", () =>
Effect.gen(function* () {
const caller = yield* Ref.make<OrchestrationV2ThreadShell>(liveThreadShell(threadId));
const updates = yield* Ref.make(0);
// Completes once the declaration's own check has read the caller.
const checked = yield* Deferred.make<void>();
const layerDependencies = Layer.mergeAll(
ThreadCommandExecutor.layer,
Layer.succeed(McpInvocationContext.McpInvocationContext, {
environmentId,
requestNamespace: "provider:preferences",
thread: {
threadId,
providerSessionId: "provider:preferences",
providerInstanceId: ProviderInstanceId.make("codex"),
},
client: undefined,
capabilities: new Set(["orchestration" as const]),
issuedAt: 0,
}),
Layer.mock(ThreadManagement.ThreadManagementService)({
getThreadShell: () =>
Ref.get(caller).pipe(Effect.tap(() => Deferred.succeed(checked, undefined))),
}),
Layer.mock(Environment.ServerEnvironment)({
getDescriptor: Effect.succeed({
it.effect.each([
{ change: "turn ends", patch: { activeRunId: null }, code: "parent_not_active" },
{
change: "runtime mode changes",
patch: { runtimeMode: "approval-required" },
code: "capability_denied",
},
{
change: "interaction mode changes",
patch: { interactionMode: "plan" },
code: "capability_denied",
},
] as const)(
"refuses a preferences update when the caller's $change while it waits",
({ patch, code }) =>
Effect.gen(function* () {
const caller = yield* Ref.make<OrchestrationV2ThreadShell>(liveThreadShell(threadId));
const updates = yield* Ref.make(0);
// Completes once the declaration's own check has read the caller.
const checked = yield* Deferred.make<void>();
const layerDependencies = Layer.mergeAll(
ThreadCommandExecutor.layer,
Layer.succeed(McpInvocationContext.McpInvocationContext, {
environmentId,
label: "Test",
platform: { os: "linux", arch: "x64" },
serverVersion: "0.0.0",
capabilities: { repositoryIdentity: false },
requestNamespace: "provider:preferences",
thread: {
threadId,
providerSessionId: "provider:preferences",
providerInstanceId: ProviderInstanceId.make("codex"),
},
client: undefined,
capabilities: new Set(["orchestration" as const]),
issuedAt: 0,
}),
}),
Layer.mock(Settings.ServerSettingsService)({
getSettings: Effect.succeed(DEFAULT_SERVER_SETTINGS),
updateSettings: () =>
Ref.update(updates, (count) => count + 1).pipe(Effect.as(DEFAULT_SERVER_SETTINGS)),
}),
);
yield* Effect.gen(function* () {
const toolkit = yield* EnvironmentToolkit;
const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor;
const held = yield* Deferred.make<void>();
const release = yield* Deferred.make<void>();
// A turn-completion command holds the thread's lock.
const holder = yield* executor
.withLock(
threadId,
Deferred.succeed(held, undefined).pipe(Effect.andThen(Deferred.await(release))),
)
.pipe(Effect.forkChild);
yield* Deferred.await(held);
const update = yield* toolkit
.handle("t3_environment_preferences_update", { newWorktreesStartFromOrigin: true })
.pipe(Stream.unwrap, Stream.runCollect, Effect.forkChild);
// The update passed its first check and waits for the lock; the turn then ends.
yield* Deferred.await(checked);
yield* Ref.update(caller, (shell) => ({ ...shell, activeRunId: null }));
yield* Deferred.succeed(release, undefined);
yield* Fiber.join(holder);
const result = yield* Fiber.join(update);
expect(result.at(-1)?.result).toMatchObject({ code: "parent_not_active" });
expect(yield* Ref.get(updates)).toBe(0);
}).pipe(
Effect.provide(
McpToolAccess.HandlersLayer.layer(EnvironmentHandlers.layer).pipe(
Layer.provideMerge(layerDependencies),
Layer.mock(ThreadManagement.ThreadManagementService)({
getThreadShell: () =>
Ref.get(caller).pipe(Effect.tap(() => Deferred.succeed(checked, undefined))),
}),
Layer.mock(Environment.ServerEnvironment)({
getDescriptor: Effect.succeed({
environmentId,
label: "Test",
platform: { os: "linux", arch: "x64" },
serverVersion: "0.0.0",
capabilities: { repositoryIdentity: false },
}),
}),
Layer.mock(Settings.ServerSettingsService)({
getSettings: Effect.succeed(DEFAULT_SERVER_SETTINGS),
updateSettings: () =>
Ref.update(updates, (count) => count + 1).pipe(Effect.as(DEFAULT_SERVER_SETTINGS)),
}),
);
yield* Effect.gen(function* () {
const toolkit = yield* EnvironmentToolkit;
const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor;
const held = yield* Deferred.make<void>();
const release = yield* Deferred.make<void>();
// A turn-completion command holds the thread's lock.
const holder = yield* executor
.withLock(
threadId,
Deferred.succeed(held, undefined).pipe(Effect.andThen(Deferred.await(release))),
)
.pipe(Effect.forkChild);
yield* Deferred.await(held);
const update = yield* toolkit
.handle("t3_environment_preferences_update", { newWorktreesStartFromOrigin: true })
.pipe(Stream.unwrap, Stream.runCollect, Effect.forkChild);
// The caller changes after the first check, while the update waits for the lock.
yield* Deferred.await(checked);
yield* Ref.update(caller, (shell) => ({ ...shell, ...patch }));
yield* Deferred.succeed(release, undefined);
yield* Fiber.join(holder);
const result = yield* Fiber.join(update);
expect(result.at(-1)?.result).toMatchObject({ code });
expect(yield* Ref.get(updates)).toBe(0);
}).pipe(
Effect.provide(
McpToolAccess.HandlersLayer.layer(EnvironmentHandlers.layer).pipe(
Layer.provideMerge(layerDependencies),
),
),
),
);
}),
);
}),
);
144 changes: 142 additions & 2 deletions apps/server/src/orchestration-v2/runtimeLayer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
CommandId,
ContextTransferId,
EventId,
EnvironmentId,
MessageId,
NodeId,
RuntimeRequestId,
Expand Down Expand Up @@ -38,6 +39,7 @@ import * as Option from "effect/Option";
import * as Queue from "effect/Queue";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";
import { McpSchema, McpServer } from "effect/ai";
import * as SqlClient from "effect/sql/SqlClient";

import * as CheckpointStore from "../checkpointing/CheckpointStore.ts";
Expand All @@ -49,6 +51,12 @@ import * as ProjectEnrichmentService from "../project/ProjectEnrichmentService.t
import * as ProjectService from "../project/ProjectService.ts";
import * as ServerSettings from "../serverSettings.ts";
import * as McpSessionRegistryTestkit from "../mcp/McpSessionRegistry.testkit.ts";
import * as McpHttpServer from "../mcp/McpHttpServer.ts";
import * as McpInvocationContext from "../mcp/McpInvocationContext.ts";
import * as McpToolAccessTestkit from "../mcp/McpToolAccess.testkit.ts";
import * as EnvironmentHandlers from "../mcp/toolkits/environment/handlers.ts";
import { EnvironmentToolkit } from "../mcp/toolkits/environment/tools.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import * as ProviderInstanceRegistry from "../provider/ProviderInstanceRegistry.ts";
import type { ProviderInstance } from "@t3tools/provider-core/server/driver";
import * as VcsDriverRegistry from "../vcs/VcsDriverRegistry.ts";
Expand Down Expand Up @@ -272,7 +280,6 @@ const layerTest = Layer.mergeAll(
ProjectStore.layer,
ProjectionStore.layer,
EffectOutbox.layer,
ThreadCommandExecutor.layer,
).pipe(
Layer.provide(McpSessionRegistryTestkit.layer),
Layer.provide(SqlitePersistence.layerMemory),
Expand Down Expand Up @@ -301,7 +308,6 @@ const layerProjectDeletionTest = Layer.mergeAll(
RuntimeLayer.layer.pipe(Layer.provide(RuntimeLayer.layerProjectService)),
RuntimeLayer.layerProjectService,
RuntimeLayer.layerEventSink,
ThreadCommandExecutor.layer,
).pipe(
Layer.provide(
Layer.mock(ProjectEnrichmentService.ProjectEnrichmentService)({
Expand Down Expand Up @@ -480,6 +486,140 @@ const layerSharedApplicationDataPlaneTest = Layer.mergeAll(
);

it.layer(layerTest)("OrchestrationV2LayerLive", (it) => {
it.effect("serves preferences with the runtime's shared thread lock and OAuth access", () =>
Effect.gen(function* () {
const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor;
const environmentId = EnvironmentId.make("environment:runtime-preferences");
const threadId = ThreadId.make("thread:runtime-preferences");
const invocation: McpInvocationContext.McpInvocationScope = {
environmentId,
requestNamespace: "provider:runtime-preferences",
thread: {
threadId,
providerSessionId: "provider:runtime-preferences",
providerInstanceId: ProviderInstanceId.make("codex"),
},
client: undefined,
capabilities: new Set(["orchestration"]),
issuedAt: 0,
};
const client = McpSchema.McpServerClient.of({
clientId: 1,
clientCapabilities: {},
clientInfo: { name: "preferences-test", version: "1.0.0" },
protocolVersion: "2025-06-18",
initializePayload: {
protocolVersion: "2025-06-18",
capabilities: {},
clientInfo: { name: "preferences-test", version: "1.0.0" },
},
getClient: Effect.die("unused"),
});
const registration = McpHttpServer.toolkitRegistration(
EnvironmentToolkit,
EnvironmentHandlers.layer,
).pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provideMerge(
ServerSettings.layerTest({
sourceControlWritingStyle: {
mode: "custom",
customInstructions: "Preserve these instructions.",
followChangeRequestTemplates: true,
},
}),
),
Layer.provide(McpToolAccessTestkit.liveThreadsLayer),
Layer.provide(
Layer.mock(ServerEnvironment.ServerEnvironment)({
getDescriptor: Effect.succeed({
environmentId,
label: "Preferences test",
platform: { os: "linux", arch: "x64" },
serverVersion: "0.0.0",
capabilities: { repositoryIdentity: false },
}),
}),
),
);
yield* Effect.gen(function* () {
const server = yield* McpServer.McpServer;
const settings = yield* ServerSettings.ServerSettingsService;
const call = (
scope: McpInvocationContext.McpInvocationScope,
name: string,
parameters: Record<string, unknown> = {},
) =>
server
.callTool({ name, arguments: parameters })
.pipe(
Effect.provideService(McpInvocationContext.McpInvocationContext, scope),
Effect.provideService(McpSchema.McpServerClient, client),
);
const before = yield* settings.getSettings;
const read = yield* call(invocation, "t3_environment_read");
assert.equal(read.isError, false);
const held = yield* Deferred.make<void>();
const release = yield* Deferred.make<void>();
const requested = yield* Deferred.make<void>();
const holder = yield* executor
.withLock(
threadId,
Deferred.succeed(held, undefined).pipe(Effect.andThen(Deferred.await(release))),
)
.pipe(Effect.forkChild);
yield* Deferred.await(held);
const withLock = executor.withLock;
const observeLock: ThreadCommandExecutor.ThreadCommandExecutor["Service"]["withLock"] = (
key,
effect,
) => Deferred.succeed(requested, undefined).pipe(Effect.andThen(withLock(key, effect)));
const spy = vi.spyOn(executor, "withLock").mockImplementation(observeLock);
yield* Effect.gen(function* () {
const update = yield* call(invocation, "t3_environment_preferences_update", {
newWorktreesStartFromOrigin: !before.newWorktreesStartFromOrigin,
}).pipe(Effect.forkChild);
yield* Deferred.await(requested);
assert.deepEqual(yield* settings.getSettings, before);
const oauth: McpInvocationContext.McpInvocationScope = {
...invocation,
requestNamespace: "client:runtime-preferences",
thread: undefined,
client: { sessionId: "runtime-preferences", label: "Test", access: "full-access" },
};
// A client has no thread lock, even while a provider's update is waiting.
const reapplied = yield* call(oauth, "t3_environment_preferences_update", {
sourceControlWritingStyle: {
customInstructions: before.sourceControlWritingStyle.customInstructions,
},
});
assert.equal(reapplied.isError, false);
assert.equal(spy.mock.calls.length, 1);
assert.deepEqual(yield* settings.getSettings, before);
for (const access of ["read-only", "approval-required", "auto"] as const) {
const denied = yield* call(
{ ...oauth, client: { ...oauth.client!, access } },
"t3_environment_preferences_update",
{ newWorktreesStartFromOrigin: true },
);
assert.equal(denied.isError, true);
assert.match(JSON.stringify(denied.content), /capability_denied/);
}
yield* Deferred.succeed(release, undefined);
yield* Fiber.join(holder);
assert.equal((yield* Fiber.join(update)).isError, false);
assert.deepEqual(yield* settings.getSettings, {
...before,
newWorktreesStartFromOrigin: !before.newWorktreesStartFromOrigin,
});
}).pipe(
Effect.ensuring(Deferred.succeed(release, undefined)),
Effect.ensuring(Effect.sync(() => spy.mockRestore())),
);
}).pipe(Effect.provide(registration));
}),
);

it.effect("emits model updates separately from provider switches", () =>
Effect.gen(function* () {
const orchestrator = yield* Orchestrator.OrchestratorV2;
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/orchestration-v2/runtimeLayer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -330,6 +330,7 @@ const layerMcpAppRequestsProvided = McpAppRequests.layer.pipe(
);

export const layer = Layer.mergeAll(
ThreadCommandExecutor.layer,
layerOrchestratorProvided,
layerMcpAppRequestsProvided,
layerThreadManagementProvided,
Expand Down
Loading