diff --git a/apps/server/src/persistence/ProviderSessionRuntime.ts b/apps/server/src/persistence/ProviderSessionRuntime.ts index a3475d2f190..2ccdd862522 100644 --- a/apps/server/src/persistence/ProviderSessionRuntime.ts +++ b/apps/server/src/persistence/ProviderSessionRuntime.ts @@ -1,7 +1,9 @@ +import * as Arr from "effect/Array"; import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; import * as Struct from "effect/Struct"; import * as SqlClient from "effect/unstable/sql/SqlClient"; @@ -280,18 +282,30 @@ export const make = Effect.gen(function* () { ), ), Effect.flatMap((rows) => + // Skip rows that no longer decode (e.g. written by an older build) + // instead of failing the whole list — one stale row must not disable + // every consumer that enumerates sessions, such as the reaper. Effect.forEach(rows, (row) => decodeRuntimeRow(row).pipe( - Effect.mapError((cause) => - PersistenceDecodeError.fromSchemaError( - "ProviderSessionRuntimeRepository.list:decodeRows", - cause, - { threadId: row.threadId }, - ), + Effect.map(Option.some), + Effect.catch((cause) => + Effect.logWarning("provider.session.runtime.row-skipped", { + threadId: row.threadId, + error: PersistenceDecodeError.fromSchemaError( + "ProviderSessionRuntimeRepository.list:decodeRows", + cause, + { threadId: row.threadId }, + ).message, + }).pipe(Effect.as(Option.none())), ), ), ), ), + Effect.map((decoded) => + Arr.filterMap(decoded, (row) => + Option.isSome(row) ? Result.succeed(row.value) : Result.failVoid, + ), + ), ); const deleteByThreadId: ProviderSessionRuntimeRepository["Service"]["deleteByThreadId"] = ( diff --git a/apps/server/src/persistence/RepositoryErrorCorrelation.test.ts b/apps/server/src/persistence/RepositoryErrorCorrelation.test.ts index f7425200fd1..379b06e2a22 100644 --- a/apps/server/src/persistence/RepositoryErrorCorrelation.test.ts +++ b/apps/server/src/persistence/RepositoryErrorCorrelation.test.ts @@ -182,7 +182,7 @@ describe("persistence error correlation", () => { }).pipe(Effect.provide(authPairingLinkLayer)), ); - it.effect("correlates provider runtime SQL and per-row decode failures by thread", () => + it.effect("skips undecodable provider runtime rows and correlates SQL failures by thread", () => Effect.gen(function* () { const runtimes = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; const sql = yield* SqlClient.SqlClient; @@ -215,16 +215,24 @@ describe("persistence error correlation", () => { ) `; - const decodeError = yield* Effect.flip(runtimes.list()); - assert.instanceOf(decodeError, PersistenceErrors.PersistenceDecodeError); - assert.deepStrictEqual(decodeError.correlation, { threadId }); - assert.equal( - decodeError.message, - `Decode error in ProviderSessionRuntimeRepository.list:decodeRows: ${decodeError.issue}`, + const validThreadId = ThreadId.make("thread-valid"); + yield* runtimes.upsert({ + threadId: validThreadId, + providerName: "codex", + providerInstanceId: null, + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt, + resumeCursor: null, + runtimePayload: null, + }); + + const listed = yield* runtimes.list(); + assert.deepStrictEqual( + listed.map((runtime) => runtime.threadId), + [validThreadId], ); - assert.notInclude(decodeError.issue, runtimePayload); - assert.notInclude(decodeError.message, runtimePayload); - assert.notInclude(decodeError.message, lastSeenAt); yield* sql`DROP TABLE provider_session_runtime`; const sqlFailure = yield* Effect.flip(