Skip to content
Draft
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
76 changes: 76 additions & 0 deletions apps/server/scripts/acp-mock-agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,19 @@ const emitActiveToolThenHang = process.env.T3_ACP_EMIT_ACTIVE_TOOL_THEN_HANG ===
const emitGrokMonitorPostTurnPoll = process.env.T3_ACP_EMIT_GROK_MONITOR_POST_TURN_POLL === "1";
const emitGrokBackgroundTaskStarted = process.env.T3_ACP_EMIT_GROK_BACKGROUND_TASK_STARTED === "1";
const emitForeignSessionUpdates = process.env.T3_ACP_EMIT_FOREIGN_SESSION_UPDATES === "1";
// Root-session replacement: emit the first chunk under the requested session id
// and then publish every later update (including completion) under a fresh id,
// mirroring an agent-side `/fresh` provider session swap. The fresh id is not
// addressable: prompts and cancellation keep targeting the requested id.
const emitRootSessionReplacement = process.env.T3_ACP_ROOT_SESSION_REPLACEMENT === "1";
const emitStaleRootAfterReplacement = process.env.T3_ACP_STALE_ROOT_AFTER_REPLACEMENT === "1";
const emitIdleForeignSession = process.env.T3_ACP_IDLE_FOREIGN_SESSION === "1";
const hangRootSessionReplacement = process.env.T3_ACP_ROOT_SESSION_REPLACEMENT_HANG === "1";
// Completes the prompt on the requested session id while content published
// under the replacement id. Lets a strict-mode client settle the turn without
// the replacement completion.
const completeDurableAfterRootReplacement =
process.env.T3_ACP_ROOT_SESSION_REPLACEMENT_COMPLETE_DURABLE === "1";
const waitForResumeRelease = process.env.T3_ACP_WAIT_FOR_RESUME_RELEASE === "1";
const completeFirstPromptOnCancel = process.env.T3_ACP_COMPLETE_FIRST_PROMPT_ON_CANCEL === "1";
const floodStderr = process.env.T3_ACP_FLOOD_STDERR === "1";
Expand Down Expand Up @@ -122,6 +135,8 @@ let currentFast = false;
let authenticated = !requiresAuthentication;
let promptCount = 0;
let overlappingFirstPromptId: string | undefined;
let liveReplacementSessionId: string | undefined;
let previousReplacementSessionId: string | undefined;
const cancelledSessions = new Set<string>();
let configuredProvider: AcpSchema.ProviderCurrentConfig | null = null;

