Compare commits

..

1 Commits

Author SHA1 Message Date
Kit Langton cc7ac73097 fix(core): scope v1 migration event deletion 2026-08-11 12:05:19 -04:00
6 changed files with 51 additions and 56 deletions
+15 -28
View File
@@ -297,7 +297,7 @@ export const classifyHttpFailure = (input: {
})
}
export const mapHttpError = (error: unknown, redactedNames: ReadonlyArray<string | RegExp>) => {
const toHttpError = (redactedNames: ReadonlyArray<string | RegExp>) => (error: unknown) => {
const transportError = (input: {
readonly message: string
readonly kind?: string | undefined
@@ -314,35 +314,23 @@ export const mapHttpError = (error: unknown, redactedNames: ReadonlyArray<string
}),
})
const cause =
HttpClientError.isHttpClientError(error) && "cause" in error.reason
? error.reason.cause
: error instanceof Error
? error.cause
: undefined
const code = [cause, error]
.map((value) => (typeof value === "object" && value !== null ? Reflect.get(value, "code") : undefined))
.find((value): value is string => typeof value === "string")
const request = HttpClientError.isHttpClientError(error) && "request" in error ? error.request : undefined
const raw = cause instanceof Error ? cause.message : error instanceof Error ? error.message : undefined
const detail = raw && request ? redactBody(raw, secretValues(request)) : raw
const message = code && detail && !detail.includes(code) ? `${code}: ${detail}` : detail
if (Cause.isTimeoutError(error) || Cause.isTimeoutError(cause))
return transportError({ message: message ?? "HTTP transport timed out", kind: code ?? "Timeout", request })
if (!HttpClientError.isHttpClientError(error)) {
return transportError({ message: message ?? "HTTP transport failed", kind: code, request })
if (Cause.isTimeoutError(error)) {
return transportError({ message: error.message, kind: "Timeout" })
}
if (!HttpClientError.isHttpClientError(error)) {
return transportError({ message: error instanceof Error ? error.message : "HTTP transport failed" })
}
const request = "request" in error ? error.request : undefined
if (error.reason._tag === "TransportError") {
return transportError({
message: message ?? error.reason.description ?? "HTTP transport failed",
kind: code ?? error.reason._tag,
message: error.reason.description ?? "HTTP transport failed",
kind: error.reason._tag,
request,
})
}
return transportError({
message: message ?? `HTTP transport failed: ${error.reason._tag}`,
kind: code ?? error.reason._tag,
message: `HTTP transport failed: ${error.reason._tag}`,
kind: error.reason._tag,
request,
})
}
@@ -355,16 +343,15 @@ export const layer: Layer.Layer<Service, never, HttpClient.HttpClient> = Layer.e
Effect.gen(function* () {
const redactedNames = yield* Headers.CurrentRedactedNames
if (!middleware)
return yield* http.execute(request).pipe(
Effect.mapError((error) => mapHttpError(error, redactedNames)),
Effect.flatMap(statusError(request, redactedNames)),
)
return yield* http
.execute(request)
.pipe(Effect.mapError(toHttpError(redactedNames)), Effect.flatMap(statusError(request, redactedNames)))
const response = yield* middleware(request, (input) =>
http
.execute(input)
.pipe(Effect.mapError((cause) => (cause instanceof Error ? cause : new Error(String(cause))))),
).pipe(Effect.mapError((error) => mapHttpError(error, redactedNames)))
).pipe(Effect.mapError(toHttpError(redactedNames)))
return yield* statusError(response.request, redactedNames)(response)
})
return Service.of({
+9 -6
View File
@@ -2,7 +2,6 @@ import { Effect, Stream } from "effect"
import { Headers, HttpClientRequest } from "effect/unstable/http"
import { Auth } from "../auth"
import { render as renderEndpoint } from "../endpoint"
import { mapHttpError } from "../executor"
import { Framing } from "../framing"
import type { HttpMiddleware, Transport, TransportPrepareInput } from "./index"
import * as ProviderShared from "../../protocols/shared"
@@ -87,16 +86,20 @@ export const httpJson = <Body, Frame>(input: HttpJsonInput<Body, Frame>): HttpJs
middleware: prepareInput.middleware,
}
}),
frames: (prepared, _request, runtime) =>
frames: (prepared, request, runtime) =>
Stream.unwrap(
runtime.http
.execute(prepared.request, prepared.middleware)
.pipe(
Effect.map((response) =>
Stream.unwrap(
Effect.map(Headers.CurrentRedactedNames, (redactedNames) =>
prepared.framing.frame(
response.stream.pipe(Stream.mapError((error) => mapHttpError(error, redactedNames))),
prepared.framing.frame(
response.stream.pipe(
Stream.mapError((error) =>
ProviderShared.eventError(
`${request.model.provider}/${request.model.route.id}`,
`Failed to read ${request.model.provider}/${request.model.route.id} stream`,
ProviderShared.errorText(error),
),
),
),
),
+2 -2
View File
@@ -63,14 +63,14 @@ export const dynamicResponse = (handler: Handler) => runtimeLayer(handlerLayer(h
* Layer that emits the supplied SSE chunks and then aborts mid-stream. Used to
* exercise transport errors that surface during parsing.
*/
export const truncatedStream = (chunks: ReadonlyArray<string>, error: Error = new Error("connection reset")) =>
export const truncatedStream = (chunks: ReadonlyArray<string>) =>
dynamicResponse((input) =>
Effect.sync(() => {
const encoder = new TextEncoder()
const stream = new ReadableStream({
start(controller) {
for (const chunk of chunks) controller.enqueue(encoder.encode(chunk))
controller.error(error)
controller.error(new Error("connection reset"))
},
})
return input.respond(stream, { headers: SSE_HEADERS })
+4 -10
View File
@@ -1221,18 +1221,12 @@ describe("OpenAI Chat route", () => {
it.effect("surfaces transport errors that occur mid-stream", () =>
Effect.gen(function* () {
const layer = truncatedStream(
[`data: ${JSON.stringify(deltaChunk({ role: "assistant", content: "Hello" }))}\n\n`],
Object.assign(new Error("socket closed unexpectedly"), { code: "ECONNRESET" }),
)
const layer = truncatedStream([
`data: ${JSON.stringify(deltaChunk({ role: "assistant", content: "Hello" }))}\n\n`,
])
const error = yield* LLMClient.generate(request).pipe(Effect.provide(layer), Effect.flip)
expect(error.reason).toMatchObject({
_tag: "Transport",
message: "ECONNRESET: socket closed unexpectedly",
kind: "ECONNRESET",
url: "https://api.openai.test/v1/chat/completions",
})
expect(error.message).toContain("Failed to read openai/openai-chat stream")
}),
)
+2 -2
View File
@@ -440,7 +440,7 @@ export function status(): Effect.Effect<Status, never, Database.Service> {
export const layer = Layer.effectDiscard(
Effect.gen(function* () {
runtimeState = { status: "running", progress: { label: "Clearing old events" } }
runtimeState = { status: "running", progress: { label: "Migrating sessions" } }
yield* run().pipe(
Effect.matchCauseEffect({
onFailure: (cause) =>
@@ -485,7 +485,6 @@ export function run(options: Options = {}): Effect.Effect<RunResult, never, Data
yield* db
.transaction((tx) =>
Effect.gen(function* () {
yield* tx.delete(EventTable).run()
yield* tx
.insert(KVTable)
.values({ key: MIGRATION_STATE_KEY, value: { phase: "sessions" } })
@@ -567,6 +566,7 @@ export function run(options: Options = {}): Effect.Effect<RunResult, never, Data
yield* Effect.forEach(transformed.warnings, (warning) =>
Effect.logWarning("Skipped V1 migration row", warning),
)
yield* tx.delete(EventTable).where(eq(EventTable.aggregate_id, next.id)).run()
yield* tx.delete(SessionMessageTable).where(eq(SessionMessageTable.session_id, next.id)).run()
yield* Effect.forEach(transformed.messages, (message) =>
tx.run(sql`
+19 -8
View File
@@ -1039,7 +1039,7 @@ describe("V1Migration database workflow", () => {
)
})
test("rolls back one session atomically and resumes from the committed cursor", async () => {
test("deletes events only in each successfully checkpointed session transaction", async () => {
await database(
Effect.gen(function* () {
const { db } = yield* Database.Service
@@ -1057,9 +1057,19 @@ describe("V1Migration database workflow", () => {
yield* db.run(
sql`INSERT INTO session_message (id, session_id, type, seq, time_created, time_updated, data) VALUES ('msg_stale_b', 'ses_b', 'user', 0, 7, 8, '{"text":"stale","time":{"created":7}}')`,
)
yield* db.run(sql`INSERT INTO event_sequence (aggregate_id, seq, owner_id) VALUES ('ses_b', 7, 'owner')`)
yield* db.run(sql`
INSERT INTO event_sequence (aggregate_id, seq, owner_id) VALUES
('ses_a', 7, 'owner'),
('ses_b', 7, 'owner'),
('ses_c', 7, 'owner'),
('ses_unrelated', 7, 'owner')
`)
yield* db.run(
sql`INSERT INTO event (id, aggregate_id, seq, created, type, data) VALUES ('event_stale_b', 'ses_b', 7, 1, 'session.renamed.1', '{}')`,
sql`INSERT INTO event (id, aggregate_id, seq, created, type, data) VALUES
('event_stale_a', 'ses_a', 7, 1, 'session.renamed.1', '{}'),
('event_stale_b', 'ses_b', 7, 1, 'session.renamed.1', '{}'),
('event_stale_c', 'ses_c', 7, 1, 'session.renamed.1', '{}'),
('event_unrelated', 'ses_unrelated', 7, 1, 'session.renamed.1', '{}')`,
)
yield* Layer.launch(V1Migration.layer).pipe(Effect.forkScoped)
const failed = yield* V1Migration.status().pipe(
@@ -1091,13 +1101,14 @@ describe("V1Migration database workflow", () => {
seq: 7,
owner_id: "owner",
})
expect(yield* db.all(sql`SELECT id FROM event WHERE aggregate_id = 'ses_b'`)).toEqual([])
expect(yield* db.all(sql`SELECT id FROM event ORDER BY id`)).toEqual([
{ id: "event_stale_a" },
{ id: "event_stale_b" },
{ id: "event_unrelated" },
])
expect(yield* db.get(sql`SELECT value FROM kv WHERE key = 'migration.v1-v2'`)).toEqual({
value: '{"phase":"sessions","cursor":"ses_c"}',
})
yield* db.run(
sql`INSERT INTO event (id, aggregate_id, seq, created, type, data) VALUES ('event_after_clear', 'ses_c', 0, 2, 'session.renamed.1', '{}')`,
)
yield* db.run(sql`DROP TRIGGER fail_b`)
yield* Layer.launch(V1Migration.layer).pipe(Effect.forkScoped)
yield* V1Migration.status().pipe(
@@ -1110,7 +1121,7 @@ describe("V1Migration database workflow", () => {
seq: -1,
owner_id: null,
})
expect(yield* db.all(sql`SELECT id FROM event`)).toEqual([{ id: "event_after_clear" }])
expect(yield* db.all(sql`SELECT id FROM event`)).toEqual([{ id: "event_unrelated" }])
}),
)
})