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
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
import {
EventId,
ProviderDriverKind,
RuntimeTaskId,
ThreadId,
type ProviderRuntimeEvent,
} from "@t3tools/contracts";
import { describe, expect, it } from "vite-plus/test";

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

const base = {
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-08-06T00:00:00.000Z",
threadId: ThreadId.make("thread-1"),
};

describe("runtimeEventToActivities task progress", () => {
it("persists usage independently from replaceable activity", () => {
const taskId = RuntimeTaskId.make("agent-1");
const usageOnly = {
...base,
type: "task.progress",
eventId: EventId.make("evt-usage"),
payload: {
taskId,
description: "Agent one",
typedUsage: { totalTokens: 73_700_000 },
},
} satisfies ProviderRuntimeEvent;
const command = {
...base,
type: "task.progress",
eventId: EventId.make("evt-command"),
payload: {
taskId,
description: "Agent one",
summary: "Running tests",
lastToolName: "exec_command",
},
} satisfies ProviderRuntimeEvent;

const usageActivities = runtimeEventToActivities(usageOnly);
const commandActivities = runtimeEventToActivities(command);

expect(usageActivities.map((activity) => activity.id)).toEqual(["task-usage:thread-1:agent-1"]);
expect(commandActivities.map((activity) => activity.id)).toEqual([
"task-progress:thread-1:agent-1",
]);
const usagePayload = usageActivities[0]?.payload as Record<string, unknown> | undefined;
expect(usagePayload?.typedUsage).toEqual({ totalTokens: 73_700_000 });
expect(usagePayload?.usageSnapshot).toBe(true);
});

it("splits combined progress and usage into their independent snapshots", () => {
const event = {
...base,
type: "task.progress",
eventId: EventId.make("evt-combined"),
payload: {
taskId: RuntimeTaskId.make("agent-2"),
description: "Agent two",
summary: "Inspecting the panel",
typedUsage: { totalTokens: 4_200, toolUses: 7 },
status: "running",
},
} satisfies ProviderRuntimeEvent;

const activities = runtimeEventToActivities(event);
const progressPayload = activities[0]?.payload as Record<string, unknown>;
const usagePayload = activities[1]?.payload as Record<string, unknown>;

expect(activities.map((activity) => activity.id)).toEqual([
"task-progress:thread-1:agent-2",
"task-usage:thread-1:agent-2",
]);
expect(progressPayload.summary).toBe("Inspecting the panel");
expect(progressPayload.status).toBe("running");
expect(progressPayload).not.toHaveProperty("typedUsage");
expect(usagePayload.typedUsage).toEqual({ totalTokens: 4_200, toolUses: 7 });
expect(usagePayload.usageSnapshot).toBe(true);
expect(usagePayload).not.toHaveProperty("status");
});
});
105 changes: 73 additions & 32 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -563,39 +563,80 @@ export function runtimeEventToActivities(
}

case "task.progress": {
const linkage = taskLinkageActivityFields(event.payload as Record<string, unknown>);
// Usage and activity are independent latest-state streams. Keeping them
// under separate stable ids prevents a command/reasoning update from
// replacing the last known token count (and prevents a usage-only tick
// from blanking the last meaningful activity).
const identityLinkage = { ...linkage };
delete identityLinkage.typedUsage;
delete identityLinkage.status;
delete identityLinkage.error;
const title =
event.payload.description.trim().length > 0
? { title: truncateDetail(event.payload.description, 120) }
: {};
const hasProgressState =
event.payload.typedUsage === undefined ||
event.payload.summary !== undefined ||
event.payload.lastToolName !== undefined ||
event.payload.status !== undefined ||
event.payload.error !== undefined;
return [
{
// Stable per-task id: progress is "latest state", not history, so
// each tick REPLACES the last via the activity upsert (PK + the
// replace-by-id apply in projector and client reducer). Keeps one
// progress row per task instead of thousands, so a large fleet's
// ticks can no longer evict its own start/terminal rows out of
// the 500-row retention window. Thread-scoped: activity_id is a
// GLOBAL primary key and Claude task ids are session-local, so a
// bare taskId could collide across threads and steal another
// thread's row (review finding).
id: EventId.make(`task-progress:${event.threadId}:${event.payload.taskId}`),
createdAt: event.createdAt,
tone: "info",
kind: "task.progress",
summary:
event.payload.description.trim().length > 0
? truncateDetail(event.payload.description, 120)
: "Reasoning update",
payload: {
taskId: event.payload.taskId,
...(event.payload.description.trim().length > 0
? { title: truncateDetail(event.payload.description, 120) }
: {}),
detail: truncateDetail(event.payload.summary ?? event.payload.description),
...(event.payload.summary ? { summary: truncateDetail(event.payload.summary) } : {}),
...(event.payload.lastToolName ? { lastToolName: event.payload.lastToolName } : {}),
...(event.payload.usage !== undefined ? { usage: event.payload.usage } : {}),
...taskLinkageActivityFields(event.payload as Record<string, unknown>),
},
turnId: toTurnId(event.turnId) ?? null,
...maybeSequence,
},
...(hasProgressState
? [
{
// Stable per-task id: activity is "latest state", not
// history, so each meaningful tick replaces the last. This
// bounds a large fleet to one activity row per task.
id: EventId.make(`task-progress:${event.threadId}:${event.payload.taskId}`),
createdAt: event.createdAt,
tone: "info" as const,
kind: "task.progress" as const,
summary:
event.payload.description.trim().length > 0
? truncateDetail(event.payload.description, 120)
: "Reasoning update",
payload: {
taskId: event.payload.taskId,
...title,
detail: truncateDetail(event.payload.summary ?? event.payload.description),
...(event.payload.summary
? { summary: truncateDetail(event.payload.summary) }
: {}),
...(event.payload.lastToolName
? { lastToolName: event.payload.lastToolName }
: {}),
...(event.payload.status ? { status: event.payload.status } : {}),
...(event.payload.error ? { error: event.payload.error } : {}),
...(event.payload.usage !== undefined ? { usage: event.payload.usage } : {}),
...identityLinkage,
},
turnId: toTurnId(event.turnId) ?? null,
...maybeSequence,
},
]
: []),
...(event.payload.typedUsage !== undefined
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
? [
{
id: EventId.make(`task-usage:${event.threadId}:${event.payload.taskId}`),
createdAt: event.createdAt,
tone: "info" as const,
kind: "task.progress" as const,
summary: "Task usage updated",
payload: {
taskId: event.payload.taskId,
...title,
...identityLinkage,
usageSnapshot: true,
typedUsage: event.payload.typedUsage,
},
turnId: toTurnId(event.turnId) ?? null,
...maybeSequence,
},
]
: []),
];
}

Expand Down
Loading
Loading