Expand Down Expand Up @@ -981,6 +996,67 @@ const program = Effect.gen(function* () {
return yield* finishPrompt(requestedSessionId, "end_turn");
}

if (emitRootSessionReplacement) {
previousReplacementSessionId = liveReplacementSessionId;
liveReplacementSessionId = `mock-session-fresh-${promptCount}`;
yield* agent.client.sessionUpdate({
sessionId: requestedSessionId,
update: {
sessionUpdate: "agent_message_chunk",
messageId: "mock-agent-message",
content: { type: "text", text: "root before fresh" },
},
});
yield* agent.client.sessionUpdate({
sessionId: liveReplacementSessionId,
update: {
sessionUpdate: "agent_message_chunk",
messageId: "mock-agent-message",
content: { type: "text", text: "replaced live root" },
},
});
if (emitStaleRootAfterReplacement) {
// Stragglers from replaced identities: the durable id stays accepted
// as root; a previously adopted live id is stale.
yield* agent.client.sessionUpdate({
sessionId: requestedSessionId,
update: {
sessionUpdate: "agent_message_chunk",
messageId: "mock-agent-message",
content: { type: "text", text: "durable root straggler" },
},
});
if (previousReplacementSessionId !== undefined) {
yield* agent.client.sessionUpdate({
sessionId: previousReplacementSessionId,
update: {
sessionUpdate: "agent_message_chunk",
messageId: "mock-agent-message",
content: { type: "text", text: "stale replaced root" },
},
});
}
}
if (hangRootSessionReplacement) {
return yield* Effect.never;
}
yield* finishPrompt(
completeDurableAfterRootReplacement ? requestedSessionId : liveReplacementSessionId,
"end_turn",
);
if (emitIdleForeignSession) {
yield* agent.client.sessionUpdate({
sessionId: "mock-session-idle-foreign",
update: {
sessionUpdate: "agent_message_chunk",
messageId: "mock-agent-message",
content: { type: "text", text: "foreign while idle" },
},
});
}
return {};
}

if (residualCallbackTriggerPath !== undefined) {
yield* Effect.gen(function* () {
while (!(yield* Effect.sync(() => NodeFS.existsSync(residualCallbackTriggerPath)))) {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,16 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { assert, describe, it } from "@effect/vitest";
import { ProviderInstanceId, ProviderSessionId, ThreadId } from "@t3tools/contracts";
import {
MessageId,
NodeId,
ProjectId,
ProviderInstanceId,
ProviderSessionId,
RunAttemptId,
RunId,
ThreadId,
} from "@t3tools/contracts";
import * as DateTime from "effect/DateTime";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Crypto from "effect/Crypto";
Expand All @@ -9,6 +19,7 @@ import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";
import { HttpClient, HttpClientResponse } from "effect/unstable/http";
import { ChildProcessSpawner } from "effect/unstable/process";

Expand Down Expand Up @@ -57,6 +68,11 @@ const registryLayer = Layer.succeed(
cmd: "fixture-agent",
args: [],
},
"darwin-x86_64": {
archive: "https://registry.test/unused",
cmd: "fixture-agent",
args: [],
},
"linux-x86_64": {
archive: "https://registry.test/unused",
cmd: "fixture-agent",
Expand Down Expand Up @@ -90,6 +106,7 @@ describe("AcpRegistryAdapterV2", () => {
authMethodId: "",
distribution: "auto",
customModels: [],
rootSessionReplacement: false,
});
});

Expand Down Expand Up @@ -236,4 +253,130 @@ describe("AcpRegistryAdapterV2", () => {
});
}).pipe(Effect.provide(testLayer), Effect.scoped),
);

it.effect("adopts a replaced root session only when the instance opts in", () =>
Effect.gen(function* () {
const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner;
const fileSystem = yield* FileSystem.FileSystem;
const idAllocator = yield* IdAllocatorV2;
const path = yield* Path.Path;
const serverConfig = yield* ServerConfig;
const mockAgentPath = yield* path.fromFileUrl(
new URL("../../../scripts/acp-mock-agent.ts", import.meta.url),
);
const resolver = yield* makeAcpRegistryCatalog({
cacheDir: serverConfig.providerStatusCacheDir,
registryUrl,
});
const settings = yield* decodeAcpRegistryAdapterSettings({
agentId: "fixture-agent",
commandPath: process.execPath,
authMethodId: "test",
rootSessionReplacement: true,
});
const instanceId = ProviderInstanceId.make("acp-registry-root-replacement");
const adapter = makeAcpRegistryAdapterV2({
crypto: yield* Crypto.Crypto,
instanceId,
settings,
environment: {
T3_ACP_SESSION_LIFECYCLE: "1",
T3_ACP_ROOT_SESSION_REPLACEMENT: "1",
},
childProcessSpawner,
fileSystem,
idAllocator,
resolver: {
resolve: (configuredSettings, cwd, environment) =>
resolver.resolve(configuredSettings, cwd, environment).pipe(
Effect.map((resolved) => ({
...resolved,
spawn: { ...resolved.spawn, args: [mockAgentPath] },
})),
),
},
serverConfig,
});
const threadId = ThreadId.make("thread-acp-registry-root-replacement");
const runtimePolicy = ProviderAdapterV2RuntimePolicy.make({
runtimeMode: "full-access",
interactionMode: "default",
cwd: process.cwd(),
});
const modelSelection = { instanceId, model: "default" } as const;
const runtime = yield* adapter.openSession({
threadId,
providerSessionId: ProviderSessionId.make("provider-session-acp-root-replacement"),
modelSelection,
runtimePolicy,
});
const providerThread = yield* runtime.ensureThread({
threadId,
modelSelection,
runtimePolicy,
});
const now = yield* DateTime.now;
const runId = RunId.make(`run:${threadId}:1`);
yield* runtime.startTurn({
appThread: {
createdBy: "user",
creationSource: "web",
id: threadId,
projectId: ProjectId.make(`project:${threadId}`),
title: "ACP registry root replacement",
providerInstanceId: instanceId,
modelSelection,
runtimeMode: "approval-required",
interactionMode: "default",
branch: null,
worktreePath: null,
activeProviderThreadId: providerThread.id,
lineage: {
parentThreadId: null,
relationshipToParent: null,
rootThreadId: threadId,
},
forkedFrom: null,
createdAt: now,
updatedAt: now,
archivedAt: null,
settledOverride: null,
settledAt: null,
lastVisitedAt: null,
deletedAt: null,
},
threadId,
runId,
runOrdinal: 1,
providerTurnOrdinal: 1,
attemptId: RunAttemptId.make(`attempt:${threadId}:1`),
rootNodeId: NodeId.make(`node:${threadId}:1`),
providerThread,
message: {
createdBy: "user",
creationSource: "web",
messageId: MessageId.make(`message:${threadId}:1`),
text: "hi",
attachments: [],
},
modelSelection,
runtimePolicy,
});
const events = Array.from(
yield* runtime.events.pipe(
Stream.takeUntil((event) => event.type === "turn.terminal"),
Stream.runCollect,
),
);
const assistantText = events
.flatMap((event) => (event.type === "turn_item.updated" ? [event.turnItem] : []))
.flatMap((item) =>
item.type === "assistant_message" && item.threadId === threadId ? [item.text] : [],
)
.join("");
assert.include(assistantText, "replaced live root");
// The durable native thread id still addresses session/load.
assert.equal(providerThread.nativeThreadRef?.nativeId, "mock-session-1");
}).pipe(Effect.provide(testLayer), Effect.scoped),
);
});
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,10 @@ function makeAcpRegistryRuntime(options: AcpRegistryAdapterV2Options) {
...resolved.spawn,
env: { ...resolved.spawn.env, ...processEnvironment },
},
// Generic per-instance opt-in: the operator declares that this agent
// replaces its root session on the same connection (see
// AcpSessionRuntimeOptions.adoptRootSessionReplacement).
...(options.settings.rootSessionReplacement ? { adoptRootSessionReplacement: true } : {}),
...(options.settings.authMethodId ? { authMethodId: options.settings.authMethodId } : {}),
}).pipe(
Layer.provide(
Expand Down
Loading