Skip to content
Open
8 changes: 8 additions & 0 deletions apps/server/integration/transferBudgetV2.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
AuthSessionId,
AuthOrchestrationReadScope,
EnvironmentHttpApi,
EnvironmentId,
EnvironmentAuthenticatedAuth,
EnvironmentAuthenticatedPrincipal,
ORCHESTRATION_V2_WS_METHODS,
Expand All @@ -35,6 +36,7 @@ import * as HttpApi from "effect/http-api/HttpApi";
import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder";
import { HttpRouter, HttpServer } from "effect/http";
import { Rpc, RpcGroup, RpcServer, RpcSerialization } from "effect/rpc";
import * as ServerEnvironment from "../src/environment/ServerEnvironment.ts";
import * as SqlitePersistence from "../src/persistence/Sqlite.ts";
import * as OrchestrationEventStore from "../src/persistence/OrchestrationEventStore.ts";
import * as EventStore from "../src/orchestration-v2/EventStore.ts";
Expand Down Expand Up @@ -111,6 +113,11 @@ const layerEnrichment = Layer.unwrap(
);
// The transfer history has no project events, so shell streams never read a project shell.
const layerServices = layerManagement.pipe(
Layer.provideMerge(
Layer.succeed(ServerEnvironment.ServerEnvironmentIdentity, {
getEnvironmentId: Effect.succeed(EnvironmentId.make("transfer-test")),
}),
),
Layer.provideMerge(ProjectStore.layer),
Layer.provideMerge(Layer.mock(ProjectService.ProjectService)({})),
Layer.provideMerge(layerEnrichment),
Expand Down Expand Up @@ -237,6 +244,7 @@ it.live(
Layer.provide(Layer.succeedContext(context)),
Layer.provideMerge(
NodeHttpServer.layer(NodeHttp.createServer, {
host: "127.0.0.1",
port: 0,
websocket: { perMessageDeflate: true },
}),
Expand Down
48 changes: 48 additions & 0 deletions apps/server/src/orchestration-v2/ThreadManagementService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,14 +21,17 @@ import {
type OrchestrationV2Run,
type OrchestrationV2ThreadShellSnapshot,
type OrchestrationV2ThreadProjection,
type OrchestrationV2ThreadTranscript,
type OrchestrationV2ThreadShell,
type OrchestrationV2TurnItem,
ProjectId,
PROVIDER_SEND_TURN_MAX_FILE_BYTES,
RunId,
type ScheduledTaskId,
ThreadId,
type TurnItemId,
} from "@t3tools/contracts";
import { threadTranscriptHeader } from "@t3tools/shared/threadTranscript";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Duration from "effect/Duration";
Expand All @@ -39,6 +42,8 @@ import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";

import * as Orchestrator from "./Orchestrator.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import { projectedRowEncodedBytes } from "./threadHistoryPaging.ts";
import { projectTurnItemForDetail } from "./WireProjection.ts";
import * as LegacyV1ThreadImporter from "./legacy/LegacyV1ThreadImporter.ts";

Expand Down Expand Up @@ -176,6 +181,15 @@ export class ThreadManagementThreadNotFoundError extends Schema.TaggedError<Thre
}
}

export class ThreadTranscriptTooLargeError extends Schema.TaggedError<ThreadTranscriptTooLargeError>()(
"ThreadTranscriptTooLargeError",
{ threadId: ThreadId },
) {
override get message(): string {
return "The thread transcript exceeds the file attachment limit.";
}
}

export class ThreadManagementRunNotFoundError extends Schema.TaggedError<ThreadManagementRunNotFoundError>()(
"ThreadManagementRunNotFoundError",
{
Expand Down Expand Up @@ -304,6 +318,13 @@ export interface ThreadManagementServiceShape {
) => Effect.Effect<OrchestrationV2ThreadProjection, Orchestrator.OrchestratorV2Error>;
readonly getCheckpointContext: Orchestrator.OrchestratorV2["Service"]["getCheckpointContext"];
readonly getThreadSnapshot: Orchestrator.OrchestratorV2["Service"]["getThreadSnapshot"];
readonly getThreadTranscript: (
threadId: ThreadId,
) => Effect.Effect<
OrchestrationV2ThreadTranscript,
Orchestrator.OrchestratorV2Error | ThreadTranscriptTooLargeError,
ServerEnvironment.ServerEnvironmentIdentity
>;
readonly getThreadSnapshotWindow: Orchestrator.OrchestratorV2["Service"]["getThreadSnapshotWindow"];
readonly getProjectThreadRecords: <K extends ProjectionRecordField>(
input: { readonly projectId: ProjectId; readonly threadId: ThreadId },
Expand Down Expand Up @@ -501,6 +522,32 @@ const make = Effect.gen(function* () {
ensureProjectionTranscript(threadId).pipe(
Effect.andThen(orchestrator.getThreadSnapshot(threadId)),
);
const getThreadTranscript = Effect.fn("orchestrationV2.threadManagement.getThreadTranscript")(
function* (threadId: ThreadId) {
const { projection } = yield* getThreadSnapshot(threadId);
const environment = yield* ServerEnvironment.ServerEnvironmentIdentity;
const transcript = {
threadId: projection.thread.id,
title: projection.thread.title,
updatedAt: projection.updatedAt,
items: projection.visibleTurnItems,
};
let sizeBytes = Buffer.byteLength(
threadTranscriptHeader(yield* environment.getEnvironmentId, transcript),
"utf8",
);
if (sizeBytes > PROVIDER_SEND_TURN_MAX_FILE_BYTES) {
return yield* new ThreadTranscriptTooLargeError({ threadId });
}
for (const row of transcript.items) {
sizeBytes += projectedRowEncodedBytes(row) + 1;
if (sizeBytes > PROVIDER_SEND_TURN_MAX_FILE_BYTES) {
return yield* new ThreadTranscriptTooLargeError({ threadId });
}
}
return transcript;
},
);
const getThreadSnapshotWindow: ThreadManagementServiceShape["getThreadSnapshotWindow"] = (
threadId,
options,
Expand Down Expand Up @@ -907,6 +954,7 @@ const make = Effect.gen(function* () {
getThreadProjection,
getCheckpointContext,
getThreadSnapshot,
getThreadTranscript,
getThreadSnapshotWindow,
getProjectThreadRecords,
getProjectThread,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import {
EnvironmentAuthenticatedAuth,
EnvironmentAuthenticatedPrincipal,
EnvironmentHttpApi,
EnvironmentId,
EventId,
MessageId,
NodeId,
Expand Down Expand Up @@ -34,6 +35,7 @@ import * as HttpApi from "effect/http-api/HttpApi";
import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder";
import { Etag, HttpRouter } from "effect/http";

import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import * as SqlitePersistence from "../persistence/Sqlite.ts";
import * as OrchestrationEventStore from "../persistence/OrchestrationEventStore.ts";
import * as ProjectEnrichmentService from "../project/ProjectEnrichmentService.ts";
Expand Down Expand Up @@ -226,6 +228,9 @@ const seed = Effect.gen(function* () {
// thread routes never touch projects.
const TestLayer = Layer.mergeAll(
management,
Layer.succeed(ServerEnvironment.ServerEnvironmentIdentity, {
getEnvironmentId: Effect.succeed(EnvironmentId.make("compact-transport")),
}),
Layer.effectDiscard(seed),
Layer.mock(OrchestrationEventStore.OrchestrationEventStore)({}),
Layer.mock(ProjectStore.ProjectStoreV2)({}),
Expand Down
238 changes: 238 additions & 0 deletions apps/server/src/orchestration-v2/http.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,238 @@
import {
AuthSessionId,
type AuthEnvironmentScope,
EnvironmentAuthenticatedAuth,
EnvironmentAuthenticatedPrincipal,
EnvironmentHttpApi,
EnvironmentId,
EventId,
ORCHESTRATION_PROTOCOL_HEADER,
ORCHESTRATION_PROTOCOL_VERSION_TEXT,
OrchestrationV2AppThread,
OrchestrationV2TurnItem,
OrchestrationV2ThreadTranscript,
PROVIDER_SEND_TURN_MAX_FILE_BYTES,
ThreadId,
} from "@t3tools/contracts";
import { expect, it } from "@effect/vitest";
import { threadTranscriptHeader } from "@t3tools/shared/threadTranscript";
import * as Effect from "effect/Effect";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Layer from "effect/Layer";
import * as Schema from "effect/Schema";
import * as HttpRouter from "effect/http/HttpRouter";
import * as HttpPlatform from "effect/http/HttpPlatform";
import * as Etag from "effect/http/Etag";
import * as NodeServices from "@effect/platform-node/NodeServices";
import * as HttpApi from "effect/http-api/HttpApi";
import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import * as SqlitePersistence from "../persistence/Sqlite.ts";

import * as OrchestrationEventStore from "../persistence/OrchestrationEventStore.ts";
import { ProjectEnrichmentService } from "../project/ProjectEnrichmentService.ts";
import * as OrchestrationHttp from "./http.ts";
import { EventSinkV2 } from "./EventSink.ts";
import * as ProjectStore from "./ProjectStore.ts";
import * as ThreadManagementService from "./ThreadManagementService.ts";
import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts";
import * as ProviderReplayHarness from "./testkit/ProviderReplayHarness.ts";

class TranscriptTestApi extends HttpApi.make("environment").add(
EnvironmentHttpApi.groups.orchestration,
) {}
const decodeTranscript = Schema.decodeUnknownEffect(
Schema.toCodecJson(OrchestrationV2ThreadTranscript),
);

const now = DateTime.makeUnsafe("2026-09-29T00:00:00.000Z");
const environmentId = EnvironmentId.make("source-بيئة");
const output = "Full tool output. ".repeat(4_000);
const thread = Schema.decodeUnknownSync(OrchestrationV2AppThread)({
id: "source-thread",
projectId: "source-project",
title: "Source 界\nconversation",
providerInstanceId: "codex",
modelSelection: { instanceId: "codex", model: "gpt-5.4" },
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
activeProviderThreadId: null,
lineage: { rootThreadId: "source-thread", parentThreadId: null, relationshipToParent: null },
forkedFrom: null,
createdBy: "user",
creationSource: "web",
createdAt: now,
updatedAt: now,
archivedAt: null,
settledOverride: null,
settledAt: null,
lastVisitedAt: null,
deletedAt: null,
});
const items = Schema.decodeUnknownSync(Schema.Array(OrchestrationV2TurnItem))(
Array.from({ length: 200 }, (_, position) => ({
id: `item-${position}`,
threadId: "source-thread",
runId: null,
nodeId: null,
providerThreadId: null,
providerTurnId: null,
nativeItemRef: null,
parentItemId: null,
ordinal: position,
type: "command_execution",
status: "completed",
title: "Read notes",
input: "cat notes.md",
output: position === 0 ? output : "Done",
startedAt: now,
completedAt: now,
updatedAt: now,
})),
);

it.effect("exports full history with size, scope, and missing-thread checks", () =>
Effect.gen(function* () {
let allowed = true;
const runtime = ProviderReplayHarness.layerWithRegistry(
{ name: "thread-transcript" },
ProviderAdapterRegistry.layerFromAdapters([]),
{ runEffectWorker: false },
);
const services = yield* Layer.build(
Layer.mergeAll(
runtime,
ThreadManagementService.layer.pipe(Layer.provide(runtime)),
OrchestrationEventStore.layer,
ProjectStore.layer,
).pipe(Layer.provideMerge(SqlitePersistence.layerMemory)),
);
const eventSink = Context.get(services, EventSinkV2);
yield* eventSink.write({
events: [
{
id: EventId.make("create-source"),
type: "thread.created",
threadId: thread.id,
occurredAt: now,
payload: thread,
},
],
});
const routes = HttpApiBuilder.layer(TranscriptTestApi).pipe(
Layer.provide(OrchestrationHttp.layer),
Layer.provide(
Layer.succeed(ServerEnvironment.ServerEnvironmentIdentity, {
getEnvironmentId: Effect.succeed(environmentId),
}),
),
Layer.provide(
Layer.succeed(EnvironmentAuthenticatedAuth, (effect) =>
Effect.suspend(() =>
effect.pipe(
Effect.provideService(EnvironmentAuthenticatedPrincipal, {
sessionId: AuthSessionId.make("test-session"),
subject: "test",
method: "bearer-access-token",
scopes: new Set<AuthEnvironmentScope>(allowed ? ["orchestration:read"] : []),
}),
),
),
),
),
Layer.provide(Layer.succeedContext(services)),
Layer.provide(Layer.mock(ProjectEnrichmentService)({})),
Layer.provide(HttpPlatform.layer),
Layer.provide(Etag.layerWeak),
Layer.provide(NodeServices.layer),
);
const { handler, dispose } = HttpRouter.toWebHandler(routes, { disableLogger: true });
yield* Effect.addFinalizer(() => Effect.promise(dispose));
const read = (threadId: ThreadId) =>
Effect.promise(() =>
handler(
new Request(`http://localhost/api/orchestration/threads/${threadId}/transcript`, {
headers: { [ORCHESTRATION_PROTOCOL_HEADER]: ORCHESTRATION_PROTOCOL_VERSION_TEXT },
}),
),
);

const emptyResponse = yield* read(thread.id);
expect(emptyResponse.status).toBe(200);
expect(
(yield* decodeTranscript(yield* Effect.promise(() => emptyResponse.json()))).items,
).toEqual([]);
yield* eventSink.write({
events: items.map((item) => ({
id: EventId.make(`create-${item.id}`),
type: "turn-item.updated" as const,
threadId: thread.id,
occurredAt: now,
payload: item,
})),
});

const response = yield* read(thread.id);
expect(response.status).toBe(200);
const transcript = yield* decodeTranscript(yield* Effect.promise(() => response.json()));
expect(transcript.items).toHaveLength(200);
expect(transcript.items[0]?.item).toMatchObject({ output });
expect(transcript.items[0]?.visibility).toBe("local");
expect(transcript.items[199]?.position).toBe(199);

const emptyOutputRows = transcript.items.map((row, position) =>
position === 0 ? { ...row, item: { ...row.item, output: "" } } : row,
);
const header = threadTranscriptHeader(environmentId, transcript);
const baseSize =
Buffer.byteLength(header, "utf8") +
emptyOutputRows.reduce(
(size, row) => size + Buffer.byteLength(`${JSON.stringify(row)}\n`, "utf8"),
0,
);
const boundaryOutput = "界" + "x".repeat(PROVIDER_SEND_TURN_MAX_FILE_BYTES - baseSize - 3);
for (const [suffix, expectedStatus] of [
["", 200],
["x", 400],
] as const) {
yield* eventSink.write({
events: [
{
id: EventId.make(`boundary-${expectedStatus}`),
type: "turn-item.updated",
threadId: thread.id,
occurredAt: now,
payload: {
...items.find((item) => item.type === "command_execution")!,
output: boundaryOutput + suffix,
},
},
],
});
const response = yield* read(thread.id);
expect(response.status).toBe(expectedStatus);
if (expectedStatus === 400) {
const error = yield* Context.get(services, ThreadManagementService.ThreadManagementService)
.getThreadTranscript(thread.id)
.pipe(
Effect.provideService(ServerEnvironment.ServerEnvironmentIdentity, {
getEnvironmentId: Effect.succeed(environmentId),
}),
Effect.flip,
);
expect(error).toBeInstanceOf(ThreadManagementService.ThreadTranscriptTooLargeError);
expect(yield* Effect.promise(() => response.json())).toMatchObject({
_tag: "EnvironmentRequestInvalidError",
reason: "thread_transcript_too_large",
});
}
}

expect((yield* read(ThreadId.make("missing"))).status).toBe(404);
allowed = false;
expect((yield* read(thread.id)).status).toBe(403);
}),
);
Loading
Loading