Skip to content
Merged
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
61 changes: 61 additions & 0 deletions apps/server/src/compadre/NativeThreadEvents.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@ import {
ProjectId,
ThreadId,
ProviderInstanceId,
ProviderDriverKind,
TurnId,
} from "@t3tools/contracts";
import * as NodeServices from "@effect/platform-node/NodeServices";
import * as Effect from "effect/Effect";
Expand All @@ -36,6 +38,7 @@ import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSna
import { ServerConfig } from "../config.ts";

import { mapNativeThreadEvent, nativeId } from "./NativeThreadEvents.ts";
import { runtimeEventToActivities } from "../orchestration/Layers/ProviderRuntimeIngestion.ts";
import { HttpRouter, HttpClient, HttpClientResponse } from "effect/unstable/http";
import * as EnvironmentAuth from "../auth/EnvironmentAuth.ts";
import * as SqlClient from "effect/unstable/sql/SqlClient";
Expand Down Expand Up @@ -116,6 +119,64 @@ async function seed(system: Awaited<ReturnType<typeof createOrchestrationSystem>
}

describe("native thread replication", () => {
it("preserves upstream compaction counts and summary through central storage and replay", async () => {
const threadId = ThreadId.make(NodeCrypto.randomUUID());
const sourceThreadId = ThreadId.make(NodeCrypto.randomUUID());
const source = await createOrchestrationSystem();
const central = await createOrchestrationSystem(true);
try {
await seed(source, sourceThreadId);
await seed(central, threadId);
await central.run(
bindNativeThreadStream({
threadId,
sourceThreadId,
epoch: 1,
sourceSequence: 0,
checkpointOffset: 0,
}),
);
const [activity] = runtimeEventToActivities({
type: "thread.state.changed",
eventId: EventId.make("compacted"),
provider: ProviderDriverKind.make("claudeAgent"),
threadId: sourceThreadId,
turnId: TurnId.make("compact-turn"),
createdAt,
payload: { state: "compacted", beforeTokens: 60_877, afterTokens: 9_651 },
});
if (!activity) throw new Error("Missing compaction activity");
await source.run(
source.engine.dispatch({
type: "thread.activity.append",
commandId: CommandId.make("compact-activity"),
threadId: sourceThreadId,
activity,
createdAt,
}),
);
const page = await source.run(readNativeEventPage(source.engine, sourceThreadId, 0));
for (const event of page.events) {
const command = mapNativeThreadEvent(sourceThreadId, threadId, event, 1);
if (!command) continue;
await central.run(central.engine.dispatch(command));
await central.run(central.engine.dispatch(command));
}
const snapshot = await central.readModel();
const activities = snapshot.threads.find((thread) => thread.id === threadId)?.activities;
expect(activities).toHaveLength(1);
expect(activities?.[0]).toMatchObject({
summary: "Compacted context 60.9K → 9.65K tokens",
kind: "context-compaction",
payload: { state: "compacted", beforeTokens: 60_877, afterTokens: 9_651 },
turnId: nativeId(sourceThreadId, "compact-turn"),
});
} finally {
await source.dispose();
await central.dispose();
}
});

it("preserves old history and replays deltas, completion, and later background output exactly once", async () => {
const threadId = ThreadId.make(NodeCrypto.randomUUID());
const sourceThreadId = ThreadId.make(NodeCrypto.randomUUID());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,28 @@ const asMessageId = (value: string): MessageId => MessageId.make(value);
const asThreadId = (value: string): ThreadId => ThreadId.make(value);
const asTurnId = (value: string): TurnId => TurnId.make(value);

it.each([
{
counts: { beforeTokens: 60_877, afterTokens: 9_651 },
label: "Compacted context 60.9K → 9.65K tokens",
},
{ counts: { beforeTokens: 120_000, afterTokens: 0 }, label: "Compacted context 120K → 0 tokens" },
{ counts: { beforeTokens: 60_877 }, label: "Context compacted" },
{ counts: { afterTokens: 9_651 }, label: "Context compacted" },
{ counts: {}, label: "Context compacted" },
])("uses the upstream compaction summary: $label", ({ counts, label }) => {
const [activity] = runtimeEventToActivities({
type: "thread.state.changed",
eventId: asEventId("compacted"),
provider: ProviderDriverKind.make("claudeAgent"),
threadId: asThreadId("thread-1"),
createdAt: "2026-09-11T00:00:00.000Z",
payload: { state: "compacted", ...counts },
});
expect(activity?.summary).toBe(label);
expect(activity?.payload).toEqual({ state: "compacted", ...counts });
});

it("persists provider stop reasons as non-visible completion metadata", () => {
const activities = runtimeEventToActivities({
type: "turn.completed",
Expand Down Expand Up @@ -3301,6 +3323,71 @@ describe("ProviderRuntimeIngestion", () => {
});
});

it.each([
{
usage: [120_000, 30_000],
previousBoundary: false,
label: "Compacted context 120K → 30K tokens",
},
{ usage: [120_000, 0], previousBoundary: false, label: "Compacted context 120K → 0 tokens" },
{ usage: [30_000, 120_000], previousBoundary: false, label: "Context compacted" },
{ usage: [120_000, 30_000], previousBoundary: true, label: "Context compacted" },
])(
"uses upstream usage fallback without reusing old boundaries: %j",
async ({ usage, previousBoundary, label }) => {
const harness = await createHarness();
const threadId = asThreadId("thread-1");
for (const [index, usedTokens] of usage.entries()) {
const createdAt = `2026-01-01T00:00:0${index + 1}.000Z`;
await harness.dispatch({
type: "thread.activity.append",
commandId: CommandId.make(`usage-${index}`),
threadId,
createdAt,
activity: {
id: asEventId(`usage-${index}`),
kind: "context-window.updated",
tone: "info",
summary: "Context window updated",
payload: { usedTokens },
turnId: null,
createdAt,
},
});
}
if (previousBoundary) {
await harness.dispatch({
type: "thread.activity.append",
commandId: CommandId.make("previous-compaction"),
threadId,
createdAt: "2026-01-01T00:00:03.000Z",
activity: {
id: asEventId("previous-compaction"),
kind: "context-compaction",
tone: "info",
summary: "Context compacted",
payload: { state: "compacted" },
turnId: null,
createdAt: "2026-01-01T00:00:03.000Z",
},
});
}
harness.emit({
type: "thread.state.changed",
eventId: asEventId("compacted-fallback"),
provider: ProviderDriverKind.make("codex"),
threadId,
createdAt: "2026-01-01T00:00:04.000Z",
payload: { state: "compacted" },
});
await harness.drain();
const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId);
expect(
thread?.activities.find((activity) => activity.id === "compacted-fallback")?.summary,
).toBe(label);
},
);

it("projects compacted thread state into context compaction activities", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
71 changes: 68 additions & 3 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,10 @@ import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Predicate from "effect/Predicate";
import * as Stream from "effect/Stream";
import { makeFairDrainableWorker } from "@t3tools/shared/DrainableWorker";
import { formatTokens } from "@t3tools/shared/usageFormat";

import { ProviderService } from "../../provider/Services/ProviderService.ts";
import { ProjectionTurnRepository } from "../../persistence/Services/ProjectionTurns.ts";
Expand Down Expand Up @@ -252,12 +254,45 @@ function assistantSegmentMessageId(baseKey: string, segmentIndex: number): Messa
function buildContextWindowActivityPayload(
event: ProviderRuntimeEvent,
): ThreadTokenUsageSnapshot | undefined {
if (event.type !== "thread.token-usage.updated" || event.payload.usage.usedTokens <= 0) {
if (event.type !== "thread.token-usage.updated" || event.payload.usage.usedTokens < 0) {
return undefined;
}
return event.payload.usage;
}

function compactedTokenCountsFromActivities(
activities: ReadonlyArray<OrchestrationThreadActivity> | undefined,
): { readonly beforeTokens: number; readonly afterTokens: number } | undefined {
const lastCompactionIndex = activities?.findLastIndex(
(activity) => activity.kind === "context-compaction",
);
const lastCompaction =
lastCompactionIndex !== undefined && lastCompactionIndex >= 0
? activities?.[lastCompactionIndex]
: undefined;
const activitiesSinceLastCompaction = activities?.slice((lastCompactionIndex ?? -1) + 1) ?? [];
const usedTokens = activitiesSinceLastCompaction.flatMap((activity) => {
if (activity.kind !== "context-window.updated") return [];
if (lastCompaction !== undefined) {
const isAfterLastCompaction =
activity.sequence !== undefined && lastCompaction.sequence !== undefined
? activity.sequence > lastCompaction.sequence
: activity.createdAt > lastCompaction.createdAt;
if (!isAfterLastCompaction) return [];
}
const payload = Predicate.isObject(activity.payload) ? activity.payload : undefined;
return Predicate.isNumber(payload?.usedTokens) && payload.usedTokens >= 0
? [payload.usedTokens]
: [];
});
const beforeTokens = usedTokens.at(-2);
const afterTokens = usedTokens.at(-1);
if (beforeTokens === undefined || afterTokens === undefined || afterTokens >= beforeTokens) {
return undefined;
}
return { beforeTokens, afterTokens };
}

function normalizeRuntimeTurnState(
value: string | undefined,
): "completed" | "failed" | "interrupted" | "cancelled" {
Expand Down Expand Up @@ -778,15 +813,24 @@ export function runtimeEventToActivities(
return [];
}

const beforeTokens = event.payload.beforeTokens;
const afterTokens = event.payload.afterTokens;
const summary =
beforeTokens !== undefined && afterTokens !== undefined
? `Compacted context ${formatTokens(beforeTokens)} → ${formatTokens(afterTokens)} tokens`
: "Context compacted";
return [
{
id: event.eventId,
createdAt: event.createdAt,
tone: "info",
kind: "context-compaction",
summary: "Context compacted",
summary,
payload: {
state: event.payload.state,
...(beforeTokens !== undefined ? { beforeTokens } : {}),
...(afterTokens !== undefined ? { afterTokens } : {}),
...(event.requestId !== undefined ? { requestId: event.requestId } : {}),
...(event.payload.detail !== undefined ? { detail: event.payload.detail } : {}),
},
turnId: toTurnId(event.turnId) ?? null,
Expand Down Expand Up @@ -2157,7 +2201,28 @@ const make = Effect.gen(function* () {
reasoningText = next.summary || next.text;
}

const activities = runtimeEventToActivities(event, taskTitle, reasoningText);
let activityEvent = event;
if (
activityEvent.type === "thread.state.changed" &&
activityEvent.payload.state === "compacted" &&
(activityEvent.payload.beforeTokens === undefined ||
activityEvent.payload.afterTokens === undefined)
) {
const threadDetail = yield* resolveThreadDetail(thread.id);
const tokenCounts = compactedTokenCountsFromActivities(threadDetail?.activities);
if (tokenCounts) {
activityEvent = {
...activityEvent,
payload: {
...activityEvent.payload,
beforeTokens: activityEvent.payload.beforeTokens ?? tokenCounts.beforeTokens,
afterTokens: activityEvent.payload.afterTokens ?? tokenCounts.afterTokens,
},
};
}
}

const activities = runtimeEventToActivities(activityEvent, taskTitle, reasoningText);
yield* Effect.forEach(activities, (activity) =>
providerCommandId(event, "thread-activity-append").pipe(
Effect.flatMap((commandId) =>
Expand Down
38 changes: 38 additions & 0 deletions apps/server/src/provider/Layers/ClaudeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -913,6 +913,44 @@ describe("ClaudeAdapterLive", () => {
});
}

for (const preTokens of [60_877, undefined]) {
it.effect(`emits upstream compaction token fields (before: ${preTokens})`, () => {
const harness = makeHarness();
return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
const eventsFiber = yield* Stream.takeUntil(
adapter.streamEvents,
(event) => event.type === "thread.state.changed" && event.payload.state === "compacted",
).pipe(Stream.runCollect, Effect.forkChild);
yield* adapter.startSession({
threadId: THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
});
yield* adapter.sendTurn({
threadId: THREAD_ID,
providerAction: { type: "compact" },
input: "/compact",
});
harness.query.emit({
type: "system",
subtype: "compact_boundary",
session_id: "sdk-session-1",
uuid: "compact-1",
compact_metadata: { trigger: "manual", pre_tokens: preTokens, post_tokens: 9_651 },
} as unknown as SDKMessage);
const events = Array.from(yield* Fiber.join(eventsFiber));
const compacted = events.find((event) => event.type === "thread.state.changed");
assert.ok(compacted?.type === "thread.state.changed");
assert.equal(compacted.payload.beforeTokens, preTokens);
assert.equal(compacted.payload.afterTokens, 9_651);
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(harness.layer),
);
});
}

it.effect("embeds image attachments in Claude user messages", () => {
const baseDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "claude-attachments-"));
const harness = makeHarness({
Expand Down
Loading
Loading