feat(v2): add session storage service
This commit is contained in:
parent
5381795844
commit
e674c242e1
7 changed files with 1182 additions and 189 deletions
|
|
@ -1,9 +1,6 @@
|
|||
import { SessionMessageTable, SessionTable } from "@/session/session.sql"
|
||||
import { SessionID } from "@/session/schema"
|
||||
import { WorkspaceID } from "@/control-plane/schema"
|
||||
import { and, asc, desc, eq, gt, gte, isNull, like, lt, or, type SQL } from "@/storage/db"
|
||||
import * as Database from "@/storage/db"
|
||||
import { Context, DateTime, Effect, Layer, Option, Schema } from "effect"
|
||||
import { Context, DateTime, Effect, Layer, Schema } from "effect"
|
||||
import { SessionMessage } from "@opencode-ai/core/session-message"
|
||||
import type { Prompt } from "@opencode-ai/core/session-prompt"
|
||||
import { ProjectID } from "@/project/schema"
|
||||
|
|
@ -13,7 +10,8 @@ import { optionalOmitUndefined } from "@opencode-ai/core/schema"
|
|||
import { EventV2 } from "@opencode-ai/core/event"
|
||||
import { EventV2Bridge } from "@/event-v2-bridge"
|
||||
import { ModelV2 } from "@opencode-ai/core/model"
|
||||
import { ProviderV2 } from "@opencode-ai/core/provider"
|
||||
import { SessionStorage } from "./session/storage"
|
||||
import { SessionStorageSql } from "./session/storage-sql"
|
||||
|
||||
export const Delivery = Schema.Literals(["immediate", "deferred"]).annotate({
|
||||
identifier: "Session.Delivery",
|
||||
|
|
@ -73,40 +71,17 @@ export interface Interface {
|
|||
workspaceID?: WorkspaceID
|
||||
}) => Effect.Effect<Info>
|
||||
readonly get: (sessionID: SessionID) => Effect.Effect<Info, NotFoundError>
|
||||
readonly list: (input: {
|
||||
limit?: number
|
||||
order?: "asc" | "desc"
|
||||
directory?: string
|
||||
path?: string
|
||||
workspaceID?: WorkspaceID
|
||||
roots?: boolean
|
||||
start?: number
|
||||
search?: string
|
||||
cursor?: {
|
||||
id: SessionID
|
||||
time: number
|
||||
direction: "previous" | "next"
|
||||
}
|
||||
}) => Effect.Effect<Info[], never>
|
||||
readonly messages: (input: {
|
||||
sessionID: SessionID
|
||||
limit?: number
|
||||
order?: "asc" | "desc"
|
||||
cursor?: {
|
||||
id: SessionMessage.ID
|
||||
time: number
|
||||
direction: "previous" | "next"
|
||||
}
|
||||
}) => Effect.Effect<SessionMessage.Message[], never>
|
||||
readonly context: (sessionID: SessionID) => Effect.Effect<SessionMessage.Message[], never>
|
||||
readonly list: (input: SessionStorage.SessionListInput) => Effect.Effect<Info[]>
|
||||
readonly messages: (input: SessionStorage.MessageListInput) => Effect.Effect<SessionMessage.Message[]>
|
||||
readonly context: (sessionID: SessionID) => Effect.Effect<SessionMessage.Message[]>
|
||||
readonly prompt: (input: {
|
||||
id?: EventV2.ID
|
||||
sessionID: SessionID
|
||||
prompt: Prompt
|
||||
delivery?: Delivery
|
||||
}) => Effect.Effect<SessionMessage.User, never>
|
||||
readonly shell: (input: { id?: EventV2.ID; sessionID: SessionID; command: string }) => Effect.Effect<void, never>
|
||||
readonly skill: (input: { id?: EventV2.ID; sessionID: SessionID; skill: string }) => Effect.Effect<void, never>
|
||||
}) => Effect.Effect<SessionMessage.User>
|
||||
readonly shell: (input: { id?: EventV2.ID; sessionID: SessionID; command: string }) => Effect.Effect<void>
|
||||
readonly skill: (input: { id?: EventV2.ID; sessionID: SessionID; skill: string }) => Effect.Effect<void>
|
||||
readonly subagent: (input: {
|
||||
id?: EventV2.ID
|
||||
parentID: SessionID
|
||||
|
|
@ -114,10 +89,10 @@ export interface Interface {
|
|||
agent: string
|
||||
model?: ModelV2.Ref
|
||||
}) => Effect.Effect<void, NotFoundError>
|
||||
readonly switchAgent: (input: { sessionID: SessionID; agent: string }) => Effect.Effect<void, never>
|
||||
readonly switchModel: (input: { sessionID: SessionID; model: ModelV2.Ref }) => Effect.Effect<void, never>
|
||||
readonly compact: (sessionID: SessionID) => Effect.Effect<void, never>
|
||||
readonly wait: (sessionID: SessionID) => Effect.Effect<void, never>
|
||||
readonly switchAgent: (input: { sessionID: SessionID; agent: string }) => Effect.Effect<void>
|
||||
readonly switchModel: (input: { sessionID: SessionID; model: ModelV2.Ref }) => Effect.Effect<void>
|
||||
readonly compact: (sessionID: SessionID) => Effect.Effect<void>
|
||||
readonly wait: (sessionID: SessionID) => Effect.Effect<void>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/v2/Session") {}
|
||||
|
|
@ -126,170 +101,28 @@ export const layer = Layer.effect(
|
|||
Service,
|
||||
Effect.gen(function* () {
|
||||
const events = yield* EventV2Bridge.Service
|
||||
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Message)
|
||||
|
||||
const decode = (row: typeof SessionMessageTable.$inferSelect) =>
|
||||
decodeMessage({ ...row.data, id: row.id, type: row.type })
|
||||
|
||||
function fromRow(row: typeof SessionTable.$inferSelect): Info {
|
||||
return new Info({
|
||||
id: SessionID.make(row.id),
|
||||
projectID: ProjectID.make(row.project_id),
|
||||
workspaceID: row.workspace_id ? WorkspaceID.make(row.workspace_id) : undefined,
|
||||
title: row.title,
|
||||
parentID: row.parent_id ? SessionID.make(row.parent_id) : undefined,
|
||||
path: row.path ?? "",
|
||||
agent: row.agent ?? undefined,
|
||||
model: row.model
|
||||
? {
|
||||
id: ModelV2.ID.make(row.model.id),
|
||||
providerID: ProviderV2.ID.make(row.model.providerID),
|
||||
variant: ModelV2.VariantID.make(row.model.variant ?? "default"),
|
||||
}
|
||||
: undefined,
|
||||
cost: row.cost,
|
||||
tokens: {
|
||||
input: row.tokens_input,
|
||||
output: row.tokens_output,
|
||||
reasoning: row.tokens_reasoning,
|
||||
cache: {
|
||||
read: row.tokens_cache_read,
|
||||
write: row.tokens_cache_write,
|
||||
},
|
||||
},
|
||||
time: {
|
||||
created: DateTime.makeUnsafe(row.time_created),
|
||||
updated: DateTime.makeUnsafe(row.time_updated),
|
||||
archived: row.time_archived ? DateTime.makeUnsafe(row.time_archived) : undefined,
|
||||
},
|
||||
})
|
||||
}
|
||||
const storage = yield* SessionStorage.Service
|
||||
|
||||
const result = Service.of({
|
||||
create: Effect.fn("V2Session.create")(function* (_input) {
|
||||
return {} as any
|
||||
return yield* Effect.die(new Error("V2Session.create is not implemented"))
|
||||
}),
|
||||
get: Effect.fn("V2Session.get")(function* (sessionID) {
|
||||
const row = Database.use((db) => db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get())
|
||||
const row = yield* storage.get(sessionID).pipe(Effect.orDie)
|
||||
if (!row) return yield* new NotFoundError({ sessionID })
|
||||
return fromRow(row)
|
||||
return new Info(row)
|
||||
}),
|
||||
list: Effect.fn("V2Session.list")(function* (input) {
|
||||
const direction = input.cursor?.direction ?? "next"
|
||||
let order = input.order ?? "desc"
|
||||
// This is a load bearing sort, desktop relies on this
|
||||
const sortColumn = SessionTable.time_updated
|
||||
// Query the adjacent rows in reverse, then flip them back into the requested order below.
|
||||
if (direction === "previous" && order === "asc") order = "desc"
|
||||
if (direction === "previous" && order === "desc") order = "asc"
|
||||
const conditions: SQL[] = []
|
||||
if (input.directory) conditions.push(eq(SessionTable.directory, input.directory))
|
||||
if (input.path)
|
||||
conditions.push(or(eq(SessionTable.path, input.path), like(SessionTable.path, `${input.path}/%`))!)
|
||||
if (input.workspaceID) conditions.push(eq(SessionTable.workspace_id, input.workspaceID))
|
||||
if (input.roots) conditions.push(isNull(SessionTable.parent_id))
|
||||
if (input.start) conditions.push(gte(sortColumn, input.start))
|
||||
if (input.search) conditions.push(like(SessionTable.title, `%${input.search}%`))
|
||||
if (input.cursor) {
|
||||
conditions.push(
|
||||
order === "asc"
|
||||
? or(
|
||||
gt(sortColumn, input.cursor.time),
|
||||
and(eq(sortColumn, input.cursor.time), gt(SessionTable.id, input.cursor.id)),
|
||||
)!
|
||||
: or(
|
||||
lt(sortColumn, input.cursor.time),
|
||||
and(eq(sortColumn, input.cursor.time), lt(SessionTable.id, input.cursor.id)),
|
||||
)!,
|
||||
)
|
||||
}
|
||||
const query = Database.Client()
|
||||
.select()
|
||||
.from(SessionTable)
|
||||
.where(conditions.length > 0 ? and(...conditions) : undefined)
|
||||
.orderBy(
|
||||
order === "asc" ? asc(sortColumn) : desc(sortColumn),
|
||||
order === "asc" ? asc(SessionTable.id) : desc(SessionTable.id),
|
||||
)
|
||||
|
||||
const rows = input.limit === undefined ? query.all() : query.limit(input.limit).all()
|
||||
return (direction === "previous" ? rows.toReversed() : rows).map((row) => fromRow(row))
|
||||
return (yield* storage.list(input).pipe(Effect.orDie)).map((row) => new Info(row))
|
||||
}),
|
||||
messages: Effect.fn("V2Session.messages")(function* (input) {
|
||||
const direction = input.cursor?.direction ?? "next"
|
||||
let order = input.order ?? "desc"
|
||||
// Query the adjacent rows in reverse, then flip them back into the requested order below.
|
||||
if (direction === "previous" && order === "asc") order = "desc"
|
||||
if (direction === "previous" && order === "desc") order = "asc"
|
||||
const boundary = input.cursor
|
||||
? order === "asc"
|
||||
? or(
|
||||
gt(SessionMessageTable.time_created, input.cursor.time),
|
||||
and(
|
||||
eq(SessionMessageTable.time_created, input.cursor.time),
|
||||
gt(SessionMessageTable.id, input.cursor.id),
|
||||
),
|
||||
)
|
||||
: or(
|
||||
lt(SessionMessageTable.time_created, input.cursor.time),
|
||||
and(
|
||||
eq(SessionMessageTable.time_created, input.cursor.time),
|
||||
lt(SessionMessageTable.id, input.cursor.id),
|
||||
),
|
||||
)
|
||||
: undefined
|
||||
const where = boundary
|
||||
? and(eq(SessionMessageTable.session_id, input.sessionID), boundary)
|
||||
: eq(SessionMessageTable.session_id, input.sessionID)
|
||||
|
||||
const rows = Database.use((db) => {
|
||||
const query = db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(where)
|
||||
.orderBy(
|
||||
order === "asc" ? asc(SessionMessageTable.time_created) : desc(SessionMessageTable.time_created),
|
||||
order === "asc" ? asc(SessionMessageTable.id) : desc(SessionMessageTable.id),
|
||||
)
|
||||
const rows = input.limit === undefined ? query.all() : query.limit(input.limit).all()
|
||||
return direction === "previous" ? rows.toReversed() : rows
|
||||
})
|
||||
return rows.map((row) => decode(row))
|
||||
return yield* storage.messages(input).pipe(Effect.orDie)
|
||||
}),
|
||||
context: Effect.fn("V2Session.context")(function* (sessionID) {
|
||||
const rows = Database.use((db) => {
|
||||
const compaction = db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(and(eq(SessionMessageTable.session_id, sessionID), eq(SessionMessageTable.type, "compaction")))
|
||||
.orderBy(desc(SessionMessageTable.time_created), desc(SessionMessageTable.id))
|
||||
.limit(1)
|
||||
.get()
|
||||
|
||||
return db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(
|
||||
and(
|
||||
eq(SessionMessageTable.session_id, sessionID),
|
||||
compaction
|
||||
? or(
|
||||
gt(SessionMessageTable.time_created, compaction.time_created),
|
||||
and(
|
||||
eq(SessionMessageTable.time_created, compaction.time_created),
|
||||
gte(SessionMessageTable.id, compaction.id),
|
||||
),
|
||||
)
|
||||
: undefined,
|
||||
),
|
||||
)
|
||||
.orderBy(asc(SessionMessageTable.time_created), asc(SessionMessageTable.id))
|
||||
.all()
|
||||
})
|
||||
return rows.map((row) => decode(row))
|
||||
return yield* storage.context(sessionID).pipe(Effect.orDie)
|
||||
}),
|
||||
prompt: Effect.fn("V2Session.prompt")(function* (_input) {
|
||||
return {} as any
|
||||
return yield* Effect.die(new Error("V2Session.prompt is not implemented"))
|
||||
}),
|
||||
shell: Effect.fn("V2Session.shell")(function* (_input) {}),
|
||||
skill: Effect.fn("V2Session.skill")(function* (_input) {}),
|
||||
|
|
@ -336,6 +169,8 @@ export const layer = Layer.effect(
|
|||
}),
|
||||
)
|
||||
|
||||
export const defaultLayer = layer.pipe(Layer.provide(EventV2Bridge.defaultLayer))
|
||||
export const defaultLayer = layer.pipe(
|
||||
Layer.provide(Layer.mergeAll(EventV2Bridge.defaultLayer, SessionStorageSql.defaultLayer)),
|
||||
)
|
||||
|
||||
export * as SessionV2 from "./session"
|
||||
|
|
|
|||
102
packages/opencode/src/v2/session/storage-memory.ts
Normal file
102
packages/opencode/src/v2/session/storage-memory.ts
Normal file
|
|
@ -0,0 +1,102 @@
|
|||
import { DateTime, Effect, Layer } from "effect"
|
||||
import { SessionMessage } from "@opencode-ai/core/session-message"
|
||||
import { SessionStorage } from "./storage"
|
||||
|
||||
export interface State {
|
||||
readonly sessions: Map<string, SessionStorage.SessionRow>
|
||||
readonly messages: Map<string, SessionMessage.Message[]>
|
||||
}
|
||||
|
||||
export const makeState = (): State => ({
|
||||
sessions: new Map(),
|
||||
messages: new Map(),
|
||||
})
|
||||
|
||||
export const layer = (state: State = makeState()) =>
|
||||
Layer.succeed(
|
||||
SessionStorage.Service,
|
||||
SessionStorage.Service.of({
|
||||
get: (sessionID) => Effect.sync(() => state.sessions.get(sessionID)),
|
||||
list: (input) =>
|
||||
Effect.sync(() => {
|
||||
const direction = input.cursor?.direction ?? "next"
|
||||
const order = SessionStorage.pageOrder(input.order ?? "desc", direction)
|
||||
const rows = Array.from(state.sessions.values())
|
||||
.filter((row) => {
|
||||
if (input.directory && row.directory !== input.directory) return false
|
||||
if (input.path && row.path !== input.path && !row.path?.startsWith(`${input.path}/`)) return false
|
||||
if (input.workspaceID && row.workspaceID !== input.workspaceID) return false
|
||||
if (input.roots && row.parentID) return false
|
||||
if (input.start && DateTime.toEpochMillis(row.time.updated) < input.start) return false
|
||||
if (input.search && !row.title.includes(input.search)) return false
|
||||
if (!input.cursor) return true
|
||||
return compareCursor(row.id, DateTime.toEpochMillis(row.time.updated), input.cursor, order)
|
||||
})
|
||||
.toSorted((a, b) =>
|
||||
compareRows(
|
||||
a.id,
|
||||
DateTime.toEpochMillis(a.time.updated),
|
||||
b.id,
|
||||
DateTime.toEpochMillis(b.time.updated),
|
||||
order,
|
||||
),
|
||||
)
|
||||
const limited = input.limit === undefined ? rows : rows.slice(0, input.limit)
|
||||
return direction === "previous" ? limited.toReversed() : limited
|
||||
}),
|
||||
messages: (input) =>
|
||||
Effect.sync(() => {
|
||||
const direction = input.cursor?.direction ?? "next"
|
||||
const order = SessionStorage.pageOrder(input.order ?? "desc", direction)
|
||||
const rows = (state.messages.get(input.sessionID) ?? [])
|
||||
.filter((message) => {
|
||||
if (!input.cursor) return true
|
||||
return compareCursor(message.id, DateTime.toEpochMillis(message.time.created), input.cursor, order)
|
||||
})
|
||||
.toSorted((a, b) =>
|
||||
compareRows(
|
||||
a.id,
|
||||
DateTime.toEpochMillis(a.time.created),
|
||||
b.id,
|
||||
DateTime.toEpochMillis(b.time.created),
|
||||
order,
|
||||
),
|
||||
)
|
||||
const limited = input.limit === undefined ? rows : rows.slice(0, input.limit)
|
||||
return direction === "previous" ? limited.toReversed() : limited
|
||||
}),
|
||||
context: (sessionID) =>
|
||||
Effect.sync(() => {
|
||||
const messages = (state.messages.get(sessionID) ?? []).toSorted((a, b) =>
|
||||
compareRows(
|
||||
a.id,
|
||||
DateTime.toEpochMillis(a.time.created),
|
||||
b.id,
|
||||
DateTime.toEpochMillis(b.time.created),
|
||||
"asc",
|
||||
),
|
||||
)
|
||||
const index = messages.findLastIndex((message) => message.type === "compaction")
|
||||
return index === -1 ? messages : messages.slice(index)
|
||||
}),
|
||||
}),
|
||||
)
|
||||
|
||||
export const defaultLayer = layer()
|
||||
|
||||
function compareCursor(
|
||||
id: string,
|
||||
time: number,
|
||||
cursor: { readonly id: string; readonly time: number },
|
||||
order: SessionStorage.SortOrder,
|
||||
) {
|
||||
if (order === "asc") return time > cursor.time || (time === cursor.time && id > cursor.id)
|
||||
return time < cursor.time || (time === cursor.time && id < cursor.id)
|
||||
}
|
||||
|
||||
function compareRows(aID: string, aTime: number, bID: string, bTime: number, order: SessionStorage.SortOrder) {
|
||||
const result = aTime === bTime ? aID.localeCompare(bID) : aTime - bTime
|
||||
return order === "asc" ? result : -result
|
||||
}
|
||||
|
||||
export * as SessionStorageMemory from "./storage-memory"
|
||||
173
packages/opencode/src/v2/session/storage-sql.ts
Normal file
173
packages/opencode/src/v2/session/storage-sql.ts
Normal file
|
|
@ -0,0 +1,173 @@
|
|||
import { SessionMessageTable, SessionTable } from "@/session/session.sql"
|
||||
import { and, asc, Database, desc, eq, gt, gte, isNull, like, lt, or, type SQL } from "@/storage/db"
|
||||
import { SessionMessage } from "@opencode-ai/core/session-message"
|
||||
import { Effect, Layer, Schema } from "effect"
|
||||
import { SessionStorage } from "./storage"
|
||||
|
||||
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Message)
|
||||
const decodeSessionRow = Schema.decodeUnknownSync(SessionStorage.SessionRow)
|
||||
|
||||
export const layer = Layer.effect(
|
||||
SessionStorage.Service,
|
||||
Effect.gen(function* () {
|
||||
const get: SessionStorage.Interface["get"] = Effect.fn("SessionStorageSql.get")((sessionID) =>
|
||||
attempt(() =>
|
||||
Database.use((db) => db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get()),
|
||||
).pipe(Effect.map((row) => (row ? fromSessionRow(row) : undefined))),
|
||||
)
|
||||
|
||||
const list: SessionStorage.Interface["list"] = Effect.fn("SessionStorageSql.list")((input) =>
|
||||
attempt(() => {
|
||||
const direction = input.cursor?.direction ?? "next"
|
||||
const order = SessionStorage.pageOrder(input.order ?? "desc", direction)
|
||||
const sortColumn = SessionTable.time_updated
|
||||
const conditions: SQL[] = []
|
||||
if (input.directory) conditions.push(eq(SessionTable.directory, input.directory))
|
||||
if (input.path)
|
||||
conditions.push(or(eq(SessionTable.path, input.path), like(SessionTable.path, `${input.path}/%`))!)
|
||||
if (input.workspaceID) conditions.push(eq(SessionTable.workspace_id, input.workspaceID))
|
||||
if (input.roots) conditions.push(isNull(SessionTable.parent_id))
|
||||
if (input.start) conditions.push(gte(sortColumn, input.start))
|
||||
if (input.search) conditions.push(like(SessionTable.title, `%${input.search}%`))
|
||||
if (input.cursor) conditions.push(sessionCursorBoundary(input.cursor, order))
|
||||
|
||||
return Database.use((db) => {
|
||||
const query = db
|
||||
.select()
|
||||
.from(SessionTable)
|
||||
.where(conditions.length > 0 ? and(...conditions) : undefined)
|
||||
.orderBy(
|
||||
order === "asc" ? asc(sortColumn) : desc(sortColumn),
|
||||
order === "asc" ? asc(SessionTable.id) : desc(SessionTable.id),
|
||||
)
|
||||
const rows = input.limit === undefined ? query.all() : query.limit(input.limit).all()
|
||||
return direction === "previous" ? rows.toReversed() : rows
|
||||
})
|
||||
}).pipe(Effect.map((rows) => rows.map(fromSessionRow))),
|
||||
)
|
||||
|
||||
const messages: SessionStorage.Interface["messages"] = Effect.fn("SessionStorageSql.messages")((input) =>
|
||||
attempt(() => {
|
||||
const direction = input.cursor?.direction ?? "next"
|
||||
const order = SessionStorage.pageOrder(input.order ?? "desc", direction)
|
||||
const boundary = input.cursor ? messageCursorBoundary(input.cursor, order) : undefined
|
||||
const where = boundary
|
||||
? and(eq(SessionMessageTable.session_id, input.sessionID), boundary)
|
||||
: eq(SessionMessageTable.session_id, input.sessionID)
|
||||
|
||||
return Database.use((db) => {
|
||||
const query = db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(where)
|
||||
.orderBy(
|
||||
order === "asc" ? asc(SessionMessageTable.time_created) : desc(SessionMessageTable.time_created),
|
||||
order === "asc" ? asc(SessionMessageTable.id) : desc(SessionMessageTable.id),
|
||||
)
|
||||
const rows = input.limit === undefined ? query.all() : query.limit(input.limit).all()
|
||||
return direction === "previous" ? rows.toReversed() : rows
|
||||
})
|
||||
}).pipe(Effect.map((rows) => rows.map((row) => decodeMessage({ ...row.data, id: row.id, type: row.type })))),
|
||||
)
|
||||
|
||||
const context: SessionStorage.Interface["context"] = Effect.fn("SessionStorageSql.context")((sessionID) =>
|
||||
attempt(() =>
|
||||
Database.use((db) => {
|
||||
const compaction = db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(and(eq(SessionMessageTable.session_id, sessionID), eq(SessionMessageTable.type, "compaction")))
|
||||
.orderBy(desc(SessionMessageTable.time_created), desc(SessionMessageTable.id))
|
||||
.limit(1)
|
||||
.get()
|
||||
|
||||
return db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(
|
||||
and(
|
||||
eq(SessionMessageTable.session_id, sessionID),
|
||||
compaction
|
||||
? or(
|
||||
gt(SessionMessageTable.time_created, compaction.time_created),
|
||||
and(
|
||||
eq(SessionMessageTable.time_created, compaction.time_created),
|
||||
gte(SessionMessageTable.id, compaction.id),
|
||||
),
|
||||
)
|
||||
: undefined,
|
||||
),
|
||||
)
|
||||
.orderBy(asc(SessionMessageTable.time_created), asc(SessionMessageTable.id))
|
||||
.all()
|
||||
}),
|
||||
).pipe(Effect.map((rows) => rows.map((row) => decodeMessage({ ...row.data, id: row.id, type: row.type })))),
|
||||
)
|
||||
|
||||
return SessionStorage.Service.of({ get, list, messages, context })
|
||||
}),
|
||||
)
|
||||
|
||||
export const defaultLayer = layer
|
||||
|
||||
function attempt<A>(body: () => A) {
|
||||
return Effect.try({
|
||||
try: body,
|
||||
catch: (cause) => new SessionStorage.StorageError({ message: "Session storage SQL operation failed", cause }),
|
||||
})
|
||||
}
|
||||
|
||||
function sessionCursorBoundary(cursor: SessionStorage.SessionCursor, order: SessionStorage.SortOrder) {
|
||||
if (order === "asc")
|
||||
return or(
|
||||
gt(SessionTable.time_updated, cursor.time),
|
||||
and(eq(SessionTable.time_updated, cursor.time), gt(SessionTable.id, cursor.id)),
|
||||
)!
|
||||
return or(
|
||||
lt(SessionTable.time_updated, cursor.time),
|
||||
and(eq(SessionTable.time_updated, cursor.time), lt(SessionTable.id, cursor.id)),
|
||||
)!
|
||||
}
|
||||
|
||||
function messageCursorBoundary(cursor: SessionStorage.MessageCursor, order: SessionStorage.SortOrder) {
|
||||
if (order === "asc")
|
||||
return or(
|
||||
gt(SessionMessageTable.time_created, cursor.time),
|
||||
and(eq(SessionMessageTable.time_created, cursor.time), gt(SessionMessageTable.id, cursor.id)),
|
||||
)!
|
||||
return or(
|
||||
lt(SessionMessageTable.time_created, cursor.time),
|
||||
and(eq(SessionMessageTable.time_created, cursor.time), lt(SessionMessageTable.id, cursor.id)),
|
||||
)!
|
||||
}
|
||||
|
||||
function fromSessionRow(row: typeof SessionTable.$inferSelect) {
|
||||
return decodeSessionRow({
|
||||
id: row.id,
|
||||
parentID: row.parent_id ?? undefined,
|
||||
projectID: row.project_id,
|
||||
workspaceID: row.workspace_id ?? undefined,
|
||||
directory: row.directory,
|
||||
title: row.title,
|
||||
path: row.path ?? "",
|
||||
agent: row.agent ?? undefined,
|
||||
model: row.model ? { ...row.model, variant: row.model.variant ?? "default" } : undefined,
|
||||
cost: row.cost,
|
||||
tokens: {
|
||||
input: row.tokens_input,
|
||||
output: row.tokens_output,
|
||||
reasoning: row.tokens_reasoning,
|
||||
cache: {
|
||||
read: row.tokens_cache_read,
|
||||
write: row.tokens_cache_write,
|
||||
},
|
||||
},
|
||||
time: {
|
||||
created: row.time_created,
|
||||
updated: row.time_updated,
|
||||
archived: row.time_archived ?? undefined,
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
export * as SessionStorageSql from "./storage-sql"
|
||||
100
packages/opencode/src/v2/session/storage.ts
Normal file
100
packages/opencode/src/v2/session/storage.ts
Normal file
|
|
@ -0,0 +1,100 @@
|
|||
import { WorkspaceID } from "@/control-plane/schema"
|
||||
import { ProjectID } from "@/project/schema"
|
||||
import { SessionID } from "@/session/schema"
|
||||
import { V2Schema } from "@opencode-ai/core/v2-schema"
|
||||
import { SessionMessage } from "@opencode-ai/core/session-message"
|
||||
import { ModelV2 } from "@opencode-ai/core/model"
|
||||
import { Context, Effect, Schema } from "effect"
|
||||
|
||||
export const SortOrder = Schema.Literals(["asc", "desc"]).annotate({
|
||||
identifier: "SortOrder",
|
||||
})
|
||||
export type SortOrder = typeof SortOrder.Type
|
||||
|
||||
export const PageDirection = Schema.Literals(["previous", "next"]).annotate({
|
||||
identifier: "PageDirection",
|
||||
})
|
||||
export type PageDirection = typeof PageDirection.Type
|
||||
|
||||
export class StorageError extends Schema.TaggedErrorClass<StorageError>()("StorageError", {
|
||||
message: Schema.String,
|
||||
cause: Schema.Defect,
|
||||
}) {}
|
||||
|
||||
export class SessionRow extends Schema.Class<SessionRow>("SessionRow")({
|
||||
id: SessionID,
|
||||
parentID: Schema.optional(SessionID),
|
||||
projectID: ProjectID,
|
||||
workspaceID: Schema.optional(WorkspaceID),
|
||||
directory: Schema.optional(Schema.String),
|
||||
path: Schema.optional(Schema.String),
|
||||
agent: Schema.optional(Schema.String),
|
||||
model: Schema.optional(ModelV2.Ref),
|
||||
cost: Schema.Finite,
|
||||
tokens: Schema.Struct({
|
||||
input: Schema.Finite,
|
||||
output: Schema.Finite,
|
||||
reasoning: Schema.Finite,
|
||||
cache: Schema.Struct({
|
||||
read: Schema.Finite,
|
||||
write: Schema.Finite,
|
||||
}),
|
||||
}),
|
||||
time: Schema.Struct({
|
||||
created: V2Schema.DateTimeUtcFromMillis,
|
||||
updated: V2Schema.DateTimeUtcFromMillis,
|
||||
archived: Schema.optional(V2Schema.DateTimeUtcFromMillis),
|
||||
}),
|
||||
title: Schema.String,
|
||||
}) {}
|
||||
|
||||
export const SessionCursor = Schema.Struct({
|
||||
id: SessionID,
|
||||
time: Schema.Finite,
|
||||
direction: PageDirection,
|
||||
}).annotate({ identifier: "SessionCursor" })
|
||||
export type SessionCursor = typeof SessionCursor.Type
|
||||
|
||||
export const SessionListInput = Schema.Struct({
|
||||
limit: Schema.optional(Schema.Finite),
|
||||
order: Schema.optional(SortOrder),
|
||||
directory: Schema.optional(Schema.String),
|
||||
path: Schema.optional(Schema.String),
|
||||
workspaceID: Schema.optional(WorkspaceID),
|
||||
roots: Schema.optional(Schema.Boolean),
|
||||
start: Schema.optional(Schema.Finite),
|
||||
search: Schema.optional(Schema.String),
|
||||
cursor: Schema.optional(SessionCursor),
|
||||
}).annotate({ identifier: "SessionListInput" })
|
||||
export type SessionListInput = typeof SessionListInput.Type
|
||||
|
||||
export const MessageCursor = Schema.Struct({
|
||||
id: SessionMessage.ID,
|
||||
time: Schema.Finite,
|
||||
direction: PageDirection,
|
||||
}).annotate({ identifier: "MessageCursor" })
|
||||
export type MessageCursor = typeof MessageCursor.Type
|
||||
|
||||
export const MessageListInput = Schema.Struct({
|
||||
sessionID: SessionID,
|
||||
limit: Schema.optional(Schema.Finite),
|
||||
order: Schema.optional(SortOrder),
|
||||
cursor: Schema.optional(MessageCursor),
|
||||
}).annotate({ identifier: "MessageListInput" })
|
||||
export type MessageListInput = typeof MessageListInput.Type
|
||||
|
||||
export interface Interface {
|
||||
readonly get: (sessionID: SessionID) => Effect.Effect<SessionRow | undefined, StorageError>
|
||||
readonly list: (input: SessionListInput) => Effect.Effect<SessionRow[], StorageError>
|
||||
readonly messages: (input: MessageListInput) => Effect.Effect<SessionMessage.Message[], StorageError>
|
||||
readonly context: (sessionID: SessionID) => Effect.Effect<SessionMessage.Message[], StorageError>
|
||||
}
|
||||
|
||||
export function pageOrder(order: SortOrder, direction: PageDirection) {
|
||||
if (direction !== "previous") return order
|
||||
return order === "asc" ? "desc" : "asc"
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/v2/session/Storage") {}
|
||||
|
||||
export * as SessionStorage from "./storage"
|
||||
|
|
@ -36,6 +36,7 @@ import { SessionRunState } from "../../src/session/run-state"
|
|||
import { MessageID, PartID, SessionID } from "../../src/session/schema"
|
||||
import { SessionStatus } from "../../src/session/status"
|
||||
import { SessionV2 } from "../../src/v2/session"
|
||||
import { SessionStorageSql } from "../../src/v2/session/storage-sql"
|
||||
import { Skill } from "../../src/skill"
|
||||
import { SystemPrompt } from "../../src/session/system"
|
||||
import { Shell } from "../../src/shell/shell"
|
||||
|
|
@ -507,6 +508,7 @@ noLLMServer.instance(
|
|||
|
||||
const messages = yield* SessionV2.Service.use((session) => session.messages({ sessionID: chat.id })).pipe(
|
||||
Effect.provide(SessionV2.layer),
|
||||
Effect.provide(SessionStorageSql.defaultLayer),
|
||||
)
|
||||
const row = Database.use((db) =>
|
||||
db.select().from(SessionMessageTable).where(Database.eq(SessionMessageTable.session_id, chat.id)).get(),
|
||||
|
|
|
|||
307
packages/opencode/test/v2/session-storage.test.ts
Normal file
307
packages/opencode/test/v2/session-storage.test.ts
Normal file
|
|
@ -0,0 +1,307 @@
|
|||
import { expect } from "bun:test"
|
||||
import { ProjectID } from "@/project/schema"
|
||||
import { ProjectTable } from "@/project/project.sql"
|
||||
import { SessionID } from "@/session/schema"
|
||||
import { SessionMessageTable, SessionTable } from "@/session/session.sql"
|
||||
import { Database } from "@/storage/db"
|
||||
import { SessionStorage } from "@/v2/session/storage"
|
||||
import { SessionStorageMemory } from "@/v2/session/storage-memory"
|
||||
import { SessionStorageSql } from "@/v2/session/storage-sql"
|
||||
import { EventV2 } from "@opencode-ai/core/event"
|
||||
import { SessionMessage } from "@opencode-ai/core/session-message"
|
||||
import { eq, or } from "@/storage/db"
|
||||
import { DateTime, Effect, Layer, Schema } from "effect"
|
||||
import { testEffect } from "../lib/effect"
|
||||
|
||||
const projectID = ProjectID.make("project-session-storage")
|
||||
const sessionA = SessionID.make("ses_storage_a")
|
||||
const sessionB = SessionID.make("ses_storage_b")
|
||||
const sessionC = SessionID.make("ses_storage_c")
|
||||
const encodeMessage = Schema.encodeSync(SessionMessage.Message)
|
||||
const memoryState = SessionStorageMemory.makeState()
|
||||
|
||||
interface Seeds<R> {
|
||||
readonly reset: Effect.Effect<void, never, R>
|
||||
readonly project: Effect.Effect<void, never, R>
|
||||
readonly session: (input: {
|
||||
id: SessionID
|
||||
title: string
|
||||
directory?: string
|
||||
path: string
|
||||
updated: number
|
||||
}) => Effect.Effect<void, never, R>
|
||||
readonly userMessage: (input: { id: SessionMessage.ID; text: string; time: number }) => Effect.Effect<void, never, R>
|
||||
readonly compaction: (input: {
|
||||
id: SessionMessage.ID
|
||||
summary: string
|
||||
time: number
|
||||
}) => Effect.Effect<void, never, R>
|
||||
}
|
||||
|
||||
function sessionStorageContract<R, E>(name: string, layer: Layer.Layer<SessionStorage.Service | R, E>, seed: Seeds<R>) {
|
||||
const it = testEffect(layer)
|
||||
|
||||
const setup = Effect.gen(function* () {
|
||||
yield* seed.reset
|
||||
yield* Effect.addFinalizer(() => seed.reset)
|
||||
yield* seed.project
|
||||
})
|
||||
|
||||
it.effect("gets and lists sessions with filters and cursors", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
yield* seed.session({
|
||||
id: sessionA,
|
||||
title: "Alpha",
|
||||
directory: "/tmp/project-session-storage",
|
||||
path: "apps/api",
|
||||
updated: 1000,
|
||||
})
|
||||
yield* seed.session({
|
||||
id: sessionB,
|
||||
title: "Beta",
|
||||
directory: "/tmp/project-session-storage",
|
||||
path: "apps/web",
|
||||
updated: 2000,
|
||||
})
|
||||
yield* seed.session({
|
||||
id: sessionC,
|
||||
title: "Gamma",
|
||||
directory: "/tmp/other-project",
|
||||
path: "docs",
|
||||
updated: 3000,
|
||||
})
|
||||
|
||||
const storage = yield* SessionStorage.Service
|
||||
const found = yield* storage.get(sessionB)
|
||||
expect(found?.title).toBe("Beta")
|
||||
expect(found ? DateTime.toEpochMillis(found.time.updated) : undefined).toBe(2000)
|
||||
|
||||
expect((yield* storage.list({ path: "apps", order: "asc" })).map((row) => row.id)).toEqual([sessionA, sessionB])
|
||||
expect(
|
||||
(yield* storage.list({ directory: "/tmp/project-session-storage", order: "asc" })).map((row) => row.id),
|
||||
).toEqual([sessionA, sessionB])
|
||||
expect(
|
||||
(yield* storage.list({ order: "asc", cursor: { id: sessionA, time: 1000, direction: "next" } })).map(
|
||||
(row) => row.id,
|
||||
),
|
||||
).toEqual([sessionB, sessionC])
|
||||
expect(
|
||||
(yield* storage.list({ order: "asc", cursor: { id: sessionC, time: 3000, direction: "previous" } })).map(
|
||||
(row) => row.id,
|
||||
),
|
||||
).toEqual([sessionA, sessionB])
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("lists session messages with cursor direction", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
yield* seed.session({ id: sessionA, title: "Alpha", path: "apps/api", updated: 1000 })
|
||||
yield* seed.userMessage({ id: EventV2.ID.make("evt_msg_1"), text: "one", time: 1000 })
|
||||
yield* seed.userMessage({ id: EventV2.ID.make("evt_msg_2"), text: "two", time: 2000 })
|
||||
yield* seed.userMessage({ id: EventV2.ID.make("evt_msg_3"), text: "three", time: 3000 })
|
||||
|
||||
const storage = yield* SessionStorage.Service
|
||||
|
||||
expect((yield* storage.messages({ sessionID: sessionA, order: "asc", limit: 2 })).map((row) => row.id)).toEqual([
|
||||
EventV2.ID.make("evt_msg_1"),
|
||||
EventV2.ID.make("evt_msg_2"),
|
||||
])
|
||||
expect(
|
||||
(yield* storage.messages({
|
||||
sessionID: sessionA,
|
||||
order: "asc",
|
||||
cursor: { id: EventV2.ID.make("evt_msg_3"), time: 3000, direction: "previous" },
|
||||
})).map((row) => row.id),
|
||||
).toEqual([EventV2.ID.make("evt_msg_1"), EventV2.ID.make("evt_msg_2")])
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("returns context from the latest compaction boundary", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
yield* seed.session({ id: sessionA, title: "Alpha", path: "apps/api", updated: 1000 })
|
||||
yield* seed.userMessage({ id: EventV2.ID.make("evt_context_1"), text: "before", time: 1000 })
|
||||
yield* seed.compaction({ id: EventV2.ID.make("evt_context_2"), summary: "compact", time: 2000 })
|
||||
yield* seed.userMessage({ id: EventV2.ID.make("evt_context_3"), text: "after", time: 3000 })
|
||||
|
||||
const storage = yield* SessionStorage.Service
|
||||
const context = yield* storage.context(sessionA)
|
||||
|
||||
expect(context.map((message) => message.id)).toEqual([
|
||||
EventV2.ID.make("evt_context_2"),
|
||||
EventV2.ID.make("evt_context_3"),
|
||||
])
|
||||
expect(context.map((message) => message.type)).toEqual(["compaction", "user"])
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
const sqlSeeds: Seeds<never> = {
|
||||
reset: Effect.sync(resetSqlSeeds),
|
||||
project: Effect.sync(seedProject),
|
||||
session: (input) => Effect.sync(() => seedSession(input)),
|
||||
userMessage: (input) => Effect.sync(() => seedUserMessage(input)),
|
||||
compaction: (input) => Effect.sync(() => seedCompaction(input)),
|
||||
}
|
||||
|
||||
sessionStorageContract("SessionStorageSql", SessionStorageSql.defaultLayer, sqlSeeds)
|
||||
|
||||
const memorySeeds: Seeds<never> = {
|
||||
reset: Effect.sync(() => {
|
||||
memoryState.sessions.clear()
|
||||
memoryState.messages.clear()
|
||||
}),
|
||||
project: Effect.void,
|
||||
session: (input) =>
|
||||
Effect.sync(() => {
|
||||
memoryState.sessions.set(input.id, makeSessionRow(input))
|
||||
}),
|
||||
userMessage: (input) =>
|
||||
Effect.sync(() => {
|
||||
appendMemoryMessage(makeUserMessage(input))
|
||||
}),
|
||||
compaction: (input) =>
|
||||
Effect.sync(() => {
|
||||
appendMemoryMessage(makeCompaction(input))
|
||||
}),
|
||||
}
|
||||
|
||||
sessionStorageContract("SessionStorageMemory", SessionStorageMemory.layer(memoryState), memorySeeds)
|
||||
|
||||
function seedProject() {
|
||||
Database.use((db) =>
|
||||
db
|
||||
.insert(ProjectTable)
|
||||
.values({
|
||||
id: projectID,
|
||||
worktree: "/tmp/project-session-storage",
|
||||
time_created: 1,
|
||||
time_updated: 1,
|
||||
sandboxes: [],
|
||||
})
|
||||
.onConflictDoNothing()
|
||||
.run(),
|
||||
)
|
||||
}
|
||||
|
||||
function resetSqlSeeds() {
|
||||
Database.use((db) => {
|
||||
db.delete(SessionMessageTable).where(eq(SessionMessageTable.session_id, sessionA)).run()
|
||||
db.delete(SessionTable)
|
||||
.where(or(eq(SessionTable.id, sessionA), eq(SessionTable.id, sessionB), eq(SessionTable.id, sessionC)))
|
||||
.run()
|
||||
db.delete(ProjectTable).where(eq(ProjectTable.id, projectID)).run()
|
||||
})
|
||||
}
|
||||
|
||||
function seedSession(input: { id: SessionID; title: string; directory?: string; path: string; updated: number }) {
|
||||
Database.use((db) =>
|
||||
db
|
||||
.insert(SessionTable)
|
||||
.values({
|
||||
id: input.id,
|
||||
project_id: projectID,
|
||||
slug: input.title.toLowerCase(),
|
||||
directory: input.directory ?? "/tmp/project-session-storage",
|
||||
path: input.path,
|
||||
title: input.title,
|
||||
version: "test",
|
||||
cost: 0,
|
||||
tokens_input: 0,
|
||||
tokens_output: 0,
|
||||
tokens_reasoning: 0,
|
||||
tokens_cache_read: 0,
|
||||
tokens_cache_write: 0,
|
||||
time_created: input.updated,
|
||||
time_updated: input.updated,
|
||||
})
|
||||
.run(),
|
||||
)
|
||||
}
|
||||
|
||||
function seedUserMessage(input: { id: SessionMessage.ID; text: string; time: number }) {
|
||||
const encoded = encodeMessage(makeUserMessage(input))
|
||||
const { id: _, type: __, ...data } = encoded
|
||||
seedMessage(input.id, "user", input.time, data)
|
||||
}
|
||||
|
||||
function seedCompaction(input: { id: SessionMessage.ID; summary: string; time: number }) {
|
||||
const encoded = encodeMessage(makeCompaction(input))
|
||||
const { id: _, type: __, ...data } = encoded
|
||||
seedMessage(input.id, "compaction", input.time, data)
|
||||
}
|
||||
|
||||
function makeSessionRow(input: { id: SessionID; title: string; directory?: string; path: string; updated: number }) {
|
||||
return new SessionStorage.SessionRow({
|
||||
id: input.id,
|
||||
projectID,
|
||||
title: input.title,
|
||||
directory: input.directory ?? "/tmp/project-session-storage",
|
||||
path: input.path,
|
||||
cost: 0,
|
||||
tokens: {
|
||||
input: 0,
|
||||
output: 0,
|
||||
reasoning: 0,
|
||||
cache: {
|
||||
read: 0,
|
||||
write: 0,
|
||||
},
|
||||
},
|
||||
time: {
|
||||
created: DateTime.makeUnsafe(input.updated),
|
||||
updated: DateTime.makeUnsafe(input.updated),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
function makeUserMessage(input: { id: SessionMessage.ID; text: string; time: number }) {
|
||||
return new SessionMessage.User({
|
||||
id: input.id,
|
||||
type: "user",
|
||||
text: input.text,
|
||||
files: [],
|
||||
agents: [],
|
||||
references: [],
|
||||
time: { created: DateTime.makeUnsafe(input.time) },
|
||||
})
|
||||
}
|
||||
|
||||
function makeCompaction(input: { id: SessionMessage.ID; summary: string; time: number }) {
|
||||
return new SessionMessage.Compaction({
|
||||
id: input.id,
|
||||
type: "compaction",
|
||||
reason: "manual",
|
||||
summary: input.summary,
|
||||
time: { created: DateTime.makeUnsafe(input.time) },
|
||||
})
|
||||
}
|
||||
|
||||
function seedMessage(
|
||||
id: SessionMessage.ID,
|
||||
type: SessionMessage.Type,
|
||||
time: number,
|
||||
data: typeof SessionMessageTable.$inferInsert.data,
|
||||
) {
|
||||
Database.use((db) =>
|
||||
db
|
||||
.insert(SessionMessageTable)
|
||||
.values({
|
||||
id,
|
||||
session_id: sessionA,
|
||||
type,
|
||||
time_created: time,
|
||||
time_updated: time,
|
||||
data,
|
||||
})
|
||||
.run(),
|
||||
)
|
||||
}
|
||||
|
||||
function appendMemoryMessage(message: SessionMessage.Message) {
|
||||
const current = memoryState.messages.get(sessionA) ?? []
|
||||
current.push(message)
|
||||
memoryState.messages.set(sessionA, current)
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue