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
Expand Up @@ -129,6 +129,21 @@ Verify:
- the legacy relay received no `/slack/events` traffic;
- the delivery log came from the intended Render instance.

### Headless native reliability verification

For request persistence, native delivery retries, duplicate acknowledgments,
and central session state, use the API-key-protected canary routes documented in
`hosted/compadre/docs/runbooks/api-reliability-verification.md`. Start through
central commands, then inspect the persisted central projection, controller run,
delivery cursor/block, and Temporal attempts together. A completed worker alone
does not prove the UI's read model received its output. GET inspection must not
wake Modal. Exercise faults only on server-registered verification threads;
never crash the shared database to prove recovery.

Verify image writes have private S3 objects and reference-only request metadata.
Record deployed instances, canary IDs and terminal evidence, then stop canary
work. Report browser rendering as untested when verification is API-only.

### Compatibility API

- authenticate with the real compatibility credential without printing it;
Expand Down
6 changes: 6 additions & 0 deletions .agents/skills/debug-compadre-workflows/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,12 @@ authenticated tool request fails.

Prefer exact identifiers and narrow time windows. Render service instance suffixes such as `web-8cv7x` are ephemeral; use the stable host/service identity and discovered IDs instead of copying an old suffix.

For headless live verification, use the API-key-protected, canary-only routes in
`hosted/compadre/docs/runbooks/api-reliability-verification.md`. They exercise
canonical commands, expose central projections and delivery cursors, and offer
expiring one-shot faults without crashing shared infrastructure. Verify that the
routes are deployed before using them; GET inspection never wakes a worker.

## Interpret the evidence

- For Modal, distinguish controller failure, sandbox lifecycle failure,
Expand Down
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 @@ -353,6 +353,67 @@ describe("native thread replication", () => {
const sequence = after.snapshotSequence;
await central.run(central.engine.dispatch({ ...command, epoch: 2 }));
expect((await central.readModel()).snapshotSequence).toBe(sequence);
await central.run(
central.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make(NodeCrypto.randomUUID()),
threadId,
session: {
threadId,
status: "starting",
providerName: "codex",
runtimeMode: "full-access",
activeTurnId: null,
lastError: null,
updatedAt: createdAt,
},
createdAt,
}),
);
await central.run(
central.engine.dispatch({
...command,
epoch: 2,
commandId: CommandId.make(NodeCrypto.randomUUID()),
status: "error",
reason: "Event delivery blocked; pending output retained",
}),
);
const blocked = (await central.readModel()).threads.find(
(thread) => thread.id === threadId,
)?.session;
expect(blocked?.status).toBe("error");
expect(blocked?.activeTurnId).toBeNull();
expect(blocked?.lastError).toBe("Event delivery blocked; pending output retained");
await central.run(
central.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make(NodeCrypto.randomUUID()),
threadId,
session: {
threadId,
status: "ready",
providerName: "codex",
runtimeMode: "full-access",
activeTurnId: null,
lastError: null,
updatedAt: createdAt,
},
createdAt,
}),
);
await central.run(
central.engine.dispatch({
...command,
epoch: 2,
commandId: CommandId.make(NodeCrypto.randomUUID()),
status: "error",
}),
);
expect(
(await central.readModel()).threads.find((thread) => thread.id === threadId)?.session
?.status,
).toBe("ready");
} finally {
await central.dispose();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ import { it as effectIt } from "@effect/vitest";
import { afterEach, describe, expect, it } from "vite-plus/test";

import { OrchestrationEventStoreLive } from "../../persistence/Layers/OrchestrationEventStore.ts";
import { PersistenceSqlError } from "../../persistence/Errors.ts";
import { OrchestrationCommandReceiptRepositoryLive } from "../../persistence/Layers/OrchestrationCommandReceipts.ts";
import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts";
import {
Expand Down Expand Up @@ -287,6 +288,7 @@ describe("ProviderRuntimeIngestion", () => {
async function createHarness(options?: {
serverSettings?: Partial<ServerSettings>;
threadTitle?: string;
failRuntimeErrorWrite?: "before" | "after";
}) {
const workspaceRoot = makeTempDir("t3-provider-project-");
NodeFS.mkdirSync(NodePath.join(workspaceRoot, ".git"));
Expand All @@ -303,8 +305,31 @@ describe("ProviderRuntimeIngestion", () => {
Layer.provide(RepositoryIdentityResolver.layer),
Layer.provide(SqlitePersistenceMemory),
);
let failedRuntimeErrorWrite = false;
const recoveringEngine = Layer.effect(
OrchestrationEngineService,
Effect.gen(function* () {
const engine = yield* OrchestrationEngineService;
return {
...engine,
dispatch: (command: OrchestrationCommand) =>
Effect.gen(function* () {
if (
options?.failRuntimeErrorWrite &&
!failedRuntimeErrorWrite &&
command.commandId.includes("runtime-error-session-set")
) {
failedRuntimeErrorWrite = true;
if (options.failRuntimeErrorWrite === "after") yield* engine.dispatch(command);
return yield* new PersistenceSqlError({ operation: "test.outage" });
}
return yield* engine.dispatch(command);
}),
};
}),
).pipe(Layer.provide(orchestrationLayer));
const layer = ProviderRuntimeIngestionLive.pipe(
Layer.provideMerge(orchestrationLayer),
Layer.provideMerge(recoveringEngine),
Layer.provideMerge(projectionSnapshotLayer),
// Single shared liveness instance across ingestion (writer), the
// engine, and the snapshot query (reader).
Expand Down Expand Up @@ -2789,6 +2814,28 @@ describe("ProviderRuntimeIngestion", () => {
expect(resolvedPayload?.requestType).toBe("command_execution_approval");
});

it.each(["before", "after"] as const)(
"recovers runtime errors after database failure %s commit",
async (failure) => {
const harness = await createHarness({ failRuntimeErrorWrite: failure });
harness.emit({
type: "runtime.error",
eventId: asEventId("evt-database-recovery"),
provider: ProviderDriverKind.make("codex"),
threadId: asThreadId("thread-1"),
createdAt: "2026-01-01T00:00:01.000Z",
payload: { message: "Start failed" },
});
await harness.drain();
const thread = (await harness.readModel()).threads.find((entry) => entry.id === "thread-1");
expect(thread?.session?.status).toBe("error");
expect(thread?.session?.lastError).toBe("Start failed");
expect(
thread?.activities.filter((activity) => activity.id === "evt-database-recovery"),
).toHaveLength(1);
},
);

it("maps runtime.error into errored session state", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
21 changes: 18 additions & 3 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ 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 Schedule from "effect/Schedule";
import * as Stream from "effect/Stream";
import { makeFairDrainableWorker } from "@t3tools/shared/DrainableWorker";
import { formatTokens } from "@t3tools/shared/usageFormat";
Expand Down Expand Up @@ -1015,9 +1016,11 @@ const make = Effect.gen(function* () {
const projectionTurnRepository = yield* ProjectionTurnRepository;
const serverSettingsService = yield* ServerSettingsService;
const providerCommandId = (event: ProviderRuntimeEvent, tag: string) =>
crypto.randomUUIDv4.pipe(
Effect.map((uuid) => CommandId.make(`provider:${event.eventId}:${tag}:${uuid}`)),
);
event.type === "runtime.error"
? Effect.succeed(CommandId.make(`provider:${event.eventId}:${tag}`))
: crypto.randomUUIDv4.pipe(
Effect.map((uuid) => CommandId.make(`provider:${event.eventId}:${tag}:${uuid}`)),
);

const turnMessageIdsByTurnKey = yield* Cache.make<string, Set<MessageId>>({
capacity: TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY,
Expand Down Expand Up @@ -2245,6 +2248,16 @@ const make = Effect.gen(function* () {

const processInputSafely = (input: RuntimeIngestionInput) =>
processInput(input).pipe(
// Error events have no worker journal to replay when start fails. Keep
// this ordered item through a brief database recovery; stable command IDs
// make a retry after a committed-but-unacknowledged write safe.
Effect.retry({
schedule: Schedule.max([Schedule.exponential("1 second"), Schedule.recurs(5)]),
while: (error) =>
input.source === "runtime" &&
input.event.type === "runtime.error" &&
error._tag === "PersistenceSqlError",
}),
Effect.catchCause((cause) => {
if (Cause.hasInterruptsOnly(cause)) {
return Effect.failCause(cause);
Expand All @@ -2253,6 +2266,8 @@ const make = Effect.gen(function* () {
source: input.source,
eventId: input.event.eventId,
eventType: input.event.type,
threadId:
input.source === "runtime" ? input.event.threadId : input.event.payload.threadId,
cause: Cause.pretty(cause),
});
}),
Expand Down
33 changes: 21 additions & 12 deletions apps/server/src/orchestration/decider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1199,6 +1199,14 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand"

case "thread.native-stream.close": {
const thread = yield* requireThread({ readModel, command, threadId: command.threadId });
// A lost acknowledgement can exhaust delivery retries after central has
// already committed completion. Do not replace that terminal state: a
// duplicate replay is receipt-deduplicated and cannot restore it later.
const alreadySettled =
command.status === "error" &&
thread.session !== null &&
thread.session.status !== "starting" &&
thread.session.status !== "running";
return {
...(yield* withEventBase({
aggregateKind: "thread",
Expand All @@ -1210,18 +1218,19 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand"
type: "thread.session-set",
payload: {
threadId: command.threadId,
session: {
...thread.session,
threadId: command.threadId,
status: "stopped",
activeTurnId: null,
providerName: thread.session?.providerName ?? null,
runtimeMode: thread.session?.runtimeMode ?? "full-access",
// Closing delivery is normal after idle worker expiry. The run
// driver reports interrupted work separately; retain that error.
lastError: thread.session?.lastError ?? null,
updatedAt: command.createdAt,
},
session: alreadySettled
? thread.session
: {
...thread.session,
threadId: command.threadId,
status: command.status ?? "stopped",
activeTurnId: null,
providerName: thread.session?.providerName ?? null,
runtimeMode: thread.session?.runtimeMode ?? "full-access",
lastError:
command.status === "error" ? command.reason : (thread.session?.lastError ?? null),
updatedAt: command.createdAt,
},
},
};
}
Expand Down
Loading
Loading