Files
kilocode/packages/core/test/event.test.ts
T
+25 ef6b152ff8 OpenCode v1.16.2 (#12088)
* feat(worktree): add managed workspace cloning (#30117)

* test(tui): skip crashing keymap textarea renderer

* fix(core): allow skipping migration execution

* fix(opencode): remove automatic full session diffs (#30127)

* chore: generate

* refactor(worktree): move project out of repository

* zen: deepseek flash

* fix(tui): remount session view on session switch (#30129)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>

* go: minimax m3

* refactor(opencode): inline local provider helpers (#30169)

* refactor(opencode): simplify provider setup flow (#30173)

* fix(app): show project sessions before path sync resolves (#30167)

Co-authored-by: LukeParkerDev <10430890+Hona@users.noreply.github.com>

* fix(core): preserve session metadata migration identity (#30176)

* refactor(session): align namespace imports and inline trivial helpers (#30180)

* opencode(run): add queued prompt management (#30103)

Direct run mode previously made submitted follow-up prompts irrevocable while a response was still running. Let users edit or remove queued prompts before dispatch without interrupting the active turn.

* chore: generate

* fix(acp): honor session/cancel by aborting the running turn (#30145)

Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(tui): prevent prompt corruption when pasting near wide characters (#29710)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>
Co-authored-by: Simon Klee <hello@simonklee.dk>

* fix(opencode): avoid nullable webfetch format schema (#30215)

* chore: generate

* fix(core): contain lsp warmup defects (#30226)

* add run --replay mode (#30239)

* chore: generate

* chore: update nix node_modules hashes

* fix(stats): restore leaderboard spacing

* fix(stats): center top models dot grid

* fix(stats): stabilize top models hover

* fix(stats): align big-pickle provider resolution (#30274)

* feat(app): v2 desktop UI improvements (#29689)

Co-authored-by: Brendan Allan <git@brendonovich.dev>
Co-authored-by: Brendan Allan <14191578+Brendonovich@users.noreply.github.com>

* chore: generate

* fix(tui): clarify inline subagent rows (#30051)

* fix(tui): handle events across workspaces (#30281)

* feat(core): update Copilot for token-based billing (#30181)

* fix(tui): keep background marker with subagent label (#30271)

* fix(tui): keep retry attempt before message (#30275)

* chore: generate

* fix(opencode): enforce storage path invariants (#29666)

* chore: generate

* feat(core): add location-based permission service (#30287)

* chore: generate

* fix(tui): preserve live parts during session hydration (#30300)

* fix(app): restore deferred MCP status updates (#30220)

* fix: export v2 stylesheets and declare core node types (#30312)

* chore: update nix node_modules hashes

* fix(app): avoid suspending on pending child path (#30314)

* fix(opencode): remove sunsetted gpt-5.2 and gpt-5.3-codex from allowed models for codex subscriptions (#30316)

* chore: generate

* feat(core): expose session location

* chore: generate

* fix(opencode): preserve websocket api errors (#30321)

* refactor(core): simplify session pagination

* feat(core): add location filesystem contract

* feat(core): add dummy location filesystem layer

* chore: generate

* feat(opencode): add filesystem read and list routes

* chore: generate

* infra: stats

* sync

* feat(app): inset new layout session panels (#30342)

* fix(app): tab title truncation and close button positioning (#30349)

* tui: show model context in run footer (#30380)

* tui: revert OpenTUI upgrade to 0.2.16 (#30383)

* chore: update nix node_modules hashes

* feat(core): add managed repository cache (#30408)

* chore: generate

* chore: generate

* sync

* feat(stats): add cache ratio section

* feat(core): add flagged project references (#30414)

* chore: generate

* feat(core): support named migrations (#30418)

* fix(stats): clean retired provider rows during sync (#30420)

* fix(stats): mention opencode go in top models copy

* feat(core): expose project reference filesystem access (#30423)

* chore: generate

* sync

* fix(tui): scope diff viewer to session directory (#30426)

* test: widen provider header timeout margin (#30427)

* fix(plugin): restore private git install fallback (#30430)

* fix(stats): remove leaderboard nav link

* chore(opencode): remove scout agent (#30435)

* chore: generate

* feat(stats): improve cache ratio chart

* chore: generate

* fix(effect-drizzle-sqlite): preserve transaction begin errors (#30448)

* chore: bump effect beta to 74 (#30449)

* Revert "tui: revert OpenTUI upgrade to 0.2.16 (#30383)" (#30452)

* chore: update nix node_modules hashes

* refactor(opencode): improve startup time by 38% (#30453)

Co-authored-by: starptech <starptech@starptechs-MBP.fritz.box>

* chore: generate

* fix(opencode): patch empty Gemini replay messages (#30463)

* chore: generate

* refactor(core): consolidate filesystem services (#30447)

* chore: generate

* run: enable interactive replay by default (#30465)

* chore: update nix node_modules hashes

* refactor(opencode): remove JSON storage migration (#30461)

* chore: generate

* chore: update nix node_modules hashes

* fix(tui): stop idle background task spinner (#30484)

* refactor(core): move v1 schemas into core (#30473)

* chore: generate

* fix: task id passed to background job for continuation (#30485)

* chore: generate

* feat(core): project copying and tracking directories (#30139)

* chore: generate

* fix(opencode): preserve signed thinking during anthropic reorder (#30182)

* Revert "fix(opencode): preserve signed thinking during anthropic reorder" (#30502)

* fix: rm tool reorder logic from old bug (#30483)

* chore: generate

* feat(app): polish home projects list UI (#30436)

* feat(app): polish select-v2 component (#30446)

Co-authored-by: Brendan Allan <git@brendonovich.dev>

* fix(github): enforce existing git author identity (#30507)

* feat(app): new update button  (#30460)

Co-authored-by: Brendan Allan <git@brendonovich.dev>

* fix(opencode): fallback to sh for curl upgrade (#30499)

Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(ui): render whole-file patches as complete diffs (#30516)

* chore: generate

* feat(app): add servers tab to settings dialog (#29675)

* refactor(core): consolidate pty service (#30537)

* chore: generate

* tui: truncate sidebar file paths (#30531)

* chore: update nix node_modules hashes

* feat(stats): add geo breakdown (#30456)

* chore: generate

* chore: update nix node_modules hashes

* fix(acp): classify apply_patch as edit (#30564)

* fix(acp): classify task as think (#30565)

* fix(acp): include external directory permission context (#30567)

* fix(acp): clean read tool display content (#30569)

* fix(tui): route question responses by session directory (#30578)

* fix(stats): serve stats og image from banner

* docs(go): add Qwen3.7 Plus model (#30594)

* fix(openai): preserve websocket idle state (#30586)

* refactor(core): remove ai sdk option fields (#30581)

* chore: generate

* test(core): cover v1 provider option lowering (#30599)

* chore: generate

* refactor(core): nest model api id (#30603)

* fix(core): expose azure openai xhigh efforts (#30620)

* feat(core): add skill registry and file agent loading (#30617)

* chore: generate

* chore: update nix node_modules hashes

* fix(stats): count all go usage

* chore: remove zed extension and automation (#30628)

* fix(opencode): preserve variant for delegated tasks (#30630)

* zen: update nvidia tos

* fix(opencode): route SAP AI Core reasoning variants through modelParams (#30482)

* chore: generate

* fix(app): hide unavailable titlebar update (#30642)

* feat(app): v2 thinking level selector (#30646)

* fix(app,ui): session review reactivity and VCS query cache (#30660)

* feat(core): add embedded v2 session runtime and tool foundation (#30632)

* chore: generate

* chore: update nix node_modules hashes

* docs: correct compaction prune default (#30670)

* fix(opencode): avoid shell cancel race (#30641)

* feat: bump bedrock and add proper mantle support for openai models through aws bedrock (#30464)

* test: wait for shell truncation readiness (#30679)

* chore: update nix node_modules hashes

* refactor(opencode): clean up task tool prompts (#30687)

* feat(core): add command registry (#30624)

* chore: generate

* fix(acp): replay loaded session transcript (#30645)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>
Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(core): reset pre-launch session projections (#30728)

* feat(tui): improve experimental session switcher (#30738)

* fix(opencode): respect disabled auto compaction on overflow (#30749)

* zen: nemotron 3 ultra

* fix(enterprise): install hono standard validator peer (#30740)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>

* fix build

* chore: update nix node_modules hashes

* make scripts executable

* fix(tui): show toast when variant_list keybind used with no variants (#30724)

* fix(opencode): `ACP.loadSession` should replay all messages (#30761)

Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(opencode): attribute task child agent on creation (#30786)

* fix(tui): add Vue syntax highlighting (#30802)

* fix: bump @openrouter/ai-sdk-provider to 2.9.0 (#30800)

* feat(core): moving sessions (#30640)

* chore: generate

* tweak: background agent prompting to avoid polling issues (#30790)

* upgrade opentui to 0.3.2 (#30748)

* chore: update nix node_modules hashes

* feat(desktop): surface local server startup failures (#30822)

* ci: publish

* refactor(core): make v2 session inputs event sourced (#30785)

* chore: generate

* fix(llm): normalize OpenAI function tool schemas

* chore: generate

* feat(stats): refresh stats routes and homepage (#30419)

* fix(stats): sort metric charts by top usage

* feat(core): add public native API (#30828)

* chore: generate

* feat(app): color themes (#30824)

Co-authored-by: LukeParkerDev <10430890+Hona@users.noreply.github.com>

* chore: generate

* sync release versions for v1.16.0

* feat(core): attach global native tools (#30832)

* chore: generate

* feat(core): add Snowflake Cortex provider (#29901)

Co-authored-by: Cortex Code <noreply@snowflake.com>

* chore: generate

* feat(core): persist v2 session context epochs (#30789)

* chore: generate

* feat(tui): allow backgrounding synchronous subagents (#30488)

* fix(app): improve tab handling (#30669)

* chore: generate

* fix(tui): prioritize models slash autocomplete (#30848)

* fix(tui): route permission replies to session directory (#30851)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>

* fix(cli): harden daemon lifecycle (#30844)

* chore: generate

* feat(app): improve desktop multi-server support (#30678)

Co-authored-by: Brendan Allan <git@brendonovich.dev>

* chore: generate

* fix(app): handle tab overflow and scrolling in titlebar (#30886)

* fix(app): tab overflow (#30894)

* tui: guard path formatting inputs (#30469)

Fixes #27726, #25216, #24856, #24294, #17071, #29164, #24837, #16865, #14279, #29895

* opencode/run: refresh themes after terminal reloads (#30917)

* chore: generate

* fix(tui): fall back to local cwd when editor spawns in attach mode (#30583)

* docs: update Go Qwen tiered pricing (#30936)

* chore: generate

* feat(tui): add diff hunk navigation (#30935)

* chore: rm fuzzy search on references (#30931)

* fix: use mapError instead of orDie for context snapshot decoding (#30905)

Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(core): recover corrupted models cache (#30947)

* chore: bun install (#30968)

* fix(opencode): resolve Bedrock hang by using node build conditions (#30873)

* fix(workflows): retry nix-hashes compute-hash on transient failure (#30743)

* fix(stats): scroll model charts to latest on mobile

* fix(opencode): prevent destructive edit matches (#30932)

* chore: generate

* fix(core): respect v2 default agents (#30969)

* chore: generate

* test(opencode): remove disposal event wait race (#30971)

* test(opencode): remove shell timeout output race (#30974)

* fix(opencode): gate reasoning summaries by provider (#30973)

* feat(core): admit v2 skill guidance (#30843)

* fix(workflows): serialize desktop release uploads (#30978)

* fix(stats): add mobile chart end spacing

* release: v1.16.2

* refactor: kilo compat for v1.16.2

* fix(opencode): address v1.16.2 merge regressions

* chore: update kilo-vscode visual regression baselines

* fix(opencode): restore Kilo behavior after v1.16.2 merge

* fix(opencode): retry Windows migration cleanup

* test(opencode): restore clone and macOS watcher coverage

* fix(opencode): address second-pass review for #12099

Preserve imported usage and retry partial JSON migrations. Refresh active dependency patches, remove the obsolete GCP patch, and regenerate Kilo HttpApi branding.

---------

Co-authored-by: Dax <mail@thdxr.com>
Co-authored-by: Dax Raad <d@ironbay.co>
Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>
Co-authored-by: Frank <frank@anoma.ly>
Co-authored-by: opencode-agent[bot] <219766164+opencode-agent[bot]@users.noreply.github.com>
Co-authored-by: Aiden Cline <63023139+rekram1-node@users.noreply.github.com>
Co-authored-by: Michael Hart <mhart@cloudflare.com>
Co-authored-by: LukeParkerDev <10430890+Hona@users.noreply.github.com>
Co-authored-by: Simon Klee <hello@simonklee.dk>
Co-authored-by: smagnuso <smagnuso@gmail.com>
Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>
Co-authored-by: Orca丶 <93272799+dauphinYan@users.noreply.github.com>
Co-authored-by: Adam <2363879+adamdotdevin@users.noreply.github.com>
Co-authored-by: Aarav Sareen <96787824+arvsrn@users.noreply.github.com>
Co-authored-by: Brendan Allan <git@brendonovich.dev>
Co-authored-by: Brendan Allan <14191578+Brendonovich@users.noreply.github.com>
Co-authored-by: Kit Langton <kit.langton@gmail.com>
Co-authored-by: James Long <longster@gmail.com>
Co-authored-by: Dustin Deus <deusdustin@gmail.com>
Co-authored-by: starptech <starptech@starptechs-MBP.fritz.box>
Co-authored-by: Ulises Jeremias <ulisescf.24@gmail.com>
Co-authored-by: Jack <jack@anoma.ly>
Co-authored-by: Jérôme Benoit <jerome.benoit@sap.com>
Co-authored-by: Ariane Emory <97994360+ariane-emory@users.noreply.github.com>
Co-authored-by: LIU Xinyu <contact@lxy.cc>
Co-authored-by: Colin McDonnell <colinmcd94@gmail.com>
Co-authored-by: Sebastian <hasta84@gmail.com>
Co-authored-by: opencode <opencode@sst.dev>
Co-authored-by: Kamesh Sampath <kamesh.sampath@hotmail.com>
Co-authored-by: Cortex Code <noreply@snowflake.com>
Co-authored-by: pcadena-lila <pcadena@lila.ai>
Co-authored-by: weiconghe <46336277+weiconghe@users.noreply.github.com>
Co-authored-by: alberto <914199+alblez@users.noreply.github.com>
Co-authored-by: kilo-maintainer[bot] <kilo-maintainer[bot]@users.noreply.github.com>
2026-07-13 18:00:35 +02:00

1141 lines
37 KiB
TypeScript

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 { 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 } from "@opencode-ai/core/schema"
import { WorkspaceV2 } from "@opencode-ai/core/workspace"
import { V2Schema } from "@opencode-ai/core/v2-schema"
import { eq } from "drizzle-orm"
import { location } from "./fixture/location"
import { testEffect } from "./lib/effect"
const locationLayer = Layer.succeed(
Location.Service,
Location.Service.of(
location({ directory: AbsolutePath.make("project"), workspaceID: WorkspaceV2.ID.make("wrk_test") }),
),
)
// kilocode_change start - keep concurrent tests isolated from process database migrations
const database = Database.layerFromPath(":memory:")
const eventLayer = Layer.mergeAll(EventV2.defaultLayer, database)
// kilocode_change end
const it = testEffect(eventLayer.pipe(Layer.provideMerge(locationLayer)))
const itWithoutLocation = testEffect(eventLayer)
const Message = EventV2.define({
type: "test.message",
schema: {
text: Schema.String,
},
})
const SyncMessage = EventV2.define({
type: "test.sync",
sync: {
version: 1,
aggregate: "id",
},
schema: {
id: Schema.String,
text: Schema.String,
},
})
const SyncSent = EventV2.define({
type: "test.sent",
sync: {
version: 1,
aggregate: "messageID",
},
schema: {
messageID: Schema.String,
text: Schema.String,
},
})
const GlobalMessage = EventV2.define({
type: "test.global",
schema: {
text: Schema.String,
},
})
const VersionedMessage = EventV2.define({
type: "test.versioned",
sync: {
version: 2,
aggregate: "id",
},
schema: {
id: Schema.String,
text: Schema.String,
},
})
const SyncTimestamp = EventV2.define({
type: "test.timestamp",
sync: {
version: 1,
aggregate: "id",
},
schema: {
id: Schema.String,
timestamp: V2Schema.DateTimeUtcFromMillis,
},
})
describe("EventV2", () => {
it.effect("derives stable namespaced external IDs", () =>
Effect.sync(() => {
const input = { namespace: "opencord.agent-input", key: "input-1" }
expect(EventV2.ID.fromExternal(input)).toBe(EventV2.ID.fromExternal(input))
expect(EventV2.ID.fromExternal(input)).toMatch(/^evt_[a-f0-9]{64}$/)
expect(EventV2.ID.fromExternal({ ...input, namespace: "another-app" })).not.toBe(EventV2.ID.fromExternal(input))
expect(EventV2.ID.fromExternal({ namespace: "a:b", key: "c" })).not.toBe(
EventV2.ID.fromExternal({ namespace: "a", key: "b:c" }),
)
}),
)
it.effect("publishes events with the current location", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const fiber = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
yield* Effect.yieldNow
const event = yield* events.publish(Message, { text: "hello" })
const received = Array.from(yield* Fiber.join(fiber))
expect(received).toEqual([event])
expect(event.type).toBe("test.message")
expect(event).not.toHaveProperty("version")
expect(event.data).toEqual({ text: "hello" })
expect(event.location).toEqual({
directory: AbsolutePath.make("project"),
workspaceID: WorkspaceV2.ID.make("wrk_test"),
})
}),
)
itWithoutLocation.effect("omits location when no location is available", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const event = yield* events.publish(GlobalMessage, { text: "hello" })
expect(event).not.toHaveProperty("location")
expect(event.type).toBe("test.global")
}),
)
it.effect("publishes definition version", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const event = yield* events.publish(VersionedMessage, { id: "one", text: "hello" })
expect(event.type).toBe("test.versioned")
expect(event.version).toBe(2)
}),
)
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", () =>
Effect.sync(() => {
const latest = EventV2.define({
type: "test.out-of-order",
sync: { version: 2, aggregate: "id" },
schema: { id: Schema.String },
})
EventV2.define({
type: "test.out-of-order",
sync: { version: 1, aggregate: "id" },
schema: { id: Schema.String },
})
expect(EventV2.registry.get("test.out-of-order")).toBe(latest)
}),
)
it.effect("publishes to typed and wildcard subscriptions", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const typed = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
const wildcard = yield* events.all().pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
yield* Effect.yieldNow
const event = yield* events.publish(Message, { text: "hello" })
expect(Array.from(yield* Fiber.join(typed))).toEqual([event])
expect(Array.from(yield* Fiber.join(wildcard))).toEqual([event])
}),
)
it.effect("runs projectors inline", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<EventV2.Payload>()
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
const event = yield* events.publish(SyncMessage, { id: "one", text: "hello" })
yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" })
expect(received[0]).toEqual(event)
expect(received[1]?.data).toEqual({ id: "one", text: "after unsubscribe" })
}),
)
it.effect("commits local operational state inside a new synchronized event transaction", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<string>()
const aggregateID = EventV2.ID.create()
yield* events.project(SyncMessage, () => Effect.sync(() => received.push("projector")))
yield* events.publish(
SyncMessage,
{ id: aggregateID, text: "hello" },
{ commit: (seq) => Effect.sync(() => received.push(`commit:${seq}`)) },
)
expect(received).toEqual(["projector", "commit:0"])
}),
)
it.effect("rolls back the synchronized event and projector when the local commit fails", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const aggregateID = EventV2.ID.create()
yield* db.run("CREATE TABLE IF NOT EXISTS event_commit_probe (value text NOT NULL)")
yield* db.run("DELETE FROM event_commit_probe")
yield* events.project(SyncMessage, () =>
db.run("INSERT INTO event_commit_probe (value) VALUES ('projected')").pipe(Effect.orDie, Effect.asVoid),
)
const exit = yield* events
.publish(SyncMessage, { id: aggregateID, text: "hello" }, { commit: () => Effect.die("commit failed") })
.pipe(Effect.exit)
expect(String(exit)).toContain("commit failed")
expect(yield* db.all("SELECT value FROM event_commit_probe")).toEqual([])
expect(yield* db.select().from(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).all()).toEqual([])
expect(
yield* db.select().from(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).all(),
).toEqual([])
}),
)
it.effect("rejects local commit hooks on live-only events", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const exit = yield* events.publish(Message, { text: "hello" }, { commit: () => Effect.void }).pipe(Effect.exit)
expect(String(exit)).toContain("Local commit hooks require a synchronized event")
}),
)
it.effect("runs projectors before publishing to streams", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<string>()
const fiber = yield* events.all().pipe(
Stream.take(1),
Stream.runForEach(() => Effect.sync(() => received.push("stream"))),
Effect.forkScoped,
)
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event.type)
}),
)
yield* Effect.yieldNow
yield* events.publish(SyncMessage, { id: "one", text: "hello" })
yield* Fiber.join(fiber)
expect(received).toEqual([SyncMessage.type, "stream"])
}),
)
it.effect("runs listeners inline after projectors", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<string>()
yield* events.project(SyncMessage, () =>
Effect.sync(() => {
received.push("projector")
}),
)
const unsubscribe = yield* events.listen(() =>
Effect.sync(() => {
received.push("listener")
}),
)
yield* events.publish(SyncMessage, { id: "one", text: "hello" })
yield* unsubscribe
yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" })
expect(received).toEqual(["projector", "listener", "projector"])
}),
)
it.effect("isolates observer defects after durable events commit", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<string>()
yield* events.sync(() => Effect.die("sync defect"))
yield* events.listen(() => {
throw new Error("listener defect")
})
yield* events.listen((event) =>
Effect.sync(() => {
received.push(event.type)
}),
)
const event = yield* events.publish(SyncMessage, { id: "one", text: "hello" })
expect(received).toEqual([SyncMessage.type])
expect(event.seq).toBeNumber()
}),
)
it.effect("preserves observer interruption", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
yield* events.listen(() => Effect.interrupt)
const exit = yield* events.publish(SyncMessage, { id: "interrupted", text: "hello" }).pipe(Effect.exit)
const committed = yield* db
.select({ id: EventTable.id })
.from(EventTable)
.where(eq(EventTable.aggregate_id, "interrupted"))
.get()
.pipe(Effect.orDie)
expect(Exit.isFailure(exit) && Cause.hasInterrupts(exit.cause)).toBeTrue()
expect(committed).toBeDefined()
}),
)
it.effect("keeps live-only listener defects fail-fast", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const defect = new Error("listener defect")
yield* events.listen(() => Effect.die(defect))
expect(yield* events.publish(Message, { text: "hello" }).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect)
}),
)
it.effect("does not synchronize live-only events", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const synchronized = new Array<string>()
const unsubscribe = yield* events.sync((event) =>
Effect.sync(() => {
synchronized.push(event.type)
}),
)
yield* Effect.addFinalizer(() => unsubscribe)
yield* events.publish(Message, { text: "live only" })
yield* events.publish(SyncMessage, { id: "one", text: "durable" })
expect(synchronized).toEqual([SyncMessage.type])
}),
)
it.effect("synchronizes only after the durable event commits", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const synchronized = new Array<boolean>()
yield* events.sync((event) =>
db
.select({ id: EventTable.id })
.from(EventTable)
.where(eq(EventTable.id, event.id))
.get()
.pipe(
Effect.orDie,
Effect.map((row) => synchronized.push(row !== undefined)),
Effect.asVoid,
),
)
yield* events.publish(SyncMessage, { id: EventV2.ID.create(), text: "durable" })
expect(synchronized).toEqual([true])
}),
)
it.effect("inserts sync event rows on publish", () =>
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: "first" })
const rows = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, aggregateID))
.all()
.pipe(Effect.orDie)
expect(rows).toHaveLength(1)
expect(rows[0]?.type).toBe(EventV2.versionedType(SyncMessage.type, 1))
expect(rows[0]?.aggregate_id).toBe(aggregateID)
}),
)
it.effect("increments sync event seq per aggregate", () =>
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: "first" })
yield* events.publish(SyncMessage, { id: aggregateID, text: "second" })
const rows = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, aggregateID))
.all()
.pipe(Effect.orDie)
expect(rows.map((row) => row.seq)).toEqual([0, 1])
}),
)
it.effect("replays durable aggregate events after a cursor 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 fiber = yield* events
.aggregateEvents({ aggregateID, after: EventV2.Cursor.make(0) })
.pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
yield* Effect.yieldNow
yield* events.publish(SyncMessage, { id: aggregateID, text: "two" })
expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.cursor, event.event.data])).toEqual([
[EventV2.Cursor.make(1), { id: aggregateID, text: "one" }],
[EventV2.Cursor.make(2), { id: aggregateID, text: "two" }],
])
}),
)
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 fiber = yield* events
.aggregateEvents({ aggregateID })
.pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
yield* events.publish(SyncMessage, { id: aggregateID, text: "one" })
expect(
Array.from(yield* Fiber.join(fiber)).map((event) => [
event.cursor,
(event.event.data as { text: string }).text,
]),
).toEqual([
[EventV2.Cursor.make(0), "zero"],
[EventV2.Cursor.make(1), "one"],
])
}),
)
it.effect("retains a durable wake committed while historical replay is paused", () =>
Effect.gen(function* () {
const readStarted = yield* Deferred.make<void>()
const continueRead = yield* Deferred.make<void>()
let pause = true
const database = Database.layerFromPath(":memory:")
const eventLayer = EventV2.layerWith({
beforeAggregateRead: () =>
pause
? Deferred.succeed(readStarted, undefined).pipe(Effect.andThen(Deferred.await(continueRead)))
: Effect.void,
}).pipe(Layer.provide(database))
yield* Effect.gen(function* () {
const events = yield* EventV2.Service
const aggregateID = EventV2.ID.create()
const fiber = yield* events
.aggregateEvents({ 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* Deferred.succeed(continueRead, undefined)
expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.cursor, event.event.data])).toEqual([
[EventV2.Cursor.make(0), { id: aggregateID, text: "during handoff" }],
])
}).pipe(Effect.provide(Layer.mergeAll(database, eventLayer)))
}),
)
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 count = 64
const fiber = yield* events
.aggregateEvents({ aggregateID })
.pipe(Stream.take(count), Stream.runCollect, Effect.forkScoped)
yield* Effect.yieldNow
for (let index = 0; index < count; index++) {
yield* events.publish(SyncMessage, { id: aggregateID, text: String(index) })
}
expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.cursor, event.event.data])).toEqual(
Array.from({ length: count }, (_, index) => [
EventV2.Cursor.make(index),
{ id: aggregateID, text: String(index) },
]),
)
}),
)
it.effect("omits live-only events from durable aggregate streams", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const aggregateID = EventV2.ID.create()
const fiber = yield* events
.aggregateEvents({ 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" })
expect(Array.from(yield* Fiber.join(fiber)).map((event) => event.event.type)).toEqual([SyncMessage.type])
}),
)
it.effect("uses custom sync aggregate field", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const aggregateID = EventV2.ID.create()
yield* events.publish(SyncSent, { messageID: aggregateID, text: "sent" })
const rows = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, aggregateID))
.all()
.pipe(Effect.orDie)
expect(rows).toHaveLength(1)
expect(rows[0]?.aggregate_id).toBe(aggregateID)
}),
)
it.effect("replays sync events through projectors", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<EventV2.Payload>()
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
const aggregateID = EventV2.ID.create()
yield* events.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "hello" },
})
expect(received[0]?.type).toBe(SyncMessage.type)
expect(received[0]?.data).toEqual({ id: aggregateID, text: "hello" })
}),
)
it.effect("replay inserts external event rows", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const aggregateID = EventV2.ID.create()
yield* events.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "replayed" },
})
const rows = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, aggregateID))
.all()
.pipe(Effect.orDie)
expect(rows).toHaveLength(1)
expect(rows[0]?.aggregate_id).toBe(aggregateID)
}),
)
it.effect(
"replay rejects an envelope aggregate that differs from its payload without mutating the payload aggregate",
() =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const envelopeAggregateID = EventV2.ID.create()
const payloadAggregateID = EventV2.ID.create()
const received = new Array<EventV2.Payload>()
yield* events.publish(SyncMessage, { id: payloadAggregateID, text: "seed" })
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
const exit = yield* events
.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID: envelopeAggregateID,
data: { id: payloadAggregateID, text: "replayed" },
})
.pipe(Effect.exit)
const rows = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, payloadAggregateID))
.all()
.pipe(Effect.orDie)
const sequence = yield* db
.select({ seq: EventSequenceTable.seq })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, payloadAggregateID))
.get()
.pipe(Effect.orDie)
expect(String(exit)).toContain("Aggregate mismatch")
expect(received).toHaveLength(0)
expect(rows).toHaveLength(1)
expect(sequence).toEqual({ seq: 0 })
}),
)
it.effect("replay defects on sequence mismatch", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const aggregateID = EventV2.ID.create()
yield* events.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "first" },
})
const exit = yield* events
.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 5,
aggregateID,
data: { id: aggregateID, text: "bad" },
})
.pipe(Effect.exit)
expect(String(exit)).toContain("Sequence mismatch")
}),
)
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) =>
Effect.sync(() => {
received.push(event)
}),
)
yield* events.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncTimestamp.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, timestamp: 0 },
})
expect(received[0]?.data.timestamp).toEqual(DateTime.makeUnsafe(0))
}),
)
it.effect("replay defects on unknown event type", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const exit = yield* events
.replay({
id: EventV2.ID.create(),
type: "unknown.event.1",
seq: 0,
aggregateID: EventV2.ID.create(),
data: {},
})
.pipe(Effect.exit)
expect(String(exit)).toContain("Unknown sync event type")
}),
)
it.effect("replayAll validates contiguous aggregate events", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const aggregateID = EventV2.ID.create()
const source = yield* events.replayAll([
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "one" },
},
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "two" },
},
])
expect(source).toBe(aggregateID)
}),
)
it.effect("replayAll accepts later chunks after the first batch", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const aggregateID = EventV2.ID.create()
const one = yield* events.replayAll([
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "one" },
},
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "two" },
},
])
const two = yield* events.replayAll([
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 2,
aggregateID,
data: { id: aggregateID, text: "three" },
},
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 3,
aggregateID,
data: { id: aggregateID, text: "four" },
},
])
const rows = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, aggregateID))
.all()
.pipe(Effect.orDie)
expect(one).toBe(aggregateID)
expect(two).toBe(aggregateID)
expect(rows.map((row) => row.seq)).toEqual([0, 1, 2, 3])
}),
)
it.effect("claim fences replay owners", () =>
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" })
yield* events.claim(aggregateID, "owner-a")
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
yield* events.replay(
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "ignored" },
},
{ ownerID: "owner-b" },
)
expect(received).toHaveLength(0)
}),
)
it.effect("strict owner fences exact replay", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const aggregateID = EventV2.ID.create()
const id = EventV2.ID.create()
const replayed = {
id,
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "owned" },
}
yield* events.replay(replayed, { ownerID: "owner-a" })
const exit = yield* events.replay(replayed, { ownerID: "owner-b", strictOwner: true }).pipe(Effect.exit)
expect(String(exit)).toContain("Replay owner mismatch")
}),
)
it.effect("exact replay claims an unowned aggregate", () =>
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 replayed = {
id: published.id,
type: EventV2.versionedType(SyncMessage.type, 1),
seq: published.seq!,
aggregateID,
data: published.data,
}
yield* events.replay(replayed, { ownerID: "owner-a", strictOwner: true })
const row = yield* db
.select({ ownerID: EventSequenceTable.owner_id })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, aggregateID))
.get()
.pipe(Effect.orDie)
expect(row?.ownerID).toBe("owner-a")
const exit = yield* events
.replay(
{ ...replayed, id: EventV2.ID.create(), seq: 1, data: { id: aggregateID, text: "conflict" } },
{ ownerID: "owner-b", strictOwner: true },
)
.pipe(Effect.exit)
expect(String(exit)).toContain("Replay owner mismatch")
}),
)
it.effect("replay with owner claims an unowned sequence", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const aggregateID = EventV2.ID.create()
yield* events.replay(
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "owned" },
},
{ ownerID: "owner-1" },
)
const row = yield* db
.select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, aggregateID))
.get()
.pipe(Effect.orDie)
expect(row).toEqual({ seq: 0, ownerID: "owner-1" })
}),
)
it.effect("replay claims an existing unowned sequence before fencing a different owner", () =>
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" })
yield* events.replay(
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "claimed" },
},
{ ownerID: "owner-1" },
)
yield* events.replay(
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 2,
aggregateID,
data: { id: aggregateID, text: "fenced" },
},
{ ownerID: "owner-2" },
)
const rows = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, aggregateID))
.all()
.pipe(Effect.orDie)
const sequence = yield* db
.select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, aggregateID))
.get()
.pipe(Effect.orDie)
expect(rows.map((row) => row.seq)).toEqual([0, 1])
expect(sequence).toEqual({ seq: 1, ownerID: "owner-1" })
}),
)
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()
yield* events.replay(
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "claimed" },
},
{ ownerID: "owner-1" },
)
const exit = yield* events
.replay(
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "conflict" },
},
{ ownerID: "owner-2", strictOwner: true },
)
.pipe(Effect.exit)
expect(String(exit)).toContain("Replay owner mismatch")
}),
)
it.effect("publishes accepted replay with its durable sequence and suppresses stale replay", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<EventV2.Payload>()
const aggregateID = EventV2.ID.create()
yield* events.listen((event) => Effect.sync(() => received.push(event)))
const replayed = {
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "replayed" },
}
yield* events.replay(replayed, { publish: true })
yield* events.replay(replayed, { publish: true })
expect(received).toMatchObject([{ id: replayed.id, seq: 0, data: replayed.data }])
}),
)
it.effect("rejects divergent stale replay without publishing it", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<EventV2.Payload>()
const aggregateID = EventV2.ID.create()
const replayed = {
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "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 })
.pipe(Effect.exit)
expect(String(exit)).toContain("Replay diverged")
expect(received).toHaveLength(1)
}),
)
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 id = EventV2.ID.create()
yield* events.replay({
id,
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "first" },
})
const exit = yield* events
.replay({
id,
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "second" },
})
.pipe(Effect.exit)
expect(String(exit)).toContain(`Event ${id} already exists`)
}),
)
it.effect("replay from a different owner leaves claimed sequence unchanged", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const aggregateID = EventV2.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),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "first" },
},
{ ownerID: "owner-1" },
)
yield* events.replay(
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "ignored" },
},
{ ownerID: "owner-2", publish: true },
)
const rows = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, aggregateID))
.all()
.pipe(Effect.orDie)
const sequence = yield* db
.select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, aggregateID))
.get()
.pipe(Effect.orDie)
expect(rows).toHaveLength(1)
expect(sequence).toEqual({ seq: 0, ownerID: "owner-1" })
expect(received).toHaveLength(0)
}),
)
it.effect("claim updates the event sequence owner", () =>
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: "claimed" })
yield* events.claim(aggregateID, "owner-1")
yield* events.claim(aggregateID, "owner-2")
const row = yield* db
.select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, aggregateID))
.get()
.pipe(Effect.orDie)
expect(row).toEqual({ seq: 0, ownerID: "owner-2" })
}),
)
it.effect("remove clears sync event sequence", () =>
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" })
yield* events.remove(aggregateID)
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
yield* events.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "replayed" },
})
expect(received[0]?.data).toEqual({ id: aggregateID, text: "replayed" })
}),
)
})