|
|
@@ -1,10 +1,14 @@
|
|
|
import { describe, expect } from "bun:test"
|
|
|
import { Cause, DateTime, Deferred, Effect, Exit, Fiber, Layer, Schema, Stream } from "effect"
|
|
|
import { EventV2 } from "@opencode-ai/core/event"
|
|
|
+import { Event } from "@opencode-ai/schema/event"
|
|
|
+import { Session } from "@opencode-ai/schema/session"
|
|
|
+import { SessionEvent } from "@opencode-ai/schema/session-event"
|
|
|
+import { SessionV1 } from "@opencode-ai/schema/session-v1"
|
|
|
import { Database } from "@opencode-ai/core/database/database"
|
|
|
import { EventSequenceTable, EventTable } from "@opencode-ai/core/event/sql"
|
|
|
import { Location } from "@opencode-ai/core/location"
|
|
|
-import { AbsolutePath, DateTimeUtcFromMillis } from "@opencode-ai/core/schema"
|
|
|
+import { AbsolutePath } from "@opencode-ai/core/schema"
|
|
|
import { WorkspaceV2 } from "@opencode-ai/core/workspace"
|
|
|
import { eq } from "drizzle-orm"
|
|
|
import { location } from "./fixture/location"
|
|
|
@@ -16,10 +20,6 @@ const locationLayer = Layer.succeed(
|
|
|
location({ directory: AbsolutePath.make("project"), workspaceID: WorkspaceV2.ID.make("wrk_test") }),
|
|
|
),
|
|
|
)
|
|
|
-const eventLayer = Layer.mergeAll(EventV2.defaultLayer, Database.defaultLayer)
|
|
|
-const it = testEffect(eventLayer.pipe(Layer.provideMerge(locationLayer)))
|
|
|
-const itWithoutLocation = testEffect(eventLayer)
|
|
|
-
|
|
|
const Message = EventV2.define({
|
|
|
type: "test.message",
|
|
|
schema: {
|
|
|
@@ -70,18 +70,16 @@ const VersionedMessage = EventV2.define({
|
|
|
},
|
|
|
})
|
|
|
|
|
|
-const SyncTimestamp = EventV2.define({
|
|
|
- type: "test.timestamp",
|
|
|
- durable: {
|
|
|
- version: 1,
|
|
|
- aggregate: "id",
|
|
|
- },
|
|
|
- schema: {
|
|
|
- id: Schema.String,
|
|
|
- timestamp: DateTimeUtcFromMillis,
|
|
|
- },
|
|
|
+const DurableMessage = SessionV1.Event.MessageRemoved
|
|
|
+const durableData = (sessionID: Session.ID, text: string) => ({
|
|
|
+ sessionID,
|
|
|
+ messageID: SessionV1.MessageID.ascending(`msg_${text}`),
|
|
|
})
|
|
|
|
|
|
+const eventLayer = Layer.mergeAll(EventV2.layerWith().pipe(Layer.provide(Database.defaultLayer)), Database.defaultLayer)
|
|
|
+const it = testEffect(eventLayer.pipe(Layer.provideMerge(locationLayer)))
|
|
|
+const itWithoutLocation = testEffect(eventLayer)
|
|
|
+
|
|
|
describe("EventV2", () => {
|
|
|
it.effect("publishes events with the current location", () =>
|
|
|
Effect.gen(function* () {
|
|
|
@@ -122,26 +120,21 @@ describe("EventV2", () => {
|
|
|
}),
|
|
|
)
|
|
|
|
|
|
- it.effect("stores definitions in the exported registry", () =>
|
|
|
- Effect.sync(() => {
|
|
|
- expect(EventV2.registry.get(Message.type)).toBe(Message)
|
|
|
- }),
|
|
|
- )
|
|
|
-
|
|
|
- it.effect("keeps the latest sync definition in the registry", () =>
|
|
|
+ it.effect("selects the latest durable definition independent of declaration order", () =>
|
|
|
Effect.sync(() => {
|
|
|
const latest = EventV2.define({
|
|
|
type: "test.out-of-order",
|
|
|
durable: { version: 2, aggregate: "id" },
|
|
|
schema: { id: Schema.String },
|
|
|
})
|
|
|
- EventV2.define({
|
|
|
+ const historical = EventV2.define({
|
|
|
type: "test.out-of-order",
|
|
|
durable: { version: 1, aggregate: "id" },
|
|
|
schema: { id: Schema.String },
|
|
|
})
|
|
|
|
|
|
- expect(EventV2.registry.get("test.out-of-order")).toBe(latest)
|
|
|
+ expect(Event.latest([latest, historical]).get("test.out-of-order")).toBe(latest)
|
|
|
+ expect(Event.latest([historical, latest]).get("test.out-of-order")).toBe(latest)
|
|
|
}),
|
|
|
)
|
|
|
|
|
|
@@ -363,19 +356,19 @@ describe("EventV2", () => {
|
|
|
it.effect("replays durable aggregate events after a sequence and tails new events", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "zero" })
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "one" })
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
|
|
|
const fiber = yield* events
|
|
|
.durable({ aggregateID, after: 0 })
|
|
|
.pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
|
|
|
yield* Effect.yieldNow
|
|
|
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "two" })
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "two"))
|
|
|
|
|
|
expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
|
|
|
- [1, { id: aggregateID, text: "one" }],
|
|
|
- [2, { id: aggregateID, text: "two" }],
|
|
|
+ [1, durableData(aggregateID, "one")],
|
|
|
+ [2, durableData(aggregateID, "two")],
|
|
|
])
|
|
|
}),
|
|
|
)
|
|
|
@@ -383,20 +376,15 @@ describe("EventV2", () => {
|
|
|
it.effect("catches durable aggregate events published during replay handoff", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "zero" })
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
|
|
|
const fiber = yield* events.durable({ aggregateID }).pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
|
|
|
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "one" })
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
|
|
|
|
|
|
- expect(
|
|
|
- Array.from(yield* Fiber.join(fiber)).map((event) => [
|
|
|
- event.durable?.seq,
|
|
|
- (event.data as { text: string }).text,
|
|
|
- ]),
|
|
|
- ).toEqual([
|
|
|
- [0, "zero"],
|
|
|
- [1, "one"],
|
|
|
+ expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
|
|
|
+ [0, durableData(aggregateID, "zero")],
|
|
|
+ [1, durableData(aggregateID, "one")],
|
|
|
])
|
|
|
}),
|
|
|
)
|
|
|
@@ -415,16 +403,16 @@ describe("EventV2", () => {
|
|
|
|
|
|
yield* Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
const fiber = yield* events.durable({ aggregateID }).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
|
|
|
yield* Deferred.await(readStarted)
|
|
|
|
|
|
pause = false
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "during handoff" })
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "during handoff"))
|
|
|
yield* Deferred.succeed(continueRead, undefined)
|
|
|
|
|
|
expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
|
|
|
- [0, { id: aggregateID, text: "during handoff" }],
|
|
|
+ [0, durableData(aggregateID, "during handoff")],
|
|
|
])
|
|
|
}).pipe(Effect.provide(Layer.mergeAll(Database.defaultLayer, eventLayer)))
|
|
|
}),
|
|
|
@@ -433,7 +421,7 @@ describe("EventV2", () => {
|
|
|
it.effect("coalesces durable aggregate wakes while draining every committed event", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
const count = 64
|
|
|
const fiber = yield* events
|
|
|
.durable({ aggregateID })
|
|
|
@@ -441,11 +429,11 @@ describe("EventV2", () => {
|
|
|
yield* Effect.yieldNow
|
|
|
|
|
|
for (let index = 0; index < count; index++) {
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: String(index) })
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, String(index)))
|
|
|
}
|
|
|
|
|
|
expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual(
|
|
|
- Array.from({ length: count }, (_, index) => [index, { id: aggregateID, text: String(index) }]),
|
|
|
+ Array.from({ length: count }, (_, index) => [index, durableData(aggregateID, String(index))]),
|
|
|
)
|
|
|
}),
|
|
|
)
|
|
|
@@ -453,14 +441,14 @@ describe("EventV2", () => {
|
|
|
it.effect("omits live-only events from durable aggregate streams", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
const fiber = yield* events.durable({ aggregateID }).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
|
|
|
yield* Effect.yieldNow
|
|
|
|
|
|
yield* events.publish(Message, { text: "live only" })
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "durable" })
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "durable"))
|
|
|
|
|
|
- expect(Array.from(yield* Fiber.join(fiber)).map((event) => event.type)).toEqual([SyncMessage.type])
|
|
|
+ expect(Array.from(yield* Fiber.join(fiber)).map((event) => event.type)).toEqual([DurableMessage.type])
|
|
|
}),
|
|
|
)
|
|
|
|
|
|
@@ -487,23 +475,23 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const received = new Array<EventV2.Payload>()
|
|
|
- yield* events.project(SyncMessage, (event) =>
|
|
|
+ yield* events.project(DurableMessage, (event) =>
|
|
|
Effect.sync(() => {
|
|
|
received.push(event)
|
|
|
}),
|
|
|
)
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
|
|
|
yield* events.replay({
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "hello" },
|
|
|
+ data: durableData(aggregateID, "hello"),
|
|
|
})
|
|
|
|
|
|
- expect(received[0]?.type).toBe(SyncMessage.type)
|
|
|
- expect(received[0]?.data).toEqual({ id: aggregateID, text: "hello" })
|
|
|
+ expect(received[0]?.type).toBe(DurableMessage.type)
|
|
|
+ expect(received[0]?.data).toEqual(durableData(aggregateID, "hello"))
|
|
|
}),
|
|
|
)
|
|
|
|
|
|
@@ -511,14 +499,14 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const { db } = yield* Database.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
|
|
|
yield* events.replay({
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "replayed" },
|
|
|
+ data: durableData(aggregateID, "replayed"),
|
|
|
})
|
|
|
const rows = yield* db
|
|
|
.select()
|
|
|
@@ -538,11 +526,11 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const { db } = yield* Database.Service
|
|
|
- const envelopeAggregateID = EventV2.ID.create()
|
|
|
- const payloadAggregateID = EventV2.ID.create()
|
|
|
+ const envelopeAggregateID = Session.ID.create()
|
|
|
+ const payloadAggregateID = Session.ID.create()
|
|
|
const received = new Array<EventV2.Payload>()
|
|
|
- yield* events.publish(SyncMessage, { id: payloadAggregateID, text: "seed" })
|
|
|
- yield* events.project(SyncMessage, (event) =>
|
|
|
+ yield* events.publish(DurableMessage, durableData(payloadAggregateID, "seed"))
|
|
|
+ yield* events.project(DurableMessage, (event) =>
|
|
|
Effect.sync(() => {
|
|
|
received.push(event)
|
|
|
}),
|
|
|
@@ -551,10 +539,10 @@ describe("EventV2", () => {
|
|
|
const exit = yield* events
|
|
|
.replay({
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 1,
|
|
|
aggregateID: envelopeAggregateID,
|
|
|
- data: { id: payloadAggregateID, text: "replayed" },
|
|
|
+ data: durableData(payloadAggregateID, "replayed"),
|
|
|
})
|
|
|
.pipe(Effect.exit)
|
|
|
const rows = yield* db
|
|
|
@@ -580,22 +568,22 @@ describe("EventV2", () => {
|
|
|
it.effect("replay defects on sequence mismatch", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
|
|
|
yield* events.replay({
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "first" },
|
|
|
+ data: durableData(aggregateID, "first"),
|
|
|
})
|
|
|
const exit = yield* events
|
|
|
.replay({
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 5,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "bad" },
|
|
|
+ data: durableData(aggregateID, "bad"),
|
|
|
})
|
|
|
.pipe(Effect.exit)
|
|
|
|
|
|
@@ -606,9 +594,9 @@ describe("EventV2", () => {
|
|
|
it.effect("replay decodes synchronized transformed values before projection", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
- const received = new Array<typeof SyncTimestamp.Type>()
|
|
|
- yield* events.project(SyncTimestamp, (event) =>
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
+ const received = new Array<typeof SessionEvent.ContextUpdated.Type>()
|
|
|
+ yield* events.project(SessionEvent.ContextUpdated, (event) =>
|
|
|
Effect.sync(() => {
|
|
|
received.push(event)
|
|
|
}),
|
|
|
@@ -616,10 +604,10 @@ describe("EventV2", () => {
|
|
|
|
|
|
yield* events.replay({
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncTimestamp.type, 1),
|
|
|
+ type: EventV2.versionedType(SessionEvent.ContextUpdated.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, timestamp: 0 },
|
|
|
+ data: { sessionID: aggregateID, messageID: "msg_context", timestamp: 0, text: "context" },
|
|
|
})
|
|
|
|
|
|
expect(received[0]?.data.timestamp).toEqual(DateTime.makeUnsafe(0))
|
|
|
@@ -646,21 +634,21 @@ describe("EventV2", () => {
|
|
|
it.effect("replayAll validates contiguous aggregate events", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
const source = yield* events.replayAll([
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "one" },
|
|
|
+ data: durableData(aggregateID, "one"),
|
|
|
},
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 1,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "two" },
|
|
|
+ data: durableData(aggregateID, "two"),
|
|
|
},
|
|
|
])
|
|
|
|
|
|
@@ -672,38 +660,38 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const { db } = yield* Database.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
|
|
|
const one = yield* events.replayAll([
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "one" },
|
|
|
+ data: durableData(aggregateID, "one"),
|
|
|
},
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 1,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "two" },
|
|
|
+ data: durableData(aggregateID, "two"),
|
|
|
},
|
|
|
])
|
|
|
const two = yield* events.replayAll([
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 2,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "three" },
|
|
|
+ data: durableData(aggregateID, "three"),
|
|
|
},
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 3,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "four" },
|
|
|
+ data: durableData(aggregateID, "four"),
|
|
|
},
|
|
|
])
|
|
|
const rows = yield* db
|
|
|
@@ -723,10 +711,10 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const received = new Array<EventV2.Payload>()
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "seed" })
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "seed"))
|
|
|
yield* events.claim(aggregateID, "owner-a")
|
|
|
- yield* events.project(SyncMessage, (event) =>
|
|
|
+ yield* events.project(DurableMessage, (event) =>
|
|
|
Effect.sync(() => {
|
|
|
received.push(event)
|
|
|
}),
|
|
|
@@ -735,10 +723,10 @@ describe("EventV2", () => {
|
|
|
yield* events.replay(
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 1,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "ignored" },
|
|
|
+ data: durableData(aggregateID, "ignored"),
|
|
|
},
|
|
|
{ ownerID: "owner-b" },
|
|
|
)
|
|
|
@@ -750,14 +738,14 @@ describe("EventV2", () => {
|
|
|
it.effect("strict owner fences exact replay", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
const id = EventV2.ID.create()
|
|
|
const replayed = {
|
|
|
id,
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "owned" },
|
|
|
+ data: durableData(aggregateID, "owned"),
|
|
|
}
|
|
|
yield* events.replay(replayed, { ownerID: "owner-a" })
|
|
|
|
|
|
@@ -771,11 +759,11 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const { db } = yield* Database.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
- const published = yield* events.publish(SyncMessage, { id: aggregateID, text: "owned" })
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
+ const published = yield* events.publish(DurableMessage, durableData(aggregateID, "owned"))
|
|
|
const replayed = {
|
|
|
id: published.id,
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: published.durable!.seq,
|
|
|
aggregateID,
|
|
|
data: published.data,
|
|
|
@@ -792,7 +780,7 @@ describe("EventV2", () => {
|
|
|
expect(row?.ownerID).toBe("owner-a")
|
|
|
const exit = yield* events
|
|
|
.replay(
|
|
|
- { ...replayed, id: EventV2.ID.create(), seq: 1, data: { id: aggregateID, text: "conflict" } },
|
|
|
+ { ...replayed, id: EventV2.ID.create(), seq: 1, data: durableData(aggregateID, "conflict") },
|
|
|
{ ownerID: "owner-b", strictOwner: true },
|
|
|
)
|
|
|
.pipe(Effect.exit)
|
|
|
@@ -804,15 +792,15 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const { db } = yield* Database.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
|
|
|
yield* events.replay(
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "owned" },
|
|
|
+ data: durableData(aggregateID, "owned"),
|
|
|
},
|
|
|
{ ownerID: "owner-1" },
|
|
|
)
|
|
|
@@ -831,26 +819,26 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const { db } = yield* Database.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "local" })
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "local"))
|
|
|
|
|
|
yield* events.replay(
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 1,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "claimed" },
|
|
|
+ data: durableData(aggregateID, "claimed"),
|
|
|
},
|
|
|
{ ownerID: "owner-1" },
|
|
|
)
|
|
|
yield* events.replay(
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 2,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "fenced" },
|
|
|
+ data: durableData(aggregateID, "fenced"),
|
|
|
},
|
|
|
{ ownerID: "owner-2" },
|
|
|
)
|
|
|
@@ -875,14 +863,14 @@ describe("EventV2", () => {
|
|
|
it.effect("strict replay rejects an owner conflict instead of silently skipping it", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
yield* events.replay(
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "claimed" },
|
|
|
+ data: durableData(aggregateID, "claimed"),
|
|
|
},
|
|
|
{ ownerID: "owner-1" },
|
|
|
)
|
|
|
@@ -891,10 +879,10 @@ describe("EventV2", () => {
|
|
|
.replay(
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 1,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "conflict" },
|
|
|
+ data: durableData(aggregateID, "conflict"),
|
|
|
},
|
|
|
{ ownerID: "owner-2", strictOwner: true },
|
|
|
)
|
|
|
@@ -908,14 +896,14 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const received = new Array<EventV2.Payload>()
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
yield* events.listen((event) => Effect.sync(() => received.push(event)))
|
|
|
const replayed = {
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "replayed" },
|
|
|
+ data: durableData(aggregateID, "replayed"),
|
|
|
}
|
|
|
|
|
|
yield* events.replay(replayed, { publish: true })
|
|
|
@@ -929,19 +917,19 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const received = new Array<EventV2.Payload>()
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
const replayed = {
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "original" },
|
|
|
+ data: durableData(aggregateID, "original"),
|
|
|
}
|
|
|
yield* events.listen((event) => Effect.sync(() => received.push(event)))
|
|
|
yield* events.replay(replayed, { publish: true })
|
|
|
|
|
|
const exit = yield* events
|
|
|
- .replay({ ...replayed, data: { id: aggregateID, text: "divergent" } }, { publish: true })
|
|
|
+ .replay({ ...replayed, data: durableData(aggregateID, "divergent") }, { publish: true })
|
|
|
.pipe(Effect.exit)
|
|
|
|
|
|
expect(String(exit)).toContain("Replay diverged")
|
|
|
@@ -952,23 +940,23 @@ describe("EventV2", () => {
|
|
|
it.effect("rejects an event ID reused at another aggregate position", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
const id = EventV2.ID.create()
|
|
|
yield* events.replay({
|
|
|
id,
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "first" },
|
|
|
+ data: durableData(aggregateID, "first"),
|
|
|
})
|
|
|
|
|
|
const exit = yield* events
|
|
|
.replay({
|
|
|
id,
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 1,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "second" },
|
|
|
+ data: durableData(aggregateID, "second"),
|
|
|
})
|
|
|
.pipe(Effect.exit)
|
|
|
|
|
|
@@ -980,27 +968,27 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const { db } = yield* Database.Service
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
const received = new Array<EventV2.Payload>()
|
|
|
yield* events.listen((event) => Effect.sync(() => received.push(event)))
|
|
|
|
|
|
yield* events.replay(
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "first" },
|
|
|
+ data: durableData(aggregateID, "first"),
|
|
|
},
|
|
|
{ ownerID: "owner-1" },
|
|
|
)
|
|
|
yield* events.replay(
|
|
|
{
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 1,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "ignored" },
|
|
|
+ data: durableData(aggregateID, "ignored"),
|
|
|
},
|
|
|
{ ownerID: "owner-2", publish: true },
|
|
|
)
|
|
|
@@ -1047,10 +1035,10 @@ describe("EventV2", () => {
|
|
|
Effect.gen(function* () {
|
|
|
const events = yield* EventV2.Service
|
|
|
const received = new Array<EventV2.Payload>()
|
|
|
- const aggregateID = EventV2.ID.create()
|
|
|
- yield* events.publish(SyncMessage, { id: aggregateID, text: "seed" })
|
|
|
+ const aggregateID = Session.ID.create()
|
|
|
+ yield* events.publish(DurableMessage, durableData(aggregateID, "seed"))
|
|
|
yield* events.remove(aggregateID)
|
|
|
- yield* events.project(SyncMessage, (event) =>
|
|
|
+ yield* events.project(DurableMessage, (event) =>
|
|
|
Effect.sync(() => {
|
|
|
received.push(event)
|
|
|
}),
|
|
|
@@ -1058,13 +1046,13 @@ describe("EventV2", () => {
|
|
|
|
|
|
yield* events.replay({
|
|
|
id: EventV2.ID.create(),
|
|
|
- type: EventV2.versionedType(SyncMessage.type, 1),
|
|
|
+ type: EventV2.versionedType(DurableMessage.type, 1),
|
|
|
seq: 0,
|
|
|
aggregateID,
|
|
|
- data: { id: aggregateID, text: "replayed" },
|
|
|
+ data: durableData(aggregateID, "replayed"),
|
|
|
})
|
|
|
|
|
|
- expect(received[0]?.data).toEqual({ id: aggregateID, text: "replayed" })
|
|
|
+ expect(received[0]?.data).toEqual(durableData(aggregateID, "replayed"))
|
|
|
}),
|
|
|
)
|
|
|
})
|