mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-06 17:19:49 -04:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1c2b2e682b | |||
| 4f5af93e44 | |||
| bdefdc2306 | |||
| ec98b656fd |
@@ -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,
|
||||||
|
|||||||
@@ -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),
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
@@ -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 })
|
||||||
|
|||||||
@@ -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])
|
||||||
|
}),
|
||||||
|
),
|
||||||
)
|
)
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
Reference in New Issue
Block a user