fix(opencode): restore global sync event compatibility
This commit is contained in:
parent
1ed76ccace
commit
5e3026b5ce
9 changed files with 766 additions and 303 deletions
|
|
@ -30,6 +30,7 @@ export type Payload<D extends Definition = Definition> = {
|
|||
readonly type: D["type"]
|
||||
readonly data: Data<D>
|
||||
readonly version?: number
|
||||
readonly sync?: SyncMetadata
|
||||
readonly location?: Location.Info
|
||||
readonly metadata?: Record<string, unknown>
|
||||
}
|
||||
|
|
@ -48,6 +49,11 @@ export type SerializedEvent = {
|
|||
readonly data: Record<string, unknown>
|
||||
}
|
||||
|
||||
export type SyncMetadata = {
|
||||
readonly seq: number
|
||||
readonly aggregateID: string
|
||||
}
|
||||
|
||||
export class InvalidSyncEventError extends Schema.TaggedErrorClass<InvalidSyncEventError>()(
|
||||
"EventV2.InvalidSyncEvent",
|
||||
{
|
||||
|
|
@ -77,6 +83,12 @@ export function define<const Type extends string, Fields extends Schema.Struct.F
|
|||
metadata: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)),
|
||||
type: Schema.Literal(input.type),
|
||||
version: Schema.optional(Schema.Number),
|
||||
sync: Schema.optional(
|
||||
Schema.Struct({
|
||||
seq: Schema.Finite,
|
||||
aggregateID: Schema.String,
|
||||
}),
|
||||
),
|
||||
location: Schema.optional(Location.Info),
|
||||
data: Data,
|
||||
}).annotate({ identifier: input.type })
|
||||
|
|
@ -186,7 +198,7 @@ export const layer = Layer.effect(
|
|||
)
|
||||
} else {
|
||||
const list = projectors.get(event.type) ?? []
|
||||
yield* db
|
||||
return yield* db
|
||||
.transaction(
|
||||
() =>
|
||||
Effect.gen(function* () {
|
||||
|
|
@ -208,8 +220,12 @@ export const layer = Layer.effect(
|
|||
}),
|
||||
)
|
||||
}
|
||||
const serialized = {
|
||||
seq,
|
||||
aggregateID,
|
||||
}
|
||||
for (const projector of list) {
|
||||
yield* projector(event as Payload)
|
||||
yield* projector({ ...event, sync: serialized })
|
||||
}
|
||||
yield* db
|
||||
.insert(EventSequenceTable)
|
||||
|
|
@ -233,6 +249,7 @@ export const layer = Layer.effect(
|
|||
])
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
return serialized
|
||||
}),
|
||||
{ behavior: "immediate" },
|
||||
)
|
||||
|
|
@ -247,14 +264,15 @@ export const layer = Layer.effect(
|
|||
for (const sync of syncHandlers) {
|
||||
yield* sync(event as Payload)
|
||||
}
|
||||
yield* commitSyncEvent(event as Payload)
|
||||
const sync = yield* commitSyncEvent(event as Payload)
|
||||
const payload = sync ? { ...event, sync } : event
|
||||
for (const listener of listeners) {
|
||||
yield* listener(event as Payload)
|
||||
yield* listener(payload as Payload)
|
||||
}
|
||||
const pubsub = typed.get(event.type)
|
||||
if (pubsub) yield* PubSub.publish(pubsub, event as Payload)
|
||||
yield* PubSub.publish(all, event as Payload)
|
||||
return event
|
||||
if (pubsub) yield* PubSub.publish(pubsub, payload as Payload)
|
||||
yield* PubSub.publish(all, payload as Payload)
|
||||
return payload
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -109,6 +109,19 @@ describe("EventV2", () => {
|
|||
}),
|
||||
)
|
||||
|
||||
it.effect("publishes sync metadata", () =>
|
||||
Effect.gen(function* () {
|
||||
const events = yield* EventV2.Service
|
||||
const aggregateID = EventV2.ID.create()
|
||||
const event = yield* events.publish(SyncMessage, { id: aggregateID, text: "hello" })
|
||||
|
||||
expect(event.sync).toEqual({
|
||||
seq: 0,
|
||||
aggregateID,
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("stores definitions in the exported registry", () =>
|
||||
Effect.sync(() => {
|
||||
expect(EventV2.registry.get(Message.type)).toBe(Message)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue