| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574 |
- import { describe, expect } from "bun:test"
- import { DateTime, Effect, Fiber, Layer, Stream } from "effect"
- import { eq } from "drizzle-orm"
- import { Database } from "@opencode-ai/core/database/database"
- import { EventV2 } from "@opencode-ai/core/event"
- import { EventTable } from "@opencode-ai/core/event/sql"
- import { SessionEvent } from "@opencode-ai/core/session/event"
- import { Project } from "@opencode-ai/core/project"
- import { ProjectTable } from "@opencode-ai/core/project/sql"
- import { AbsolutePath } from "@opencode-ai/core/schema"
- import { SessionV2 } from "@opencode-ai/core/session"
- import { LocationServiceMap } from "@opencode-ai/core/location-layer"
- import { Prompt } from "@opencode-ai/core/session/prompt"
- import { SessionMessage } from "@opencode-ai/core/session/message"
- import { SessionProjector } from "@opencode-ai/core/session/projector"
- import { SessionExecution } from "@opencode-ai/core/session/execution"
- import { SessionInput } from "@opencode-ai/core/session/input"
- import { SessionInputTable, SessionMessageTable, SessionTable } from "@opencode-ai/core/session/sql"
- import { SessionStore } from "@opencode-ai/core/session/store"
- import { testEffect } from "./lib/effect"
- const executionCalls: SessionV2.ID[] = []
- const interruptCalls: SessionV2.ID[] = []
- const wakeCalls: SessionV2.ID[] = []
- const activeSessions = new Set<SessionV2.ID>()
- const execution = Layer.succeed(
- SessionExecution.Service,
- SessionExecution.Service.of({
- active: Effect.sync(() => new Set(activeSessions)),
- resume: (sessionID) =>
- Effect.sync(() => {
- executionCalls.push(sessionID)
- }),
- interrupt: (sessionID) =>
- Effect.sync(() => {
- interruptCalls.push(sessionID)
- }),
- wake: (sessionID) =>
- Effect.sync(() => {
- wakeCalls.push(sessionID)
- }),
- }),
- )
- const sessions = SessionV2.layer.pipe(
- Layer.provide(LocationServiceMap.layer),
- Layer.provide(EventV2.defaultLayer),
- Layer.provide(Database.defaultLayer),
- Layer.provide(SessionStore.defaultLayer),
- Layer.provide(Project.defaultLayer),
- Layer.provide(execution),
- )
- const it = testEffect(
- Layer.mergeAll(
- Database.defaultLayer,
- EventV2.defaultLayer,
- SessionProjector.defaultLayer,
- SessionStore.defaultLayer,
- execution,
- sessions,
- ),
- )
- const sessionID = SessionV2.ID.make("ses_prompt_test")
- const messageID = SessionMessage.ID.create()
- const setup = Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .onConflictDoNothing()
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .onConflictDoNothing()
- .run()
- .pipe(Effect.orDie)
- })
- const admitted = (id: SessionMessage.ID) => Database.Service.use(({ db }) => SessionInput.find(db, id))
- const admittedCount = Database.Service.use(({ db }) =>
- db
- .select()
- .from(SessionInputTable)
- .all()
- .pipe(
- Effect.orDie,
- Effect.map((rows) => rows.length),
- ),
- )
- const eventCount = (type: string) =>
- Database.Service.use(({ db }) =>
- db
- .select()
- .from(EventTable)
- .where(eq(EventTable.type, type))
- .all()
- .pipe(
- Effect.orDie,
- Effect.map((rows) => rows.length),
- ),
- )
- describe("SessionV2.prompt", () => {
- it.effect("exposes the execution registry", () =>
- Effect.gen(function* () {
- activeSessions.add(sessionID)
- expect(Array.from(yield* (yield* SessionV2.Service).active)).toEqual([sessionID])
- }).pipe(Effect.ensuring(Effect.sync(() => activeSessions.clear()))),
- )
- it.effect("delegates execution continuation through SessionExecution", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- executionCalls.length = 0
- wakeCalls.length = 0
- yield* session.resume(sessionID)
- expect(executionCalls).toEqual([sessionID])
- expect(wakeCalls).toEqual([])
- }),
- )
- it.effect("delegates process-local interruption through SessionExecution", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- interruptCalls.length = 0
- yield* session.interrupt(sessionID)
- expect(interruptCalls).toEqual([sessionID])
- expect(yield* session.messages({ sessionID })).toEqual([])
- }),
- )
- it.effect("delegates interruption without requiring a recorded Session", () =>
- Effect.gen(function* () {
- const session = yield* SessionV2.Service
- interruptCalls.length = 0
- yield* session.interrupt(SessionV2.ID.make("ses_missing"))
- expect(interruptCalls).toEqual([SessionV2.ID.make("ses_missing")])
- }),
- )
- it.effect("durably admits one user message before transcript promotion", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const message = yield* session.prompt({
- sessionID,
- prompt: Prompt.make({ text: "Fix the failing tests" }),
- resume: false,
- })
- expect(message.prompt.text).toBe("Fix the failing tests")
- expect(yield* session.messages({ sessionID })).toEqual([])
- expect(yield* admitted(message.id)).toMatchObject({
- id: message.id,
- sessionID,
- prompt: { text: "Fix the failing tests" },
- delivery: "steer",
- })
- }),
- )
- it.effect("streams durable Session events after an aggregate sequence", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const events = yield* EventV2.Service
- const { db } = yield* Database.Service
- const fiber = yield* session.events({ sessionID }).pipe(Stream.take(4), Stream.runCollect, Effect.forkScoped)
- yield* Effect.yieldNow
- yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "First" }), resume: false })
- yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Second" }), resume: false })
- yield* SessionInput.promoteSteers(db, events, sessionID, Number.MAX_SAFE_INTEGER)
- const streamed = Array.from(yield* Fiber.join(fiber))
- expect(streamed.map((event) => [event.durable?.seq, event.type])).toEqual([
- [0, "session.next.prompt.admitted"],
- [1, "session.next.prompt.admitted"],
- [2, "session.next.prompted"],
- [3, "session.next.prompted"],
- ])
- expect(
- Array.from(
- yield* session
- .events({ sessionID, after: streamed[0]!.durable?.seq })
- .pipe(Stream.take(1), Stream.runCollect),
- ).map((event) => [event.durable?.seq, event.type]),
- ).toEqual([[1, "session.next.prompt.admitted"]])
- }),
- )
- it.effect("resumes through a recorded message without appending another prompt", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const message = yield* session.prompt({
- sessionID,
- prompt: Prompt.make({ text: "Fix the failing tests" }),
- resume: false,
- })
- executionCalls.length = 0
- wakeCalls.length = 0
- yield* session.resume(sessionID)
- expect(yield* session.messages({ sessionID })).toEqual([])
- expect(yield* admitted(message.id)).not.toHaveProperty("promotedSeq")
- expect(executionCalls).toEqual([sessionID])
- expect(wakeCalls).toEqual([])
- }),
- )
- it.effect("records distinct messages when the ID is omitted", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const input = { sessionID, prompt: Prompt.make({ text: "Fix the failing tests" }), resume: false }
- const first = yield* session.prompt(input)
- const second = yield* session.prompt(input)
- expect(second.id).not.toBe(first.id)
- expect(yield* session.messages({ sessionID })).toEqual([])
- expect(yield* admittedCount).toBe(2)
- }),
- )
- it.effect("returns the original recorded message when the ID is retried", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const input = {
- sessionID,
- id: messageID,
- prompt: Prompt.make({ text: "Fix the failing tests" }),
- resume: false,
- }
- const first = yield* session.prompt(input)
- const retried = yield* session.prompt(input)
- expect(retried).toEqual(first)
- expect(yield* session.messages({ sessionID })).toEqual([])
- expect(yield* admittedCount).toBe(1)
- }),
- )
- it.effect("wakes execution when an exact prompt retry recovers a committed message", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const input = {
- sessionID,
- id: messageID,
- prompt: Prompt.make({ text: "Recover committed prompt" }),
- resume: false,
- }
- const first = yield* session.prompt(input)
- wakeCalls.length = 0
- const retried = yield* session.prompt({ ...input, resume: true })
- expect(retried).toEqual(first)
- expect(wakeCalls).toEqual([sessionID])
- }),
- )
- it.effect("rejects reuse of one ID with a different prompt", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- yield* session.prompt({
- sessionID,
- id: messageID,
- prompt: Prompt.make({ text: "Fix the failing tests" }),
- })
- const failure = yield* session
- .prompt({
- sessionID,
- id: messageID,
- prompt: Prompt.make({ text: "Delete the failing tests" }),
- resume: false,
- })
- .pipe(Effect.flip)
- expect(failure._tag).toBe("Session.PromptConflictError")
- expect(yield* session.messages({ sessionID })).toHaveLength(0)
- expect(yield* admittedCount).toBe(1)
- }),
- )
- it.effect("rejects reuse of one ID with a different delivery mode", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- yield* session.prompt({
- id: messageID,
- sessionID,
- prompt: Prompt.make({ text: "Fix the failing tests" }),
- resume: false,
- })
- const failure = yield* session
- .prompt({
- id: messageID,
- sessionID,
- prompt: Prompt.make({ text: "Fix the failing tests" }),
- delivery: "queue",
- resume: false,
- })
- .pipe(Effect.flip)
- expect(failure._tag).toBe("Session.PromptConflictError")
- }),
- )
- it.effect("returns one recorded message to concurrent exact retries", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const input = {
- sessionID,
- id: messageID,
- prompt: Prompt.make({ text: "Fix the failing tests" }),
- resume: false,
- }
- const messages = yield* Effect.all([session.prompt(input), session.prompt(input)], { concurrency: "unbounded" })
- expect(messages[1]).toEqual(messages[0])
- expect(yield* session.messages({ sessionID })).toEqual([])
- expect(yield* admittedCount).toBe(1)
- expect(yield* eventCount(EventV2.versionedType(SessionEvent.PromptAdmitted.type, 1))).toBe(1)
- }),
- )
- it.effect("promotes one message once under concurrent promotion attempts", () =>
- Effect.gen(function* () {
- yield* setup
- const { db } = yield* Database.Service
- const session = yield* SessionV2.Service
- const events = yield* EventV2.Service
- yield* session.prompt({ id: messageID, sessionID, prompt: Prompt.make({ text: "Promote once" }), resume: false })
- yield* Effect.all(
- [
- SessionInput.promoteSteers(db, events, sessionID, Number.MAX_SAFE_INTEGER),
- SessionInput.promoteSteers(db, events, sessionID, Number.MAX_SAFE_INTEGER),
- ],
- { concurrency: "unbounded" },
- )
- expect(yield* eventCount(EventV2.versionedType(SessionEvent.Prompted.type, 1))).toBe(1)
- expect(yield* admitted(messageID)).toMatchObject({ promotedSeq: 1 })
- expect(yield* session.messages({ sessionID })).toMatchObject([
- { id: messageID, type: "user", text: "Promote once" },
- ])
- }),
- )
- it.effect("promotes steers only through the captured inbox cutoff", () =>
- Effect.gen(function* () {
- yield* setup
- const { db } = yield* Database.Service
- const session = yield* SessionV2.Service
- const events = yield* EventV2.Service
- const first = yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Before cutoff" }), resume: false })
- const cutoff = first.admittedSeq
- const second = yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "After cutoff" }), resume: false })
- yield* SessionInput.promoteSteers(db, events, sessionID, cutoff)
- expect(yield* admitted(first.id)).toHaveProperty("promotedSeq")
- expect(yield* admitted(second.id)).not.toHaveProperty("promotedSeq")
- }),
- )
- it.effect("reprojects pending inbox input without scheduling execution", () =>
- Effect.gen(function* () {
- yield* setup
- const { db } = yield* Database.Service
- const session = yield* SessionV2.Service
- const events = yield* EventV2.Service
- wakeCalls.length = 0
- yield* session.prompt({
- id: messageID,
- sessionID,
- prompt: Prompt.make({ text: "Replay pending" }),
- resume: false,
- })
- const recorded = yield* db
- .select()
- .from(EventTable)
- .where(eq(EventTable.aggregate_id, sessionID))
- .all()
- .pipe(Effect.orDie)
- yield* events.remove(sessionID)
- yield* db.delete(SessionInputTable).where(eq(SessionInputTable.session_id, sessionID)).run().pipe(Effect.orDie)
- yield* db
- .delete(SessionMessageTable)
- .where(eq(SessionMessageTable.session_id, sessionID))
- .run()
- .pipe(Effect.orDie)
- yield* events.replayAll(
- recorded.map((event) => ({
- id: event.id,
- aggregateID: event.aggregate_id,
- seq: event.seq,
- type: event.type,
- data: event.data,
- })),
- )
- expect(yield* admitted(messageID)).toMatchObject({ id: messageID, prompt: { text: "Replay pending" } })
- expect(yield* session.messages({ sessionID })).toEqual([])
- expect(wakeCalls).toEqual([])
- }),
- )
- it.effect("returns an exact retry of a legacy projected prompt", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const events = yield* EventV2.Service
- const prompt = Prompt.make({ text: "Historical prompt" })
- yield* events.publish(SessionEvent.Prompted, {
- sessionID,
- messageID,
- timestamp: yield* DateTime.now,
- prompt,
- delivery: "steer",
- })
- const retried = yield* session.prompt({ id: messageID, sessionID, prompt, resume: false })
- expect(retried).toMatchObject({ id: messageID, prompt: { text: "Historical prompt" } })
- expect(yield* admitted(messageID)).toHaveProperty("promotedSeq")
- }),
- )
- it.effect("returns an exact retry of a legacy projected queued prompt", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const events = yield* EventV2.Service
- const prompt = Prompt.make({ text: "Historical queued prompt" })
- yield* events.publish(SessionEvent.Prompted, {
- sessionID,
- messageID,
- timestamp: yield* DateTime.now,
- prompt,
- delivery: "queue",
- })
- const retried = yield* session.prompt({ id: messageID, sessionID, prompt, delivery: "queue", resume: false })
- expect(retried).toMatchObject({ id: messageID, prompt: { text: "Historical queued prompt" } })
- expect(yield* admitted(messageID)).toMatchObject({ delivery: "queue" })
- }),
- )
- it.effect("rejects reuse of one globally unique message ID across sessions", () =>
- Effect.gen(function* () {
- yield* setup
- const { db } = yield* Database.Service
- const session = yield* SessionV2.Service
- const other = SessionV2.ID.make("ses_prompt_other")
- yield* db
- .insert(SessionTable)
- .values({
- id: other,
- project_id: Project.ID.global,
- slug: "other",
- directory: "/project",
- title: "other",
- version: "test",
- })
- .onConflictDoNothing()
- .run()
- .pipe(Effect.orDie)
- const prompt = Prompt.make({ text: "Fix the failing tests" })
- yield* session.prompt({ id: messageID, sessionID, prompt, resume: false })
- const failure = yield* session
- .prompt({ id: messageID, sessionID: other, prompt, resume: false })
- .pipe(Effect.flip)
- expect(failure).toMatchObject({ _tag: "Session.PromptConflictError", sessionID: other, messageID })
- }),
- )
- it.effect("rejects a prompt ID already used by visible Session history", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- const events = yield* EventV2.Service
- yield* events.publish(SessionEvent.Synthetic, {
- sessionID,
- messageID,
- timestamp: yield* DateTime.now,
- text: "Existing history",
- })
- const failure = yield* session
- .prompt({ id: messageID, sessionID, prompt: Prompt.make({ text: "Conflicting prompt" }), resume: false })
- .pipe(Effect.flip)
- expect(failure).toMatchObject({ _tag: "Session.PromptConflictError", sessionID, messageID })
- expect(yield* admitted(messageID)).toBeUndefined()
- }),
- )
- it.effect("starts execution by default after recording the prompt", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- executionCalls.length = 0
- wakeCalls.length = 0
- yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Run by default" }) })
- expect(executionCalls).toEqual([])
- expect(wakeCalls).toEqual([sessionID])
- }),
- )
- it.effect("starts execution when resume is explicitly true", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- executionCalls.length = 0
- wakeCalls.length = 0
- yield* session.prompt({
- sessionID,
- prompt: Prompt.make({ text: "Run explicitly" }),
- resume: true,
- })
- expect(executionCalls).toEqual([])
- expect(wakeCalls).toEqual([sessionID])
- }),
- )
- it.effect("only records the prompt when resume is false", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* SessionV2.Service
- executionCalls.length = 0
- wakeCalls.length = 0
- yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Do not run" }), resume: false })
- expect(executionCalls).toEqual([])
- expect(wakeCalls).toEqual([])
- }),
- )
- })
|