refactor(opencode): move session reads to core database
This commit is contained in:
@@ -235,10 +235,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
|
||||
if (currentAssistant) {
|
||||
yield* adapter.updateAssistant(
|
||||
produce(currentAssistant, (draft) => {
|
||||
draft.content.push({
|
||||
type: "text",
|
||||
text: "",
|
||||
})
|
||||
draft.content.push(new SessionMessage.AssistantText({ type: "text", text: "" }) as DraftText)
|
||||
}),
|
||||
)
|
||||
}
|
||||
@@ -276,18 +273,15 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
|
||||
if (currentAssistant) {
|
||||
yield* adapter.updateAssistant(
|
||||
produce(currentAssistant, (draft) => {
|
||||
draft.content.push({
|
||||
type: "tool",
|
||||
id: event.data.callID,
|
||||
name: event.data.name,
|
||||
time: {
|
||||
created: event.data.timestamp,
|
||||
},
|
||||
state: {
|
||||
status: "pending",
|
||||
input: "",
|
||||
},
|
||||
})
|
||||
draft.content.push(
|
||||
new SessionMessage.AssistantTool({
|
||||
type: "tool",
|
||||
id: event.data.callID,
|
||||
name: event.data.name,
|
||||
time: { created: event.data.timestamp },
|
||||
state: { status: "pending", input: "" },
|
||||
}) as DraftTool,
|
||||
)
|
||||
}),
|
||||
)
|
||||
}
|
||||
@@ -397,11 +391,13 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
|
||||
if (currentAssistant) {
|
||||
yield* adapter.updateAssistant(
|
||||
produce(currentAssistant, (draft) => {
|
||||
draft.content.push({
|
||||
type: "reasoning",
|
||||
id: event.data.reasoningID,
|
||||
text: "",
|
||||
})
|
||||
draft.content.push(
|
||||
new SessionMessage.AssistantReasoning({
|
||||
type: "reasoning",
|
||||
id: event.data.reasoningID,
|
||||
text: "",
|
||||
}) as DraftReasoning,
|
||||
)
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ import {
|
||||
import { NamedError } from "@opencode-ai/core/util/error"
|
||||
import { APICallError, convertToModelMessages, LoadAPIKeyError, type ModelMessage, type UIMessage } from "ai"
|
||||
import { SyncEvent } from "../sync"
|
||||
import { Database } from "@/storage/db"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { NotFoundError } from "@/storage/storage"
|
||||
import { and } from "drizzle-orm"
|
||||
import { desc } from "drizzle-orm"
|
||||
@@ -150,30 +150,31 @@ const part = (row: typeof PartTable.$inferSelect) =>
|
||||
const older = (row: Cursor) =>
|
||||
or(lt(MessageTable.time_created, row.time), and(eq(MessageTable.time_created, row.time), lt(MessageTable.id, row.id)))
|
||||
|
||||
function hydrate(rows: (typeof MessageTable.$inferSelect)[]) {
|
||||
function hydrate(db: Database.Interface["db"], rows: (typeof MessageTable.$inferSelect)[]) {
|
||||
const ids = rows.map((row) => row.id)
|
||||
const partByMessage = new Map<string, Part[]>()
|
||||
if (ids.length > 0) {
|
||||
const partRows = Database.use((db) =>
|
||||
db
|
||||
return Effect.gen(function* () {
|
||||
if (ids.length > 0) {
|
||||
const partRows = yield* db
|
||||
.select()
|
||||
.from(PartTable)
|
||||
.where(inArray(PartTable.message_id, ids))
|
||||
.orderBy(PartTable.message_id, PartTable.id)
|
||||
.all(),
|
||||
)
|
||||
for (const row of partRows) {
|
||||
const next = part(row)
|
||||
const list = partByMessage.get(row.message_id)
|
||||
if (list) list.push(next)
|
||||
else partByMessage.set(row.message_id, [next])
|
||||
.all()
|
||||
.pipe(Effect.orDie)
|
||||
for (const row of partRows) {
|
||||
const next = part(row)
|
||||
const list = partByMessage.get(row.message_id)
|
||||
if (list) list.push(next)
|
||||
else partByMessage.set(row.message_id, [next])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return rows.map((row) => ({
|
||||
info: info(row),
|
||||
parts: partByMessage.get(row.id) ?? [],
|
||||
}))
|
||||
return rows.map((row) => ({
|
||||
info: info(row),
|
||||
parts: partByMessage.get(row.id) ?? [],
|
||||
}))
|
||||
})
|
||||
}
|
||||
|
||||
function providerMeta(metadata: Record<string, any> | undefined) {
|
||||
@@ -480,23 +481,26 @@ export const page = Effect.fn("MessageV2.page")(function* (input: {
|
||||
limit: number
|
||||
before?: string
|
||||
}) {
|
||||
const { db } = yield* Database.Service
|
||||
const before = input.before ? cursor.decode(input.before) : undefined
|
||||
const where = before
|
||||
? and(eq(MessageTable.session_id, input.sessionID), older(before))
|
||||
: eq(MessageTable.session_id, input.sessionID)
|
||||
const rows = Database.use((db) =>
|
||||
db
|
||||
.select()
|
||||
.from(MessageTable)
|
||||
.where(where)
|
||||
.orderBy(desc(MessageTable.time_created), desc(MessageTable.id))
|
||||
.limit(input.limit + 1)
|
||||
.all(),
|
||||
)
|
||||
const rows = yield* db
|
||||
.select()
|
||||
.from(MessageTable)
|
||||
.where(where)
|
||||
.orderBy(desc(MessageTable.time_created), desc(MessageTable.id))
|
||||
.limit(input.limit + 1)
|
||||
.all()
|
||||
.pipe(Effect.orDie)
|
||||
if (rows.length === 0) {
|
||||
const row = Database.use((db) =>
|
||||
db.select({ id: SessionTable.id }).from(SessionTable).where(eq(SessionTable.id, input.sessionID)).get(),
|
||||
)
|
||||
const row = yield* db
|
||||
.select({ id: SessionTable.id })
|
||||
.from(SessionTable)
|
||||
.where(eq(SessionTable.id, input.sessionID))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (!row) return yield* new NotFoundError({ message: `Session not found: ${input.sessionID}` })
|
||||
return {
|
||||
items: [] as WithParts[],
|
||||
@@ -506,7 +510,7 @@ export const page = Effect.fn("MessageV2.page")(function* (input: {
|
||||
|
||||
const more = rows.length > input.limit
|
||||
const slice = more ? rows.slice(0, input.limit) : rows
|
||||
const items = hydrate(slice)
|
||||
const items = yield* hydrate(db, slice)
|
||||
items.reverse()
|
||||
const tail = slice.at(-1)
|
||||
return {
|
||||
@@ -516,53 +520,55 @@ export const page = Effect.fn("MessageV2.page")(function* (input: {
|
||||
}
|
||||
})
|
||||
|
||||
export function* stream(sessionID: SessionID) {
|
||||
export function stream(sessionID: SessionID) {
|
||||
const size = 50
|
||||
let before: string | undefined
|
||||
while (true) {
|
||||
const next = Effect.runSync(
|
||||
page({ sessionID, limit: size, before }).pipe(
|
||||
return Effect.gen(function* () {
|
||||
const result = [] as WithParts[]
|
||||
let before: string | undefined
|
||||
while (true) {
|
||||
const next = yield* page({ sessionID, limit: size, before }).pipe(
|
||||
Effect.catchIf(NotFoundError.isInstance, () =>
|
||||
Effect.succeed({ items: [] as WithParts[], more: false, cursor: undefined }),
|
||||
),
|
||||
),
|
||||
)
|
||||
if (next.items.length === 0) break
|
||||
for (let i = next.items.length - 1; i >= 0; i--) {
|
||||
yield next.items[i]
|
||||
)
|
||||
if (next.items.length === 0) break
|
||||
for (let i = next.items.length - 1; i >= 0; i--) {
|
||||
const item = next.items[i]
|
||||
if (item) result.push(item)
|
||||
}
|
||||
if (!next.more || !next.cursor) break
|
||||
before = next.cursor
|
||||
}
|
||||
if (!next.more || !next.cursor) break
|
||||
before = next.cursor
|
||||
}
|
||||
return result
|
||||
})
|
||||
}
|
||||
|
||||
export function parts(message_id: MessageID) {
|
||||
const rows = Database.use((db) =>
|
||||
db.select().from(PartTable).where(eq(PartTable.message_id, message_id)).orderBy(PartTable.id).all(),
|
||||
)
|
||||
return rows.map(
|
||||
(row) =>
|
||||
({
|
||||
...row.data,
|
||||
id: row.id,
|
||||
sessionID: row.session_id,
|
||||
messageID: row.message_id,
|
||||
}) as Part,
|
||||
)
|
||||
export function parts(messageID: MessageID) {
|
||||
return Effect.gen(function* () {
|
||||
const { db } = yield* Database.Service
|
||||
const rows = yield* db
|
||||
.select()
|
||||
.from(PartTable)
|
||||
.where(eq(PartTable.message_id, messageID))
|
||||
.orderBy(PartTable.id)
|
||||
.all()
|
||||
.pipe(Effect.orDie)
|
||||
return rows.map(part)
|
||||
})
|
||||
}
|
||||
|
||||
export const get = Effect.fn("MessageV2.get")(function* (input: { sessionID: SessionID; messageID: MessageID }) {
|
||||
const row = Database.use((db) =>
|
||||
db
|
||||
.select()
|
||||
.from(MessageTable)
|
||||
.where(and(eq(MessageTable.id, input.messageID), eq(MessageTable.session_id, input.sessionID)))
|
||||
.get(),
|
||||
)
|
||||
const { db } = yield* Database.Service
|
||||
const row = yield* db
|
||||
.select()
|
||||
.from(MessageTable)
|
||||
.where(and(eq(MessageTable.id, input.messageID), eq(MessageTable.session_id, input.sessionID)))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (!row) return yield* new NotFoundError({ message: `Message not found: ${input.messageID}` })
|
||||
return {
|
||||
info: info(row),
|
||||
parts: parts(input.messageID),
|
||||
parts: yield* parts(input.messageID),
|
||||
}
|
||||
})
|
||||
|
||||
@@ -620,7 +626,7 @@ export function filterCompacted(msgs: Iterable<WithParts>) {
|
||||
}
|
||||
|
||||
export const filterCompactedEffect = Effect.fnUntraced(function* (sessionID: SessionID) {
|
||||
return filterCompacted(stream(sessionID))
|
||||
return filterCompacted(yield* stream(sessionID))
|
||||
})
|
||||
|
||||
// filterCompacted reorders messages for model consumption
|
||||
|
||||
@@ -23,6 +23,7 @@ import { errorMessage } from "@/util/error"
|
||||
import * as Log from "@opencode-ai/core/util/log"
|
||||
import { isRecord } from "@/util/record"
|
||||
import { EventV2Bridge } from "@/event-v2-bridge"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { SessionEvent } from "@opencode-ai/core/session/event"
|
||||
import { ModelV2 } from "@opencode-ai/core/model"
|
||||
import { ProviderV2 } from "@opencode-ai/core/provider"
|
||||
@@ -102,6 +103,7 @@ export const layer = Layer.effect(
|
||||
const image = yield* Image.Service
|
||||
const events = yield* EventV2Bridge.Service
|
||||
const flags = yield* RuntimeFlags.Service
|
||||
const database = yield* Database.Service
|
||||
|
||||
const create = Effect.fn("SessionProcessor.create")(function* (input: Input) {
|
||||
// Pre-capture snapshot before the LLM stream starts. The AI SDK
|
||||
@@ -422,7 +424,9 @@ export const layer = Layer.effect(
|
||||
: value.providerMetadata,
|
||||
}))
|
||||
|
||||
const parts = MessageV2.parts(ctx.assistantMessage.id)
|
||||
const parts = yield* MessageV2.parts(ctx.assistantMessage.id).pipe(
|
||||
Effect.provideService(Database.Service, database),
|
||||
)
|
||||
const recentParts = parts.slice(-DOOM_LOOP_THRESHOLD)
|
||||
|
||||
if (
|
||||
@@ -877,6 +881,7 @@ export const defaultLayer = Layer.suspend(() =>
|
||||
Layer.provide(Bus.layer),
|
||||
Layer.provide(Config.defaultLayer),
|
||||
Layer.provide(RuntimeFlags.defaultLayer),
|
||||
Layer.provide(Database.defaultLayer),
|
||||
Layer.provide(EventV2Bridge.defaultLayer),
|
||||
),
|
||||
)
|
||||
|
||||
@@ -49,14 +49,14 @@ import { TaskTool, type TaskPromptOps } from "@/tool/task"
|
||||
import { SessionRunState } from "./run-state"
|
||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||
import { EventV2Bridge } from "@/event-v2-bridge"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { SessionEvent } from "@opencode-ai/core/session/event"
|
||||
import { ModelV2 } from "@opencode-ai/core/model"
|
||||
import { ProviderV2 } from "@opencode-ai/core/provider"
|
||||
import { AgentAttachment, FileAttachment, ReferenceAttachment, Source } from "@opencode-ai/core/session/prompt"
|
||||
import { Reference } from "@/reference/reference"
|
||||
import * as DateTime from "effect/DateTime"
|
||||
import { eq } from "@/storage/db"
|
||||
import * as Database from "@/storage/db"
|
||||
import { eq } from "drizzle-orm"
|
||||
import { SessionTable } from "@opencode-ai/core/session/sql"
|
||||
import { referencePromptMetadata, referenceTextPart } from "./prompt/reference"
|
||||
import { SessionReminders } from "./reminders"
|
||||
@@ -124,6 +124,8 @@ export const layer = Layer.effect(
|
||||
const references = yield* Reference.Service
|
||||
const events = yield* EventV2Bridge.Service
|
||||
const flags = yield* RuntimeFlags.Service
|
||||
const database = yield* Database.Service
|
||||
const { db } = database
|
||||
const ops = Effect.fn("SessionPrompt.ops")(function* () {
|
||||
return {
|
||||
cancel: (sessionID: SessionID) => cancel(sessionID),
|
||||
@@ -669,9 +671,12 @@ export const layer = Layer.effect(
|
||||
})
|
||||
|
||||
const currentModel = Effect.fnUntraced(function* (sessionID: SessionID) {
|
||||
const current = Database.use((db) =>
|
||||
db.select({ model: SessionTable.model }).from(SessionTable).where(eq(SessionTable.id, sessionID)).get(),
|
||||
)
|
||||
const current = yield* db
|
||||
.select({ model: SessionTable.model })
|
||||
.from(SessionTable)
|
||||
.where(eq(SessionTable.id, sessionID))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (current?.model) {
|
||||
return {
|
||||
providerID: ProviderID.make(current.model.providerID),
|
||||
@@ -697,13 +702,12 @@ export const layer = Layer.effect(
|
||||
throw error
|
||||
}
|
||||
|
||||
const current = Database.use((db) =>
|
||||
db
|
||||
.select({ agent: SessionTable.agent, model: SessionTable.model })
|
||||
.from(SessionTable)
|
||||
.where(eq(SessionTable.id, input.sessionID))
|
||||
.get(),
|
||||
)
|
||||
const current = yield* db
|
||||
.select({ agent: SessionTable.agent, model: SessionTable.model })
|
||||
.from(SessionTable)
|
||||
.where(eq(SessionTable.id, input.sessionID))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
const model = input.model ?? ag.model ?? (yield* currentModel(input.sessionID))
|
||||
const same = ag.model && model.providerID === ag.model.providerID && model.modelID === ag.model.modelID
|
||||
const full =
|
||||
@@ -1249,7 +1253,9 @@ export const layer = Layer.effect(
|
||||
yield* status.set(sessionID, { type: "busy" })
|
||||
yield* slog.info("loop", { step })
|
||||
|
||||
let msgs = yield* MessageV2.filterCompactedEffect(sessionID)
|
||||
let msgs = yield* MessageV2.filterCompactedEffect(sessionID).pipe(
|
||||
Effect.provideService(Database.Service, database),
|
||||
)
|
||||
|
||||
const { user: lastUser, assistant: lastAssistant, finished: lastFinished, tasks } = MessageV2.latest(msgs)
|
||||
|
||||
@@ -1648,6 +1654,7 @@ export const defaultLayer = Layer.suspend(() =>
|
||||
Layer.mergeAll(
|
||||
EventV2Bridge.defaultLayer,
|
||||
Agent.defaultLayer,
|
||||
Database.defaultLayer,
|
||||
SystemPrompt.defaultLayer,
|
||||
LLM.defaultLayer,
|
||||
Reference.defaultLayer,
|
||||
|
||||
@@ -8,8 +8,9 @@ import { Bus } from "@/bus"
|
||||
import { Decimal } from "decimal.js"
|
||||
import type { ProviderMetadata, Usage } from "@opencode-ai/llm"
|
||||
import { InstallationVersion } from "@opencode-ai/core/installation/version"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { makeRuntime } from "@opencode-ai/core/effect/runtime"
|
||||
|
||||
import { Database } from "@/storage/db"
|
||||
import { NotFoundError } from "@/storage/storage"
|
||||
import { eq } from "drizzle-orm"
|
||||
import { and } from "drizzle-orm"
|
||||
@@ -43,6 +44,7 @@ import { NonNegativeInt, optionalOmitUndefined } from "@opencode-ai/core/schema"
|
||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||
|
||||
const log = Log.create({ service: "session" })
|
||||
const runtime = makeRuntime(Database.Service, Database.defaultLayer)
|
||||
|
||||
const parentTitlePrefix = "New session - "
|
||||
const childTitlePrefix = "Child session - "
|
||||
@@ -505,16 +507,15 @@ export const use = serviceUse(Service)
|
||||
|
||||
export type Patch = Types.DeepMutable<SyncEvent.Event<typeof Event.Updated>["data"]["info"]>
|
||||
|
||||
const db = <T>(fn: (d: Parameters<typeof Database.use>[0] extends (trx: infer D) => any ? D : never) => T) =>
|
||||
Effect.sync(() => Database.use(fn))
|
||||
|
||||
export const layer: Layer.Layer<
|
||||
Service,
|
||||
never,
|
||||
BackgroundJob.Service | Bus.Service | Storage.Service | SyncEvent.Service | RuntimeFlags.Service
|
||||
BackgroundJob.Service | Bus.Service | Storage.Service | SyncEvent.Service | RuntimeFlags.Service | Database.Service
|
||||
> = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const { db } = yield* Database.Service
|
||||
const database = yield* Database.Service
|
||||
const background = yield* BackgroundJob.Service
|
||||
const bus = yield* Bus.Service
|
||||
const storage = yield* Storage.Service
|
||||
@@ -570,7 +571,7 @@ export const layer: Layer.Layer<
|
||||
})
|
||||
|
||||
const get = Effect.fn("Session.get")(function* (id: SessionID) {
|
||||
const row = yield* db((d) => d.select().from(SessionTable).where(eq(SessionTable.id, id)).get())
|
||||
const row = yield* db.select().from(SessionTable).where(eq(SessionTable.id, id)).get().pipe(Effect.orDie)
|
||||
if (!row) return yield* Effect.fail(new NotFoundError({ message: `Session not found: ${id}` }))
|
||||
return fromRow(row)
|
||||
})
|
||||
@@ -583,13 +584,12 @@ export const layer: Layer.Layer<
|
||||
})
|
||||
|
||||
const children = Effect.fn("Session.children")(function* (parentID: SessionID) {
|
||||
const rows = yield* db((d) =>
|
||||
d
|
||||
.select()
|
||||
.from(SessionTable)
|
||||
.where(and(eq(SessionTable.parent_id, parentID)))
|
||||
.all(),
|
||||
)
|
||||
const rows = yield* db
|
||||
.select()
|
||||
.from(SessionTable)
|
||||
.where(and(eq(SessionTable.parent_id, parentID)))
|
||||
.all()
|
||||
.pipe(Effect.orDie)
|
||||
return rows.map(fromRow)
|
||||
})
|
||||
|
||||
@@ -633,19 +633,18 @@ export const layer: Layer.Layer<
|
||||
}).pipe(Effect.withSpan("Session.updatePart"))
|
||||
|
||||
const getPart: Interface["getPart"] = Effect.fn("Session.getPart")(function* (input) {
|
||||
const row = Database.use((db) =>
|
||||
db
|
||||
.select()
|
||||
.from(PartTable)
|
||||
.where(
|
||||
and(
|
||||
eq(PartTable.session_id, input.sessionID),
|
||||
eq(PartTable.message_id, input.messageID),
|
||||
eq(PartTable.id, input.partID),
|
||||
),
|
||||
)
|
||||
.get(),
|
||||
)
|
||||
const row = yield* db
|
||||
.select()
|
||||
.from(PartTable)
|
||||
.where(
|
||||
and(
|
||||
eq(PartTable.session_id, input.sessionID),
|
||||
eq(PartTable.message_id, input.messageID),
|
||||
eq(PartTable.id, input.partID),
|
||||
),
|
||||
)
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (!row) return
|
||||
return {
|
||||
...row.data,
|
||||
@@ -767,14 +766,18 @@ export const layer: Layer.Layer<
|
||||
|
||||
const messages: Interface["messages"] = Effect.fn("Session.messages")(function* (input) {
|
||||
if (input.limit) {
|
||||
return (yield* MessageV2.page({ sessionID: input.sessionID, limit: input.limit })).items
|
||||
return (yield* MessageV2.page({ sessionID: input.sessionID, limit: input.limit }).pipe(
|
||||
Effect.provideService(Database.Service, database),
|
||||
)).items
|
||||
}
|
||||
|
||||
const size = 50
|
||||
const result = [] as SessionLegacy.WithParts[]
|
||||
let before: string | undefined
|
||||
while (true) {
|
||||
const page = yield* MessageV2.page({ sessionID: input.sessionID, limit: size, before })
|
||||
const page = yield* MessageV2.page({ sessionID: input.sessionID, limit: size, before }).pipe(
|
||||
Effect.provideService(Database.Service, database),
|
||||
)
|
||||
if (page.items.length === 0) break
|
||||
for (let i = page.items.length - 1; i >= 0; i--) {
|
||||
const item = page.items[i]
|
||||
@@ -825,7 +828,9 @@ export const layer: Layer.Layer<
|
||||
const size = 50
|
||||
let before: string | undefined
|
||||
while (true) {
|
||||
const page = yield* MessageV2.page({ sessionID, limit: size, before })
|
||||
const page = yield* MessageV2.page({ sessionID, limit: size, before }).pipe(
|
||||
Effect.provideService(Database.Service, database),
|
||||
)
|
||||
if (page.items.length === 0) break
|
||||
for (let i = page.items.length - 1; i >= 0; i--) {
|
||||
const item = page.items[i]
|
||||
@@ -869,6 +874,7 @@ export const defaultLayer = layer.pipe(
|
||||
Layer.provide(Bus.layer),
|
||||
Layer.provide(Storage.defaultLayer),
|
||||
Layer.provide(SyncEvent.defaultLayer),
|
||||
Layer.provide(Database.defaultLayer),
|
||||
Layer.provide(RuntimeFlags.defaultLayer),
|
||||
)
|
||||
|
||||
@@ -927,14 +933,15 @@ function* listByProject(
|
||||
|
||||
const limit = input.limit ?? 100
|
||||
|
||||
const rows = Database.use((db) =>
|
||||
const rows = runtime.runSync(({ db }) =>
|
||||
db
|
||||
.select()
|
||||
.from(SessionTable)
|
||||
.where(and(...conditions))
|
||||
.orderBy(desc(SessionTable.time_updated))
|
||||
.limit(limit)
|
||||
.all(),
|
||||
.all()
|
||||
.pipe(Effect.orDie),
|
||||
)
|
||||
for (const row of rows) {
|
||||
yield fromRow(row)
|
||||
@@ -973,7 +980,7 @@ export function* listGlobal(input?: {
|
||||
|
||||
const limit = input?.limit ?? 100
|
||||
|
||||
const rows = Database.use((db) => {
|
||||
const rows = runtime.runSync(({ db }) => {
|
||||
const query =
|
||||
conditions.length > 0
|
||||
? db
|
||||
@@ -981,19 +988,20 @@ export function* listGlobal(input?: {
|
||||
.from(SessionTable)
|
||||
.where(and(...conditions))
|
||||
: db.select().from(SessionTable)
|
||||
return query.orderBy(desc(SessionTable.time_updated), desc(SessionTable.id)).limit(limit).all()
|
||||
return query.orderBy(desc(SessionTable.time_updated), desc(SessionTable.id)).limit(limit).all().pipe(Effect.orDie)
|
||||
})
|
||||
|
||||
const ids = [...new Set(rows.map((row) => row.project_id))]
|
||||
const projects = new Map<string, ProjectInfo>()
|
||||
|
||||
if (ids.length > 0) {
|
||||
const items = Database.use((db) =>
|
||||
const items = runtime.runSync(({ db }) =>
|
||||
db
|
||||
.select({ id: ProjectTable.id, name: ProjectTable.name, worktree: ProjectTable.worktree })
|
||||
.from(ProjectTable)
|
||||
.where(inArray(ProjectTable.id, ids))
|
||||
.all(),
|
||||
.all()
|
||||
.pipe(Effect.orDie),
|
||||
)
|
||||
for (const item of items) {
|
||||
projects.set(item.id, {
|
||||
|
||||
@@ -2,7 +2,7 @@ import { BusEvent } from "@/bus/bus-event"
|
||||
import { Bus } from "@/bus"
|
||||
import { SessionID } from "./schema"
|
||||
import { Effect, Layer, Context, Schema } from "effect"
|
||||
import { Database } from "@/storage/db"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { eq } from "drizzle-orm"
|
||||
import { asc } from "drizzle-orm"
|
||||
import { TodoTable } from "@opencode-ai/core/session/sql"
|
||||
@@ -37,34 +37,40 @@ export const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const bus = yield* Bus.Service
|
||||
const { db } = yield* Database.Service
|
||||
|
||||
const update = Effect.fn("Todo.update")(function* (input: { sessionID: SessionID; todos: Info[] }) {
|
||||
yield* Effect.sync(() =>
|
||||
Database.transaction((db) => {
|
||||
db.delete(TodoTable).where(eq(TodoTable.session_id, input.sessionID)).run()
|
||||
if (input.todos.length === 0) return
|
||||
db.insert(TodoTable)
|
||||
.values(
|
||||
input.todos.map((todo, position) => ({
|
||||
session_id: input.sessionID,
|
||||
content: todo.content,
|
||||
status: todo.status,
|
||||
priority: todo.priority,
|
||||
position,
|
||||
})),
|
||||
)
|
||||
.run()
|
||||
}),
|
||||
)
|
||||
yield* db
|
||||
.transaction((tx) =>
|
||||
Effect.gen(function* () {
|
||||
yield* tx.delete(TodoTable).where(eq(TodoTable.session_id, input.sessionID)).run()
|
||||
if (input.todos.length === 0) return
|
||||
yield* tx
|
||||
.insert(TodoTable)
|
||||
.values(
|
||||
input.todos.map((todo, position) => ({
|
||||
session_id: input.sessionID,
|
||||
content: todo.content,
|
||||
status: todo.status,
|
||||
priority: todo.priority,
|
||||
position,
|
||||
})),
|
||||
)
|
||||
.run()
|
||||
}),
|
||||
)
|
||||
.pipe(Effect.orDie)
|
||||
yield* bus.publish(Event.Updated, input)
|
||||
})
|
||||
|
||||
const get = Effect.fn("Todo.get")(function* (sessionID: SessionID) {
|
||||
const rows = yield* Effect.sync(() =>
|
||||
Database.use((db) =>
|
||||
db.select().from(TodoTable).where(eq(TodoTable.session_id, sessionID)).orderBy(asc(TodoTable.position)).all(),
|
||||
),
|
||||
)
|
||||
const rows = yield* db
|
||||
.select()
|
||||
.from(TodoTable)
|
||||
.where(eq(TodoTable.session_id, sessionID))
|
||||
.orderBy(asc(TodoTable.position))
|
||||
.all()
|
||||
.pipe(Effect.orDie)
|
||||
return rows.map((row) => ({
|
||||
content: row.content,
|
||||
status: row.status,
|
||||
@@ -76,6 +82,6 @@ export const layer = Layer.effect(
|
||||
}),
|
||||
)
|
||||
|
||||
export const defaultLayer = layer.pipe(Layer.provide(Bus.layer))
|
||||
export const defaultLayer = layer.pipe(Layer.provide(Bus.layer), Layer.provide(Database.defaultLayer))
|
||||
|
||||
export * as Todo from "./todo"
|
||||
|
||||
@@ -7,6 +7,7 @@ import { GlobTool } from "./glob"
|
||||
import { GrepTool } from "./grep"
|
||||
import { ReadTool } from "./read"
|
||||
import { TaskTool } from "./task"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { TaskStatusTool } from "./task_status"
|
||||
import { TodoWriteTool } from "./todo"
|
||||
import { WebFetchTool } from "./webfetch"
|
||||
@@ -107,6 +108,7 @@ export const layer: Layer.Layer<
|
||||
| Format.Service
|
||||
| Truncate.Service
|
||||
| RuntimeFlags.Service
|
||||
| Database.Service
|
||||
> = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
@@ -399,7 +401,7 @@ export const defaultLayer = Layer.suspend(() =>
|
||||
Layer.provide(Ripgrep.defaultLayer),
|
||||
Layer.provide(Truncate.defaultLayer),
|
||||
)
|
||||
.pipe(Layer.provide(RuntimeFlags.defaultLayer)),
|
||||
.pipe(Layer.provide(Database.defaultLayer), Layer.provide(RuntimeFlags.defaultLayer)),
|
||||
)
|
||||
|
||||
function isZodType(value: unknown): value is z.ZodType {
|
||||
|
||||
@@ -16,6 +16,7 @@ import { TuiEvent } from "@/cli/cmd/tui/event"
|
||||
import { Cause, Effect, Exit, Option, Schema, Scope } from "effect"
|
||||
import { EffectBridge } from "@/effect/bridge"
|
||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
|
||||
export interface TaskPromptOps {
|
||||
cancel(sessionID: SessionID): Effect.Effect<void>
|
||||
@@ -112,6 +113,7 @@ export const TaskTool = Tool.define(
|
||||
const scope = yield* Scope.Scope
|
||||
const status = yield* SessionStatus.Service
|
||||
const flags = yield* RuntimeFlags.Service
|
||||
const database = yield* Database.Service
|
||||
|
||||
const run = Effect.fn("TaskTool.execute")(function* (
|
||||
params: Schema.Schema.Type<typeof Parameters>,
|
||||
@@ -169,7 +171,10 @@ export const TaskTool = Tool.define(
|
||||
],
|
||||
}))
|
||||
|
||||
const msg = yield* MessageV2.get({ sessionID: ctx.sessionID, messageID: ctx.messageID }).pipe(Effect.orDie)
|
||||
const msg = yield* MessageV2.get({ sessionID: ctx.sessionID, messageID: ctx.messageID }).pipe(
|
||||
Effect.provideService(Database.Service, database),
|
||||
Effect.orDie,
|
||||
)
|
||||
if (msg.info.role !== "assistant") return yield* Effect.fail(new Error("Not an assistant message"))
|
||||
|
||||
const model = next.model ?? {
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
import { afterEach, describe, expect } from "bun:test"
|
||||
import { Effect, Layer } from "effect"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { Session as SessionNs } from "@/session/session"
|
||||
import * as Log from "@opencode-ai/core/util/log"
|
||||
import { disposeAllInstances, provideInstance, TestInstance } from "../fixture/fixture"
|
||||
import { mkdir } from "fs/promises"
|
||||
import path from "path"
|
||||
import { Database } from "@/storage/db"
|
||||
import { SessionTable } from "@opencode-ai/core/session/sql"
|
||||
import { eq } from "drizzle-orm"
|
||||
import { testEffect } from "../lib/effect"
|
||||
@@ -17,12 +17,16 @@ import { BackgroundJob } from "@/background/job"
|
||||
|
||||
void Log.init({ print: false })
|
||||
const it = testEffect(
|
||||
SessionNs.layer.pipe(
|
||||
Layer.provide(Bus.layer),
|
||||
Layer.provide(Storage.defaultLayer),
|
||||
Layer.provide(SyncEvent.defaultLayer),
|
||||
Layer.provide(RuntimeFlags.layer({ experimentalWorkspaces: false })),
|
||||
Layer.provide(BackgroundJob.defaultLayer),
|
||||
Layer.mergeAll(
|
||||
Database.defaultLayer,
|
||||
SessionNs.layer.pipe(
|
||||
Layer.provide(Bus.layer),
|
||||
Layer.provide(Storage.defaultLayer),
|
||||
Layer.provide(SyncEvent.defaultLayer),
|
||||
Layer.provide(Database.defaultLayer),
|
||||
Layer.provide(RuntimeFlags.layer({ experimentalWorkspaces: false })),
|
||||
Layer.provide(BackgroundJob.defaultLayer),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
@@ -148,16 +152,9 @@ describe("session.list", () => {
|
||||
provideInstance(path.join(test.directory, "packages", "app")),
|
||||
)
|
||||
|
||||
yield* Effect.sync(() =>
|
||||
Database.use((db) =>
|
||||
db.update(SessionTable).set({ path: null }).where(eq(SessionTable.id, current.id)).run(),
|
||||
),
|
||||
)
|
||||
yield* Effect.sync(() =>
|
||||
Database.use((db) =>
|
||||
db.update(SessionTable).set({ path: null }).where(eq(SessionTable.id, sibling.id)).run(),
|
||||
),
|
||||
)
|
||||
const { db } = yield* Database.Service
|
||||
yield* db.update(SessionTable).set({ path: null }).where(eq(SessionTable.id, current.id)).run().pipe(Effect.orDie)
|
||||
yield* db.update(SessionTable).set({ path: null }).where(eq(SessionTable.id, sibling.id)).run().pipe(Effect.orDie)
|
||||
|
||||
const pathIDs = (yield* SessionNs.Service.use((session) =>
|
||||
session.list({
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { afterEach, describe, expect, mock, test } from "bun:test"
|
||||
import { SessionLegacy } from "@opencode-ai/core/session/legacy"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { APICallError } from "ai"
|
||||
import { Cause, Deferred, Effect, Exit, Fiber, Layer, Schema } from "effect"
|
||||
import * as Stream from "effect/Stream"
|
||||
@@ -235,17 +236,19 @@ const deps = Layer.mergeAll(
|
||||
SyncEvent.defaultLayer,
|
||||
RuntimeFlags.layer({ experimentalEventSystem: true }),
|
||||
EventV2Bridge.defaultLayer,
|
||||
Database.defaultLayer,
|
||||
)
|
||||
|
||||
const env = Layer.mergeAll(
|
||||
SessionNs.defaultLayer,
|
||||
Database.defaultLayer,
|
||||
CrossSpawnSpawner.defaultLayer,
|
||||
SessionCompaction.layer.pipe(Layer.provide(SessionNs.defaultLayer), Layer.provideMerge(deps)),
|
||||
)
|
||||
|
||||
const it = testEffect(env)
|
||||
|
||||
const compactionEnv = Layer.mergeAll(SessionNs.defaultLayer, CrossSpawnSpawner.defaultLayer)
|
||||
const compactionEnv = Layer.mergeAll(SessionNs.defaultLayer, Database.defaultLayer, CrossSpawnSpawner.defaultLayer)
|
||||
const itCompaction = testEffect(compactionEnv)
|
||||
|
||||
type CompactionProcessOptions = {
|
||||
@@ -1065,7 +1068,7 @@ describe("session.compaction.process", () => {
|
||||
expect(captured).toContain("zzzz")
|
||||
expect(captured).not.toContain("keep tail")
|
||||
|
||||
const filtered = MessageV2.filterCompacted(MessageV2.stream(session.id))
|
||||
const filtered = MessageV2.filterCompacted(yield* MessageV2.stream(session.id))
|
||||
expect(filtered.map((msg) => msg.info.id).slice(0, 3)).toEqual([parent!, expect.any(String), keep.id])
|
||||
expect(filtered[1]?.info.role).toBe("assistant")
|
||||
expect(filtered[1]?.info.role === "assistant" ? filtered[1].info.summary : false).toBe(true)
|
||||
@@ -1406,7 +1409,7 @@ describe("session.compaction.process", () => {
|
||||
yield* createUserMessage(session.id, "latest turn")
|
||||
yield* createCompactionMarker(session.id)
|
||||
|
||||
msgs = MessageV2.filterCompacted(MessageV2.stream(session.id))
|
||||
msgs = MessageV2.filterCompacted(yield* MessageV2.stream(session.id))
|
||||
parent = msgs.at(-1)?.info.id
|
||||
expect(parent).toBeTruthy()
|
||||
yield* SessionCompaction.use.process({ parentID: parent!, messages: msgs, sessionID: session.id, auto: false })
|
||||
@@ -1442,12 +1445,12 @@ describe("session.compaction.process", () => {
|
||||
const u4 = yield* createUserMessage(session.id, "four")
|
||||
yield* createCompactionMarker(session.id)
|
||||
|
||||
msgs = MessageV2.filterCompacted(MessageV2.stream(session.id))
|
||||
msgs = MessageV2.filterCompacted(yield* MessageV2.stream(session.id))
|
||||
parent = msgs.at(-1)?.info.id
|
||||
expect(parent).toBeTruthy()
|
||||
yield* SessionCompaction.use.process({ parentID: parent!, messages: msgs, sessionID: session.id, auto: false })
|
||||
|
||||
const filtered = MessageV2.filterCompacted(MessageV2.stream(session.id))
|
||||
const filtered = MessageV2.filterCompacted(yield* MessageV2.stream(session.id))
|
||||
const ids = filtered.map((msg) => msg.info.id)
|
||||
|
||||
expect(ids).not.toContain(u1.id)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { describe, expect, test } from "bun:test"
|
||||
import { SessionLegacy } from "@opencode-ai/core/session/legacy"
|
||||
import { Effect, Option } from "effect"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { Effect, Layer, Option } from "effect"
|
||||
import { Session as SessionNs } from "@/session/session"
|
||||
import { MessageV2 } from "../../src/session/message-v2"
|
||||
import { MessageID, PartID, type SessionID } from "../../src/session/schema"
|
||||
@@ -11,7 +12,7 @@ import { testEffect } from "../lib/effect"
|
||||
|
||||
void Log.init({ print: false })
|
||||
|
||||
const it = testEffect(SessionNs.defaultLayer)
|
||||
const it = testEffect(Layer.mergeAll(SessionNs.defaultLayer, Database.defaultLayer))
|
||||
|
||||
const withSession = <A, E, R>(
|
||||
fn: (input: { session: SessionNs.Interface; sessionID: SessionID }) => Effect.Effect<A, E, R>,
|
||||
@@ -311,7 +312,7 @@ describe("MessageV2.stream", () => {
|
||||
Effect.gen(function* () {
|
||||
const ids = yield* fill(sessionID, 5)
|
||||
|
||||
const items = Array.from(MessageV2.stream(sessionID))
|
||||
const items = yield* MessageV2.stream(sessionID)
|
||||
expect(items.map((item) => item.info.id)).toEqual(ids.slice().reverse())
|
||||
}),
|
||||
),
|
||||
@@ -320,7 +321,7 @@ describe("MessageV2.stream", () => {
|
||||
it.instance("yields nothing for empty session", () =>
|
||||
withSession(({ sessionID }) =>
|
||||
Effect.gen(function* () {
|
||||
const items = Array.from(MessageV2.stream(sessionID))
|
||||
const items = yield* MessageV2.stream(sessionID)
|
||||
expect(items).toHaveLength(0)
|
||||
}),
|
||||
),
|
||||
@@ -331,7 +332,7 @@ describe("MessageV2.stream", () => {
|
||||
Effect.gen(function* () {
|
||||
const ids = yield* fill(sessionID, 1)
|
||||
|
||||
const items = Array.from(MessageV2.stream(sessionID))
|
||||
const items = yield* MessageV2.stream(sessionID)
|
||||
expect(items).toHaveLength(1)
|
||||
expect(items[0].info.id).toBe(ids[0])
|
||||
}),
|
||||
@@ -343,7 +344,7 @@ describe("MessageV2.stream", () => {
|
||||
Effect.gen(function* () {
|
||||
yield* fill(sessionID, 3)
|
||||
|
||||
const items = Array.from(MessageV2.stream(sessionID))
|
||||
const items = yield* MessageV2.stream(sessionID)
|
||||
for (const item of items) {
|
||||
expect(item.parts).toHaveLength(1)
|
||||
expect(item.parts[0].type).toBe("text")
|
||||
@@ -357,7 +358,7 @@ describe("MessageV2.stream", () => {
|
||||
Effect.gen(function* () {
|
||||
const ids = yield* fill(sessionID, 60)
|
||||
|
||||
const items = Array.from(MessageV2.stream(sessionID))
|
||||
const items = yield* MessageV2.stream(sessionID)
|
||||
expect(items).toHaveLength(60)
|
||||
expect(items[0].info.id).toBe(ids[ids.length - 1])
|
||||
expect(items[59].info.id).toBe(ids[0])
|
||||
@@ -365,17 +366,13 @@ describe("MessageV2.stream", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.instance("is a sync generator", () =>
|
||||
it.instance("returns an Effect", () =>
|
||||
withSession(({ sessionID }) =>
|
||||
Effect.gen(function* () {
|
||||
yield* fill(sessionID, 1)
|
||||
|
||||
const gen = MessageV2.stream(sessionID)
|
||||
const first = gen.next()
|
||||
// sync generator returns { value, done } directly, not a Promise
|
||||
expect(first).toHaveProperty("value")
|
||||
expect(first).toHaveProperty("done")
|
||||
expect(first.done).toBe(false)
|
||||
const result = yield* MessageV2.stream(sessionID)
|
||||
expect(result).toHaveLength(1)
|
||||
}),
|
||||
),
|
||||
)
|
||||
@@ -387,7 +384,7 @@ describe("MessageV2.parts", () => {
|
||||
Effect.gen(function* () {
|
||||
const [id] = yield* fill(sessionID, 1)
|
||||
|
||||
const result = MessageV2.parts(id)
|
||||
const result = yield* MessageV2.parts(id)
|
||||
expect(result).toHaveLength(1)
|
||||
expect(result[0].type).toBe("text")
|
||||
expect((result[0] as SessionLegacy.TextPart).text).toBe("m0")
|
||||
@@ -400,7 +397,7 @@ describe("MessageV2.parts", () => {
|
||||
Effect.gen(function* () {
|
||||
const id = yield* addUser(sessionID)
|
||||
|
||||
const result = MessageV2.parts(id)
|
||||
const result = yield* MessageV2.parts(id)
|
||||
expect(result).toEqual([])
|
||||
}),
|
||||
),
|
||||
@@ -426,7 +423,7 @@ describe("MessageV2.parts", () => {
|
||||
text: "third",
|
||||
})
|
||||
|
||||
const result = MessageV2.parts(id)
|
||||
const result = yield* MessageV2.parts(id)
|
||||
expect(result).toHaveLength(3)
|
||||
expect((result[0] as SessionLegacy.TextPart).text).toBe("m0")
|
||||
expect((result[1] as SessionLegacy.TextPart).text).toBe("second")
|
||||
@@ -438,7 +435,7 @@ describe("MessageV2.parts", () => {
|
||||
it.instance("returns empty for non-existent message id", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* SessionNs.Service
|
||||
const result = MessageV2.parts(MessageID.ascending())
|
||||
const result = yield* MessageV2.parts(MessageID.ascending())
|
||||
expect(result).toEqual([])
|
||||
}),
|
||||
)
|
||||
@@ -448,7 +445,7 @@ describe("MessageV2.parts", () => {
|
||||
Effect.gen(function* () {
|
||||
const [id] = yield* fill(sessionID, 1)
|
||||
|
||||
const result = MessageV2.parts(id)
|
||||
const result = yield* MessageV2.parts(id)
|
||||
expect(result[0].sessionID).toBe(sessionID)
|
||||
expect(result[0].messageID).toBe(id)
|
||||
}),
|
||||
@@ -605,7 +602,7 @@ describe("MessageV2.filterCompacted", () => {
|
||||
Effect.gen(function* () {
|
||||
const ids = yield* fill(sessionID, 5)
|
||||
|
||||
const result = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
|
||||
expect(result).toHaveLength(5)
|
||||
// reversed from newest-first to chronological
|
||||
expect(result.map((item) => item.info.id)).toEqual(ids)
|
||||
@@ -639,7 +636,7 @@ describe("MessageV2.filterCompacted", () => {
|
||||
text: "new response",
|
||||
})
|
||||
|
||||
const result = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
|
||||
// Includes compaction boundary: u1, a1, u2, a2
|
||||
expect(result[0].info.id).toBe(u1)
|
||||
expect(result.length).toBe(4)
|
||||
@@ -661,7 +658,7 @@ describe("MessageV2.filterCompacted", () => {
|
||||
yield* addCompactionPart(sessionID, u1)
|
||||
yield* addUser(sessionID, "world")
|
||||
|
||||
const result = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
|
||||
expect(result).toHaveLength(2)
|
||||
}),
|
||||
),
|
||||
@@ -680,7 +677,7 @@ describe("MessageV2.filterCompacted", () => {
|
||||
yield* addAssistant(sessionID, u1, { summary: true, finish: "end_turn", error })
|
||||
yield* addUser(sessionID, "retry")
|
||||
|
||||
const result = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
|
||||
// Error assistant doesn't add to completed, so compaction boundary never triggers
|
||||
expect(result).toHaveLength(3)
|
||||
}),
|
||||
@@ -697,7 +694,7 @@ describe("MessageV2.filterCompacted", () => {
|
||||
yield* addAssistant(sessionID, u1, { summary: true })
|
||||
yield* addUser(sessionID, "next")
|
||||
|
||||
const result = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
|
||||
expect(result).toHaveLength(3)
|
||||
}),
|
||||
),
|
||||
@@ -747,7 +744,7 @@ describe("MessageV2.filterCompacted", () => {
|
||||
text: "third reply",
|
||||
})
|
||||
|
||||
const result = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
|
||||
|
||||
expect(result.map((item) => item.info.id)).toEqual([c1, s1, u2, a2, u3, a3])
|
||||
}),
|
||||
@@ -800,11 +797,11 @@ describe("MessageV2.filterCompacted", () => {
|
||||
text: "third reply",
|
||||
})
|
||||
|
||||
const parentFiltered = MessageV2.filterCompacted(MessageV2.stream(created.id))
|
||||
const parentFiltered = MessageV2.filterCompacted(yield* MessageV2.stream(created.id))
|
||||
expect(parentFiltered.map((item) => item.info.id)).toEqual([c1, s1, u2, a2, u3, a3])
|
||||
|
||||
const forked = yield* session.fork({ sessionID: created.id })
|
||||
const childFiltered = MessageV2.filterCompacted(MessageV2.stream(forked.id))
|
||||
const childFiltered = MessageV2.filterCompacted(yield* MessageV2.stream(forked.id))
|
||||
expect(childFiltered).toHaveLength(parentFiltered.length)
|
||||
|
||||
const tailPart = childFiltered.flatMap((m) => m.parts).find((p) => p.type === "compaction")
|
||||
@@ -870,7 +867,7 @@ describe("MessageV2.filterCompacted", () => {
|
||||
text: "third reply",
|
||||
})
|
||||
|
||||
const result = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
|
||||
|
||||
expect(result.map((item) => item.info.id)).toEqual([c1, s1, a3, u3, a4])
|
||||
}),
|
||||
@@ -942,7 +939,7 @@ describe("MessageV2.filterCompacted", () => {
|
||||
text: "fourth reply",
|
||||
})
|
||||
|
||||
const result = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
|
||||
|
||||
expect(result.map((item) => item.info.id)).toEqual([c2, s2, u3, a3, u4, a4])
|
||||
}),
|
||||
@@ -1015,7 +1012,7 @@ describe("MessageV2 consistency", () => {
|
||||
const [id] = yield* fill(sessionID, 1)
|
||||
|
||||
const got = yield* MessageV2.get({ sessionID, messageID: id })
|
||||
const standalone = MessageV2.parts(id)
|
||||
const standalone = yield* MessageV2.parts(id)
|
||||
expect(got.parts).toEqual(standalone)
|
||||
}),
|
||||
),
|
||||
@@ -1026,7 +1023,7 @@ describe("MessageV2 consistency", () => {
|
||||
Effect.gen(function* () {
|
||||
yield* fill(sessionID, 7)
|
||||
|
||||
const streamed = Array.from(MessageV2.stream(sessionID))
|
||||
const streamed = yield* MessageV2.stream(sessionID)
|
||||
|
||||
const paged = [] as SessionLegacy.WithParts[]
|
||||
let cursor: string | undefined
|
||||
@@ -1049,8 +1046,9 @@ describe("MessageV2 consistency", () => {
|
||||
Effect.gen(function* () {
|
||||
yield* fill(sessionID, 4)
|
||||
|
||||
const filtered = MessageV2.filterCompacted(MessageV2.stream(sessionID))
|
||||
const all = Array.from(MessageV2.stream(sessionID)).reverse()
|
||||
const stream = yield* MessageV2.stream(sessionID)
|
||||
const filtered = MessageV2.filterCompacted(stream)
|
||||
const all = stream.toReversed()
|
||||
|
||||
expect(filtered.map((m) => m.info.id)).toEqual(all.map((m) => m.info.id))
|
||||
}),
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { NodeFileSystem } from "@effect/platform-node"
|
||||
import { SessionLegacy } from "@opencode-ai/core/session/legacy"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { expect } from "bun:test"
|
||||
import { tool } from "ai"
|
||||
import { Cause, Effect, Exit, Fiber, Layer } from "effect"
|
||||
@@ -185,6 +186,7 @@ const deps = Layer.mergeAll(
|
||||
status,
|
||||
SyncEvent.defaultLayer,
|
||||
EventV2Bridge.defaultLayer,
|
||||
Database.defaultLayer,
|
||||
).pipe(Layer.provideMerge(infra))
|
||||
const env = Layer.mergeAll(
|
||||
TestLLMServer.layer,
|
||||
@@ -213,6 +215,7 @@ it.live("session.processor effect tests capture llm input cleanly", () =>
|
||||
provideTmpdirServer(
|
||||
({ dir, llm }) =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const { processors, session, provider } = yield* boot()
|
||||
|
||||
yield* llm.text("hello")
|
||||
@@ -245,7 +248,7 @@ it.live("session.processor effect tests capture llm input cleanly", () =>
|
||||
} satisfies LLM.StreamInput
|
||||
|
||||
const value = yield* handle.process(input)
|
||||
const parts = MessageV2.parts(msg.id)
|
||||
const parts = yield* MessageV2.parts(msg.id)
|
||||
const calls = yield* llm.calls
|
||||
|
||||
expect(value).toBe("continue")
|
||||
@@ -260,6 +263,7 @@ it.live("session.processor effect tests preserve text start time", () =>
|
||||
provideTmpdirServer(
|
||||
({ dir, llm }) =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const gate = defer<void>()
|
||||
const { processors, session, provider } = yield* boot()
|
||||
|
||||
@@ -318,14 +322,19 @@ it.live("session.processor effect tests preserve text start time", () =>
|
||||
.pipe(Effect.forkChild)
|
||||
|
||||
yield* waitFor(
|
||||
Effect.sync(() => MessageV2.parts(msg.id).find((part): part is SessionLegacy.TextPart => part.type === "text")),
|
||||
MessageV2.parts(msg.id).pipe(
|
||||
Effect.map((parts) => parts.find((part): part is SessionLegacy.TextPart => part.type === "text")),
|
||||
Effect.provideService(Database.Service, database),
|
||||
),
|
||||
"timed out waiting for text part",
|
||||
)
|
||||
yield* Effect.sleep("20 millis")
|
||||
gate.resolve()
|
||||
|
||||
const exit = yield* Fiber.await(run)
|
||||
const text = MessageV2.parts(msg.id).find((part): part is SessionLegacy.TextPart => part.type === "text")
|
||||
const text = (yield* MessageV2.parts(msg.id)).find(
|
||||
(part): part is SessionLegacy.TextPart => part.type === "text",
|
||||
)
|
||||
|
||||
expect(Exit.isSuccess(exit)).toBe(true)
|
||||
expect(text?.text).toBe("hello")
|
||||
@@ -342,6 +351,7 @@ it.live("session.processor effect tests stop after token overflow requests compa
|
||||
provideTmpdirServer(
|
||||
({ dir, llm }) =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const { processors, session, provider } = yield* boot()
|
||||
|
||||
yield* llm.text("after", { usage: { input: 100, output: 0 } })
|
||||
@@ -374,7 +384,7 @@ it.live("session.processor effect tests stop after token overflow requests compa
|
||||
tools: {},
|
||||
})
|
||||
|
||||
const parts = MessageV2.parts(msg.id)
|
||||
const parts = yield* MessageV2.parts(msg.id)
|
||||
|
||||
expect(value).toBe("compact")
|
||||
expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
|
||||
@@ -388,6 +398,7 @@ it.live("session.processor effect tests capture reasoning from http mock", () =>
|
||||
provideTmpdirServer(
|
||||
({ dir, llm }) =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const { processors, session, provider } = yield* boot()
|
||||
|
||||
yield* llm.push(reply().reason("think").text("done").stop())
|
||||
@@ -419,7 +430,7 @@ it.live("session.processor effect tests capture reasoning from http mock", () =>
|
||||
tools: {},
|
||||
})
|
||||
|
||||
const parts = MessageV2.parts(msg.id)
|
||||
const parts = yield* MessageV2.parts(msg.id)
|
||||
const reasoning = parts.find((part): part is SessionLegacy.ReasoningPart => part.type === "reasoning")
|
||||
const text = parts.find((part): part is SessionLegacy.TextPart => part.type === "text")
|
||||
|
||||
@@ -467,7 +478,7 @@ it.live("session.processor effect tests reset reasoning state across retries", (
|
||||
tools: {},
|
||||
})
|
||||
|
||||
const parts = MessageV2.parts(msg.id)
|
||||
const parts = yield* MessageV2.parts(msg.id)
|
||||
const reasoning = parts.filter((part): part is SessionLegacy.ReasoningPart => part.type === "reasoning")
|
||||
|
||||
expect(value).toBe("continue")
|
||||
@@ -558,7 +569,7 @@ it.live("session.processor effect tests retry recognized structured json errors"
|
||||
tools: {},
|
||||
})
|
||||
|
||||
const parts = MessageV2.parts(msg.id)
|
||||
const parts = yield* MessageV2.parts(msg.id)
|
||||
|
||||
expect(value).toBe("continue")
|
||||
expect(yield* llm.calls).toBe(2)
|
||||
@@ -709,7 +720,7 @@ it.live("session.processor effect tests complete AI SDK tool calls when native f
|
||||
},
|
||||
})
|
||||
|
||||
const parts = MessageV2.parts(msg.id)
|
||||
const parts = yield* MessageV2.parts(msg.id)
|
||||
const call = parts.find((part): part is SessionLegacy.ToolPart => part.type === "tool")
|
||||
|
||||
expect(value).toBe("continue")
|
||||
@@ -733,6 +744,7 @@ it.live("session.processor effect tests mark pending tools as aborted on cleanup
|
||||
provideTmpdirServer(
|
||||
({ dir, llm }) =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const { processors, session, provider } = yield* boot()
|
||||
|
||||
yield* llm.toolHang("bash", { cmd: "pwd" })
|
||||
@@ -768,13 +780,16 @@ it.live("session.processor effect tests mark pending tools as aborted on cleanup
|
||||
|
||||
yield* llm.wait(1)
|
||||
yield* waitFor(
|
||||
Effect.sync(() => MessageV2.parts(msg.id).find((part): part is SessionLegacy.ToolPart => part.type === "tool")),
|
||||
MessageV2.parts(msg.id).pipe(
|
||||
Effect.map((parts) => parts.find((part): part is SessionLegacy.ToolPart => part.type === "tool")),
|
||||
Effect.provideService(Database.Service, database),
|
||||
),
|
||||
"timed out waiting for tool part",
|
||||
)
|
||||
yield* Fiber.interrupt(run)
|
||||
|
||||
const exit = yield* Fiber.await(run)
|
||||
const parts = MessageV2.parts(msg.id)
|
||||
const parts = yield* MessageV2.parts(msg.id)
|
||||
const call = parts.find((part): part is SessionLegacy.ToolPart => part.type === "tool")
|
||||
|
||||
expect(Exit.isFailure(exit)).toBe(true)
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { NodeFileSystem } from "@effect/platform-node"
|
||||
import { SessionLegacy } from "@opencode-ai/core/session/legacy"
|
||||
import { Database as CoreDatabase } from "@opencode-ai/core/database/database"
|
||||
import { FetchHttpClient } from "effect/unstable/http"
|
||||
import { expect } from "bun:test"
|
||||
import { Cause, Deferred, Duration, Effect, Exit, Fiber, Layer } from "effect"
|
||||
@@ -184,6 +185,7 @@ function makePrompt(input?: { processor?: "blocking" }) {
|
||||
status,
|
||||
SyncEvent.defaultLayer,
|
||||
EventV2Bridge.defaultLayer,
|
||||
CoreDatabase.defaultLayer,
|
||||
).pipe(Layer.provideMerge(infra))
|
||||
const question = Question.layer.pipe(Layer.provideMerge(deps))
|
||||
const todo = Todo.layer.pipe(Layer.provideMerge(deps))
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { describe, expect } from "bun:test"
|
||||
import { SessionLegacy } from "@opencode-ai/core/session/legacy"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { Deferred, Effect, Exit, Layer } from "effect"
|
||||
import { Session as SessionNs } from "@/session/session"
|
||||
import { GlobalBus, type GlobalEvent } from "../../src/bus/global"
|
||||
@@ -23,6 +24,7 @@ const it = testEffect(
|
||||
Layer.provide(Bus.layer),
|
||||
Layer.provide(Storage.defaultLayer),
|
||||
Layer.provide(SyncEvent.defaultLayer),
|
||||
Layer.provide(Database.defaultLayer),
|
||||
Layer.provide(RuntimeFlags.layer({ experimentalWorkspaces: false })),
|
||||
Layer.provide(BackgroundJob.defaultLayer),
|
||||
),
|
||||
|
||||
@@ -30,6 +30,7 @@ import { TestLLMServer } from "../lib/llm-server"
|
||||
|
||||
// Same layer setup as prompt-effect.test.ts
|
||||
import { NodeFileSystem } from "@effect/platform-node"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { Agent as AgentSvc } from "../../src/agent/agent"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { Git } from "../../src/git"
|
||||
@@ -133,6 +134,7 @@ function makeHttp() {
|
||||
status,
|
||||
SyncEvent.defaultLayer,
|
||||
EventV2Bridge.defaultLayer,
|
||||
Database.defaultLayer,
|
||||
).pipe(Layer.provideMerge(infra))
|
||||
const question = Question.layer.pipe(Layer.provideMerge(deps))
|
||||
const todo = Todo.layer.pipe(Layer.provideMerge(deps))
|
||||
|
||||
@@ -4,6 +4,7 @@ import fs from "fs/promises"
|
||||
import { fileURLToPath, pathToFileURL } from "url"
|
||||
import { Effect, Layer, Result, Schema } from "effect"
|
||||
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { ToolRegistry } from "@/tool/registry"
|
||||
import { Tool } from "@/tool/tool"
|
||||
import { disposeAllInstances, TestInstance } from "../fixture/fixture"
|
||||
@@ -65,7 +66,7 @@ const registryLayer = (opts: RegistryLayerOptions = {}) =>
|
||||
Layer.provide(Bus.layer),
|
||||
Layer.provide(FetchHttpClient.layer),
|
||||
Layer.provide(Format.defaultLayer),
|
||||
Layer.provide(node),
|
||||
Layer.provide(Layer.mergeAll(node, Database.defaultLayer)),
|
||||
Layer.provide(Ripgrep.defaultLayer),
|
||||
Layer.provide(Truncate.defaultLayer),
|
||||
)
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { afterEach, describe, expect } from "bun:test"
|
||||
import { SessionLegacy } from "@opencode-ai/core/session/legacy"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { Effect, Exit, Fiber, Layer } from "effect"
|
||||
import { Agent } from "../../src/agent/agent"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
@@ -41,6 +42,7 @@ const layer = (flags: Partial<RuntimeFlags.Info> = {}) =>
|
||||
SessionStatus.defaultLayer,
|
||||
Truncate.defaultLayer,
|
||||
ToolRegistry.defaultLayer,
|
||||
Database.defaultLayer,
|
||||
RuntimeFlags.layer(flags),
|
||||
)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user