mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-24 14:43:37 -04:00
feat(core): Session.snapshot
This commit is contained in:
@@ -53,6 +53,7 @@ import { Shell as ShellSchema } from "@opencode-ai/schema/shell"
|
||||
import { KeyedMutex } from "./effect/keyed-mutex.js"
|
||||
import { fileURLToPath } from "url"
|
||||
import { SessionEnvironment } from "./session/environment.js"
|
||||
import { EventSequenceTable } from "./event/sql.js"
|
||||
|
||||
// get project -> project.locations
|
||||
//
|
||||
@@ -194,6 +195,19 @@ export interface Interface {
|
||||
* unhandled compaction barriers.
|
||||
*/
|
||||
readonly inbox: (sessionID: SessionSchema.ID) => Effect.Effect<SessionInbox.Info[], NotFoundError>
|
||||
readonly snapshot: (input: {
|
||||
sessionID: SessionSchema.ID
|
||||
recent?: number
|
||||
}) => Effect.Effect<
|
||||
{
|
||||
readonly session: SessionSchema.Info
|
||||
readonly children: SessionSchema.Info[]
|
||||
readonly inbox: SessionInbox.Info[]
|
||||
readonly messages: SessionMessage.Info[]
|
||||
readonly seq: Event.Seq
|
||||
},
|
||||
NotFoundError | MessageDecodeError
|
||||
>
|
||||
readonly cancelInbox: (input: InboxItemRef) => Effect.Effect<void, NotFoundError | InboxConflictError>
|
||||
readonly steerInbox: (input: InboxItemRef) => Effect.Effect<void, NotFoundError | InboxConflictError>
|
||||
readonly queueInbox: (input: InboxItemRef) => Effect.Effect<void, NotFoundError | InboxConflictError>
|
||||
@@ -557,6 +571,51 @@ const layer = Layer.effect(
|
||||
yield* result.get(sessionID)
|
||||
return yield* SessionInbox.list(db, sessionID)
|
||||
}),
|
||||
snapshot: Effect.fn("Session.snapshot")(function* (input) {
|
||||
return yield* db
|
||||
.transaction(() =>
|
||||
Effect.gen(function* () {
|
||||
const row = yield* db
|
||||
.select()
|
||||
.from(SessionTable)
|
||||
.where(eq(SessionTable.id, input.sessionID))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (!row) return yield* new NotFoundError({ sessionID: input.sessionID })
|
||||
const children = yield* db
|
||||
.select()
|
||||
.from(SessionTable)
|
||||
.where(eq(SessionTable.parent_id, input.sessionID))
|
||||
.orderBy(desc(SessionTable.time_updated), desc(SessionTable.id))
|
||||
.all()
|
||||
.pipe(Effect.orDie)
|
||||
const inbox = yield* SessionInbox.list(db, input.sessionID)
|
||||
const messages = yield* db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(eq(SessionMessageTable.session_id, input.sessionID))
|
||||
.orderBy(desc(SessionMessageTable.seq))
|
||||
.limit(input.recent ?? 200)
|
||||
.all()
|
||||
.pipe(Effect.orDie)
|
||||
const sequence = yield* db
|
||||
.select({ seq: EventSequenceTable.seq })
|
||||
.from(EventSequenceTable)
|
||||
.where(eq(EventSequenceTable.aggregate_id, input.sessionID))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (!sequence) return yield* Effect.die(new Error(`Session ${input.sessionID} has no event sequence`))
|
||||
return {
|
||||
session: fromRow(row),
|
||||
children: children.map(fromRow),
|
||||
inbox,
|
||||
messages: yield* Effect.forEach(messages.toReversed(), decode),
|
||||
seq: Event.Seq.make(sequence.seq),
|
||||
}
|
||||
}),
|
||||
)
|
||||
.pipe(Effect.catchTag("SqlError", Effect.die))
|
||||
}),
|
||||
cancelInbox: Effect.fn("Session.cancelInbox")((input) => mutatePending(input, SessionInbox.cancel)),
|
||||
steerInbox: Effect.fn("Session.steerInbox")((input) => mutatePending(input, SessionInbox.steer, true)),
|
||||
queueInbox: Effect.fn("Session.queueInbox")((input) => mutatePending(input, SessionInbox.queue)),
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Effect } from "effect"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
|
||||
import { Bus } from "@opencode-ai/core/bus"
|
||||
import { Location } from "@opencode-ai/core/location"
|
||||
import { Project } from "@opencode-ai/core/project"
|
||||
import { AbsolutePath } from "@opencode-ai/core/schema"
|
||||
import { Session } from "@opencode-ai/core/session"
|
||||
import { SessionProjector } from "@opencode-ai/core/session/projector"
|
||||
import { SessionExecution } from "@opencode-ai/core/session/execution"
|
||||
import { SessionStore } from "@opencode-ai/core/session/store"
|
||||
import { SessionEvent } from "@opencode-ai/core/session/event"
|
||||
import { Event } from "@opencode-ai/schema/event"
|
||||
import { testEffect } from "./lib/effect"
|
||||
import { globalProjectLayer } from "./lib/project"
|
||||
|
||||
const it = testEffect(
|
||||
AppNodeBuilder.build(
|
||||
LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node, Session.node]),
|
||||
[
|
||||
[Bus.node, Bus.configured({ persist: true })],
|
||||
[Project.node, globalProjectLayer],
|
||||
[SessionExecution.node, SessionExecution.noopLayer],
|
||||
],
|
||||
),
|
||||
)
|
||||
const location = Location.Ref.make({ directory: AbsolutePath.make("/project") })
|
||||
|
||||
describe("Session.snapshot", () => {
|
||||
it.effect("returns an empty projected session at its aggregate watermark", () =>
|
||||
Effect.gen(function* () {
|
||||
const sessions = yield* Session.Service
|
||||
const created = yield* sessions.create({ location })
|
||||
|
||||
expect(yield* sessions.snapshot({ sessionID: created.id })).toEqual({
|
||||
session: created,
|
||||
children: [],
|
||||
inbox: [],
|
||||
messages: [],
|
||||
seq: Event.Seq.make(0),
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("returns the most recent messages in aggregate order", () =>
|
||||
Effect.gen(function* () {
|
||||
const sessions = yield* Session.Service
|
||||
const bus = yield* Bus.Service
|
||||
const created = yield* sessions.create({ location })
|
||||
yield* Effect.forEach(["first", "second", "third"], (text) =>
|
||||
bus.publish(SessionEvent.Synthetic, { sessionID: created.id, text }),
|
||||
)
|
||||
|
||||
const snapshot = yield* sessions.snapshot({ sessionID: created.id, recent: 2 })
|
||||
|
||||
expect(snapshot.messages.map((message) => (message.type === "synthetic" ? message.text : message.type))).toEqual([
|
||||
"second",
|
||||
"third",
|
||||
])
|
||||
expect(snapshot.seq).toBe(Event.Seq.make(3))
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("keeps rows and watermark consistent during concurrent publication", () =>
|
||||
Effect.gen(function* () {
|
||||
const sessions = yield* Session.Service
|
||||
const bus = yield* Bus.Service
|
||||
const created = yield* sessions.create({ location })
|
||||
const publish = Effect.forEach(
|
||||
Array.from({ length: 40 }, (_, index) => index + 1),
|
||||
(index) => bus.publish(SessionEvent.Synthetic, { sessionID: created.id, text: String(index) }),
|
||||
)
|
||||
const read = Effect.forEach(Array.from({ length: 40 }), () => sessions.snapshot({ sessionID: created.id }))
|
||||
|
||||
const [, snapshots] = yield* Effect.all([publish, read], { concurrency: "unbounded" })
|
||||
|
||||
snapshots.forEach((snapshot) => {
|
||||
expect(snapshot.messages).toHaveLength(snapshot.seq)
|
||||
expect(
|
||||
snapshot.messages.map((message) => (message.type === "synthetic" ? Number(message.text) : -1)),
|
||||
).toEqual(Array.from({ length: snapshot.seq }, (_, index) => index + 1))
|
||||
})
|
||||
}),
|
||||
)
|
||||
})
|
||||
Reference in New Issue
Block a user