From 9e92cbe9be751cb2fbef3c654d30ef378aed2879 Mon Sep 17 00:00:00 2001 From: Aiden Date: Fri, 21 Aug 2026 04:28:17 +0000 Subject: [PATCH] fix(core): reuse durable event codecs Co-authored-by: Hona <10430890+Hona@users.noreply.github.com> --- packages/core/src/bus.ts | 33 +++++++++++++++++++++++---------- 1 file changed, 23 insertions(+), 10 deletions(-) diff --git a/packages/core/src/bus.ts b/packages/core/src/bus.ts index 837b5a4df50..60974322848 100644 --- a/packages/core/src/bus.ts +++ b/packages/core/src/bus.ts @@ -67,6 +67,25 @@ const envelope = (aggregateID: string, seq: number, version: number) => ({ version: Event.Version.make(version), }) +const encoders = new WeakMap unknown>() +const decoders = new WeakMap unknown>() + +function encodeData(definition: Event.Definition, data: unknown) { + const cached = encoders.get(definition) + if (cached) return cached(data) + const encode = Schema.encodeUnknownSync(definition.data) + encoders.set(definition, encode) + return encode(data) +} + +function decodeData(definition: Event.Definition, data: unknown) { + const cached = decoders.get(definition) + if (cached) return cached(data) + const decode = Schema.decodeUnknownSync(definition.data) + decoders.set(definition, decode) + return decode(data) +} + const decodeSerializedEvent = (event: SerializedEvent): Event.Payload => { const definition = Durable.get(event.type) if (!definition?.durable) { @@ -77,7 +96,7 @@ const decodeSerializedEvent = (event: SerializedEvent): Event.Payload => { created: event.created ?? 0, type: definition.type, durable: envelope(event.aggregateID, event.seq, definition.durable.version), - data: Schema.decodeUnknownSync(definition.data)(event.data), + data: decodeData(definition, event.data), } } @@ -260,10 +279,7 @@ export function configured(options?: Options) { .get() .pipe(Effect.orDie) const latest = row?.seq ?? -1 - const encoded = Schema.encodeUnknownSync(definition.data)(event.data) as Record< - string, - unknown - > + const encoded = encodeData(definition, event.data) as Record if (input?.strictOwner && row?.ownerID && row.ownerID !== input.ownerID) { yield* Effect.die( new InvalidDurableEventError({ @@ -529,10 +545,7 @@ export function configured(options?: Options) { const ids = new Set() for (const [index, item] of payloads.entries()) { const seq = firstSeq + index - const encoded = Schema.encodeUnknownSync(item.definition.data)(item.event.data) as Record< - string, - unknown - > + const encoded = encodeData(item.definition, item.event.data) as Record if (persist) { if (ids.has(item.event.id)) yield* Effect.die( @@ -621,7 +634,7 @@ export function configured(options?: Options) { id: event.id, created: event.created ?? 0, type: definition.type, - data: Schema.decodeUnknownSync(definition.data)(event.data), + data: decodeData(definition, event.data), } as Event.Payload const committed = yield* commitDurableEvent(definition, payload, { seq: event.seq,