Compare commits

...

1 Commits

Author SHA1 Message Date
James Long b50cd706ce fix(core): gate durable event persistence 2026-08-04 23:35:14 +00:00
4 changed files with 32 additions and 18 deletions
+21 -15
View File
@@ -165,6 +165,7 @@ export const allBounded = (events: Interface, capacity: number) =>
export interface LayerOptions {
readonly beforeAggregateRead?: (aggregateID: string) => Effect.Effect<void>
readonly persistDurableEvents?: boolean
}
export const layerWith = (options?: LayerOptions) =>
@@ -180,6 +181,7 @@ export const layerWith = (options?: LayerOptions) =>
// TODO: Bind durable projectors to exact type+version before supporting incompatible historical payloads.
const listeners = new Array<Subscriber>()
const { db } = yield* Database.Service
const persistDurableEvents = options?.persistDurableEvents ?? true
const getOrCreate = (definition: Definition) =>
Effect.gen(function* () {
@@ -333,19 +335,20 @@ export const layerWith = (options?: LayerOptions) =>
})
.run()
.pipe(Effect.orDie)
yield* db
.insert(EventTable)
.values([
{
id: event.id,
aggregate_id: aggregateID,
seq,
type: versionedType(definition.type, durable.version),
data: encoded,
},
])
.run()
.pipe(Effect.orDie)
if (persistDurableEvents)
yield* db
.insert(EventTable)
.values([
{
id: event.id,
aggregate_id: aggregateID,
seq,
type: versionedType(definition.type, durable.version),
data: encoded,
},
])
.run()
.pipe(Effect.orDie)
return { aggregateID, seq }
}),
{ behavior: "immediate" },
@@ -634,5 +637,8 @@ export const layerWith = (options?: LayerOptions) =>
}),
)
const layer = layerWith()
export const node = makeGlobalNode({ service: Service, layer: layer, deps: [Database.node] })
export const nodeWith = (options?: LayerOptions) =>
makeGlobalNode({ service: Service, layer: layerWith(options), deps: [Database.node] })
export const node = nodeWith()
export const sequenceOnlyNode = nodeWith({ persistDurableEvents: false })
@@ -51,6 +51,7 @@ import { memoMap } from "@opencode-ai/core/effect/memo-map"
import { BackgroundJob } from "@/background/job"
import { RuntimeFlags } from "@/effect/runtime-flags"
import { EventV2Bridge } from "@/event-v2-bridge"
import { EventV2 } from "@opencode-ai/core/event"
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { AppNodeBuilderV1 } from "./app-node-builder-v1"
import { SessionProjector } from "@opencode-ai/core/session/projector"
@@ -106,6 +107,7 @@ export const AppLayer = AppNodeBuilderV1.build(
ShareNext.node,
SessionShare.node,
]),
[[EventV2.node, EventV2.sequenceOnlyNode]],
).pipe(Layer.provideMerge(AppNodeBuilderV1.build(Ripgrep.node)), Layer.provideMerge(Observability.layer))
const rt = ManagedRuntime.make(AppLayer, { memoMap })
@@ -270,8 +270,10 @@ const app = LayerNode.group([
export function createRoutes(
corsOptions?: CorsOptions,
persistDurableEvents = false,
): Layer.Layer<never, EffectConfig.ConfigError, RouteRequirements> {
const locationServiceMapV2 = buildLocationServiceMap()
const eventNode = persistDurableEvents ? EventV2.node : EventV2.sequenceOnlyNode
return Layer.mergeAll(
rootApiRoutes,
@@ -288,7 +290,10 @@ export function createRoutes(
corsVaryFix,
fenceLayer,
cors(corsOptions),
AppNodeBuilderV1.build(MoveSession.node, [[LocationServiceMap.node, locationServiceMapV2]]),
AppNodeBuilderV1.build(MoveSession.node, [
[LocationServiceMap.node, locationServiceMapV2],
[EventV2.node, eventNode],
]),
HttpServer.layerServices,
]),
Layer.provide(Layer.succeed(CorsConfig)(corsOptions)),
@@ -299,11 +304,12 @@ export function createRoutes(
AppNodeBuilderV1.build(SessionV2.node, [
[LocationServiceMap.node, locationServiceMapV2],
[SessionExecution.node, SessionExecutionLocal.node],
[EventV2.node, eventNode],
]),
),
Layer.provide(locationServiceMapV2),
Layer.provide(AppNodeBuilderV1.build(app)),
Layer.provide(AppNodeBuilderV1.build(app, [[EventV2.node, eventNode]])),
// Must stay last: layers provided later in this pipe build beneath earlier ones,
// so Observability must come after every service graph. Otherwise eagerly forked
// fibers (e.g. the ModelsDev background refresh) capture Effect's default stdout
@@ -48,7 +48,7 @@ const appLayer = AppNodeBuilder.build(
[[InstanceStore.bootstrapNode, noopBootstrapLayer]],
)
const servedRoutes: Layer.Layer<never, Config.ConfigError, HttpServer.HttpServer> = HttpRouter.serve(
HttpApiApp.routes,
HttpApiApp.createRoutes(undefined, true),
{
disableListenLog: true,
disableLogger: true,