Compare commits

...

1 Commits

Author SHA1 Message Date
Dax Raad 402e5999da feat(core): implement v2 session forking
Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>
2026-06-28 20:59:01 +00:00
7 changed files with 454 additions and 86 deletions
+85 -73
View File
@@ -86,45 +86,56 @@ const Endpoint3_3 = (raw: RawClient["server.session"]) => (input: Endpoint3_3Inp
Effect.map((value) => value.data),
)
type Endpoint3_4Request = Parameters<RawClient["server.session"]["session.switchAgent"]>[0]
type Endpoint3_4Request = Parameters<RawClient["server.session"]["session.fork"]>[0]
type Endpoint3_4Input = {
readonly sessionID: Endpoint3_4Request["params"]["sessionID"]
readonly agent: Endpoint3_4Request["payload"]["agent"]
readonly messageID?: Endpoint3_4Request["payload"]["messageID"]
}
const Endpoint3_4 = (raw: RawClient["server.session"]) => (input: Endpoint3_4Input) =>
raw["session.fork"]({ params: { sessionID: input["sessionID"] }, payload: { messageID: input["messageID"] } }).pipe(
Effect.mapError(mapClientError),
Effect.map((value) => value.data),
)
type Endpoint3_5Request = Parameters<RawClient["server.session"]["session.switchAgent"]>[0]
type Endpoint3_5Input = {
readonly sessionID: Endpoint3_5Request["params"]["sessionID"]
readonly agent: Endpoint3_5Request["payload"]["agent"]
}
const Endpoint3_5 = (raw: RawClient["server.session"]) => (input: Endpoint3_5Input) =>
raw["session.switchAgent"]({ params: { sessionID: input["sessionID"] }, payload: { agent: input["agent"] } }).pipe(
Effect.mapError(mapClientError),
)
type Endpoint3_5Request = Parameters<RawClient["server.session"]["session.switchModel"]>[0]
type Endpoint3_5Input = {
readonly sessionID: Endpoint3_5Request["params"]["sessionID"]
readonly model: Endpoint3_5Request["payload"]["model"]
type Endpoint3_6Request = Parameters<RawClient["server.session"]["session.switchModel"]>[0]
type Endpoint3_6Input = {
readonly sessionID: Endpoint3_6Request["params"]["sessionID"]
readonly model: Endpoint3_6Request["payload"]["model"]
}
const Endpoint3_5 = (raw: RawClient["server.session"]) => (input: Endpoint3_5Input) =>
const Endpoint3_6 = (raw: RawClient["server.session"]) => (input: Endpoint3_6Input) =>
raw["session.switchModel"]({ params: { sessionID: input["sessionID"] }, payload: { model: input["model"] } }).pipe(
Effect.mapError(mapClientError),
)
type Endpoint3_6Request = Parameters<RawClient["server.session"]["session.rename"]>[0]
type Endpoint3_6Input = {
readonly sessionID: Endpoint3_6Request["params"]["sessionID"]
readonly title: Endpoint3_6Request["payload"]["title"]
type Endpoint3_7Request = Parameters<RawClient["server.session"]["session.rename"]>[0]
type Endpoint3_7Input = {
readonly sessionID: Endpoint3_7Request["params"]["sessionID"]
readonly title: Endpoint3_7Request["payload"]["title"]
}
const Endpoint3_6 = (raw: RawClient["server.session"]) => (input: Endpoint3_6Input) =>
const Endpoint3_7 = (raw: RawClient["server.session"]) => (input: Endpoint3_7Input) =>
raw["session.rename"]({ params: { sessionID: input["sessionID"] }, payload: { title: input["title"] } }).pipe(
Effect.mapError(mapClientError),
)
type Endpoint3_7Request = Parameters<RawClient["server.session"]["session.prompt"]>[0]
type Endpoint3_7Input = {
readonly sessionID: Endpoint3_7Request["params"]["sessionID"]
readonly id?: Endpoint3_7Request["payload"]["id"]
readonly prompt: Endpoint3_7Request["payload"]["prompt"]
readonly delivery?: Endpoint3_7Request["payload"]["delivery"]
readonly resume?: Endpoint3_7Request["payload"]["resume"]
type Endpoint3_8Request = Parameters<RawClient["server.session"]["session.prompt"]>[0]
type Endpoint3_8Input = {
readonly sessionID: Endpoint3_8Request["params"]["sessionID"]
readonly id?: Endpoint3_8Request["payload"]["id"]
readonly prompt: Endpoint3_8Request["payload"]["prompt"]
readonly delivery?: Endpoint3_8Request["payload"]["delivery"]
readonly resume?: Endpoint3_8Request["payload"]["resume"]
}
const Endpoint3_7 = (raw: RawClient["server.session"]) => (input: Endpoint3_7Input) =>
const Endpoint3_8 = (raw: RawClient["server.session"]) => (input: Endpoint3_8Input) =>
raw["session.prompt"]({
params: { sessionID: input["sessionID"] },
payload: { id: input["id"], prompt: input["prompt"], delivery: input["delivery"], resume: input["resume"] },
@@ -133,23 +144,23 @@ const Endpoint3_7 = (raw: RawClient["server.session"]) => (input: Endpoint3_7Inp
Effect.map((value) => value.data),
)
type Endpoint3_8Request = Parameters<RawClient["server.session"]["session.compact"]>[0]
type Endpoint3_8Input = { readonly sessionID: Endpoint3_8Request["params"]["sessionID"] }
const Endpoint3_8 = (raw: RawClient["server.session"]) => (input: Endpoint3_8Input) =>
raw["session.compact"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError))
type Endpoint3_9Request = Parameters<RawClient["server.session"]["session.wait"]>[0]
type Endpoint3_9Request = Parameters<RawClient["server.session"]["session.compact"]>[0]
type Endpoint3_9Input = { readonly sessionID: Endpoint3_9Request["params"]["sessionID"] }
const Endpoint3_9 = (raw: RawClient["server.session"]) => (input: Endpoint3_9Input) =>
raw["session.compact"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError))
type Endpoint3_10Request = Parameters<RawClient["server.session"]["session.wait"]>[0]
type Endpoint3_10Input = { readonly sessionID: Endpoint3_10Request["params"]["sessionID"] }
const Endpoint3_10 = (raw: RawClient["server.session"]) => (input: Endpoint3_10Input) =>
raw["session.wait"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError))
type Endpoint3_10Request = Parameters<RawClient["server.session"]["session.revert.stage"]>[0]
type Endpoint3_10Input = {
readonly sessionID: Endpoint3_10Request["params"]["sessionID"]
readonly messageID: Endpoint3_10Request["payload"]["messageID"]
readonly files?: Endpoint3_10Request["payload"]["files"]
type Endpoint3_11Request = Parameters<RawClient["server.session"]["session.revert.stage"]>[0]
type Endpoint3_11Input = {
readonly sessionID: Endpoint3_11Request["params"]["sessionID"]
readonly messageID: Endpoint3_11Request["payload"]["messageID"]
readonly files?: Endpoint3_11Request["payload"]["files"]
}
const Endpoint3_10 = (raw: RawClient["server.session"]) => (input: Endpoint3_10Input) =>
const Endpoint3_11 = (raw: RawClient["server.session"]) => (input: Endpoint3_11Input) =>
raw["session.revert.stage"]({
params: { sessionID: input["sessionID"] },
payload: { messageID: input["messageID"], files: input["files"] },
@@ -158,42 +169,42 @@ const Endpoint3_10 = (raw: RawClient["server.session"]) => (input: Endpoint3_10I
Effect.map((value) => value.data),
)
type Endpoint3_11Request = Parameters<RawClient["server.session"]["session.revert.clear"]>[0]
type Endpoint3_11Input = { readonly sessionID: Endpoint3_11Request["params"]["sessionID"] }
const Endpoint3_11 = (raw: RawClient["server.session"]) => (input: Endpoint3_11Input) =>
raw["session.revert.clear"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError))
type Endpoint3_12Request = Parameters<RawClient["server.session"]["session.revert.commit"]>[0]
type Endpoint3_12Request = Parameters<RawClient["server.session"]["session.revert.clear"]>[0]
type Endpoint3_12Input = { readonly sessionID: Endpoint3_12Request["params"]["sessionID"] }
const Endpoint3_12 = (raw: RawClient["server.session"]) => (input: Endpoint3_12Input) =>
raw["session.revert.commit"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError))
raw["session.revert.clear"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError))
type Endpoint3_13Request = Parameters<RawClient["server.session"]["session.context"]>[0]
type Endpoint3_13Request = Parameters<RawClient["server.session"]["session.revert.commit"]>[0]
type Endpoint3_13Input = { readonly sessionID: Endpoint3_13Request["params"]["sessionID"] }
const Endpoint3_13 = (raw: RawClient["server.session"]) => (input: Endpoint3_13Input) =>
raw["session.revert.commit"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError))
type Endpoint3_14Request = Parameters<RawClient["server.session"]["session.context"]>[0]
type Endpoint3_14Input = { readonly sessionID: Endpoint3_14Request["params"]["sessionID"] }
const Endpoint3_14 = (raw: RawClient["server.session"]) => (input: Endpoint3_14Input) =>
raw["session.context"]({ params: { sessionID: input["sessionID"] } }).pipe(
Effect.mapError(mapClientError),
Effect.map((value) => value.data),
)
type Endpoint3_14Request = Parameters<RawClient["server.session"]["session.history"]>[0]
type Endpoint3_14Input = {
readonly sessionID: Endpoint3_14Request["params"]["sessionID"]
readonly limit?: Endpoint3_14Request["query"]["limit"]
readonly after?: Endpoint3_14Request["query"]["after"]
type Endpoint3_15Request = Parameters<RawClient["server.session"]["session.history"]>[0]
type Endpoint3_15Input = {
readonly sessionID: Endpoint3_15Request["params"]["sessionID"]
readonly limit?: Endpoint3_15Request["query"]["limit"]
readonly after?: Endpoint3_15Request["query"]["after"]
}
const Endpoint3_14 = (raw: RawClient["server.session"]) => (input: Endpoint3_14Input) =>
const Endpoint3_15 = (raw: RawClient["server.session"]) => (input: Endpoint3_15Input) =>
raw["session.history"]({
params: { sessionID: input["sessionID"] },
query: { limit: input["limit"], after: input["after"] },
}).pipe(Effect.mapError(mapClientError))
type Endpoint3_15Request = Parameters<RawClient["server.session"]["session.events"]>[0]
type Endpoint3_15Input = {
readonly sessionID: Endpoint3_15Request["params"]["sessionID"]
readonly after?: Endpoint3_15Request["query"]["after"]
type Endpoint3_16Request = Parameters<RawClient["server.session"]["session.events"]>[0]
type Endpoint3_16Input = {
readonly sessionID: Endpoint3_16Request["params"]["sessionID"]
readonly after?: Endpoint3_16Request["query"]["after"]
}
const Endpoint3_15 = (raw: RawClient["server.session"]) => (input: Endpoint3_15Input) =>
const Endpoint3_16 = (raw: RawClient["server.session"]) => (input: Endpoint3_16Input) =>
Stream.unwrap(
raw["session.events"]({ params: { sessionID: input["sessionID"] }, query: { after: input["after"] } }).pipe(
Effect.mapError(mapClientError),
@@ -201,17 +212,17 @@ const Endpoint3_15 = (raw: RawClient["server.session"]) => (input: Endpoint3_15I
),
)
type Endpoint3_16Request = Parameters<RawClient["server.session"]["session.interrupt"]>[0]
type Endpoint3_16Input = { readonly sessionID: Endpoint3_16Request["params"]["sessionID"] }
const Endpoint3_16 = (raw: RawClient["server.session"]) => (input: Endpoint3_16Input) =>
type Endpoint3_17Request = Parameters<RawClient["server.session"]["session.interrupt"]>[0]
type Endpoint3_17Input = { readonly sessionID: Endpoint3_17Request["params"]["sessionID"] }
const Endpoint3_17 = (raw: RawClient["server.session"]) => (input: Endpoint3_17Input) =>
raw["session.interrupt"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError))
type Endpoint3_17Request = Parameters<RawClient["server.session"]["session.message"]>[0]
type Endpoint3_17Input = {
readonly sessionID: Endpoint3_17Request["params"]["sessionID"]
readonly messageID: Endpoint3_17Request["params"]["messageID"]
type Endpoint3_18Request = Parameters<RawClient["server.session"]["session.message"]>[0]
type Endpoint3_18Input = {
readonly sessionID: Endpoint3_18Request["params"]["sessionID"]
readonly messageID: Endpoint3_18Request["params"]["messageID"]
}
const Endpoint3_17 = (raw: RawClient["server.session"]) => (input: Endpoint3_17Input) =>
const Endpoint3_18 = (raw: RawClient["server.session"]) => (input: Endpoint3_18Input) =>
raw["session.message"]({ params: { sessionID: input["sessionID"], messageID: input["messageID"] } }).pipe(
Effect.mapError(mapClientError),
Effect.map((value) => value.data),
@@ -222,20 +233,21 @@ const adaptGroup3 = (raw: RawClient["server.session"]) => ({
create: Endpoint3_1(raw),
active: Endpoint3_2(raw),
get: Endpoint3_3(raw),
switchAgent: Endpoint3_4(raw),
switchModel: Endpoint3_5(raw),
rename: Endpoint3_6(raw),
prompt: Endpoint3_7(raw),
compact: Endpoint3_8(raw),
wait: Endpoint3_9(raw),
stage: Endpoint3_10(raw),
clear: Endpoint3_11(raw),
commit: Endpoint3_12(raw),
context: Endpoint3_13(raw),
history: Endpoint3_14(raw),
events: Endpoint3_15(raw),
interrupt: Endpoint3_16(raw),
message: Endpoint3_17(raw),
fork: Endpoint3_4(raw),
switchAgent: Endpoint3_5(raw),
switchModel: Endpoint3_6(raw),
rename: Endpoint3_7(raw),
prompt: Endpoint3_8(raw),
compact: Endpoint3_9(raw),
wait: Endpoint3_10(raw),
stage: Endpoint3_11(raw),
clear: Endpoint3_12(raw),
commit: Endpoint3_13(raw),
context: Endpoint3_14(raw),
history: Endpoint3_15(raw),
events: Endpoint3_16(raw),
interrupt: Endpoint3_17(raw),
message: Endpoint3_18(raw),
})
type Endpoint4_0Request = Parameters<RawClient["server.message"]["session.messages"]>[0]
+14
View File
@@ -11,6 +11,8 @@ import type {
SessionsActiveOutput,
SessionsGetInput,
SessionsGetOutput,
SessionsForkInput,
SessionsForkOutput,
SessionsSwitchAgentInput,
SessionsSwitchAgentOutput,
SessionsSwitchModelInput,
@@ -345,6 +347,18 @@ export function make(options: ClientOptions) {
},
requestOptions,
).then((value) => value.data),
fork: (input: SessionsForkInput, requestOptions?: RequestOptions) =>
request<{ readonly data: SessionsForkOutput }>(
{
method: "POST",
path: `/api/session/${encodeURIComponent(input.sessionID)}/fork`,
body: { messageID: input["messageID"] },
successStatus: 200,
declaredStatuses: [404, 400, 401],
empty: false,
},
requestOptions,
).then((value) => value.data),
switchAgent: (input: SessionsSwitchAgentInput, requestOptions?: RequestOptions) =>
request<SessionsSwitchAgentOutput>(
{
+48 -9
View File
@@ -33,6 +33,15 @@ export type SessionNotFoundError = {
export const isSessionNotFoundError = (value: unknown): value is SessionNotFoundError =>
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "SessionNotFoundError"
export type MessageNotFoundError = {
readonly _tag: "MessageNotFoundError"
readonly sessionID: string
readonly messageID: string
readonly message: string
}
export const isMessageNotFoundError = (value: unknown): value is MessageNotFoundError =>
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "MessageNotFoundError"
export type ConflictError = {
readonly _tag: "ConflictError"
readonly message: string
@@ -49,15 +58,6 @@ export type ServiceUnavailableError = {
export const isServiceUnavailableError = (value: unknown): value is ServiceUnavailableError =>
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "ServiceUnavailableError"
export type MessageNotFoundError = {
readonly _tag: "MessageNotFoundError"
readonly sessionID: string
readonly messageID: string
readonly message: string
}
export const isMessageNotFoundError = (value: unknown): value is MessageNotFoundError =>
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "MessageNotFoundError"
export type SessionBusyError = {
readonly _tag: "SessionBusyError"
readonly sessionID: string
@@ -373,6 +373,45 @@ export type SessionsGetOutput = {
}
}["data"]
export type SessionsForkInput = {
readonly sessionID: { readonly sessionID: string }["sessionID"]
readonly messageID?: { readonly messageID?: string | undefined }["messageID"]
}
export type SessionsForkOutput = {
readonly data: {
readonly id: string
readonly parentID?: string
readonly projectID: string
readonly agent?: string
readonly model?: { readonly id: string; readonly providerID: string; readonly variant?: string }
readonly cost: number
readonly tokens: {
readonly input: number
readonly output: number
readonly reasoning: number
readonly cache: { readonly read: number; readonly write: number }
}
readonly time: { readonly created: number; readonly updated: number; readonly archived?: number }
readonly title: string
readonly location: { readonly directory: string; readonly workspaceID?: string }
readonly subpath?: string
readonly revert?: {
readonly messageID: string
readonly partID?: string
readonly snapshot?: string
readonly diff?: string
readonly files?: ReadonlyArray<{
readonly path: string
readonly status: "added" | "modified" | "deleted"
readonly additions: number
readonly deletions: number
readonly patch: string
}>
}
}
}["data"]
export type SessionsSwitchAgentInput = {
readonly sessionID: { readonly sessionID: string }["sessionID"]
readonly agent: { readonly agent: string }["agent"]
+188 -3
View File
@@ -3,7 +3,7 @@ export * from "./session/schema"
import { DateTime, Effect, Layer, Schema, Context, Stream } from "effect"
import { ListAnchor } from "@opencode-ai/schema/session"
import { and, asc, desc, eq, gt, like, lt, or, type SQL } from "drizzle-orm"
import { and, asc, desc, eq, gt, inArray, like, lt, or, type SQL } from "drizzle-orm"
import { ProjectV2 } from "./project"
import { WorkspaceV2 } from "./workspace"
import { ModelV2 } from "./model"
@@ -12,9 +12,10 @@ import { SessionMessage } from "./session/message"
import { Prompt } from "./session/prompt"
import { PromptInput } from "@opencode-ai/schema/prompt-input"
import { EventV2 } from "./event"
import { EventSequenceTable } from "./event/sql"
import { Database } from "./database/database"
import { SessionProjector } from "./session/projector"
import { SessionMessageTable, SessionTable } from "./session/sql"
import { SessionInputTable, SessionMessageTable, SessionTable } from "./session/sql"
import { SessionSchema } from "./session/schema"
import { AbsolutePath, PositiveInt, RelativePath } from "./schema"
import { AgentV2 } from "./agent"
@@ -90,6 +91,11 @@ type CompactInput = {
prompt?: Prompt
}
type ForkInput = {
sessionID: SessionSchema.ID
messageID?: SessionMessage.ID
}
export class NotFoundError extends Schema.TaggedErrorClass<NotFoundError>()("Session.NotFoundError", {
sessionID: SessionSchema.ID,
}) {}
@@ -113,11 +119,17 @@ export class BusyError extends Schema.TaggedErrorClass<BusyError>()("Session.Bus
export const MessageNotFoundError = SessionRevert.MessageNotFoundError
export type MessageNotFoundError = SessionRevert.MessageNotFoundError
export type Error = NotFoundError | MessageDecodeError | OperationUnavailableError | PromptConflictError
export type Error =
| NotFoundError
| MessageDecodeError
| OperationUnavailableError
| PromptConflictError
| MessageNotFoundError
export interface Interface {
readonly list: (input?: ListInput) => Effect.Effect<SessionSchema.Info[]>
readonly create: (input: CreateInput) => Effect.Effect<SessionSchema.Info, NotFoundError>
readonly fork: (input: ForkInput) => Effect.Effect<SessionSchema.Info, NotFoundError | MessageNotFoundError>
readonly get: (sessionID: SessionSchema.ID) => Effect.Effect<SessionSchema.Info, NotFoundError>
readonly messages: (input: {
sessionID: SessionSchema.ID
@@ -270,6 +282,33 @@ export const layer = Layer.effect(
// TODO: Restore recorded sessions onto replacement synchronized workspaces in a future API slice.
return yield* result.get(sessionID).pipe(Effect.orDie)
}),
fork: Effect.fn("V2Session.fork")(function* (input) {
const parent = yield* result.get(input.sessionID)
const boundary = input.messageID
? yield* db
.select({ seq: SessionMessageTable.seq })
.from(SessionMessageTable)
.where(
and(eq(SessionMessageTable.session_id, input.sessionID), eq(SessionMessageTable.id, input.messageID)),
)
.get()
.pipe(Effect.orDie)
: undefined
if (input.messageID && !boundary)
return yield* new MessageNotFoundError({ sessionID: input.sessionID, messageID: input.messageID })
const child = yield* result.create({
parentID: parent.id,
title: forkTitle(parent.title),
agent: parent.agent,
model: parent.model,
})
yield* copyForkRows(db, {
parentID: parent.id,
childID: child.id,
beforeSeq: boundary?.seq,
})
return yield* result.get(child.id)
}),
get: Effect.fn("V2Session.get")(function* (sessionID) {
const session = yield* store.get(sessionID)
if (!session) return yield* new NotFoundError({ sessionID })
@@ -492,6 +531,152 @@ export const defaultLayer = layer.pipe(
Layer.orDie,
)
const ForkBatchSize = 500
const forkTitle = (value: string) => {
const match = value.match(/^(.+) \(fork #(\d+)\)$/)
if (match) return `${match[1]} (fork #${Number.parseInt(match[2], 10) + 1})`
return `${value} (fork #1)`
}
const emptyUsage = () => ({
cost: 0,
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
})
const copyForkRows = Effect.fn("V2Session.copyForkRows")(function* (
db: Database.Interface["db"],
input: { parentID: SessionSchema.ID; childID: SessionSchema.ID; beforeSeq?: number },
) {
return yield* db
.transaction(
() =>
Effect.gen(function* () {
let cursor = -1
let maxSeq = 0
const usage = emptyUsage()
while (true) {
const rows = yield* db
.select()
.from(SessionMessageTable)
.where(
and(
eq(SessionMessageTable.session_id, input.parentID),
gt(SessionMessageTable.seq, cursor),
input.beforeSeq === undefined ? undefined : lt(SessionMessageTable.seq, input.beforeSeq),
),
)
.orderBy(asc(SessionMessageTable.seq))
.limit(ForkBatchSize)
.all()
.pipe(Effect.orDie)
if (rows.length === 0) break
const idMap = new Map(rows.map((row) => [row.id, SessionMessage.ID.create()]))
yield* db
.insert(SessionMessageTable)
.values(
rows.map((row) => ({
id: idMap.get(row.id)!,
session_id: input.childID,
type: row.type,
seq: row.seq,
time_created: row.time_created,
time_updated: row.time_updated,
data: row.type === "synthetic" ? { ...row.data, sessionID: input.childID } : row.data,
})),
)
.run()
.pipe(Effect.orDie)
const inputRows = yield* db
.select()
.from(SessionInputTable)
.where(
and(
eq(SessionInputTable.session_id, input.parentID),
inArray(
SessionInputTable.id,
rows.map((row) => row.id),
),
),
)
.all()
.pipe(Effect.orDie)
if (inputRows.length > 0) {
yield* db
.insert(SessionInputTable)
.values(
inputRows.flatMap((row) => {
const id = idMap.get(row.id)
return id
? [
{
id,
session_id: input.childID,
prompt: row.prompt,
delivery: row.delivery,
admitted_seq: row.admitted_seq,
promoted_seq: row.promoted_seq,
time_created: row.time_created,
},
]
: []
}),
)
.run()
.pipe(Effect.orDie)
}
for (const row of rows) addUsage(usage, row)
cursor = rows.at(-1)!.seq
maxSeq = cursor
}
yield* db
.update(SessionTable)
.set({
cost: usage.cost,
tokens_input: usage.tokens.input,
tokens_output: usage.tokens.output,
tokens_reasoning: usage.tokens.reasoning,
tokens_cache_read: usage.tokens.cache.read,
tokens_cache_write: usage.tokens.cache.write,
})
.where(eq(SessionTable.id, input.childID))
.run()
.pipe(Effect.orDie)
if (maxSeq > 0) {
yield* db
.update(EventSequenceTable)
.set({ seq: maxSeq })
.where(eq(EventSequenceTable.aggregate_id, input.childID))
.run()
.pipe(Effect.orDie)
}
}),
{ behavior: "immediate" },
)
.pipe(Effect.orDie)
})
function addUsage(usage: ReturnType<typeof emptyUsage>, row: typeof SessionMessageTable.$inferSelect) {
if (row.type !== "assistant") return
const data = row.data as Record<string, unknown>
if (typeof data.cost === "number") usage.cost += data.cost
if (typeof data.tokens !== "object" || data.tokens === null) return
const tokens = data.tokens as Record<string, unknown>
if (typeof tokens.input === "number") usage.tokens.input += tokens.input
if (typeof tokens.output === "number") usage.tokens.output += tokens.output
if (typeof tokens.reasoning === "number") usage.tokens.reasoning += tokens.reasoning
if (typeof tokens.cache !== "object" || tokens.cache === null) return
const cache = tokens.cache as Record<string, unknown>
if (typeof cache.read === "number") usage.tokens.cache.read += cache.read
if (typeof cache.write === "number") usage.tokens.cache.write += cache.write
}
const resolvePrompt = (input: PromptInput.Prompt) =>
Prompt.make({
text: input.text,
+76 -1
View File
@@ -1,6 +1,6 @@
import { describe, expect } from "bun:test"
import path from "path"
import { Effect, Layer, Stream } from "effect"
import { DateTime, Effect, Layer, Stream } from "effect"
import { AgentV2 } from "@opencode-ai/core/agent"
import { asc, eq } from "drizzle-orm"
import { Database } from "@opencode-ai/core/database/database"
@@ -20,6 +20,7 @@ import { SessionProjector } from "@opencode-ai/core/session/projector"
import { SessionExecution } from "@opencode-ai/core/session/execution"
import { SessionInput } from "@opencode-ai/core/session/input"
import { SessionEvent } from "@opencode-ai/core/session/event"
import { SessionMessage } from "@opencode-ai/core/session/message"
import { SessionTable } from "@opencode-ai/core/session/sql"
import { SessionStore } from "@opencode-ai/core/session/store"
import { WorkspaceV2 } from "@opencode-ai/core/workspace"
@@ -131,6 +132,80 @@ describe("SessionV2.create", () => {
}),
)
it.effect("forks a session by copying projected rows with fresh message IDs", () =>
Effect.gen(function* () {
const session = yield* SessionV2.Service
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const parent = yield* session.create({ location, title: "Parent" })
const admitted = yield* session.prompt({
sessionID: parent.id,
prompt: Prompt.make({ text: "First" }),
resume: false,
})
yield* SessionInput.promoteSteers(db, events, parent.id, Number.MAX_SAFE_INTEGER)
yield* events.publish(SessionEvent.Synthetic, {
sessionID: parent.id,
messageID: SessionMessage.ID.create(),
timestamp: yield* DateTime.now,
text: "parent note",
})
const forked = yield* session.fork({ sessionID: parent.id })
const parentContext = yield* session.context(parent.id)
const forkContext = yield* session.context(forked.id)
expect(forked).toMatchObject({ parentID: parent.id, title: "Parent (fork #1)" })
expect(forkContext).toMatchObject([
{ type: "user", text: "First" },
{ type: "synthetic", text: "parent note", sessionID: forked.id },
])
expect(forkContext.map((message) => message.id)).not.toEqual(parentContext.map((message) => message.id))
expect(yield* SessionInput.find(db, forkContext[0]!.id)).toMatchObject({
sessionID: forked.id,
prompt: { text: "First" },
promotedSeq: 2,
})
yield* session.prompt({ sessionID: parent.id, prompt: Prompt.make({ text: "Parent changed" }), resume: false })
yield* SessionInput.promoteSteers(db, events, parent.id, Number.MAX_SAFE_INTEGER)
yield* session.prompt({ sessionID: forked.id, prompt: Prompt.make({ text: "Child continues" }), resume: false })
yield* SessionInput.promoteSteers(db, events, forked.id, Number.MAX_SAFE_INTEGER)
expect((yield* session.context(parent.id)).map((message) => message.type)).toEqual(["user", "synthetic", "user"])
expect((yield* session.context(forked.id)).map((message) => message.type)).toEqual(["user", "synthetic", "user"])
expect((yield* session.context(forked.id)).at(-1)).toMatchObject({ text: "Child continues" })
expect(yield* SessionInput.find(db, admitted.id)).toMatchObject({ sessionID: parent.id })
}),
)
it.effect("forks before the selected boundary message", () =>
Effect.gen(function* () {
const session = yield* SessionV2.Service
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const parent = yield* session.create({ location })
const first = yield* session.prompt({
sessionID: parent.id,
prompt: Prompt.make({ text: "First" }),
resume: false,
})
yield* SessionInput.promoteSteers(db, events, parent.id, Number.MAX_SAFE_INTEGER)
const second = yield* session.prompt({
sessionID: parent.id,
prompt: Prompt.make({ text: "Second" }),
resume: false,
})
yield* SessionInput.promoteSteers(db, events, parent.id, Number.MAX_SAFE_INTEGER)
const forked = yield* session.fork({ sessionID: parent.id, messageID: second.id })
const context = yield* session.context(forked.id)
expect(context).toMatchObject([{ text: "First" }])
expect(context[0]?.id).not.toBe(first.id)
}),
)
it.effect("returns the existing Session when one ID is reused with different create arguments", () =>
Effect.gen(function* () {
const session = yield* SessionV2.Service
+17
View File
@@ -170,6 +170,23 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(sessionLo
}),
),
)
.add(
HttpApiEndpoint.post("session.fork", "/api/session/:sessionID/fork", {
params: { sessionID: Session.ID },
payload: Schema.Struct({ messageID: SessionMessage.ID.pipe(Schema.optional) }),
success: Schema.Struct({ data: Session.Info }),
error: [SessionNotFoundError, MessageNotFoundError],
})
.middleware(sessionLocationMiddleware)
.annotateMerge(
OpenApi.annotations({
identifier: "v2.session.fork",
summary: "Fork session",
description:
"Create a child session by copying projected history from the parent. When messageID is supplied, copy messages before that boundary.",
}),
),
)
.add(
HttpApiEndpoint.post("session.switchAgent", "/api/session/:sessionID/agent", {
params: { sessionID: Session.ID },
+26
View File
@@ -107,6 +107,32 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
}
}),
)
.handle(
"session.fork",
Effect.fn(function* (ctx) {
return {
data: yield* session.fork({ sessionID: ctx.params.sessionID, messageID: ctx.payload.messageID }).pipe(
Effect.catchTag(
"Session.NotFoundError",
(error) =>
new SessionNotFoundError({
sessionID: error.sessionID,
message: `Session not found: ${error.sessionID}`,
}),
),
Effect.catchTag(
"Session.MessageNotFoundError",
(error) =>
new MessageNotFoundError({
sessionID: error.sessionID,
messageID: error.messageID,
message: `Message not found: ${error.messageID}`,
}),
),
),
}
}),
)
.handle(
"session.switchAgent",
Effect.fn(function* (ctx) {