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
118 changes: 118 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { assert, it, vi } from "@effect/vitest";
import {
EnvironmentId,
EventId,
CommandId,
CheckpointId,
Expand All @@ -25,6 +26,7 @@ import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Schema from "effect/Schema";
import * as SqlClient from "effect/unstable/sql/SqlClient";
import { projectThreadAwarenessV2 } from "@t3tools/shared/agentAwareness";

import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts";
import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts";
Expand Down Expand Up @@ -4179,4 +4181,120 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => {
);
}),
);

it.effect("dates a completion released by held background work at the release", () =>
Effect.gen(function* () {
const store = yield* ProjectionStoreV2;
const threadId = ThreadId.make("thread:held-completion");
const runId = RunId.make("run:held-completion");
const providerThreadId = ProviderThreadId.make("provider-thread:held-completion");
const completedAt = DateTime.makeUnsafe("2026-09-28T12:00:00.000Z");
const releasedAt = DateTime.makeUnsafe("2026-09-28T12:05:00.000Z");
const providerThread = {
id: providerThreadId,
driver,
providerInstanceId,
providerSessionId: null,
appThreadId: threadId,
ownerNodeId: null,
nativeThreadRef: null,
nativeConversationHeadRef: null,
status: "idle" as const,
firstRunOrdinal: 1,
lastRunOrdinal: 1,
handoffIds: [],
forkedFrom: null,
pendingBackgroundTasks: [{ taskId: "watch", kind: "monitor" as const }],
createdAt: completedAt,
updatedAt: completedAt,
};
yield* store.apply({
id: EventId.make("event:held-completion:thread"),
type: "thread.created",
threadId,
occurredAt: completedAt,
payload: {
createdBy: "user",
creationSource: "web",
id: threadId,
projectId: ProjectId.make("project:held-completion"),
title: "Watch the build",
providerInstanceId,
modelSelection,
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
activeProviderThreadId: null,
lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: threadId },
forkedFrom: null,
createdAt: completedAt,
updatedAt: completedAt,
archivedAt: null,
settledOverride: null,
settledAt: null,
lastVisitedAt: null,
deletedAt: null,
},
});
yield* store.apply({
id: EventId.make("event:held-completion:run"),
type: "run.created",
threadId,
runId,
driver,
providerInstanceId,
occurredAt: completedAt,
payload: {
id: runId,
threadId,
ordinal: 1,
providerInstanceId,
modelSelection,
providerThreadId,
userMessageId: MessageId.make("message:held-completion"),
rootNodeId: null,
activeAttemptId: null,
status: "completed",
requestedAt: completedAt,
startedAt: completedAt,
completedAt,
checkpointId: null,
contextHandoffId: null,
},
});
yield* store.apply({
id: EventId.make("event:held-completion:monitor-running"),
type: "provider-thread.updated",
threadId,
driver,
providerInstanceId,
occurredAt: completedAt,
payload: providerThread,
});
const project = { title: "Project" };
const environmentId = EnvironmentId.make("environment:held-completion");
const held = yield* store.getThreadShell(threadId);
assert.equal(
held && projectThreadAwarenessV2({ environmentId, project, thread: held })?.phase,
"running",
);

// The monitor ends five minutes after the run: the push must not look stale.
yield* store.apply({
id: EventId.make("event:held-completion:monitor-ended"),
type: "provider-thread.updated",
threadId,
driver,
providerInstanceId,
occurredAt: releasedAt,
payload: { ...providerThread, pendingBackgroundTasks: [], updatedAt: releasedAt },
});
const released = yield* store.getThreadShell(threadId);
const state =
released && projectThreadAwarenessV2({ environmentId, project, thread: released });
assert.equal(state?.phase, "completed");
assert.equal(state?.updatedAt, DateTime.formatIso(releasedAt));
}),
);
});
184 changes: 182 additions & 2 deletions apps/server/src/relay/AgentAwarenessRelay.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,17 +2,25 @@ import { assert, describe, it } from "@effect/vitest";
import * as NodeCrypto from "@effect/platform-node/NodeCrypto";
import {
EnvironmentId,
EventId,
MessageId,
NodeId,
type OrchestrationV2DomainEvent,
type OrchestrationV2ThreadShell,
type OrchestrationV2TurnItem,
ProjectId,
ProviderInstanceId,
RunId,
ThreadId,
TurnItemId,
} from "@t3tools/contracts";
import { RelayAgentActivityState } from "@t3tools/contracts/relay";
import * as DateTime from "effect/DateTime";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Option from "effect/Option";
import * as Queue from "effect/Queue";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";
Expand All @@ -27,6 +35,10 @@ import {
RELAY_URL_SECRET,
} from "../cloud/config.ts";
import { ServerEnvironment } from "../environment/ServerEnvironment.ts";
import {
layerMemory as projectionStoreMemoryLayer,
ProjectionStoreV2,
} from "../orchestration-v2/ProjectionStore.ts";
import { ThreadManagementService } from "../orchestration-v2/ThreadManagementService.ts";
import { ProjectService } from "../project/ProjectService.ts";
import {
Expand Down Expand Up @@ -137,6 +149,9 @@ const makeTestRelay = Effect.fnUntraced(function* (
readonly failSecretRead?: (name: string) => boolean;
/** Starts unlinked with publishing off when false. */
readonly linked?: boolean;
/** Serves shells from this source instead of `currentShell`. */
readonly readShell?: (threadId: ThreadId) => Effect.Effect<OrchestrationV2ThreadShell | null>;
readonly domainEvents?: Stream.Stream<OrchestrationV2DomainEvent>;
} = {},
) {
const values = new Map<string, Uint8Array>(
Expand Down Expand Up @@ -172,7 +187,9 @@ const makeTestRelay = Effect.fnUntraced(function* (
const catchUp = { shellSnapshotReads: 0 };
const threads = ThreadManagementService.of({
getThreadShell: (threadId) =>
Effect.sync(() => shellReads.push(threadId)).pipe(Effect.andThen(Ref.get(currentShell))),
Effect.sync(() => shellReads.push(threadId)).pipe(
Effect.andThen(options.readShell?.(threadId) ?? Ref.get(currentShell)),
),
getShellSnapshot: () =>
Effect.sync(() => {
catchUp.shellSnapshotReads += 1;
Expand All @@ -196,7 +213,7 @@ const makeTestRelay = Effect.fnUntraced(function* (
getThreadEventSequence: unused,
streamStoredEvents: Stream.empty,
streamStoredEventsFrom: () => Stream.empty,
streamDomainEvents: Stream.empty,
streamDomainEvents: options.domainEvents ?? Stream.empty,
});
const publications: Array<{
readonly url: string;
Expand Down Expand Up @@ -277,6 +294,9 @@ describe("AgentAwarenessRelay", () => {
"run.created",
"run.updated",
"runtime-request.updated",
// Pending background work changes can release a held completion.
"subagent.updated",
"provider-thread.updated",
"thread.metadata-updated",
"thread.model-selection-updated",
"thread.provider-switched",
Expand Down Expand Up @@ -665,6 +685,166 @@ describe("AgentAwarenessRelay", () => {
assert.equal(publications.length, 1);
}),
);
it.effect("publishes a held completion when its background item ends, not on tool output", () =>
Effect.gen(function* () {
const store = yield* ProjectionStoreV2;
const events = yield* Queue.unbounded<OrchestrationV2DomainEvent>();
// The relay pulls the next event only after it has handled the previous one.
const pulls = yield* Queue.unbounded<void>();
const { relay, shellReads, publications } = yield* makeTestRelay({
readShell: (threadId) => store.getThreadShell(threadId).pipe(Effect.orDie),
domainEvents: Stream.fromEffectRepeat(
Queue.offer(pulls, undefined).pipe(Effect.andThen(Queue.take(events))),
),
});
yield* relay.start();
yield* Queue.take(pulls);
const deliver = Effect.fnUntraced(function* (event: OrchestrationV2DomainEvent) {
yield* store.apply(event);
yield* Queue.offer(events, event);
yield* Queue.take(pulls);
yield* relay.drain;
});

const now = yield* DateTime.now;
const runId = RunId.make("run:held-item");
const run = {
id: runId,
threadId: THREAD_ID,
ordinal: 1,
providerInstanceId: ProviderInstanceId.make("codex"),
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "test-model" },
providerThreadId: null,
userMessageId: MessageId.make("message:held-item"),
rootNodeId: null,
activeAttemptId: null,
status: "running" as const,
requestedAt: now,
startedAt: now,
completedAt: null,
checkpointId: null,
contextHandoffId: null,
};
const item = (
id: string,
fields: Pick<OrchestrationV2TurnItem, "status"> &
(
| { readonly type: "dynamic_tool"; readonly toolName: string; readonly input: unknown }
| {
readonly type: "assistant_message";
readonly messageId: MessageId;
readonly text: string;
readonly streaming: boolean;
}
),
): OrchestrationV2TurnItem => ({
id: TurnItemId.make(id),
threadId: THREAD_ID,
runId,
nodeId: NodeId.make("node:held-item"),
providerThreadId: null,
providerTurnId: null,
nativeItemRef: null,
parentItemId: null,
ordinal: 1,
title: null,
startedAt: now,
completedAt: null,
updatedAt: now,
...fields,
});
const itemEvent = (id: string, payload: OrchestrationV2TurnItem) =>
({
id: EventId.make(id),
type: "turn-item.updated",
threadId: THREAD_ID,
runId,
occurredAt: now,
payload,
}) satisfies OrchestrationV2DomainEvent;
const background = { type: "dynamic_tool" as const, toolName: "watch", input: {} };

yield* deliver({
id: EventId.make("event:held-item:thread"),
type: "thread.created",
threadId: THREAD_ID,
occurredAt: now,
payload: {
createdBy: "user",
creationSource: "web",
id: THREAD_ID,
projectId: PROJECT_ID,
title: "Thread",
providerInstanceId: ProviderInstanceId.make("codex"),
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "test-model" },
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
activeProviderThreadId: null,
lineage: { rootThreadId: THREAD_ID, parentThreadId: null, relationshipToParent: null },
forkedFrom: null,
createdAt: now,
updatedAt: now,
archivedAt: null,
settledOverride: null,
settledAt: null,
lastVisitedAt: null,
deletedAt: null,
},
});
yield* deliver({
id: EventId.make("event:held-item:run"),
type: "run.created",
threadId: THREAD_ID,
runId,
occurredAt: now,
payload: run,
});
yield* deliver(
itemEvent("event:held-item:tool", item("item:tool", { ...background, status: "running" })),
);
yield* deliver({
id: EventId.make("event:held-item:run-completed"),
type: "run.updated",
threadId: THREAD_ID,
runId,
occurredAt: now,
payload: { ...run, status: "completed", completedAt: now },
});
assert.deepEqual(
publications.map((publication) => publication.state?.phase),
["running"],
);

// Streaming output never reaches the shell read.
const readsBeforeOutput = shellReads.length;
yield* deliver(
itemEvent(
"event:held-item:reply",
item("item:reply", {
type: "assistant_message",
messageId: MessageId.make("message:held-item:reply"),
text: "Still watching",
streaming: true,
status: "running",
}),
),
);
yield* deliver(
itemEvent("event:held-item:tick", item("item:tool", { ...background, status: "running" })),
);
assert.equal(shellReads.length, readsBeforeOutput);

yield* deliver(
itemEvent("event:held-item:end", item("item:tool", { ...background, status: "completed" })),
);
assert.deepEqual(
publications.map((publication) => publication.state?.phase),
["running", "completed"],
);
}).pipe(Effect.scoped, Effect.provide(projectionStoreMemoryLayer)),
);
});

describe.sequential("startup catch-up", () => {
Expand Down
Loading
Loading