refactor(sync): publish via EffectBridge.fork for codebase consistency (#28187)
This commit is contained in:
parent
762850dfe5
commit
159d271e1e
2 changed files with 79 additions and 55 deletions
|
|
@ -7,9 +7,7 @@ import { eq } from "drizzle-orm"
|
||||||
import { GlobalBus } from "@/bus/global"
|
import { GlobalBus } from "@/bus/global"
|
||||||
import { Bus as ProjectBus } from "@/bus"
|
import { Bus as ProjectBus } from "@/bus"
|
||||||
import { BusEvent } from "@/bus/bus-event"
|
import { BusEvent } from "@/bus/bus-event"
|
||||||
import type { InstanceContext } from "@/project/instance-context"
|
|
||||||
import { EventSequenceTable, EventTable } from "./event.sql"
|
import { EventSequenceTable, EventTable } from "./event.sql"
|
||||||
import type { WorkspaceID } from "@/control-plane/schema"
|
|
||||||
import { EventID } from "./schema"
|
import { EventID } from "./schema"
|
||||||
import { Context, Effect, Layer, Schema as EffectSchema } from "effect"
|
import { Context, Effect, Layer, Schema as EffectSchema } from "effect"
|
||||||
import type { DeepMutable } from "@opencode-ai/core/schema"
|
import type { DeepMutable } from "@opencode-ai/core/schema"
|
||||||
|
|
@ -17,7 +15,7 @@ import { EventV2 } from "@opencode-ai/core/event"
|
||||||
import { serviceUse } from "@/effect/service-use"
|
import { serviceUse } from "@/effect/service-use"
|
||||||
import { InstanceState } from "@/effect/instance-state"
|
import { InstanceState } from "@/effect/instance-state"
|
||||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||||
import { attachWith } from "@/effect/run-service"
|
import { EffectBridge } from "@/effect/bridge"
|
||||||
|
|
||||||
// Keep `Event["data"]` mutable because projectors mutate the persisted shape
|
// Keep `Event["data"]` mutable because projectors mutate the persisted shape
|
||||||
// when writing to the database. Bus payloads (`Properties`) stay readonly —
|
// when writing to the database. Bus payloads (`Properties`) stay readonly —
|
||||||
|
|
@ -51,10 +49,6 @@ export type SerializedEvent<Def extends Definition = Definition> = Event<Def> &
|
||||||
|
|
||||||
type ProjectorFunc = (db: Database.TxOrDb, data: unknown, event: Event) => void
|
type ProjectorFunc = (db: Database.TxOrDb, data: unknown, event: Event) => void
|
||||||
type ConvertEvent = (type: string, data: Event["data"]) => unknown | Promise<unknown>
|
type ConvertEvent = (type: string, data: Event["data"]) => unknown | Promise<unknown>
|
||||||
type PublishContext = {
|
|
||||||
instance?: InstanceContext
|
|
||||||
workspace?: WorkspaceID
|
|
||||||
}
|
|
||||||
|
|
||||||
export interface Interface {
|
export interface Interface {
|
||||||
readonly run: <Def extends Definition>(
|
readonly run: <Def extends Definition>(
|
||||||
|
|
@ -107,16 +101,14 @@ export const layer = Layer.effect(Service)(
|
||||||
}
|
}
|
||||||
|
|
||||||
const publish = !!options?.publish
|
const publish = !!options?.publish
|
||||||
const context = publish
|
// Bridge captures handler-fiber refs (InstanceRef/WorkspaceRef) and the
|
||||||
? {
|
// full Effect context, so the forked publish + GlobalBus emit run with
|
||||||
instance: yield* InstanceState.context,
|
// the right state without a per-call attachWith.
|
||||||
workspace: yield* InstanceState.workspaceID,
|
const bridge = yield* EffectBridge.make()
|
||||||
}
|
|
||||||
: undefined
|
|
||||||
process(def, event, {
|
process(def, event, {
|
||||||
bus,
|
bus,
|
||||||
|
bridge,
|
||||||
publish,
|
publish,
|
||||||
context,
|
|
||||||
ownerID: options?.ownerID,
|
ownerID: options?.ownerID,
|
||||||
experimentalWorkspaces: flags.experimentalWorkspaces,
|
experimentalWorkspaces: flags.experimentalWorkspaces,
|
||||||
})
|
})
|
||||||
|
|
@ -154,12 +146,7 @@ export const layer = Layer.effect(Service)(
|
||||||
}
|
}
|
||||||
|
|
||||||
const { publish = true } = options || {}
|
const { publish = true } = options || {}
|
||||||
const context = publish
|
const bridge = yield* EffectBridge.make()
|
||||||
? {
|
|
||||||
instance: yield* InstanceState.context,
|
|
||||||
workspace: yield* InstanceState.workspaceID,
|
|
||||||
}
|
|
||||||
: undefined
|
|
||||||
|
|
||||||
// Note that this is an "immediate" transaction which is critical.
|
// Note that this is an "immediate" transaction which is critical.
|
||||||
// We need to make sure we can safely read and write with nothing
|
// We need to make sure we can safely read and write with nothing
|
||||||
|
|
@ -175,7 +162,7 @@ export const layer = Layer.effect(Service)(
|
||||||
const seq = row?.seq != null ? row.seq + 1 : 0
|
const seq = row?.seq != null ? row.seq + 1 : 0
|
||||||
|
|
||||||
const event = { id, seq, aggregateID: agg, data }
|
const event = { id, seq, aggregateID: agg, data }
|
||||||
process(def, event, { bus, publish, context, experimentalWorkspaces: flags.experimentalWorkspaces })
|
process(def, event, { bus, bridge, publish, experimentalWorkspaces: flags.experimentalWorkspaces })
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
behavior: "immediate",
|
behavior: "immediate",
|
||||||
|
|
@ -308,8 +295,8 @@ function process<Def extends Definition>(
|
||||||
event: Event<Def>,
|
event: Event<Def>,
|
||||||
options: {
|
options: {
|
||||||
bus: ProjectBus.Interface
|
bus: ProjectBus.Interface
|
||||||
|
bridge: EffectBridge.Shape
|
||||||
publish: boolean
|
publish: boolean
|
||||||
context?: PublishContext
|
|
||||||
ownerID?: string
|
ownerID?: string
|
||||||
experimentalWorkspaces: boolean
|
experimentalWorkspaces: boolean
|
||||||
},
|
},
|
||||||
|
|
@ -351,37 +338,36 @@ function process<Def extends Definition>(
|
||||||
}
|
}
|
||||||
|
|
||||||
Database.effect(() => {
|
Database.effect(() => {
|
||||||
if (options?.publish) {
|
if (!options.publish) return
|
||||||
if (!options.context?.instance) {
|
const result = convertEvent(def.type, event.data)
|
||||||
throw new Error("SyncEvent.process: publish requires instance context")
|
// The bridge was built inside the caller's fiber so it already carries
|
||||||
}
|
// InstanceRef/WorkspaceRef and the full Effect context. Both the bus
|
||||||
|
// publish and the GlobalBus emit run inside the forked Effect so they
|
||||||
const result = convertEvent(def.type, event.data)
|
// share the same instance/workspace lookup.
|
||||||
const publish = (data: unknown) =>
|
const publish = (data: unknown) =>
|
||||||
Effect.runPromise(
|
options.bridge.fork(
|
||||||
attachWith(options.bus.publish(def, data as Properties<Def>, { id: event.id }), {
|
Effect.gen(function* () {
|
||||||
instance: options.context?.instance,
|
yield* options.bus.publish(def, data as Properties<Def>, { id: event.id })
|
||||||
workspace: options.context?.workspace,
|
const instance = yield* InstanceState.context
|
||||||
}),
|
const workspace = yield* InstanceState.workspaceID
|
||||||
)
|
GlobalBus.emit("event", {
|
||||||
if (result instanceof Promise) {
|
directory: instance.directory,
|
||||||
void result.then(publish)
|
project: instance.project.id,
|
||||||
} else {
|
workspace,
|
||||||
void publish(result)
|
payload: {
|
||||||
}
|
type: "sync",
|
||||||
|
syncEvent: {
|
||||||
GlobalBus.emit("event", {
|
type: versionedType(def.type, def.version),
|
||||||
directory: options.context.instance.directory,
|
...event,
|
||||||
project: options.context.instance.project.id,
|
},
|
||||||
workspace: options.context.workspace,
|
},
|
||||||
payload: {
|
})
|
||||||
type: "sync",
|
}),
|
||||||
syncEvent: {
|
)
|
||||||
type: versionedType(def.type, def.version),
|
if (result instanceof Promise) {
|
||||||
...event,
|
void result.then(publish)
|
||||||
},
|
} else {
|
||||||
},
|
publish(result)
|
||||||
})
|
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
|
||||||
|
|
@ -1,14 +1,15 @@
|
||||||
import { describe, expect, beforeEach, afterAll } from "bun:test"
|
import { describe, expect, beforeEach, afterAll } from "bun:test"
|
||||||
import { provideTmpdirInstance } from "../fixture/fixture"
|
import { provideTmpdirInstance } from "../fixture/fixture"
|
||||||
import { Effect, Layer, Schema } from "effect"
|
import { Deferred, Effect, Layer, Schema } from "effect"
|
||||||
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
|
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
|
||||||
import { Bus } from "../../src/bus"
|
import { Bus } from "../../src/bus"
|
||||||
|
import { GlobalBus, type GlobalEvent } from "../../src/bus/global"
|
||||||
import { SyncEvent } from "../../src/sync"
|
import { SyncEvent } from "../../src/sync"
|
||||||
import { Database, eq } from "@/storage/db"
|
import { Database, eq } from "@/storage/db"
|
||||||
import { EventSequenceTable, EventTable } from "../../src/sync/event.sql"
|
import { EventSequenceTable, EventTable } from "../../src/sync/event.sql"
|
||||||
import { MessageID } from "../../src/session/schema"
|
import { MessageID } from "../../src/session/schema"
|
||||||
import { initProjectors } from "../../src/server/projectors"
|
import { initProjectors } from "../../src/server/projectors"
|
||||||
import { testEffect } from "../lib/effect"
|
import { awaitWithTimeout, testEffect } from "../lib/effect"
|
||||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||||
|
|
||||||
const it = testEffect(
|
const it = testEffect(
|
||||||
|
|
@ -139,6 +140,43 @@ describe("SyncEvent", () => {
|
||||||
}),
|
}),
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// Regression for the EffectBridge migration. GlobalBus.emit used to fire
|
||||||
|
// synchronously inside the Database.effect post-commit callback. After the
|
||||||
|
// migration it fires inside the forked publish Effect, AFTER bus.publish
|
||||||
|
// completes. Consumers don't care about microsecond-level ordering, but
|
||||||
|
// we still need to prove the emit actually fires.
|
||||||
|
it.live(
|
||||||
|
"emits sync events to GlobalBus after publishing to ProjectBus",
|
||||||
|
provideTmpdirInstance(() =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const { Created } = setup()
|
||||||
|
// Filter for OUR specific event in the handler so we ignore any
|
||||||
|
// stray sync events from other tests' lingering forks.
|
||||||
|
const received = yield* Deferred.make<GlobalEvent>()
|
||||||
|
const handler = (evt: GlobalEvent) => {
|
||||||
|
if (evt.payload?.type === "sync" && evt.payload?.syncEvent?.type === "item.created.1") {
|
||||||
|
Deferred.doneUnsafe(received, Effect.succeed(evt))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
GlobalBus.on("event", handler)
|
||||||
|
try {
|
||||||
|
yield* SyncEvent.use.run(Created, { id: "evt_global_1", name: "global" })
|
||||||
|
const event = yield* awaitWithTimeout(
|
||||||
|
Deferred.await(received),
|
||||||
|
"timed out waiting for sync event on GlobalBus",
|
||||||
|
"2 seconds",
|
||||||
|
)
|
||||||
|
expect(event.payload).toMatchObject({
|
||||||
|
type: "sync",
|
||||||
|
syncEvent: { type: "item.created.1", data: { id: "evt_global_1", name: "global" } },
|
||||||
|
})
|
||||||
|
} finally {
|
||||||
|
GlobalBus.off("event", handler)
|
||||||
|
}
|
||||||
|
}),
|
||||||
|
),
|
||||||
|
)
|
||||||
})
|
})
|
||||||
|
|
||||||
describe("replay", () => {
|
describe("replay", () => {
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue