From 5ff6bb87cfe8e19787c39fe57219505caaa393f7 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Tue, 18 Aug 2026 17:04:02 -0400 Subject: [PATCH] fix(core): coalesce queued compactions (#43292) --- packages/core/src/session/inbox.ts | 28 ++++++++++++++++------ packages/core/test/session-compact.test.ts | 21 +++++++++++++--- 2 files changed, 39 insertions(+), 10 deletions(-) diff --git a/packages/core/src/session/inbox.ts b/packages/core/src/session/inbox.ts index 1233ebf6858..2ba69fbe889 100644 --- a/packages/core/src/session/inbox.ts +++ b/packages/core/src/session/inbox.ts @@ -180,13 +180,27 @@ export const admitCompaction = Effect.fn("SessionInbox.admitCompaction")(functio bus: Bus.Interface, input: { readonly id: SessionMessage.ID; readonly sessionID: SessionSchema.ID; readonly delivery: Delivery }, ) { - const admitted = yield* admit(db, bus, { - id: input.id, - sessionID: input.sessionID, - item: Item.make({ type: "compaction", payload: {}, delivery: input.delivery }), - }) - if (admitted.type === "compaction") return admitted - return yield* Effect.die(new LifecycleConflict({ id: input.id })) + return yield* serialized( + input.sessionID, + Effect.gen(function* () { + const exact = yield* find(db, input.id) + if (exact) { + if (exact.type === "compaction" && exact.sessionID === input.sessionID) return exact + return yield* Effect.die(new LifecycleConflict({ id: input.id })) + } + if (yield* promotedFromMessage(db, input.sessionID, input.id, input.delivery)) + return yield* Effect.die(new LifecycleConflict({ id: input.id })) + const pending = (yield* list(db, input.sessionID)).find((item) => item.type === "compaction") + if (pending) return pending + const admitted = yield* admit(db, bus, { + id: input.id, + sessionID: input.sessionID, + item: Item.make({ type: "compaction", payload: {}, delivery: input.delivery }), + }) + if (admitted.type === "compaction") return admitted + return yield* Effect.die(new LifecycleConflict({ id: input.id })) + }), + ) }) export const projectAdmitted = Effect.fn("SessionInbox.projectAdmitted")(function* ( diff --git a/packages/core/test/session-compact.test.ts b/packages/core/test/session-compact.test.ts index 7f4d7c11466..d23a3f02228 100644 --- a/packages/core/test/session-compact.test.ts +++ b/packages/core/test/session-compact.test.ts @@ -73,7 +73,7 @@ const it = testEffect( ) describe("Session.compact", () => { - it.effect("durably stacks manual compaction", () => + it.effect("durably coalesces manual compaction", () => Effect.gen(function* () { requests = [] const session = yield* Session.Service @@ -102,11 +102,10 @@ describe("Session.compact", () => { const first = yield* session.compact({ sessionID: created.id }) const second = yield* session.compact({ sessionID: created.id }) - expect(second.id).not.toBe(first.id) + expect(second.id).toBe(first.id) expect(requests).toHaveLength(0) expect(yield* session.inbox(created.id)).toEqual([ expect.objectContaining({ id: first.id, type: "compaction", delivery: "queue" }), - expect.objectContaining({ id: second.id, type: "compaction", delivery: "queue" }), ]) expect((yield* session.context(created.id)).find((message) => message.id === first.id)).toBeUndefined() @@ -115,4 +114,20 @@ describe("Session.compact", () => { expect(steer).toMatchObject({ type: "compaction", delivery: "steer" }) }), ) + + it.effect("coalesces concurrent manual compaction", () => + Effect.gen(function* () { + const session = yield* Session.Service + const created = yield* session.create({ location }) + const admitted = yield* Effect.all( + [SessionMessage.ID.create(), SessionMessage.ID.create()].map((id) => + session.compact({ id, sessionID: created.id }), + ), + { concurrency: "unbounded" }, + ) + + expect(admitted[1]?.id).toBe(admitted[0]?.id) + expect(yield* session.inbox(created.id)).toHaveLength(1) + }), + ) })