Compare commits

...

4 Commits

Author SHA1 Message Date
Kit Langton 1c2b2e682b Remove covered workspace websocket todo 2026-04-30 16:46:43 -04:00
Kit Langton 4f5af93e44 Simplify SyncEvent service access 2026-04-30 16:39:15 -04:00
Kit Langton bdefdc2306 Make SyncEvent layer canonical 2026-04-30 16:28:30 -04:00
Kit Langton ec98b656fd Add SyncEvent service 2026-04-30 16:12:01 -04:00
12 changed files with 431 additions and 353 deletions
@@ -169,6 +169,7 @@ export const layer = Layer.effect(
const auth = yield* Auth.Service const auth = yield* Auth.Service
const session = yield* Session.Service const session = yield* Session.Service
const http = yield* HttpClient.HttpClient const http = yield* HttpClient.HttpClient
const sync = yield* SyncEvent.Service
const connections = new Map<WorkspaceID, ConnectionStatus>() const connections = new Map<WorkspaceID, ConnectionStatus>()
const syncFibers = yield* FiberMap.make<WorkspaceID, void, SyncLoopError>() const syncFibers = yield* FiberMap.make<WorkspaceID, void, SyncLoopError>()
@@ -307,25 +308,30 @@ export const layer = Layer.effect(
events: events.length, events: events.length,
}) })
yield* Effect.sync(() => yield* Effect.promise(async () => {
WorkspaceContext.provide({ await WorkspaceContext.provide({
workspaceID: space.id, workspaceID: space.id,
fn: () => { async fn() {
for (const event of events) { await Effect.runPromise(
SyncEvent.replay( Effect.forEach(
{ events,
id: event.id, (event) =>
aggregateID: event.aggregate_id, sync.replay(
seq: event.seq, {
type: event.type, id: event.id,
data: event.data, aggregateID: event.aggregate_id,
}, seq: event.seq,
{ publish: true }, type: event.type,
) data: event.data,
} },
{ publish: true },
),
{ discard: true },
),
)
}, },
}), })
) })
}) })
const syncWorkspaceLoop = Effect.fn("Workspace.syncWorkspaceLoop")(function* (space: Info) { const syncWorkspaceLoop = Effect.fn("Workspace.syncWorkspaceLoop")(function* (space: Info) {
@@ -361,16 +367,28 @@ export const layer = Layer.effect(
setStatus(space.id, "connected") setStatus(space.id, "connected")
yield* parseSSE(stream, (evt) => yield* parseSSE(stream, (evt) =>
Effect.sync(() => { Effect.gen(function* () {
if (!evt || typeof evt !== "object" || !("payload" in evt)) return
const payload = evt.payload as { type?: string; syncEvent?: SyncEvent.SerializedEvent }
if (payload.type === "server.heartbeat") return
if (payload.type === "sync" && payload.syncEvent) {
const failed = yield* sync.replay(payload.syncEvent).pipe(
Effect.as(false),
Effect.catchCause((error) =>
Effect.sync(() => {
log.info("failed to replay global event", {
workspaceID: space.id,
error,
})
return true
}),
),
)
if (failed) return
}
try { try {
if (!evt || typeof evt !== "object" || !("payload" in evt)) return
const payload = evt.payload as { type?: string; syncEvent?: SyncEvent.SerializedEvent }
if (payload.type === "server.heartbeat") return
if (payload.type === "sync" && payload.syncEvent) {
SyncEvent.replay(payload.syncEvent)
}
const event = evt as { directory?: string; project?: string; payload: unknown } const event = evt as { directory?: string; project?: string; payload: unknown }
GlobalBus.emit("event", { GlobalBus.emit("event", {
directory: event.directory, directory: event.directory,
@@ -378,10 +396,10 @@ export const layer = Layer.effect(
workspace: space.id, workspace: space.id,
payload: event.payload, payload: event.payload,
}) })
} catch (err) { } catch (error) {
log.info("failed to replay global event", { log.info("failed to replay global event", {
workspaceID: space.id, workspaceID: space.id,
error: err, error,
}) })
} }
}), }),
@@ -516,14 +534,12 @@ export const layer = Layer.effect(
const adaptor = getAdaptor(space.projectID, space.type) const adaptor = getAdaptor(space.projectID, space.type)
const target = yield* Effect.promise(() => Promise.resolve(adaptor.target(space))) const target = yield* Effect.promise(() => Promise.resolve(adaptor.target(space)))
yield* Effect.sync(() => yield* sync.run(Session.Event.Updated, {
SyncEvent.run(Session.Event.Updated, { sessionID: input.sessionID,
sessionID: input.sessionID, info: {
info: { workspaceID: input.workspaceID,
workspaceID: input.workspaceID, },
}, })
}),
)
const rows = yield* db((db) => const rows = yield* db((db) =>
db db
@@ -593,7 +609,7 @@ export const layer = Layer.effect(
}) })
if (target.type === "local") { if (target.type === "local") {
SyncEvent.replayAll(events) yield* sync.replayAll(events)
log.info("session restore batch replayed locally", { log.info("session restore batch replayed locally", {
workspaceID: input.workspaceID, workspaceID: input.workspaceID,
sessionID: input.sessionID, sessionID: input.sessionID,
@@ -812,6 +828,7 @@ export const layer = Layer.effect(
export const defaultLayer = layer.pipe( export const defaultLayer = layer.pipe(
Layer.provide(Auth.defaultLayer), Layer.provide(Auth.defaultLayer),
Layer.provide(Session.defaultLayer), Layer.provide(Session.defaultLayer),
Layer.provide(SyncEvent.defaultLayer),
Layer.provide(FetchHttpClient.layer), Layer.provide(FetchHttpClient.layer),
) )
@@ -47,6 +47,7 @@ import { Pty } from "@/pty"
import { Installation } from "@/installation" import { Installation } from "@/installation"
import { ShareNext } from "@/share/share-next" import { ShareNext } from "@/share/share-next"
import { SessionShare } from "@/share/session" import { SessionShare } from "@/share/session"
import { SyncEvent } from "@/sync"
import { Npm } from "@opencode-ai/core/npm" import { Npm } from "@opencode-ai/core/npm"
import { memoMap } from "@opencode-ai/core/effect/memo-map" import { memoMap } from "@opencode-ai/core/effect/memo-map"
@@ -97,6 +98,7 @@ export const AppLayer = Layer.mergeAll(
Installation.defaultLayer, Installation.defaultLayer,
ShareNext.defaultLayer, ShareNext.defaultLayer,
SessionShare.defaultLayer, SessionShare.defaultLayer,
SyncEvent.defaultLayer,
).pipe(Layer.provideMerge(Observability.layer)) ).pipe(Layer.provideMerge(Observability.layer))
const rt = ManagedRuntime.make(AppLayer, { memoMap }) const rt = ManagedRuntime.make(AppLayer, { memoMap })
@@ -21,6 +21,7 @@ export const syncHandlers = HttpApiBuilder.group(InstanceHttpApi, "sync", (handl
Effect.gen(function* () { Effect.gen(function* () {
const workspace = yield* Workspace.Service const workspace = yield* Workspace.Service
const scope = yield* Scope.Scope const scope = yield* Scope.Scope
const sync = yield* SyncEvent.Service
const start = Effect.fn("SyncHttpApi.start")(function* () { const start = Effect.fn("SyncHttpApi.start")(function* () {
yield* workspace yield* workspace
@@ -45,7 +46,7 @@ export const syncHandlers = HttpApiBuilder.group(InstanceHttpApi, "sync", (handl
last: events.at(-1)?.seq, last: events.at(-1)?.seq,
directory: ctx.payload.directory, directory: ctx.payload.directory,
}) })
SyncEvent.replayAll(events) yield* sync.replayAll(events)
log.info("sync replay complete", { log.info("sync replay complete", {
sessionID: source, sessionID: source,
events: events.length, events: events.length,
@@ -31,6 +31,7 @@ import { SessionSummary } from "@/session/summary"
import { Todo } from "@/session/todo" import { Todo } from "@/session/todo"
import { SessionShare } from "@/share/session" import { SessionShare } from "@/share/session"
import { Skill } from "@/skill" import { Skill } from "@/skill"
import { SyncEvent } from "@/sync"
import { ToolRegistry } from "@/tool/registry" import { ToolRegistry } from "@/tool/registry"
import { lazy } from "@/util/lazy" import { lazy } from "@/util/lazy"
import { Vcs } from "@/project/vcs" import { Vcs } from "@/project/vcs"
@@ -147,6 +148,7 @@ export const routes = Layer.mergeAll(rootApiRoutes, instanceRoutes).pipe(
SessionRunState.defaultLayer, SessionRunState.defaultLayer,
SessionStatus.defaultLayer, SessionStatus.defaultLayer,
SessionSummary.defaultLayer, SessionSummary.defaultLayer,
SyncEvent.defaultLayer,
Skill.defaultLayer, Skill.defaultLayer,
Todo.defaultLayer, Todo.defaultLayer,
ToolRegistry.defaultLayer, ToolRegistry.defaultLayer,
@@ -94,7 +94,7 @@ export const SyncRoutes = lazy(() =>
last: events.at(-1)?.seq, last: events.at(-1)?.seq,
directory: body.directory, directory: body.directory,
}) })
SyncEvent.replayAll(events) await AppRuntime.runPromise(SyncEvent.use.replayAll(events))
log.info("sync replay complete", { log.info("sync replay complete", {
sessionID: source, sessionID: source,
+4 -2
View File
@@ -38,6 +38,7 @@ export const layer = Layer.effect(
const bus = yield* Bus.Service const bus = yield* Bus.Service
const summary = yield* SessionSummary.Service const summary = yield* SessionSummary.Service
const state = yield* SessionRunState.Service const state = yield* SessionRunState.Service
const sync = yield* SyncEvent.Service
const revert = Effect.fn("SessionRevert.revert")(function* (input: RevertInput) { const revert = Effect.fn("SessionRevert.revert")(function* (input: RevertInput) {
yield* state.assertNotBusy(input.sessionID) yield* state.assertNotBusy(input.sessionID)
@@ -121,7 +122,7 @@ export const layer = Layer.effect(
remove.push(msg) remove.push(msg)
} }
for (const msg of remove) { for (const msg of remove) {
SyncEvent.run(MessageV2.Event.Removed, { yield* sync.run(MessageV2.Event.Removed, {
sessionID, sessionID,
messageID: msg.info.id, messageID: msg.info.id,
}) })
@@ -133,7 +134,7 @@ export const layer = Layer.effect(
const removeParts = target.parts.slice(idx) const removeParts = target.parts.slice(idx)
target.parts = target.parts.slice(0, idx) target.parts = target.parts.slice(0, idx)
for (const part of removeParts) { for (const part of removeParts) {
SyncEvent.run(MessageV2.Event.PartRemoved, { yield* sync.run(MessageV2.Event.PartRemoved, {
sessionID, sessionID,
messageID: target.info.id, messageID: target.info.id,
partID: part.id, partID: part.id,
@@ -156,6 +157,7 @@ export const defaultLayer = Layer.suspend(() =>
Layer.provide(Storage.defaultLayer), Layer.provide(Storage.defaultLayer),
Layer.provide(Bus.layer), Layer.provide(Bus.layer),
Layer.provide(SessionSummary.defaultLayer), Layer.provide(SessionSummary.defaultLayer),
Layer.provide(SyncEvent.defaultLayer),
), ),
) )
+26 -30
View File
@@ -443,11 +443,12 @@ export type Patch = Types.DeepMutable<SyncEvent.Event<typeof Event.Updated>["dat
const db = <T>(fn: (d: Parameters<typeof Database.use>[0] extends (trx: infer D) => any ? D : never) => T) => const db = <T>(fn: (d: Parameters<typeof Database.use>[0] extends (trx: infer D) => any ? D : never) => T) =>
Effect.sync(() => Database.use(fn)) Effect.sync(() => Database.use(fn))
export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> = Layer.effect( export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service | SyncEvent.Service> = Layer.effect(
Service, Service,
Effect.gen(function* () { Effect.gen(function* () {
const bus = yield* Bus.Service const bus = yield* Bus.Service
const storage = yield* Storage.Service const storage = yield* Storage.Service
const sync = yield* SyncEvent.Service
const createNext = Effect.fn("Session.createNext")(function* (input: { const createNext = Effect.fn("Session.createNext")(function* (input: {
id?: SessionID id?: SessionID
@@ -477,7 +478,7 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
} }
log.info("created", result) log.info("created", result)
yield* Effect.sync(() => SyncEvent.run(Event.Created, { sessionID: result.id, info: result })) yield* sync.run(Event.Created, { sessionID: result.id, info: result })
if (!Flag.OPENCODE_EXPERIMENTAL_WORKSPACES) { if (!Flag.OPENCODE_EXPERIMENTAL_WORKSPACES) {
// This only exist for backwards compatibility. We should not be // This only exist for backwards compatibility. We should not be
@@ -525,10 +526,8 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
Effect.catchCause(() => Effect.succeed(false)), Effect.catchCause(() => Effect.succeed(false)),
) )
yield* Effect.sync(() => { yield* sync.run(Event.Deleted, { sessionID, info: session }, { publish: hasInstance })
SyncEvent.run(Event.Deleted, { sessionID, info: session }, { publish: hasInstance }) yield* sync.remove(sessionID)
SyncEvent.remove(sessionID)
})
} catch (e) { } catch (e) {
log.error(e) log.error(e)
} }
@@ -536,19 +535,17 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
const updateMessage = <T extends MessageV2.Info>(msg: T): Effect.Effect<T> => const updateMessage = <T extends MessageV2.Info>(msg: T): Effect.Effect<T> =>
Effect.gen(function* () { Effect.gen(function* () {
yield* Effect.sync(() => SyncEvent.run(MessageV2.Event.Updated, { sessionID: msg.sessionID, info: msg })) yield* sync.run(MessageV2.Event.Updated, { sessionID: msg.sessionID, info: msg })
return msg return msg
}).pipe(Effect.withSpan("Session.updateMessage")) }).pipe(Effect.withSpan("Session.updateMessage"))
const updatePart = <T extends MessageV2.Part>(part: T): Effect.Effect<T> => const updatePart = <T extends MessageV2.Part>(part: T): Effect.Effect<T> =>
Effect.gen(function* () { Effect.gen(function* () {
yield* Effect.sync(() => yield* sync.run(MessageV2.Event.PartUpdated, {
SyncEvent.run(MessageV2.Event.PartUpdated, { sessionID: part.sessionID,
sessionID: part.sessionID, part: structuredClone(part),
part: structuredClone(part), time: Date.now(),
time: Date.now(), })
}),
)
return part return part
}).pipe(Effect.withSpan("Session.updatePart")) }).pipe(Effect.withSpan("Session.updatePart"))
@@ -635,8 +632,7 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
return session return session
}) })
const patch = (sessionID: SessionID, info: Patch) => const patch = (sessionID: SessionID, info: Patch) => sync.run(Event.Updated, { sessionID, info })
Effect.sync(() => SyncEvent.run(Event.Updated, { sessionID, info }))
const touch = Effect.fn("Session.touch")(function* (sessionID: SessionID) { const touch = Effect.fn("Session.touch")(function* (sessionID: SessionID) {
yield* patch(sessionID, { time: { updated: Date.now() } }) yield* patch(sessionID, { time: { updated: Date.now() } })
@@ -693,12 +689,10 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
sessionID: SessionID sessionID: SessionID
messageID: MessageID messageID: MessageID
}) { }) {
yield* Effect.sync(() => yield* sync.run(MessageV2.Event.Removed, {
SyncEvent.run(MessageV2.Event.Removed, { sessionID: input.sessionID,
sessionID: input.sessionID, messageID: input.messageID,
messageID: input.messageID, })
}),
)
return input.messageID return input.messageID
}) })
@@ -707,13 +701,11 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
messageID: MessageID messageID: MessageID
partID: PartID partID: PartID
}) { }) {
yield* Effect.sync(() => yield* sync.run(MessageV2.Event.PartRemoved, {
SyncEvent.run(MessageV2.Event.PartRemoved, { sessionID: input.sessionID,
sessionID: input.sessionID, messageID: input.messageID,
messageID: input.messageID, partID: input.partID,
partID: input.partID, })
}),
)
return input.partID return input.partID
}) })
@@ -764,7 +756,11 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
}), }),
) )
export const defaultLayer = layer.pipe(Layer.provide(Bus.layer), Layer.provide(Storage.defaultLayer)) export const defaultLayer = layer.pipe(
Layer.provide(Bus.layer),
Layer.provide(Storage.defaultLayer),
Layer.provide(SyncEvent.defaultLayer),
)
export function* list(input?: { export function* list(input?: {
directory?: string directory?: string
+4 -4
View File
@@ -21,20 +21,19 @@ export const layer = Layer.effect(
const session = yield* Session.Service const session = yield* Session.Service
const shareNext = yield* ShareNext.Service const shareNext = yield* ShareNext.Service
const scope = yield* Scope.Scope const scope = yield* Scope.Scope
const sync = yield* SyncEvent.Service
const share = Effect.fn("SessionShare.share")(function* (sessionID: SessionID) { const share = Effect.fn("SessionShare.share")(function* (sessionID: SessionID) {
const conf = yield* cfg.get() const conf = yield* cfg.get()
if (conf.share === "disabled") throw new Error("Sharing is disabled in configuration") if (conf.share === "disabled") throw new Error("Sharing is disabled in configuration")
const result = yield* shareNext.create(sessionID) const result = yield* shareNext.create(sessionID)
yield* Effect.sync(() => yield* sync.run(Session.Event.Updated, { sessionID, info: { share: { url: result.url } } })
SyncEvent.run(Session.Event.Updated, { sessionID, info: { share: { url: result.url } } }),
)
return result return result
}) })
const unshare = Effect.fn("SessionShare.unshare")(function* (sessionID: SessionID) { const unshare = Effect.fn("SessionShare.unshare")(function* (sessionID: SessionID) {
yield* shareNext.remove(sessionID) yield* shareNext.remove(sessionID)
yield* Effect.sync(() => SyncEvent.run(Session.Event.Updated, { sessionID, info: { share: { url: null } } })) yield* sync.run(Session.Event.Updated, { sessionID, info: { share: { url: null } } })
}) })
const create = Effect.fn("SessionShare.create")(function* (input?: Session.CreateInput) { const create = Effect.fn("SessionShare.create")(function* (input?: Session.CreateInput) {
@@ -54,6 +53,7 @@ export const defaultLayer = layer.pipe(
Layer.provide(ShareNext.defaultLayer), Layer.provide(ShareNext.defaultLayer),
Layer.provide(Session.defaultLayer), Layer.provide(Session.defaultLayer),
Layer.provide(Config.defaultLayer), Layer.provide(Config.defaultLayer),
Layer.provide(SyncEvent.defaultLayer),
) )
export * as SessionShare from "./session" export * as SessionShare from "./session"
+126 -78
View File
@@ -9,9 +9,11 @@ import { EventSequenceTable, EventTable } from "./event.sql"
import { WorkspaceContext } from "@/control-plane/workspace-context" import { WorkspaceContext } from "@/control-plane/workspace-context"
import { EventID } from "./schema" import { EventID } from "./schema"
import { Flag } from "@opencode-ai/core/flag/flag" import { Flag } from "@opencode-ai/core/flag/flag"
import { Schema as EffectSchema } from "effect" import { Context, Effect, Layer, Schema as EffectSchema } from "effect"
import { zodObject } from "@/util/effect-zod" import { zodObject } from "@/util/effect-zod"
import type { DeepMutable } from "@/util/schema" import type { DeepMutable } from "@/util/schema"
import { makeRuntime } from "@/effect/run-service"
import { serviceUse } from "@/effect/service-use"
// 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 —
@@ -46,6 +48,125 @@ export type SerializedEvent<Def extends Definition = Definition> = Event<Def> &
type ProjectorFunc = (db: Database.TxOrDb, data: unknown) => void type ProjectorFunc = (db: Database.TxOrDb, data: unknown) => void
type ConvertEvent = (type: string, data: Event["data"]) => unknown | Promise<unknown> type ConvertEvent = (type: string, data: Event["data"]) => unknown | Promise<unknown>
export interface Interface {
readonly run: <Def extends Definition>(
def: Def,
data: Event<Def>["data"],
options?: { publish?: boolean },
) => Effect.Effect<void>
readonly replay: (event: SerializedEvent, options?: { publish: boolean }) => Effect.Effect<void>
readonly replayAll: (events: SerializedEvent[], options?: { publish: boolean }) => Effect.Effect<string | undefined>
readonly remove: (aggregateID: string) => Effect.Effect<void>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/SyncEvent") {}
export const layer = Layer.effect(Service)(
Effect.gen(function* () {
const replay: Interface["replay"] = Effect.fn("SyncEvent.replay")(function* (event, options) {
const def = registry.get(event.type)
if (!def) {
throw new Error(`Unknown event type: ${event.type}`)
}
const row = Database.use((db) =>
db
.select({ seq: EventSequenceTable.seq })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, event.aggregateID))
.get(),
)
const latest = row?.seq ?? -1
if (event.seq <= latest) return
const expected = latest + 1
if (event.seq !== expected) {
throw new Error(
`Sequence mismatch for aggregate "${event.aggregateID}": expected ${expected}, got ${event.seq}`,
)
}
process(def, event, { publish: !!options?.publish })
})
const replayAll: Interface["replayAll"] = Effect.fn("SyncEvent.replayAll")(function* (events, options) {
const source = events[0]?.aggregateID
if (!source) return undefined
if (events.some((item) => item.aggregateID !== source)) {
throw new Error("Replay events must belong to the same session")
}
const start = events[0].seq
for (const [i, item] of events.entries()) {
const seq = start + i
if (item.seq !== seq) {
throw new Error(`Replay sequence mismatch at index ${i}: expected ${seq}, got ${item.seq}`)
}
}
for (const item of events) {
yield* replay(item, options)
}
return source
})
const run: Interface["run"] = Effect.fn("SyncEvent.run")(function* (def, data, options) {
const agg = (data as Record<string, string>)[def.aggregate]
// This should never happen: we've enforced it via typescript in
// the definition
if (agg == null) {
throw new Error(`SyncEvent.run: "${def.aggregate}" required but not found: ${JSON.stringify(data)}`)
}
if (def.version !== versions.get(def.type)) {
throw new Error(`SyncEvent.run: running old versions of events is not allowed: ${def.type}`)
}
const { publish = true } = options || {}
// Note that this is an "immediate" transaction which is critical.
// We need to make sure we can safely read and write with nothing
// else changing the data from under us
Database.transaction(
(tx) => {
const id = EventID.ascending()
const row = tx
.select({ seq: EventSequenceTable.seq })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, agg))
.get()
const seq = row?.seq != null ? row.seq + 1 : 0
const event = { id, seq, aggregateID: agg, data }
process(def, event, { publish })
},
{
behavior: "immediate",
},
)
})
const remove: Interface["remove"] = Effect.fn("SyncEvent.remove")(function* (aggregateID) {
Database.transaction((tx) => {
tx.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).run()
tx.delete(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).run()
})
})
return Service.of({
run,
replay,
replayAll,
remove,
})
}),
)
export const defaultLayer = layer
export const use = serviceUse(Service)
const runtime = makeRuntime(Service, defaultLayer)
export const registry = new Map<string, Definition>() export const registry = new Map<string, Definition>()
let projectors: Map<Definition, ProjectorFunc> | undefined let projectors: Map<Definition, ProjectorFunc> | undefined
const versions = new Map<string, number>() const versions = new Map<string, number>()
@@ -186,92 +307,19 @@ function process<Def extends Definition>(def: Def, event: Event<Def>, options: {
} }
export function replay(event: SerializedEvent, options?: { publish: boolean }) { export function replay(event: SerializedEvent, options?: { publish: boolean }) {
const def = registry.get(event.type) return runtime.runSync((sync) => sync.replay(event, options))
if (!def) {
throw new Error(`Unknown event type: ${event.type}`)
}
const row = Database.use((db) =>
db
.select({ seq: EventSequenceTable.seq })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, event.aggregateID))
.get(),
)
const latest = row?.seq ?? -1
if (event.seq <= latest) {
return
}
const expected = latest + 1
if (event.seq !== expected) {
throw new Error(`Sequence mismatch for aggregate "${event.aggregateID}": expected ${expected}, got ${event.seq}`)
}
process(def, event, { publish: !!options?.publish })
} }
export function replayAll(events: SerializedEvent[], options?: { publish: boolean }) { export function replayAll(events: SerializedEvent[], options?: { publish: boolean }) {
const source = events[0]?.aggregateID return runtime.runSync((sync) => sync.replayAll(events, options))
if (!source) return
if (events.some((item) => item.aggregateID !== source)) {
throw new Error("Replay events must belong to the same session")
}
const start = events[0].seq
for (const [i, item] of events.entries()) {
const seq = start + i
if (item.seq !== seq) {
throw new Error(`Replay sequence mismatch at index ${i}: expected ${seq}, got ${item.seq}`)
}
}
for (const item of events) {
replay(item, options)
}
return source
} }
export function run<Def extends Definition>(def: Def, data: Event<Def>["data"], options?: { publish?: boolean }) { export function run<Def extends Definition>(def: Def, data: Event<Def>["data"], options?: { publish?: boolean }) {
const agg = (data as Record<string, string>)[def.aggregate] return runtime.runSync((sync) => sync.run(def, data, options))
// This should never happen: we've enforced it via typescript in
// the definition
if (agg == null) {
throw new Error(`SyncEvent.run: "${def.aggregate}" required but not found: ${JSON.stringify(data)}`)
}
if (def.version !== versions.get(def.type)) {
throw new Error(`SyncEvent.run: running old versions of events is not allowed: ${def.type}`)
}
const { publish = true } = options || {}
// Note that this is an "immediate" transaction which is critical.
// We need to make sure we can safely read and write with nothing
// else changing the data from under us
Database.transaction(
(tx) => {
const id = EventID.ascending()
const row = tx
.select({ seq: EventSequenceTable.seq })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, agg))
.get()
const seq = row?.seq != null ? row.seq + 1 : 0
const event = { id, seq, aggregateID: agg, data }
process(def, event, { publish })
},
{
behavior: "immediate",
},
)
} }
export function remove(aggregateID: string) { export function remove(aggregateID: string) {
Database.transaction((tx) => { return runtime.runSync((sync) => sync.remove(aggregateID))
tx.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).run()
tx.delete(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).run()
})
} }
export function payloads() { export function payloads() {
@@ -1,4 +1,4 @@
import { afterEach, beforeEach, describe, expect, mock, spyOn, test } from "bun:test" import { afterEach, beforeEach, describe, expect, mock, test } from "bun:test"
import fs from "node:fs/promises" import fs from "node:fs/promises"
import Http from "node:http" import Http from "node:http"
import path from "node:path" import path from "node:path"
@@ -1426,48 +1426,41 @@ describe("workspace-old sessionRestore", () => {
}) })
}) })
test("local restore replays batches without fetch and emits progress", async () => { it.live("local restore replays batches and emits progress", () =>
await withInstance(async (dir) => { provideTmpdirInstance(
const captured = captureGlobalEvents() (dir) =>
let fetchCallCount = 0 Effect.gen(function* () {
const replayAll = spyOn(SyncEvent, "replayAll") const workspace = yield* WorkspaceOld.Service
try { const sessionSvc = yield* SessionNs.Service
using server = Bun.serve({ const captured = captureGlobalEvents()
port: 0, try {
fetch() { const type = unique("restore-local")
fetchCallCount++ const info = workspaceInfo(Instance.project.id, type, { directory: dir })
return Response.json({ ok: true }) insertWorkspace(info)
}, registerAdaptor(Instance.project.id, type, localAdaptor(dir).adaptor)
}) const session = yield* sessionSvc.create({ title: "restore local" })
const type = unique("restore-local") replaceSessionEvents(session.id, 20)
const info = workspaceInfo(Instance.project.id, type, { directory: dir })
insertWorkspace(info)
registerAdaptor(Instance.project.id, type, localAdaptor(dir).adaptor)
const session = await AppRuntime.runPromise(
SessionNs.Service.use((svc) => svc.create({ title: "restore local" })),
)
replaceSessionEvents(session.id, 20)
expect(await restoreWorkspaceSession({ workspaceID: info.id, sessionID: session.id })).toEqual({ total: 3 }) expect(yield* workspace.sessionRestore({ workspaceID: info.id, sessionID: session.id })).toEqual({
total: 3,
expect(fetchCallCount).toBe(0) })
expect(replayAll).toHaveBeenCalledTimes(3) expect((yield* sessionSvc.get(session.id)).workspaceID).toBe(info.id)
expect(replayAll.mock.calls.map((call) => call[0].length)).toEqual([10, 10, 1]) expect(eventRows(session.id).map((row) => row.seq)).toEqual(Array.from({ length: 21 }, (_, i) => i))
expect((await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.get(session.id)))).workspaceID).toBe( expect(
info.id, captured.events
) .filter(
expect(eventRows(session.id).map((row) => row.seq)).toEqual(Array.from({ length: 21 }, (_, i) => i)) (event) => event.workspace === info.id && event.payload.type === WorkspaceOld.Event.Restore.type,
expect( )
captured.events .map((event) => event.payload.properties.step),
.filter((event) => event.workspace === info.id && event.payload.type === WorkspaceOld.Event.Restore.type) ).toEqual([0, 1, 2, 3])
.map((event) => event.payload.properties.step), yield* workspace.remove(info.id)
).toEqual([0, 1, 2, 3]) } finally {
await removeWorkspace(info.id) captured.dispose()
} finally { }
captured.dispose() }),
} { git: true },
}) ),
}) )
it.live("session restore includes real message and part events in sequence order", () => { it.live("session restore includes real message and part events in sequence order", () => {
const replay: FetchCall[] = [] const replay: FetchCall[] = []
@@ -1,4 +1,4 @@
import { afterEach, describe, expect, mock, test } from "bun:test" import { afterEach, describe, expect, mock } from "bun:test"
import { NodeServices } from "@effect/platform-node" import { NodeServices } from "@effect/platform-node"
import { mkdir } from "node:fs/promises" import { mkdir } from "node:fs/promises"
import path from "node:path" import path from "node:path"
@@ -133,8 +133,6 @@ afterEach(async () => {
}) })
describe("workspace HttpApi", () => { describe("workspace HttpApi", () => {
test.todo("proxies remote workspace websocket through real Effect listener", () => {})
it.live("serves read endpoints", () => it.live("serves read endpoints", () =>
Effect.gen(function* () { Effect.gen(function* () {
const dir = yield* tmpdirScoped({ git: true }) const dir = yield* tmpdirScoped({ git: true })
+174 -155
View File
@@ -1,16 +1,18 @@
import { describe, test, expect, beforeEach, afterEach, afterAll } from "bun:test" import { describe, expect, beforeEach, afterEach, afterAll } from "bun:test"
import { tmpdir } from "../fixture/fixture" import { provideTmpdirInstance } from "../fixture/fixture"
import { Schema } from "effect" import { Effect, Layer, Schema } from "effect"
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
import { Bus } from "../../src/bus" import { Bus } from "../../src/bus"
import { Instance } from "../../src/project/instance"
import { SyncEvent } from "../../src/sync" import { SyncEvent } from "../../src/sync"
import { Database } from "@/storage/db" import { Database } from "@/storage/db"
import { EventTable } from "../../src/sync/event.sql" import { EventTable } from "../../src/sync/event.sql"
import { Identifier } from "../../src/id/id" import { MessageID } from "../../src/session/schema"
import { Flag } from "@opencode-ai/core/flag/flag" import { Flag } from "@opencode-ai/core/flag/flag"
import { initProjectors } from "../../src/server/projectors" import { initProjectors } from "../../src/server/projectors"
import { testEffect } from "../lib/effect"
const original = Flag.OPENCODE_EXPERIMENTAL_WORKSPACES const original = Flag.OPENCODE_EXPERIMENTAL_WORKSPACES
const it = testEffect(Layer.mergeAll(SyncEvent.defaultLayer, CrossSpawnSpawner.defaultLayer))
beforeEach(() => { beforeEach(() => {
Database.close() Database.close()
@@ -22,19 +24,6 @@ afterEach(() => {
Flag.OPENCODE_EXPERIMENTAL_WORKSPACES = original Flag.OPENCODE_EXPERIMENTAL_WORKSPACES = original
}) })
function withInstance(fn: () => void | Promise<void>) {
return async () => {
await using tmp = await tmpdir()
await Instance.provide({
directory: tmp.path,
fn: async () => {
await fn()
},
})
}
}
describe("SyncEvent", () => { describe("SyncEvent", () => {
function setup() { function setup() {
SyncEvent.reset() SyncEvent.reset()
@@ -59,179 +48,209 @@ describe("SyncEvent", () => {
return { Created, Sent } return { Created, Sent }
} }
function expectDefect<A, E, R>(effect: Effect.Effect<A, E, R>, pattern: RegExp) {
return Effect.gen(function* () {
const exit = yield* Effect.exit(effect)
if (exit._tag === "Success") throw new Error("Expected effect to fail")
expect(String(exit.cause)).toMatch(pattern)
})
}
afterAll(() => { afterAll(() => {
SyncEvent.reset() SyncEvent.reset()
initProjectors() initProjectors()
}) })
describe("run", () => { describe("run", () => {
test( it.live(
"inserts event row", "inserts event row",
withInstance(() => { provideTmpdirInstance(() =>
const { Created } = setup() Effect.gen(function* () {
SyncEvent.run(Created, { id: "evt_1", name: "first" }) const { Created } = setup()
const rows = Database.use((db) => db.select().from(EventTable).all()) yield* SyncEvent.use.run(Created, { id: "evt_1", name: "first" })
expect(rows).toHaveLength(1) const rows = Database.use((db) => db.select().from(EventTable).all())
expect(rows[0].type).toBe("item.created.1") expect(rows).toHaveLength(1)
expect(rows[0].aggregate_id).toBe("evt_1") expect(rows[0].type).toBe("item.created.1")
}), expect(rows[0].aggregate_id).toBe("evt_1")
}),
),
) )
test( it.live(
"increments seq per aggregate", "increments seq per aggregate",
withInstance(() => { provideTmpdirInstance(() =>
const { Created } = setup() Effect.gen(function* () {
SyncEvent.run(Created, { id: "evt_1", name: "first" }) const { Created } = setup()
SyncEvent.run(Created, { id: "evt_1", name: "second" }) yield* SyncEvent.use.run(Created, { id: "evt_1", name: "first" })
const rows = Database.use((db) => db.select().from(EventTable).all()) yield* SyncEvent.use.run(Created, { id: "evt_1", name: "second" })
expect(rows).toHaveLength(2) const rows = Database.use((db) => db.select().from(EventTable).all())
expect(rows[1].seq).toBe(rows[0].seq + 1) expect(rows).toHaveLength(2)
}), expect(rows[1].seq).toBe(rows[0].seq + 1)
}),
),
) )
test( it.live(
"uses custom aggregate field from agg()", "uses custom aggregate field from agg()",
withInstance(() => { provideTmpdirInstance(() =>
const { Sent } = setup() Effect.gen(function* () {
SyncEvent.run(Sent, { item_id: "evt_1", to: "james" }) const { Sent } = setup()
const rows = Database.use((db) => db.select().from(EventTable).all()) yield* SyncEvent.use.run(Sent, { item_id: "evt_1", to: "james" })
expect(rows).toHaveLength(1) const rows = Database.use((db) => db.select().from(EventTable).all())
expect(rows[0].aggregate_id).toBe("evt_1") expect(rows).toHaveLength(1)
}), expect(rows[0].aggregate_id).toBe("evt_1")
}),
),
) )
test( it.live(
"emits events", "emits events",
withInstance(async () => { provideTmpdirInstance(() =>
const { Created } = setup() Effect.gen(function* () {
const events: Array<{ const { Created } = setup()
type: string const events: Array<{
properties: { id: string; name: string } type: string
}> = [] properties: { id: string; name: string }
const received = new Promise<void>((resolve) => { }> = []
Bus.subscribeAll((event) => { let resolve = () => {}
const received = new Promise<void>((done) => {
resolve = done
})
const dispose = Bus.subscribeAll((event) => {
events.push(event) events.push(event)
resolve() resolve()
}) })
}) try {
yield* SyncEvent.use.run(Created, { id: "evt_1", name: "test" })
SyncEvent.run(Created, { id: "evt_1", name: "test" }) yield* Effect.promise(() => received)
expect(events).toHaveLength(1)
await received expect(events[0]).toEqual({
expect(events).toHaveLength(1) type: "item.created",
expect(events[0]).toEqual({ properties: {
type: "item.created", id: "evt_1",
properties: { name: "test",
id: "evt_1", },
name: "test", })
}, } finally {
}) dispose()
}), }
}),
),
) )
}) })
describe("replay", () => { describe("replay", () => {
test( it.live(
"inserts event from external payload", "inserts event from external payload",
withInstance(() => { provideTmpdirInstance(() =>
const id = Identifier.descending("message") Effect.gen(function* () {
SyncEvent.replay({ const id = MessageID.ascending()
id: "evt_1", yield* SyncEvent.use.replay({
type: "item.created.1",
seq: 0,
aggregateID: id,
data: { id, name: "replayed" },
})
const rows = Database.use((db) => db.select().from(EventTable).all())
expect(rows).toHaveLength(1)
expect(rows[0].aggregate_id).toBe(id)
}),
)
test(
"throws on sequence mismatch",
withInstance(() => {
const id = Identifier.descending("message")
SyncEvent.replay({
id: "evt_1",
type: "item.created.1",
seq: 0,
aggregateID: id,
data: { id, name: "first" },
})
expect(() =>
SyncEvent.replay({
id: "evt_1", id: "evt_1",
type: "item.created.1", type: "item.created.1",
seq: 5,
aggregateID: id,
data: { id, name: "bad" },
}),
).toThrow(/Sequence mismatch/)
}),
)
test(
"throws on unknown event type",
withInstance(() => {
expect(() =>
SyncEvent.replay({
id: "evt_1",
type: "unknown.event.1",
seq: 0, seq: 0,
aggregateID: "x", aggregateID: id,
data: {}, data: { id, name: "replayed" },
}), })
).toThrow(/Unknown event type/) const rows = Database.use((db) => db.select().from(EventTable).all())
}), expect(rows).toHaveLength(1)
expect(rows[0].aggregate_id).toBe(id)
}),
),
) )
test( it.live(
"replayAll accepts later chunks after the first batch", "throws on sequence mismatch",
withInstance(() => { provideTmpdirInstance(() =>
const { Created } = setup() Effect.gen(function* () {
const id = Identifier.descending("message") const id = MessageID.ascending()
yield* SyncEvent.use.replay({
const one = SyncEvent.replayAll([
{
id: "evt_1", id: "evt_1",
type: SyncEvent.versionedType(Created.type, Created.version), type: "item.created.1",
seq: 0, seq: 0,
aggregateID: id, aggregateID: id,
data: { id, name: "first" }, data: { id, name: "first" },
}, })
{ yield* expectDefect(
id: "evt_2", SyncEvent.use.replay({
type: SyncEvent.versionedType(Created.type, Created.version), id: "evt_1",
seq: 1, type: "item.created.1",
aggregateID: id, seq: 5,
data: { id, name: "second" }, aggregateID: id,
}, data: { id, name: "bad" },
]) }),
/Sequence mismatch/,
)
}),
),
)
const two = SyncEvent.replayAll([ it.live(
{ "throws on unknown event type",
id: "evt_3", provideTmpdirInstance(() =>
type: SyncEvent.versionedType(Created.type, Created.version), Effect.gen(function* () {
seq: 2, yield* expectDefect(
aggregateID: id, SyncEvent.use.replay({
data: { id, name: "third" }, id: "evt_1",
}, type: "unknown.event.1",
{ seq: 0,
id: "evt_4", aggregateID: "x",
type: SyncEvent.versionedType(Created.type, Created.version), data: {},
seq: 3, }),
aggregateID: id, /Unknown event type/,
data: { id, name: "fourth" }, )
}, }),
]) ),
)
expect(one).toBe(id) it.live(
expect(two).toBe(id) "replayAll accepts later chunks after the first batch",
provideTmpdirInstance(() =>
Effect.gen(function* () {
const { Created } = setup()
const id = MessageID.ascending()
const rows = Database.use((db) => db.select().from(EventTable).all()) const one = yield* SyncEvent.use.replayAll([
expect(rows.map((row) => row.seq)).toEqual([0, 1, 2, 3]) {
}), id: "evt_1",
type: SyncEvent.versionedType(Created.type, Created.version),
seq: 0,
aggregateID: id,
data: { id, name: "first" },
},
{
id: "evt_2",
type: SyncEvent.versionedType(Created.type, Created.version),
seq: 1,
aggregateID: id,
data: { id, name: "second" },
},
])
const two = yield* SyncEvent.use.replayAll([
{
id: "evt_3",
type: SyncEvent.versionedType(Created.type, Created.version),
seq: 2,
aggregateID: id,
data: { id, name: "third" },
},
{
id: "evt_4",
type: SyncEvent.versionedType(Created.type, Created.version),
seq: 3,
aggregateID: id,
data: { id, name: "fourth" },
},
])
expect(one).toBe(id)
expect(two).toBe(id)
const rows = Database.use((db) => db.select().from(EventTable).all())
expect(rows.map((row) => row.seq)).toEqual([0, 1, 2, 3])
}),
),
) )
}) })
}) })