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
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import {
} from "@t3tools/contracts";
import { describe, expect, it } from "vite-plus/test";

import { runtimeEventToActivities } from "./ProviderRuntimeIngestion.ts";
import { mergeBackgroundIngestion, runtimeEventToActivities } from "./ProviderRuntimeIngestion.ts";

const base = {
provider: ProviderDriverKind.make("codex"),
Expand Down Expand Up @@ -82,3 +82,231 @@ describe("runtimeEventToActivities task progress", () => {
expect(usagePayload).not.toHaveProperty("status");
});
});

describe("mergeBackgroundIngestion", () => {
it("preserves usage when a later task progress tick only updates activity", () => {
const taskId = RuntimeTaskId.make("agent-partial-progress");
const usageUpdate = {
source: "runtime",
event: {
...base,
type: "task.progress",
eventId: EventId.make("evt-task-usage"),
payload: {
taskId,
description: "Agent partial progress",
typedUsage: { totalTokens: 4_200 },
},
} satisfies ProviderRuntimeEvent,
} as const;
const activityUpdate = {
source: "runtime",
event: {
...base,
type: "task.progress",
eventId: EventId.make("evt-task-activity"),
payload: {
taskId,
description: "Agent partial progress",
summary: "Running tests",
lastToolName: "exec_command",
},
} satisfies ProviderRuntimeEvent,
} as const;

const merged = mergeBackgroundIngestion([usageUpdate], [activityUpdate]);

expect(merged).toHaveLength(1);
expect(merged[0]?.event.type).toBe("task.progress");
if (merged[0]?.event.type === "task.progress") {
expect(merged[0].event.payload).toEqual({
taskId,
description: "Agent partial progress",
typedUsage: { totalTokens: 4_200 },
summary: "Running tests",
lastToolName: "exec_command",
});
}
});

it("preserves activity and max-merges a later sparse usage update", () => {
const taskId = RuntimeTaskId.make("agent-later-usage");
const activityUpdate = {
source: "runtime",
event: {
...base,
type: "task.progress",
eventId: EventId.make("evt-earlier-activity"),
payload: {
taskId,
description: "Agent later usage",
summary: "Running tests",
lastToolName: "exec_command",
typedUsage: {
totalTokens: 100,
inputTokens: 80,
cachedInputTokens: 30,
outputTokens: 20,
reasoningOutputTokens: 10,
toolUses: 3,
durationMs: 1_000,
},
},
} satisfies ProviderRuntimeEvent,
} as const;
const usageUpdate = {
source: "runtime",
event: {
...base,
type: "task.progress",
eventId: EventId.make("evt-later-usage"),
payload: {
taskId,
description: "Agent later usage",
typedUsage: {
totalTokens: 200,
inputTokens: 70,
outputTokens: 40,
durationMs: 900,
},
},
} satisfies ProviderRuntimeEvent,
} as const;

const merged = mergeBackgroundIngestion([activityUpdate], [usageUpdate]);

expect(merged).toHaveLength(1);
expect(merged[0]?.event.type).toBe("task.progress");
if (merged[0]?.event.type === "task.progress") {
expect(merged[0].event.payload).toEqual({
taskId,
description: "Agent later usage",
summary: "Running tests",
lastToolName: "exec_command",
typedUsage: {
totalTokens: 200,
inputTokens: 80,
cachedInputTokens: 30,
outputTokens: 40,
reasoningOutputTokens: 10,
toolUses: 3,
durationMs: 1_000,
},
});
}
});

it("replaces stale activity fields while preserving the last usage snapshot", () => {
const taskId = RuntimeTaskId.make("agent-replaced-progress");
const previous = {
source: "runtime",
event: {
...base,
type: "task.progress",
eventId: EventId.make("evt-previous-progress"),
payload: {
taskId,
description: "Agent replaced progress",
summary: "Old summary",
lastToolName: "old_tool",
error: "Old error",
phases: [{ index: 0, title: "Old phase" }],
typedUsage: { totalTokens: 4_200 },
},
} satisfies ProviderRuntimeEvent,
} as const;
const next = {
source: "runtime",
event: {
...base,
type: "task.progress",
eventId: EventId.make("evt-next-progress"),
payload: {
taskId,
description: "Agent replaced progress",
phases: [{ index: 1, title: "New phase" }],
typedUsage: { totalTokens: 5_000 },
},
} satisfies ProviderRuntimeEvent,
} as const;

const merged = mergeBackgroundIngestion([previous], [next]);

expect(merged).toHaveLength(1);
expect(merged[0]?.event.type).toBe("task.progress");
if (merged[0]?.event.type === "task.progress") {
expect(merged[0].event.payload).toEqual({
taskId,
description: "Agent replaced progress",
phases: [{ index: 1, title: "New phase" }],
typedUsage: { totalTokens: 5_000 },
});
}
});

it("preserves explicit task reactivation order", () => {
const taskId = RuntimeTaskId.make("agent-reactivated");
const completed = {
source: "runtime",
event: {
...base,
type: "task.completed",
eventId: EventId.make("evt-task-completed"),
payload: { taskId, status: "completed" },
} satisfies ProviderRuntimeEvent,
} as const;
const reactivated = {
source: "runtime",
event: {
...base,
type: "task.updated",
eventId: EventId.make("evt-task-reactivated"),
payload: { taskId, status: "running" },
} satisfies ProviderRuntimeEvent,
} as const;

const merged = mergeBackgroundIngestion([completed], [reactivated]);

expect(merged.map((input) => input.event.type)).toEqual(["task.completed", "task.updated"]);
});

it("preserves fields from partial task status patches", () => {
const taskId = RuntimeTaskId.make("agent-partial-update");
const statusUpdate = {
source: "runtime",
event: {
...base,
type: "task.updated",
eventId: EventId.make("evt-task-status"),
payload: {
taskId,
status: "completed",
},
} satisfies ProviderRuntimeEvent,
} as const;
const completionTimeUpdate = {
source: "runtime",
event: {
...base,
type: "task.updated",
eventId: EventId.make("evt-task-ended-at"),
payload: {
taskId,
endedAt: "2026-08-06T00:00:01.000Z",
},
} satisfies ProviderRuntimeEvent,
} as const;

const merged = mergeBackgroundIngestion([statusUpdate], [completionTimeUpdate]);

expect(merged).toHaveLength(1);
expect(merged[0]?.event.type).toBe("task.updated");
if (merged[0]?.event.type === "task.updated") {
expect(merged[0].event.payload).toEqual({
taskId,
status: "completed",
endedAt: "2026-08-06T00:00:01.000Z",
});
}
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -225,6 +225,11 @@ describe("ProviderRuntimeIngestion", () => {
const workspaceRoot = makeTempDir("t3-provider-project-");
NodeFS.mkdirSync(NodePath.join(workspaceRoot, ".git"));
const provider = createProviderServiceHarness();
const backgroundLiveness = ThreadBackgroundLiveness.make();
const backgroundLivenessLayer = Layer.succeed(
ThreadBackgroundLiveness.ThreadBackgroundLivenessService,
backgroundLiveness,
);
const orchestrationLayer = OrchestrationEngineLive.pipe(
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(OrchestrationProjectionPipelineLive),
Expand All @@ -242,7 +247,7 @@ describe("ProviderRuntimeIngestion", () => {
Layer.provideMerge(projectionSnapshotLayer),
// Single shared liveness instance across ingestion (writer), the
// engine, and the snapshot query (reader).
Layer.provideMerge(ThreadBackgroundLiveness.layer),
Layer.provideMerge(backgroundLivenessLayer),
Layer.provideMerge(ThreadPlanProgress.layer),
Layer.provideMerge(SqlitePersistenceMemory),
Layer.provideMerge(Layer.succeed(ProviderService, provider.service)),
Expand Down Expand Up @@ -318,6 +323,8 @@ describe("ProviderRuntimeIngestion", () => {
readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()),
emit: provider.emit,
setProviderSession: provider.setSession,
getBackgroundLiveness: (threadId: ThreadId) =>
backgroundLiveness.getThreadBackgroundLiveness(threadId),
drain,
};
}
Expand Down Expand Up @@ -3254,6 +3261,37 @@ describe("ProviderRuntimeIngestion", () => {
).toBe("# Plan title");
});

it("does not revive background liveness after the provider session exits", async () => {
const harness = await createHarness();
const threadId = asThreadId("thread-1");

harness.emit({
type: "session.exited",
eventId: asEventId("evt-session-exited-before-stale-task"),
provider: ProviderDriverKind.make("claudeAgent"),
createdAt: "2026-01-01T00:00:01.000Z",
threadId,
payload: {},
});
await harness.drain();
expect(harness.getBackgroundLiveness(threadId)).toBeNull();

harness.emit({
type: "task.progress",
eventId: asEventId("evt-stale-task-after-session-exit"),
provider: ProviderDriverKind.make("claudeAgent"),
createdAt: "2026-01-01T00:00:02.000Z",
threadId,
payload: {
taskId: "stale-task-after-exit",
description: "This task belongs to the exited session",
},
});
await harness.drain();

expect(harness.getBackgroundLiveness(threadId)).toBeNull();
});

it("titles task activities with the task description, including on completion", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
Loading
Loading