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
3 changes: 2 additions & 1 deletion apps/mobile/src/connection/environment-cache-store.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import {
ORCHESTRATION_CACHE_SCHEMA_VERSION,
ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
StoredOrchestrationShellSnapshot,
} from "@t3tools/client-runtime/platform";
import {
Expand Down Expand Up @@ -251,7 +252,7 @@ describe("mobile SQLite environment cache store", () => {
ORCHESTRATION_CACHE_SCHEMA_VERSION,
);
expect(memory.schemaVersions.get(cacheId(ENVIRONMENT_ID, "thread", THREAD_ID))).toBe(
ORCHESTRATION_CACHE_SCHEMA_VERSION,
ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
);
}),
);
Expand Down
11 changes: 9 additions & 2 deletions apps/mobile/src/connection/environment-cache-store.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import {
ORCHESTRATION_CACHE_SCHEMA_VERSION,
ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
StoredOrchestrationThreadSnapshot,
Persistence,
} from "@t3tools/client-runtime/platform";
Expand Down Expand Up @@ -148,13 +149,19 @@ export const make = Effect.fn("MobileEnvironmentCacheStore.make")(function* () {
saveThread: Effect.fn("MobileEnvironmentCache.saveThread")(function* (environmentId, snapshot) {
const threadId = snapshot.projection.thread.id;
const payload = yield* encodeStoredThreadSnapshot({
schemaVersion: ORCHESTRATION_CACHE_SCHEMA_VERSION,
schemaVersion: ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
environmentId,
threadId,
snapshot,
}).pipe(Effect.mapError((cause) => persistenceError("save-thread", cause)));
yield* database
.saveCache(environmentId, "thread", threadId, ORCHESTRATION_CACHE_SCHEMA_VERSION, payload)
.saveCache(
environmentId,
"thread",
threadId,
ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
payload,
)
.pipe(Effect.mapError(mapDatabaseError("save-thread")));
}),
removeThread: Effect.fn("MobileEnvironmentCache.removeThread")((environmentId, threadId) =>
Expand Down
3 changes: 2 additions & 1 deletion apps/web/src/connection/storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import {
type ConnectionCatalogDocument as ConnectionCatalogDocumentType,
EMPTY_CONNECTION_CATALOG_DOCUMENT,
ORCHESTRATION_CACHE_SCHEMA_VERSION,
ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
StoredOrchestrationShellSnapshot,
StoredOrchestrationThreadSnapshot,
decodeOrDiscardOrchestrationCache,
Expand Down Expand Up @@ -810,7 +811,7 @@ export const layer = Layer.effectContext(
saveThread: (environmentId, snapshot) =>
Effect.gen(function* () {
const encoded = yield* encodeStoredThreadSnapshot({
schemaVersion: ORCHESTRATION_CACHE_SCHEMA_VERSION,
schemaVersion: ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
environmentId,
threadId: snapshot.projection.thread.id,
snapshot,
Expand Down
22 changes: 20 additions & 2 deletions packages/client-runtime/src/platform/orchestrationCache.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import * as Schema from "effect/Schema";

import {
ORCHESTRATION_CACHE_SCHEMA_VERSION,
ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
StoredOrchestrationShellSnapshot,
StoredOrchestrationThreadSnapshot,
decodeOrDiscardOrchestrationCache,
Expand Down Expand Up @@ -102,7 +103,7 @@ describe("orchestration cache envelopes", () => {
expect(projection).not.toBeNull();
if (projection === null) throw new Error("Expected live notice projection");
const encoded = encodeStoredThreadSnapshotJson({
schemaVersion: ORCHESTRATION_CACHE_SCHEMA_VERSION,
schemaVersion: ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
environmentId,
threadId: v2ThreadId,
snapshot: { snapshotSequence: 5, projection },
Expand All @@ -121,7 +122,7 @@ describe("orchestration cache envelopes", () => {
});
const shell = decodeStoredShellSnapshotJson(encodedShell);
const encodedThread = encodeStoredThreadSnapshotJson({
schemaVersion: ORCHESTRATION_CACHE_SCHEMA_VERSION,
schemaVersion: ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
environmentId,
threadId: v2ThreadId,
snapshot: { snapshotSequence: 4, projection: v2Projection },
Expand Down Expand Up @@ -159,6 +160,23 @@ describe("orchestration cache envelopes", () => {
}),
);

it("rejects thread snapshots saved before queued-run items were kept, but not shells", () => {
// Version 3 thread windows may have dropped a queued run's items (#16987).
const staleThread = encodeStoredThreadSnapshotJson({
schemaVersion: ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION,
environmentId,
threadId: v2ThreadId,
snapshot: { snapshotSequence: 4, projection: v2Projection },
}).replace(`"schemaVersion":${ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION}`, `"schemaVersion":3`);
expect(() => decodeStoredThreadSnapshotJson(staleThread)).toThrow();
const shell = encodeStoredShellSnapshot({
schemaVersion: ORCHESTRATION_CACHE_SCHEMA_VERSION,
environmentId,
snapshot: v2ShellSnapshot,
});
expect(decodeStoredShellSnapshotSync(shell).schemaVersion).toBe(3);
});

it.effect("discards V1-versioned cache envelopes after a decode failure", () =>
Effect.gen(function* () {
let discardCount = 0;
Expand Down
7 changes: 6 additions & 1 deletion packages/client-runtime/src/platform/orchestrationCache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,11 @@ import * as Option from "effect/Option";
import * as Schema from "effect/Schema";

export const ORCHESTRATION_CACHE_SCHEMA_VERSION = 3 as const;
/**
* Thread snapshots version on their own so a thread-only invalidation keeps
* shell caches warm. 4 drops bounded windows that lost queued runs' items (#16987).
*/
export const ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION = 4 as const;

export const StoredOrchestrationShellSnapshot = Schema.Struct({
schemaVersion: Schema.Literal(ORCHESTRATION_CACHE_SCHEMA_VERSION),
Expand All @@ -18,7 +23,7 @@ export const StoredOrchestrationShellSnapshot = Schema.Struct({
});

export const StoredOrchestrationThreadSnapshot = Schema.Struct({
schemaVersion: Schema.Literal(ORCHESTRATION_CACHE_SCHEMA_VERSION),
schemaVersion: Schema.Literal(ORCHESTRATION_THREAD_CACHE_SCHEMA_VERSION),
environmentId: EnvironmentId,
threadId: ThreadId,
snapshot: OrchestrationV2ThreadDetailSnapshot.mapFields((fields) => ({
Expand Down
127 changes: 127 additions & 0 deletions packages/client-runtime/src/state/orchestrationV2Projection.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -275,6 +275,133 @@ describe("applyOrchestrationV2ProjectionEvent", () => {
expect(next?.visibleTurnItems.map((row) => row.position)).toEqual([0, 1]);
});

describe("partial timeline", () => {
const queuedRunId = RunId.make("run-queued");
const steerRunId = RunId.make("run-steer");
const queuedRun = { ...run, id: queuedRunId, ordinal: 2, status: "queued" } as const;
const steerRun = { ...run, id: steerRunId, ordinal: 3, status: "completed" } as const;
const runItem = (id: string, itemRunId: RunId, ordinal: number) => ({
...commandItem(id, id, ordinal),
runId: itemRunId,
});
const windowOf = (
runs: ReadonlyArray<OrchestrationV2Run>,
items: ReadonlyArray<OrchestrationV2TurnItem>,
): OrchestrationV2ThreadProjection => ({
...emptyProjection,
runs,
turnItems: items,
visibleTurnItems: items.map((item, position) => ({
position,
visibility: "local" as const,
sourceThreadId: threadId,
sourceItemId: item.id,
item,
})),
});
const itemEvent = (payload: OrchestrationV2TurnItem) =>
({
id: `event-${payload.id}`,
type: "turn-item.updated",
threadId,
occurredAt: now,
payload,
}) as OrchestrationV2DomainEvent;
const runEvent = (payload: OrchestrationV2Run) =>
({
id: `event-${payload.id}-${payload.status}`,
type: "run.updated",
threadId,
occurredAt: now,
payload,
}) as OrchestrationV2DomainEvent;

it("shows a queued run that starts after a newer steer raised the watermark", () => {
// Run 2 was queued first, run 3 was steered in and ran first. Its items
// sit in run 3's band and already raised the watermark past run 2's band.
const steerItems = [
runItem("item-steer-user", steerRunId, 3_000_001),
runItem("item-steer-reply", steerRunId, 3_000_002),
];
let projection: OrchestrationV2ThreadProjection | null = windowOf(
[run, queuedRun, steerRun],
steerItems,
);
let latestLocalTurnOrdinal = 3_000_002;
// Same order the orchestrator commits a queued start: the user message
// lands while the run is still queued, then the run starts and replies.
const queuedUser = runItem("item-queued-user", queuedRunId, 2_000_001);
const queuedReply = runItem("item-queued-reply", queuedRunId, 2_000_002);
for (const event of [
itemEvent(queuedUser),
runEvent({ ...queuedRun, status: "running" }),
itemEvent(queuedReply),
runEvent({ ...queuedRun, status: "completed" }),
]) {
projection = applyOrchestrationV2ProjectionEvent(projection, event, {
partialTimeline: true,
latestLocalTurnOrdinal,
});
if (event.type === "turn-item.updated") {
latestLocalTurnOrdinal = Math.max(latestLocalTurnOrdinal, event.payload.ordinal);
}
}

expect(projection?.visibleTurnItems.map((row) => row.item.id)).toEqual([
queuedUser.id,
queuedReply.id,
...steerItems.map((item) => item.id),
]);
expect(projection?.turnItems.map((item) => item.id)).toContain(queuedReply.id);
});

it("keeps a running run's item when the window starts after its earlier items", () => {
// The bounded window begins inside run 3; run 2 started and kept going,
// but none of its rows made it into the window.
const runningRun = { ...queuedRun, status: "running" } as const;
const steerItems = [
runItem("item-steer-user", steerRunId, 3_000_001),
runItem("item-steer-reply", steerRunId, 3_000_002),
];
const projection = windowOf([run, runningRun, steerRun], steerItems);
const reply = runItem("item-running-reply", queuedRunId, 2_000_003);

const next = applyOrchestrationV2ProjectionEvent(projection, itemEvent(reply), {
partialTimeline: true,
latestLocalTurnOrdinal: 3_000_002,
});

expect(next?.visibleTurnItems.map((row) => row.item.id)).toEqual([
reply.id,
...steerItems.map((item) => item.id),
]);
});

it.each([
{ name: "keeps an item of a run already in the window", itemRunId: steerRunId, kept: true },
{ name: "drops an item of a finished run outside the window", itemRunId: runId, kept: false },
])("$name that arrives below the watermark", ({ itemRunId, kept }) => {
const first = runItem("item-first", steerRunId, 3_000_001);
const later = runItem("item-later", steerRunId, 3_000_003);
const late = runItem("item-late", itemRunId, 3_000_002);
const projection = windowOf([run, steerRun], [first, later]);

const next = applyOrchestrationV2ProjectionEvent(projection, itemEvent(late), {
partialTimeline: true,
latestLocalTurnOrdinal: later.ordinal,
});
if (kept) {
expect(next?.visibleTurnItems.map((row) => row.item.id)).toEqual([
first.id,
late.id,
later.id,
]);
} else {
expect(next).toBe(projection);
}
});
});

it("removes only hidden local items while preserving inherited rows", () => {
const inherited = commandItem("item-inherited");
const local = commandItem("item-local");
Expand Down
40 changes: 40 additions & 0 deletions packages/client-runtime/src/state/orchestrationV2Projection.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import type {
OrchestrationV2DomainEvent,
OrchestrationV2Run,
OrchestrationV2ThreadProjection,
OrchestrationV2TurnItem,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -86,6 +87,14 @@ function oldestLocalTurnOrdinal(
return oldest;
}

const UNFINISHED_RUN_STATUSES: ReadonlySet<OrchestrationV2Run["status"]> = new Set([
"preparing",
"queued",
"starting",
"running",
"waiting",
]);

function shouldDropMissingPartialTurnItem(
projection: OrchestrationV2ThreadProjection,
item: OrchestrationV2TurnItem,
Expand All @@ -94,6 +103,37 @@ function shouldDropMissingPartialTurnItem(
if (projection.visibleTurnItems.some((row) => row.sourceItemId === item.id)) {
return false;
}
if (item.runId === null) {
return isBelowPartialWindow(projection, item, latestLocalTurnOrdinal);
}
// A queued run keeps the ordinal it got when queued, so a steer sent later
// can raise the watermark past it. Items of a run that has not finished are
// new, not unloaded history: the first lands while the run is still queued,
// and a bounded window can start after the ones before it.
if (
projection.runs.some((run) => run.id === item.runId && UNFINISHED_RUN_STATUSES.has(run.status))
) {
return false;
}
// Later items of a run already in the window can raise the watermark first.
if (
projection.visibleTurnItems.some(
(row) =>
row.visibility === "local" &&
row.item.runId === item.runId &&
row.item.ordinal <= item.ordinal,
)
) {
return false;
}
return isBelowPartialWindow(projection, item, latestLocalTurnOrdinal);
}

function isBelowPartialWindow(
projection: OrchestrationV2ThreadProjection,
item: OrchestrationV2TurnItem,
latestLocalTurnOrdinal: number | null | undefined,
): boolean {
if (
latestLocalTurnOrdinal !== null &&
latestLocalTurnOrdinal !== undefined &&
Expand Down
Loading