diff --git a/packages/client/src/effect/api/api.ts b/packages/client/src/effect/api/api.ts index c45a1d803665..3c652a1b15fe 100644 --- a/packages/client/src/effect/api/api.ts +++ b/packages/client/src/effect/api/api.ts @@ -17,6 +17,7 @@ import type { PromptInput } from "@opencode-ai/schema/prompt-input" import type { AgentAttachment } from "@opencode-ai/schema/prompt" import type { Skill } from "@opencode-ai/schema/skill" import type { Event } from "@opencode-ai/schema/event" +import type { FileDiff } from "@opencode-ai/schema/file-diff" import type { InstructionEntry } from "@opencode-ai/schema/instruction-entry" import type { Schema } from "effect" import type { EventLog } from "@opencode-ai/schema/event-log" @@ -36,7 +37,6 @@ import type { PtyTicket } from "@opencode-ai/schema/pty-ticket" import type { Reference } from "@opencode-ai/schema/reference" import type { Worktree } from "@opencode-ai/schema/worktree" import type { Vcs } from "@opencode-ai/schema/vcs" -import type { FileDiff } from "@opencode-ai/schema/file-diff" import type { WebSearch } from "@opencode-ai/schema/websearch" import type { Config } from "@opencode-ai/schema/config" @@ -360,6 +360,15 @@ export type SessionContextInput = { readonly sessionID: Session.ID } export type SessionContextOutput = ReadonlyArray export type SessionContextOperation = (input: SessionContextInput) => Effect.Effect +export type SessionDiffInput = { + readonly sessionID: Session.ID + readonly messageID?: SessionMessage.ID | undefined + readonly to?: SessionMessage.ID | undefined + readonly context?: number | undefined +} +export type SessionDiffOutput = ReadonlyArray +export type SessionDiffOperation = (input: SessionDiffInput) => Effect.Effect + export type SessionInboxListInput = { readonly sessionID: Session.ID } export type SessionInboxListOutput = ReadonlyArray export type SessionInboxListOperation = ( @@ -1133,6 +1142,7 @@ export interface SessionApi { readonly commit: SessionRevertCommitOperation } readonly context: SessionContextOperation + readonly diff: SessionDiffOperation readonly inbox: { readonly list: SessionInboxListOperation readonly cancel: SessionInboxCancelOperation diff --git a/packages/client/src/effect/generated/client.ts b/packages/client/src/effect/generated/client.ts index 3a0616854ddc..dbc3a22beb7e 100644 --- a/packages/client/src/effect/generated/client.ts +++ b/packages/client/src/effect/generated/client.ts @@ -68,6 +68,8 @@ import type { SessionRevertCommitOutput, SessionContextInput, SessionContextOutput, + SessionDiffInput, + SessionDiffOutput, SessionInboxListInput, SessionInboxListOutput, SessionInboxCancelInput, @@ -594,6 +596,17 @@ const EndpointSessionContext = (raw: RawClient["server.session"]) => (input: Ses ), ) +const EndpointSessionDiff = (raw: RawClient["server.session"]) => (input: SessionDiffInput) => + preserveEffect()( + raw["session.diff"]({ + params: { sessionID: input["sessionID"] }, + query: { messageID: input["messageID"], to: input["to"], context: input["context"] }, + }).pipe( + Effect.mapError(mapClientError), + Effect.map((value) => value.data), + ), + ) + const EndpointSessionInboxList = (raw: RawClient["server.session"]) => (input: SessionInboxListInput) => preserveEffect()( raw["session.inbox.list"]({ params: { sessionID: input["sessionID"] } }).pipe( @@ -744,6 +757,7 @@ const adaptGroupSession = (raw: RawClient["server.session"]) => ({ commit: EndpointSessionRevertCommit(raw), }, context: EndpointSessionContext(raw), + diff: EndpointSessionDiff(raw), inbox: { list: EndpointSessionInboxList(raw), cancel: EndpointSessionInboxCancel(raw), diff --git a/packages/client/src/promise/generated/client.ts b/packages/client/src/promise/generated/client.ts index 03ea39536d7b..6d9db46f5262 100644 --- a/packages/client/src/promise/generated/client.ts +++ b/packages/client/src/promise/generated/client.ts @@ -62,6 +62,8 @@ import type { SessionRevertCommitOutput, SessionContextInput, SessionContextOutput, + SessionDiffInput, + SessionDiffOutput, SessionInboxListInput, SessionInboxListOutput, SessionInboxCancelInput, @@ -844,6 +846,18 @@ export function make(options: ClientOptions) { }, requestOptions, ).then((value) => value.data), + diff: (input: SessionDiffInput, requestOptions?: RequestOptions) => + request<{ readonly data: SessionDiffOutput }>( + { + method: "GET", + path: `/api/session/${encodeURIComponent(input.sessionID)}/diff`, + query: { messageID: input["messageID"], to: input["to"], context: input["context"] }, + successStatus: 200, + declaredStatuses: [400, 401, 404, 500], + empty: false, + }, + requestOptions, + ).then((value) => value.data), inbox: { list: (input: SessionInboxListInput, requestOptions?: RequestOptions) => request<{ readonly data: SessionInboxListOutput }>( diff --git a/packages/client/src/promise/generated/types.ts b/packages/client/src/promise/generated/types.ts index 7f8c229e9ff8..00dc64e0fd4b 100644 --- a/packages/client/src/promise/generated/types.ts +++ b/packages/client/src/promise/generated/types.ts @@ -147,6 +147,14 @@ export type SessionProviderContextProvenance = { endpoint: string } +export type SessionMessageIdle = { + id: string + metadata?: { [x: string]: JsonValue } + time: { created: number } + type: "idle" + outcome: "succeeded" | "failed" | "interrupted" +} + export type SessionActive = { type: "running" } export type SessionInboxDelivery = "steer" | "queue" @@ -2177,6 +2185,7 @@ export type SessionMessageInfo = | SessionMessageShell | SessionMessageAssistant | SessionMessageCompaction + | SessionMessageIdle export type SessionMessageContentUpdated = { id: string @@ -3122,6 +3131,13 @@ export type SessionImportInput = { readonly error: { readonly type: string; readonly message: string; readonly status?: number } } ) + | { + readonly id: string + readonly metadata?: { readonly [x: string]: JsonValue } + readonly time: { readonly created: number } + readonly type: "idle" + readonly outcome: "succeeded" | "failed" | "interrupted" + } > readonly location?: { readonly directory: string; readonly workspaceID?: string } | null }["info"] @@ -3413,6 +3429,13 @@ export type SessionImportInput = { readonly error: { readonly type: string; readonly message: string; readonly status?: number } } ) + | { + readonly id: string + readonly metadata?: { readonly [x: string]: JsonValue } + readonly time: { readonly created: number } + readonly type: "idle" + readonly outcome: "succeeded" | "failed" | "interrupted" + } > readonly location?: { readonly directory: string; readonly workspaceID?: string } | null }["messages"] @@ -3704,6 +3727,13 @@ export type SessionImportInput = { readonly error: { readonly type: string; readonly message: string; readonly status?: number } } ) + | { + readonly id: string + readonly metadata?: { readonly [x: string]: JsonValue } + readonly time: { readonly created: number } + readonly type: "idle" + readonly outcome: "succeeded" | "failed" | "interrupted" + } > readonly location?: { readonly directory: string; readonly workspaceID?: string } | null }["location"] @@ -4193,6 +4223,27 @@ export type SessionContextInput = { readonly sessionID: { readonly sessionID: st export type SessionContextOutput = { data: Array }["data"] +export type SessionDiffInput = { + readonly sessionID: { readonly sessionID: string }["sessionID"] + readonly messageID?: { + readonly messageID?: string | undefined + readonly to?: string | undefined + readonly context?: number | undefined + }["messageID"] + readonly to?: { + readonly messageID?: string | undefined + readonly to?: string | undefined + readonly context?: number | undefined + }["to"] + readonly context?: { + readonly messageID?: string | undefined + readonly to?: string | undefined + readonly context?: number | undefined + }["context"] +} + +export type SessionDiffOutput = { data: Array }["data"] + export type SessionInboxListInput = { readonly sessionID: { readonly sessionID: string }["sessionID"] } export type SessionInboxListOutput = { data: Array }["data"] diff --git a/packages/client/src/solid/data.ts b/packages/client/src/solid/data.ts index 6cea279b906a..3d6d12fbcd6b 100644 --- a/packages/client/src/solid/data.ts +++ b/packages/client/src/solid/data.ts @@ -1032,6 +1032,18 @@ export function createData(config: CreateDataInput) { if (currentAssistant) currentAssistant.retry = undefined }) if (event.type === "session.execution.interrupted" && event.data.reason === "shutdown") return + // Mirror the projected idle marker so turn boundaries match before the next message read. + message.insert(event.data.sessionID, { + id: messageIDFromEvent(event.id), + type: "idle", + outcome: + event.type === "session.execution.succeeded" + ? "succeeded" + : event.type === "session.execution.failed" + ? "failed" + : "interrupted", + time: { created: event.created }, + }) // An event can overtake the first read; queue a revalidation when that read is still active. if (!store.session.info[event.data.sessionID] && !sync.has(`session:${event.data.sessionID}`)) return result.session.invalidate(event.data.sessionID) diff --git a/packages/core/src/git.ts b/packages/core/src/git.ts index da745e5f3697..59909940d58e 100644 --- a/packages/core/src/git.ts +++ b/packages/core/src/git.ts @@ -9,6 +9,7 @@ import { AppProcess } from "@opencode-ai/util/process" import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" import { File } from "./file.js" import { KeyedMutex } from "./effect/keyed-mutex.js" +import { VcsPatch } from "./vcs/patch.js" export class Repository extends Schema.Class("Git.Repository")({ worktree: AbsolutePath, @@ -308,7 +309,7 @@ const layer = Layer.effect( operationName: OperationError["operation"], repository: Repository, args: string[], - options?: { stdin?: string; env?: Record }, + options?: { stdin?: string; env?: Record; maxOutputBytes?: number }, ) { const result = yield* proc .run( @@ -317,7 +318,7 @@ const layer = Layer.effect( env: options?.env, extendEnv: true, }), - { stdin: options?.stdin }, + { stdin: options?.stdin, maxOutputBytes: options?.maxOutputBytes }, ) .pipe( Effect.mapError( @@ -331,7 +332,8 @@ const layer = Layer.effect( ), ) const text = result.stdout.toString("utf8") - if (result.exitCode === 0) return { text, stderr: result.stderr.toString("utf8") } + if (result.exitCode === 0) + return { text, stderr: result.stderr.toString("utf8"), truncated: result.stdoutTruncated } return yield* new OperationError({ operation: operationName, directory: repository.worktree, @@ -385,9 +387,7 @@ const layer = Layer.effect( maximumUntrackedFileBytes?: number }) { const list = (args: string[]) => - repositoryOperation("refresh", input.repository, args).pipe( - Effect.map((result) => result.text.split("\0").filter(Boolean)), - ) + repositoryOperation("refresh", input.repository, args).pipe(Effect.map((result) => nuls(result.text))) const [tracked, untracked] = yield* Effect.all( [ list(["diff-files", "--name-only", "-z", "--", input.scope]), @@ -464,13 +464,7 @@ const layer = Layer.effect( directory: input.repository.worktree, message: result.stderr.toString("utf8").trim() || "Failed to check ignored paths", }) - return new Set( - result.stdout - .toString("utf8") - .split("\0") - .filter(Boolean) - .map((file) => RelativePath.make(file)), - ) + return new Set(nuls(result.stdout.toString("utf8")).map((file) => RelativePath.make(file))) }) const writeTree = Effect.fn("Git.tree.write")(function* (repository: Repository) { @@ -499,19 +493,23 @@ const layer = Layer.effect( to: TreeID }) { // Undo needs both paths of a rename, not only its destination. - return (yield* repositoryOperation("list_files", input.repository, [ - "diff", - "--name-only", - "--no-renames", - "-z", - input.from, - input.to, - ])).text - .split("\0") - .filter(Boolean) - .map((file) => RelativePath.make(file)) + return nuls( + (yield* repositoryOperation("list_files", input.repository, [ + "diff", + "--name-only", + "--no-renames", + "-z", + input.from, + input.to, + ])).text, + ).map((file) => RelativePath.make(file)) }) + /** + * Three batched invocations over the tree pair instead of three per file. An + * explicit empty selection diffs nothing; an absent one diffs every changed path. + * Patch output is capped like VCS diffs: files past the cap get an empty patch. + */ const treeDiff = Effect.fn("Git.tree.diff")(function* (input: { repository: Repository from: TreeID @@ -519,49 +517,57 @@ const layer = Layer.effect( context?: number paths?: readonly RelativePath[] }) { - const paths = input.paths ?? (yield* treeFiles(input)) - return yield* Effect.forEach(paths, (file) => - Effect.gen(function* () { - const statusText = (yield* repositoryOperation("diff", input.repository, [ - "diff", - "--name-status", - "--no-renames", - input.from, - input.to, - "--", - file, - ])).text.trim() - const status = statusText.startsWith("A") ? "added" : statusText.startsWith("D") ? "deleted" : "modified" - const stats = (yield* repositoryOperation("diff", input.repository, [ + if (input.paths?.length === 0) return [] + const args = ["--no-renames", input.from, input.to, "--", ...(input.paths ?? [])] + // Patch headers have no -z form: unquoted paths keep chunksByFile matching non-ASCII names. + const [names, numbers, patch] = yield* Effect.all( + [ + repositoryOperation("diff", input.repository, ["diff", "--name-status", "-z", ...args]), + repositoryOperation("diff", input.repository, ["diff", "--numstat", "-z", ...args]), + repositoryOperation( "diff", - "--numstat", - "--no-renames", - input.from, - input.to, - "--", - file, - ])).text.split("\t") - const binary = stats[0] === "-" || stats[1] === "-" - const patch = binary - ? "" - : (yield* repositoryOperation("diff", input.repository, [ - "diff", - `--unified=${input.context ?? 3}`, - "--no-renames", - input.from, - input.to, - "--", - file, - ])).text - return { - file, - status, - additions: binary ? 0 : Number(stats[0] ?? 0), - deletions: binary ? 0 : Number(stats[1] ?? 0), - patch, - } satisfies File.Diff + input.repository, + ["-c", "core.quotepath=false", "diff", "--no-ext-diff", `--unified=${input.context ?? 3}`, ...args], + { maxOutputBytes: VcsPatch.MAX_TOTAL_PATCH_BYTES }, + ), + ], + { concurrency: 3 }, + ) + const statuses = nuls(names.text) + const files = statuses.flatMap((code, index) => { + const file = statuses[index + 1] + if (index % 2 !== 0 || !file) return [] + return [ + { + file: RelativePath.make(file), + status: code.startsWith("A") ? "added" : code.startsWith("D") ? "deleted" : "modified", + } as const, + ] + }) + const stats = new Map( + nuls(numbers.text).flatMap((line) => { + const [additions, deletions, ...file] = line.split("\t") + if (!additions || !deletions || file.length === 0) return [] + return [ + [ + file.join("\t"), + additions === "-" || deletions === "-" + ? { binary: true, additions: 0, deletions: 0 } + : { binary: false, additions: Number(additions), deletions: Number(deletions) }, + ] as const, + ] }), ) + const patches = VcsPatch.chunksByFile(patch, (index) => files[index]?.file) + return files.map((entry) => { + const stat = stats.get(entry.file) + return { + ...entry, + additions: stat?.additions ?? 0, + deletions: stat?.deletions ?? 0, + patch: stat?.binary ? "" : (patches.get(entry.file) ?? VcsPatch.emptyPatch(entry.file)), + } satisfies File.Diff + }) }) const hasEntry = Effect.fnUntraced(function* (repository: Repository, tree: TreeID, file: RelativePath) { @@ -733,6 +739,11 @@ function execute(cwd: string, proc: AppProcess.Interface, args: string[]) { ) } +/** Split NUL-terminated git output into its records. */ +function nuls(text: string) { + return text.split("\0").filter(Boolean) +} + function resolvePath(cwd: string, value: string) { const trimmed = value.replace(/[\r\n]+$/, "") if (!trimmed) return cwd diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index 3e8fdb577779..2b38d3ad4ff6 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -57,8 +57,11 @@ import { SessionModelTransport } from "./session/model-transport.js" import { llmClient } from "./effect/app-node-platform.js" import { Snapshot } from "./snapshot.js" import { Session } from "./session/session.js" +import { SessionDiff, TurnRangeError } from "./session/diff.js" +import { LocationServiceMap } from "./location-service-map.js" import { FSUtil } from "@opencode-ai/util/fs-util" import type { EventLog } from "@opencode-ai/schema/event-log" +import type { FileDiff } from "@opencode-ai/schema/file-diff" import { Job } from "./job.js" import type { Command } from "./command.js" import { SessionEnvironment } from "./session/environment.js" @@ -113,6 +116,7 @@ export { type InboxItemRef = { readonly sessionID: SessionSchema.ID; readonly inboxID: SessionMessage.ID } export { DestinationNotFoundError, DestinationNotDirectoryError, DestinationUnavailableError } +export { TurnRangeError } export interface Interface { readonly list: (input?: ListInput) => Effect.Effect<{ @@ -142,6 +146,13 @@ export interface Interface { readonly context: ( sessionID: SessionSchema.ID, ) => Effect.Effect + /** Structured diffs of the files changed by a turn or range of turns; see `SessionDiff.turn`. */ + readonly diff: (input: { + readonly sessionID: SessionSchema.ID + readonly messageID?: SessionMessage.ID + readonly to?: SessionMessage.ID + readonly context?: number + }) => Effect.Effect /** * Durable admitted session work not yet visible in projected history, * ordered by admission. Includes unpromoted user and synthetic inputs and @@ -230,6 +241,7 @@ const layer = Layer.effect( const moves = yield* SessionMove.Service const jobs = yield* Job.Service const environments = yield* SessionEnvironment.Service + const locations = yield* LocationServiceMap.Service const sessions = yield* Session.make() const isDurableSessionEvent = Schema.is(SessionEvent.Durable) @@ -362,6 +374,17 @@ const layer = Layer.effect( yield* result.get(sessionID) return yield* store.context(sessionID) }), + diff: Effect.fn("Session.diff")(function* (input) { + const session = yield* result.get(input.sessionID) + const active = yield* execution.isActive(input.sessionID) + return yield* SessionDiff.turn(db, locations, { + session, + active, + messageID: input.messageID, + to: input.to, + context: input.context, + }) + }), inbox: (sessionID) => sessions.forSession(sessionID).inbox(), cancelInbox: (input) => sessions.forSession(input.sessionID).cancelInbox(input.inboxID), steerInbox: (input) => sessions.forSession(input.sessionID).steerInbox(input.inboxID), @@ -450,6 +473,7 @@ export const node: LayerNode.Provider()("Session.TurnRangeError", { + sessionID: SessionSchema.ID, + field: Schema.Literals(["messageID", "to"]), + message: Schema.String, +}) {} + +const decodeLocation = Schema.decodeUnknownSync(Schema.fromJsonString(Location.Ref)) + +/** + * Diff the files changed by the turn containing a user message. A turn runs from + * the first prompt after the Session was last idle until the next idle marker, so + * prompts steered in while it was busy belong to the same turn; `to` extends the + * range through the turn containing a later user message. Compares the range's + * first recorded start snapshot with its last recorded end snapshot; only a step + * still running in the active Session compares against the working copy. Like VCS + * diffs, an omitted `context` yields full-file patches. + * + * A Session without any idle marker predates them, so its prompts span until the + * next user message instead. + * + * Snapshot trees live in the repository of the Location that captured them, so a + * range spanning a location switch is rejected rather than diffed wrongly. + */ +export const turn = Effect.fn("SessionDiff.turn")(function* ( + db: Database.Interface["db"], + locations: Context.Service.Shape, + input: { + readonly session: SessionSchema.Info + /** The process is currently executing this Session. */ + readonly active: boolean + readonly messageID?: SessionMessage.ID + readonly to?: SessionMessage.ID + readonly context?: number + }, +) { + const sessionID = input.session.id + const rows = yield* db + .select({ id: SessionMessageTable.id, type: SessionMessageTable.type, seq: SessionMessageTable.seq }) + .from(SessionMessageTable) + .where( + and( + eq(SessionMessageTable.session_id, sessionID), + or( + inArray(SessionMessageTable.type, ["user", "idle"]), + input.messageID ? eq(SessionMessageTable.id, input.messageID) : undefined, + input.to ? eq(SessionMessageTable.id, input.to) : undefined, + ), + ), + ) + .orderBy(asc(SessionMessageTable.seq)) + .all() + .pipe(Effect.orDie) + const users = rows.filter((row) => row.type === "user") + const markers = rows.filter((row) => row.type === "idle") + const resolve = Effect.fn(function* (field: "messageID" | "to", id: SessionMessage.ID) { + const row = rows.find((row) => row.id === id) + if (!row) return yield* new MessageNotFoundError({ sessionID, messageID: id }) + if (row.type !== "user") + return yield* new TurnRangeError({ sessionID, field, message: `Message ${id} is not a user message` }) + return row + }) + const anchor = input.messageID ? yield* resolve("messageID", input.messageID) : users[users.length - 1] + if (!anchor) return [] + const last = input.to ? yield* resolve("to", input.to) : anchor + if (last.seq < anchor.seq) + return yield* new TurnRangeError({ sessionID, field: "to", message: `Message ${last.id} precedes ${anchor.id}` }) + // Without any marker, history predates idle markers and a prompt's turn ends at the next prompt. + const legacy = markers.length === 0 + // The turn opens with the first prompt after the previous idle marker; the anchor itself is the latest candidate. + const opened = markers.findLast((row) => row.seq < anchor.seq)?.seq ?? -1 + const start = legacy ? anchor.seq : (users.find((row) => row.seq > opened)?.seq ?? anchor.seq) + const end = legacy ? users.find((row) => row.seq > last.seq)?.seq : markers.find((row) => row.seq > last.seq)?.seq + const steps = yield* db + .select({ + seq: SessionMessageTable.seq, + start: sql`json_extract(${SessionMessageTable.data}, '$.snapshot.start')`, + end: sql`json_extract(${SessionMessageTable.data}, '$.snapshot.end')`, + completed: sql`json_extract(${SessionMessageTable.data}, '$.time.completed')`, + }) + .from(SessionMessageTable) + .where( + and( + eq(SessionMessageTable.session_id, sessionID), + eq(SessionMessageTable.type, "assistant"), + gt(SessionMessageTable.seq, start), + end === undefined ? undefined : lt(SessionMessageTable.seq, end), + ), + ) + .orderBy(asc(SessionMessageTable.seq)) + .all() + .pipe(Effect.orDie) + const first = steps[0] + const final = steps[steps.length - 1] + const from = steps.find((step) => step.start)?.start + if (!first || !final || !from) return [] + const switches = yield* db + .select({ + seq: SessionMessageTable.seq, + location: sql`json_extract(${SessionMessageTable.data}, '$.location')`, + previous: sql`json_extract(${SessionMessageTable.data}, '$.previous.location')`, + }) + .from(SessionMessageTable) + .where(and(eq(SessionMessageTable.session_id, sessionID), eq(SessionMessageTable.type, "location-switched"))) + .orderBy(asc(SessionMessageTable.seq)) + .all() + .pipe(Effect.orDie) + if (switches.some((row) => row.seq > first.seq && row.seq < final.seq)) + return yield* new TurnRangeError({ sessionID, field: "to", message: "Turn range spans a location change" }) + const before = switches.findLast((row) => row.seq < first.seq)?.location + const after = switches.find((row) => row.seq > first.seq)?.previous + const location = before ? decodeLocation(before) : after ? decodeLocation(after) : input.session.location + const recorded = steps.findLast((step) => step.end)?.end + return yield* Effect.gen(function* () { + const snapshot = yield* Snapshot.Service + const running = input.active && final.completed === null + const to = running ? ((yield* snapshot.capture()) ?? recorded) : recorded + if (!to) return [] + return yield* snapshot.diff({ + from: Snapshot.ID.make(from), + to: Snapshot.ID.make(to), + context: input.context ?? PATCH_CONTEXT_LINES, + }) + }).pipe(Effect.provide(locations.get(location))) +}) diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index 8675690f33b2..380377761d4d 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -60,6 +60,21 @@ export function update(adapter: Adapter, event: SessionEvent.DurableEvent) { ) }) + const idle = (outcome: SessionMessage.Idle["outcome"]) => + clearCurrentRetry.pipe( + Effect.andThen( + adapter.appendMessage( + SessionMessage.Idle.make({ + id: SessionMessage.ID.fromEvent(event.id), + type: "idle", + outcome, + metadata: event.metadata, + time: { created }, + }), + ), + ), + ) + const project = pipe( Match.type(), Match.discriminatorsExhaustive("type")({ @@ -123,9 +138,11 @@ export function update(adapter: Adapter, event: SessionEvent.DurableEvent) { "session.inbox.cancelled": () => Effect.void, "session.inbox.delivery.changed": () => Effect.void, "session.execution.started": () => Effect.void, - "session.execution.succeeded": () => clearCurrentRetry, - "session.execution.failed": () => clearCurrentRetry, - "session.execution.interrupted": () => clearCurrentRetry, + "session.execution.succeeded": () => idle("succeeded"), + "session.execution.failed": () => idle("failed"), + // Shutdown keeps the execution claim and the resumed drain continues the turn. + "session.execution.interrupted": (event) => + event.data.reason === "shutdown" ? clearCurrentRetry : idle("interrupted"), "session.instructions.updated": (event) => { if (event.data.text === undefined) return Effect.void return adapter.appendMessage( diff --git a/packages/core/src/session/runner/to-llm-message.ts b/packages/core/src/session/runner/to-llm-message.ts index 17fb423c5b42..65c197412f25 100644 --- a/packages/core/src/session/runner/to-llm-message.ts +++ b/packages/core/src/session/runner/to-llm-message.ts @@ -226,6 +226,7 @@ function toLLMMessage(message: SessionMessage.Info, model: Model.Ref, providerMe switch (message.type) { case "agent-switched": case "model-switched": + case "idle": return [] case "location-switched": return [ diff --git a/packages/core/src/snapshot.ts b/packages/core/src/snapshot.ts index 25e79d29e095..887f920fdb6c 100644 --- a/packages/core/src/snapshot.ts +++ b/packages/core/src/snapshot.ts @@ -131,38 +131,55 @@ const layer = Layer.effect( ) }) - const compare = Effect.fnUntraced(function* (operation: "files" | "diff", input: CompareInput) { + const comparison = Effect.fnUntraced(function* (operation: "files" | "diff", input: CompareInput) { const repo = yield* repository.pipe(Effect.mapError((cause) => failure(operation, cause))) - const comparison = { + return { + source: repo.source, repository: repo.snapshotRepository, from: Git.TreeID.make(input.from), to: Git.TreeID.make(input.to), } - const files = yield* git.tree.files(comparison).pipe(Effect.mapError((cause) => failure(operation, cause))) - const ignored = yield* git.index - .ignored({ repository: repo.source, paths: files }) + }) + + // Snapshots track every scoped file; the source repository's ignore rules decide what callers see. + const ignored = Effect.fnUntraced(function* ( + operation: "files" | "diff", + source: Git.Repository, + paths: readonly RelativePath[], + ) { + return yield* git.index + .ignored({ repository: source, paths }) .pipe(Effect.mapError((cause) => failure(operation, cause))) - return { - input: comparison, - files, - ignored, - } }) const files = Effect.fn("Snapshot.files")(function* (input: CompareInput) { - const comparison = yield* compare("files", input) - return comparison.files.filter((file) => !comparison.ignored.has(file)) + const compared = yield* comparison("files", input) + const changed = yield* git.tree + .files({ repository: compared.repository, from: compared.from, to: compared.to }) + .pipe(Effect.mapError((cause) => failure("files", cause))) + const skipped = yield* ignored("files", compared.source, changed) + return changed.filter((file) => !skipped.has(file)) }) const diff = Effect.fn("Snapshot.diff")(function* (input: DiffInput) { - const comparison = yield* compare("diff", input) - return yield* git.tree + if (input.paths?.length === 0) return [] + const compared = yield* comparison("diff", input) + // Only an explicit selection becomes a pathspec; ignored paths are dropped from the result instead. + const diffs = yield* git.tree .diff({ - ...comparison.input, + repository: compared.repository, + from: compared.from, + to: compared.to, context: input.context, - paths: (input.paths ?? comparison.files).filter((file) => !comparison.ignored.has(file)), + paths: input.paths, }) .pipe(Effect.mapError((cause) => failure("diff", cause))) + const skipped = yield* ignored( + "diff", + compared.source, + diffs.map((file) => RelativePath.make(file.file)), + ) + return diffs.filter((file) => !skipped.has(RelativePath.make(file.file))) }) const plan = Effect.fnUntraced(function* (worktree: AbsolutePath, input: RestoreInput) { diff --git a/packages/core/test/git.test.ts b/packages/core/test/git.test.ts index 49b3e59db292..9482a030c20b 100644 --- a/packages/core/test/git.test.ts +++ b/packages/core/test/git.test.ts @@ -6,6 +6,7 @@ import { Effect } from "effect" import { LayerNode } from "@opencode-ai/util/effect/layer-node" import { Git } from "@opencode-ai/core/git" import { AbsolutePath, RelativePath } from "@opencode-ai/core/schema" +import { VcsPatch } from "@opencode-ai/core/vcs/patch" import { branch, commit, initRepo, read, withRemote } from "./fixture/git" import { tmpdir } from "./fixture/tmpdir" import { testEffect } from "./lib/effect" @@ -196,6 +197,42 @@ describe("Git trees", () => { }), ) + it.live("caps batched tree patches, keeps per-file stats past the cap, and matches non-ASCII names", () => + Effect.gen(function* () { + const root = yield* Effect.acquireRelease( + Effect.promise(() => tmpdir()), + (dir) => Effect.promise(() => dir[Symbol.asyncDispose]()), + ) + yield* Effect.promise(() => initRepo(root.path)) + const git = yield* Git.Service + const repository = yield* git.repo.discover(AbsolutePath.make(root.path)) + if (!repository) throw new Error("Repository not found") + const before = yield* git.tree.capture({ repository, scopes: [RelativePath.make(".")] }) + const lines = Math.ceil(VcsPatch.MAX_TOTAL_PATCH_BYTES / 80) + 1 + yield* Effect.promise(async () => { + await Bun.write(path.join(root.path, "a-small.txt"), "small\n") + await Bun.write(path.join(root.path, "b-large.txt"), `${"x".repeat(79)}\n`.repeat(lines)) + await Bun.write(path.join(root.path, "c-binary.bin"), new Uint8Array([0, 1, 2, 3])) + await Bun.write(path.join(root.path, "a-caf\u00e9.txt"), "caf\u00e9\n") + }) + const after = yield* git.tree.capture({ repository, scopes: [RelativePath.make(".")] }) + + const diffs = yield* git.tree.diff({ repository, from: before, to: after, context: 0 }) + expect(diffs.map((item) => [item.file, item.status, item.additions, item.deletions])).toEqual([ + ["a-caf\u00e9.txt", "added", 1, 0], + ["a-small.txt", "added", 1, 0], + ["b-large.txt", "added", lines, 0], + ["c-binary.bin", "added", 0, 0], + ]) + // Patch headers are not NUL-delimited; a quoted (octal-escaped) header would orphan this chunk. + expect(diffs[0]?.patch).toContain("+caf\u00e9\n") + expect(diffs[1]?.patch).toContain("+small\n") + expect(diffs[2]?.patch).toBe(VcsPatch.emptyPatch("b-large.txt")) + expect(diffs[3]?.patch).toBe("") + expect(yield* git.tree.diff({ repository, from: before, to: after, paths: [] })).toEqual([]) + }), + ) + it.live("captures, compares, previews, and restores scoped trees", () => Effect.gen(function* () { const root = yield* Effect.acquireRelease( diff --git a/packages/core/test/session-diff.test.ts b/packages/core/test/session-diff.test.ts new file mode 100644 index 000000000000..5c528482aa95 --- /dev/null +++ b/packages/core/test/session-diff.test.ts @@ -0,0 +1,198 @@ +import { $ } from "bun" +import { describe, expect } from "bun:test" +import fs from "fs/promises" +import path from "path" +import { Effect } from "effect" +import { Agent } from "@opencode-ai/core/agent" +import { Bus } from "@opencode-ai/core/bus" +import { Database } from "@opencode-ai/core/database/database" +import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" +import { LocationServiceMap } from "@opencode-ai/core/location-service-map" +import { Model } from "@opencode-ai/core/model" +import { Plugin } from "@opencode-ai/core/plugin" +import { Provider } from "@opencode-ai/core/provider" +import { AbsolutePath } from "@opencode-ai/core/schema" +import { Session } from "@opencode-ai/core/session" +import { SessionDiff } from "@opencode-ai/core/session/diff" +import { SessionEvent } from "@opencode-ai/core/session/event" +import { SessionExecution } from "@opencode-ai/core/session/execution" +import { SessionInbox } from "@opencode-ai/core/session/inbox" +import { SessionMessage } from "@opencode-ai/core/session/message" +import { SessionProjector } from "@opencode-ai/core/session/projector" +import { Snapshot } from "@opencode-ai/core/snapshot" +import { Money } from "@opencode-ai/schema/money" +import { LayerNode } from "@opencode-ai/util/effect/layer-node" +import { Global } from "@opencode-ai/util/global" +import { tempGlobalLayer } from "./fixture/global" +import { offlineModels } from "./fixture/models" +import { tmpdirScoped } from "./fixture/tmpdir" +import { testEffect } from "./lib/effect" + +const it = testEffect( + AppNodeBuilder.build( + LayerNode.group([Database.node, Bus.node, SessionProjector.node, Session.node, LocationServiceMap.node]), + [Global.node.replace(tempGlobalLayer), SessionExecution.node.replace(SessionExecution.noopLayer), offlineModels], + ), +) + +const summarize = (file: { file: string; status: string; additions: number; deletions: number }) => [ + file.file, + file.status, + file.additions, + file.deletions, +] + +describe("Session.diff", () => { + it.live( + "diffs the busy period containing a user message and ranges across later turns", + () => + Effect.gen(function* () { + const tmp = yield* tmpdirScoped() + const directory = path.join(tmp.path, "project") + const write = (name: string, content: string) => () => Bun.write(path.join(directory, name), content) + yield* Effect.promise(async () => { + await fs.mkdir(directory) + await write("first.txt", "first\n")() + await write("second.txt", "second\n")() + await write("manual.txt", "manual\n")() + await $`git init -q`.cwd(directory).quiet() + await $`git -c core.fsmonitor=false add .`.cwd(directory).quiet() + }) + const sessions = yield* Session.Service + const database = yield* Database.Service + const bus = yield* Bus.Service + const locations = yield* LocationServiceMap.Service + const created = yield* sessions.create({ location: { directory: AbsolutePath.make(directory) } }) + const diff = (input?: { messageID?: SessionMessage.ID; to?: SessionMessage.ID }) => + sessions + .diff({ sessionID: created.id, context: 0, ...input }) + .pipe(Effect.map((files) => files.map(summarize))) + expect(yield* diff()).toEqual([]) + + yield* Effect.gen(function* () { + const plugins = yield* Plugin.Service + yield* plugins.awaitActivation + const snapshot = yield* Snapshot.Service + const usage = { + cost: Money.USD.zero, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + } + const prompt = Effect.fn(function* (text: string) { + const admitted = yield* sessions.prompt({ sessionID: created.id, text, resume: false }) + yield* SessionInbox.promote(database.db, bus, created.id, "steer") + return admitted.id + }) + const step = Effect.fn(function* (edit: () => Promise, end: "recorded" | "unrecorded" | "running") { + const before = yield* snapshot.capture() + if (!before) throw new Error("Start snapshot missing") + const assistantMessageID = SessionMessage.ID.create() + yield* bus.publish(SessionEvent.Step.Started, { + sessionID: created.id, + assistantMessageID, + agent: Agent.defaultID, + model: { id: Model.ID.make("test-model"), providerID: Provider.ID.make("test-provider") }, + snapshot: before, + }) + yield* Effect.promise(edit) + if (end === "running") return assistantMessageID + const after = end === "recorded" ? yield* snapshot.capture() : undefined + yield* bus.publish(SessionEvent.Step.Ended, { + sessionID: created.id, + assistantMessageID, + finish: "stop", + ...usage, + snapshot: after, + files: after && before ? yield* snapshot.files({ from: before, to: after }) : undefined, + }) + return assistantMessageID + }) + + const idle = (outcome: "succeeded" | "failed") => + outcome === "succeeded" + ? bus.publish(SessionEvent.Execution.Succeeded, { sessionID: created.id }) + : bus.publish(SessionEvent.Execution.Failed, { + sessionID: created.id, + error: { type: "unknown", message: "failed" }, + }) + + // Before any idle marker exists, a prompt's turn ends at the next prompt. + const first = yield* prompt("Edit the first file") + const firstStep = yield* step(write("first.txt", "first edited\n"), "recorded") + // Edits made while idle are not a turn's work, but a range spanning them still sees them. + yield* Effect.promise(write("manual.txt", "manual edited\n")) + const second = yield* prompt("Edit the second file") + yield* step(write("second.txt", "second edited\n"), "recorded") + expect(yield* diff()).toEqual([["second.txt", "modified", 1, 1]]) + expect(yield* diff({ messageID: first })).toEqual([["first.txt", "modified", 1, 1]]) + + // Once markers exist, a turn spans a whole busy period, steers included; earlier history merges into the first one. + yield* idle("succeeded") + const third = yield* prompt("Add a third file") + yield* step(write("third.txt", "third\n"), "recorded") + const steer = yield* prompt("Also add a fourth file") + yield* step(write("fourth.txt", "fourth\n"), "recorded") + yield* idle("failed") + const busy = [ + ["fourth.txt", "added", 1, 0], + ["third.txt", "added", 1, 0], + ] + expect(yield* diff()).toEqual(busy) + expect(yield* diff({ messageID: steer })).toEqual(busy) + expect(yield* diff({ messageID: second })).toEqual([ + ["first.txt", "modified", 1, 1], + ["manual.txt", "modified", 1, 1], + ["second.txt", "modified", 1, 1], + ]) + expect(yield* diff({ messageID: first, to: third })).toEqual([ + ["first.txt", "modified", 1, 1], + ["fourth.txt", "added", 1, 0], + ["manual.txt", "modified", 1, 1], + ["second.txt", "modified", 1, 1], + ["third.txt", "added", 1, 0], + ]) + const full = yield* sessions.diff({ sessionID: created.id, messageID: first }) + expect(full[0]?.patch).toContain("-first\n+first edited\n") + expect(yield* diff({ messageID: steer, to: second }).pipe(Effect.flip)).toMatchObject({ + _tag: "Session.TurnRangeError", + field: "to", + }) + expect(yield* diff({ messageID: firstStep }).pipe(Effect.flip)).toMatchObject({ + _tag: "Session.TurnRangeError", + field: "messageID", + }) + expect(yield* diff({ messageID: SessionMessage.ID.create() }).pipe(Effect.flip)).toMatchObject({ + _tag: "Session.MessageNotFoundError", + }) + + // A completed step without an end snapshot falls back to the last recorded end. + yield* prompt("Edit both files again") + yield* step(write("first.txt", "first edited twice\n"), "recorded") + yield* step(write("second.txt", "second edited twice\n"), "unrecorded") + yield* idle("succeeded") + expect(yield* diff()).toEqual([["first.txt", "modified", 1, 1]]) + + // Only a step still running in the active session compares against the working copy. + yield* prompt("Delete the manual file") + yield* step(() => fs.rm(path.join(directory, "manual.txt")), "running") + expect(yield* diff()).toEqual([]) + const session = yield* sessions.get(created.id) + const live = yield* SessionDiff.turn(database.db, locations, { session, active: true, context: 0 }) + expect(live.map(summarize)).toEqual([["manual.txt", "deleted", 0, 1]]) + + // Reverting removes later history, markers included; a fork keeps the copied turns. + yield* sessions.revert.stage({ sessionID: created.id, messageID: steer, files: false }) + yield* sessions.revert.commit(created.id) + expect(yield* diff()).toEqual([["third.txt", "added", 1, 0]]) + expect(yield* diff({ messageID: steer }).pipe(Effect.flip)).toMatchObject({ + _tag: "Session.MessageNotFoundError", + }) + const forked = yield* sessions.fork({ sessionID: created.id, boundary: { type: "through" } }) + expect((yield* sessions.diff({ sessionID: forked.id, context: 0 })).map(summarize)).toEqual([ + ["third.txt", "added", 1, 0], + ]) + }).pipe(Effect.provide(LocationServiceMap.Service.get(created.location))) + }), + // Real Location/plugin startup and Git snapshots can exceed five seconds under CI load. + { timeout: 30_000 }, + ) +}) diff --git a/packages/core/test/session-execution.test.ts b/packages/core/test/session-execution.test.ts index 3d34e97f2c08..c6f473f464cc 100644 --- a/packages/core/test/session-execution.test.ts +++ b/packages/core/test/session-execution.test.ts @@ -561,7 +561,9 @@ describe("SessionRestart background recovery", () => { expect(yield* restarted.pendingBackground).toEqual([]) expect(yield* SessionInbox.list(database.db, sessionID)).toHaveLength(delivered ? 0 : 1) yield* SessionInbox.promote(database.db, bus, sessionID, "steer") - expect(yield* sessions.messages({ sessionID })).toMatchObject([ + // Recovery ends a busy period, so an idle marker follows the notification. + const messages = (yield* sessions.messages({ sessionID })).filter((message) => message.type !== "idle") + expect(messages).toMatchObject([ { id: background.notificationID, type: "synthetic", @@ -569,7 +571,6 @@ describe("SessionRestart background recovery", () => { metadata: { state: "completed" }, }, ]) - expect(yield* sessions.messages({ sessionID })).toHaveLength(1) }), ) } diff --git a/packages/protocol/openapi.json b/packages/protocol/openapi.json index 2b36ba43a68f..35b69f6e2bef 100644 --- a/packages/protocol/openapi.json +++ b/packages/protocol/openapi.json @@ -1621,14 +1621,7 @@ "content": { "application/json": { "schema": { - "anyOf": [ - { - "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" - }, - { - "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" - } - ] + "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" } } } @@ -3249,6 +3242,152 @@ "summary": "Get session context" } }, + "/api/session/{sessionID}/diff": { + "get": { + "tags": ["session"], + "operationId": "v2.session.diff", + "parameters": [ + { + "name": "sessionID", + "in": "path", + "schema": { + "type": "string", + "pattern": "^ses" + }, + "required": true + }, + { + "name": "messageID", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string", + "pattern": "^msg_" + }, + { + "type": "null" + } + ], + "description": "User message whose turn to diff. Defaults to the turn of the newest user message." + }, + "required": false + }, + { + "name": "to", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string", + "pattern": "^msg_" + }, + { + "type": "null" + } + ], + "description": "Later user message whose turn ends the range. Defaults to the turn of `messageID` alone." + }, + "required": false + }, + { + "name": "context", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "description": "Unchanged lines around each hunk. Omit for full-file patches." + }, + "required": false + } + ], + "security": [], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "data": { + "type": "array", + "items": { + "$ref": "#/components/schemas/FileDiff.Info" + } + } + }, + "required": ["data"], + "additionalProperties": false + } + } + } + }, + "400": { + "description": "InvalidRequestError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + }, + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + } + ] + } + } + } + }, + "401": { + "description": "UnauthorizedError", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UnauthorizedErrorEncoded" + } + } + } + }, + "404": { + "description": "MessageNotFoundError | SessionNotFoundError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/MessageNotFoundErrorEncoded" + }, + { + "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" + } + ] + } + } + } + }, + "500": { + "description": "UnknownError", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UnknownErrorEncoded" + } + } + } + } + }, + "description": "Structured per-file diffs of the files a turn changed. A turn runs from the first prompt after the session was last idle until its next idle marker, so prompts steered in while it was busy belong to the same turn; `to` extends the range through a later turn. Compares the range's first recorded snapshot with its last; a step still running in the active session compares against the working copy. Ranges that span a location change are rejected. In sessions without any idle marker, a prompt's turn spans until the next user message.", + "summary": "Diff session turns" + } + }, "/api/session/{sessionID}/inbox": { "get": { "tags": ["session"], @@ -18486,6 +18625,38 @@ "required": ["type", "id", "time", "status", "reason", "summary", "recent"], "additionalProperties": false }, + "Session.Message.Idle": { + "type": "object", + "properties": { + "id": { + "type": "string", + "pattern": "^msg_" + }, + "metadata": { + "type": "object" + }, + "time": { + "type": "object", + "properties": { + "created": { + "type": "number" + } + }, + "required": ["created"], + "additionalProperties": false + }, + "type": { + "type": "string", + "enum": ["idle"] + }, + "outcome": { + "type": "string", + "enum": ["succeeded", "failed", "interrupted"] + } + }, + "required": ["id", "time", "type", "outcome"], + "additionalProperties": false + }, "Session.Message.Info": { "anyOf": [ { @@ -18517,6 +18688,9 @@ }, { "$ref": "#/components/schemas/Session.Message.Compaction" + }, + { + "$ref": "#/components/schemas/Session.Message.Idle" } ] }, diff --git a/packages/protocol/src/groups/session.ts b/packages/protocol/src/groups/session.ts index b42d13b0e8fb..3fb8c3d54470 100644 --- a/packages/protocol/src/groups/session.ts +++ b/packages/protocol/src/groups/session.ts @@ -30,6 +30,7 @@ import { Model } from "@opencode-ai/schema/model" import { Location } from "@opencode-ai/schema/location" import { SessionEvent } from "@opencode-ai/schema/session-event" import { EventLog } from "@opencode-ai/schema/event-log" +import { FileDiff } from "@opencode-ai/schema/file-diff" const ParentIDFilter = Schema.Union([ Session.ID, @@ -521,6 +522,31 @@ export const makeSessionGroup = (sessionLo }), ), ) + .add( + HttpApiEndpoint.get("session.diff", "/api/session/:sessionID/diff", { + params: { sessionID: Session.ID }, + query: Schema.Struct({ + messageID: Schema.optional(SessionMessage.ID).annotate({ + description: "User message whose turn to diff. Defaults to the turn of the newest user message.", + }), + to: Schema.optional(SessionMessage.ID).annotate({ + description: "Later user message whose turn ends the range. Defaults to the turn of `messageID` alone.", + }), + context: Schema.NumberFromString.pipe(Schema.decodeTo(NonNegativeInt), Schema.optional).annotate({ + description: "Unchanged lines around each hunk. Omit for full-file patches.", + }), + }), + success: Schema.Struct({ data: Schema.Array(FileDiff.Info) }), + error: [InvalidRequestError, MessageNotFoundError, SessionNotFoundError, UnknownError], + }).annotateMerge( + OpenApi.annotations({ + identifier: "v2.session.diff", + summary: "Diff session turns", + description: + "Structured per-file diffs of the files a turn changed. A turn runs from the first prompt after the session was last idle until its next idle marker, so prompts steered in while it was busy belong to the same turn; `to` extends the range through a later turn. Compares the range's first recorded snapshot with its last; a step still running in the active session compares against the working copy. Ranges that span a location change are rejected. In sessions without any idle marker, a prompt's turn spans until the next user message.", + }), + ), + ) .add( HttpApiEndpoint.get("session.inbox.list", "/api/session/:sessionID/inbox", { params: { sessionID: Session.ID }, diff --git a/packages/schema/src/session-message.ts b/packages/schema/src/session-message.ts index 6bd37eb51746..7b8cf97a15ff 100644 --- a/packages/schema/src/session-message.ts +++ b/packages/schema/src/session-message.ts @@ -272,6 +272,18 @@ export const Compaction = Schema.Union([CompactionRunning, CompactionCompleted, ) export type Compaction = CompactionRunning | CompactionCompleted | CompactionFailed +/** + * Marks the Session going idle: every step since the previous marker belongs to + * one turn, including prompts steered in while it was busy. A shutdown does not + * record one, since the resumed execution continues the same turn. + */ +export interface Idle extends Schema.Schema.Type {} +export const Idle = Schema.Struct({ + ...Base, + type: Schema.tag("idle"), + outcome: Schema.Literals(["succeeded", "failed", "interrupted"]), +}).annotate({ identifier: "Session.Message.Idle" }) + export const Info = Schema.Union([ AgentSelected, ModelSelected, @@ -283,6 +295,7 @@ export const Info = Schema.Union([ Shell, Assistant, Compaction, + Idle, ]).annotate({ identifier: "Session.Message.Info" }) export type Info = | AgentSelected @@ -295,4 +308,5 @@ export type Info = | Shell | Assistant | Compaction + | Idle export type Type = Info["type"] diff --git a/packages/server/src/handlers/session-error.ts b/packages/server/src/handlers/session-error.ts index 65be0f03cee1..4d5e859c9f2d 100644 --- a/packages/server/src/handlers/session-error.ts +++ b/packages/server/src/handlers/session-error.ts @@ -1,5 +1,6 @@ import { Session } from "@opencode-ai/core/session" -import { SessionNotFoundError, UnknownError } from "@opencode-ai/protocol/errors" +import type { Snapshot } from "@opencode-ai/core/snapshot" +import { MessageNotFoundError, SessionNotFoundError, UnknownError } from "@opencode-ai/protocol/errors" import { Effect } from "effect" export function missingSession(error: Session.NotFoundError) { @@ -9,6 +10,14 @@ export function missingSession(error: Session.NotFoundError) { }) } +export function missingMessage(error: Session.MessageNotFoundError) { + return new MessageNotFoundError({ + sessionID: error.sessionID, + messageID: error.messageID, + message: `Message not found: ${error.messageID}`, + }) +} + export function failedMessageDecode(error: Session.MessageDecodeError) { const ref = `err_${crypto.randomUUID().slice(0, 8)}` return Effect.logError("failed to decode session message").pipe( @@ -18,3 +27,16 @@ export function failedMessageDecode(error: Session.MessageDecodeError) { ), ) } + +/** Snapshot repositories are host state clients cannot repair, so surface only a log reference. */ +export function failedSnapshot(operation: string, sessionID: Session.ID) { + return (error: Snapshot.Error) => { + const ref = `err_${crypto.randomUUID().slice(0, 8)}` + return Effect.logError(`failed to ${operation}`, { cause: error }).pipe( + Effect.annotateLogs({ ref, sessionID }), + Effect.andThen( + Effect.fail(new UnknownError({ message: "Unexpected server error. Check server logs for details.", ref })), + ), + ) + } +} diff --git a/packages/server/src/handlers/session.ts b/packages/server/src/handlers/session.ts index a8e8db9ba314..f5902d5d1fe0 100644 --- a/packages/server/src/handlers/session.ts +++ b/packages/server/src/handlers/session.ts @@ -17,10 +17,9 @@ import { ServiceUnavailableError, SessionBusyError, SkillNotFoundError, - UnknownError, } from "@opencode-ai/protocol/errors" import { AbsolutePath } from "@opencode-ai/core/schema" -import { failedMessageDecode, missingSession } from "./session-error" +import { failedMessageDecode, failedSnapshot, missingMessage, missingSession } from "./session-error" const DefaultSessionsLimit = 50 @@ -212,15 +211,7 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl return { data: yield* session.fork({ sessionID: ctx.params.sessionID, boundary: ctx.payload.boundary }).pipe( Effect.catchTag("Session.NotFoundError", missingSession), - Effect.catchTag( - "Session.MessageNotFoundError", - (error) => - new MessageNotFoundError({ - sessionID: error.sessionID, - messageID: error.messageID, - message: `Message not found: ${error.messageID}`, - }), - ), + Effect.catchTag("Session.MessageNotFoundError", missingMessage), Effect.catchTag( "Session.ForkEmptyError", (error) => new InvalidRequestError({ message: error.message, kind: "empty_session" }), @@ -448,32 +439,14 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl files: ctx.payload.files, }) return { - data: yield* session.revert.stage({ ...ctx.params, ...ctx.payload }).pipe( - Effect.catchTag("Session.NotFoundError", missingSession), - Effect.catchTag( - "Session.MessageNotFoundError", - (error) => - new MessageNotFoundError({ - sessionID: error.sessionID, - messageID: error.messageID, - message: `Message not found: ${error.messageID}`, - }), + data: yield* session.revert + .stage({ ...ctx.params, ...ctx.payload }) + .pipe( + Effect.catchTag("Session.NotFoundError", missingSession), + Effect.catchTag("Session.MessageNotFoundError", missingMessage), + Effect.catchTag("Session.BusyError", busySession), + Effect.catchTag("Snapshot.Error", failedSnapshot("stage session revert", ctx.params.sessionID)), ), - Effect.catchTag("Session.BusyError", busySession), - Effect.catchTag("Snapshot.Error", (error) => { - const ref = `err_${crypto.randomUUID().slice(0, 8)}` - return Effect.logError("failed to stage session revert", { cause: error }).pipe( - Effect.andThen( - Effect.fail( - new UnknownError({ - message: "Unexpected server error. Check server logs for details.", - ref, - }), - ), - ), - ) - }), - ), } }), ) @@ -481,23 +454,13 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl "session.revert.clear", Effect.fn(function* (ctx) { yield* Effect.log("session.revert.clear", { sessionID: ctx.params.sessionID }) - yield* session.revert.clear(ctx.params.sessionID).pipe( - Effect.catchTag("Session.NotFoundError", missingSession), - Effect.catchTag("Session.BusyError", busySession), - Effect.catchTag("Snapshot.Error", (error) => { - const ref = `err_${crypto.randomUUID().slice(0, 8)}` - return Effect.logError("failed to clear session revert", { cause: error }).pipe( - Effect.andThen( - Effect.fail( - new UnknownError({ - message: "Unexpected server error. Check server logs for details.", - ref, - }), - ), - ), - ) - }), - ) + yield* session.revert + .clear(ctx.params.sessionID) + .pipe( + Effect.catchTag("Session.NotFoundError", missingSession), + Effect.catchTag("Session.BusyError", busySession), + Effect.catchTag("Snapshot.Error", failedSnapshot("clear session revert", ctx.params.sessionID)), + ) return HttpApiSchema.NoContent.make() }), ) @@ -527,6 +490,22 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl } }), ) + .handle( + "session.diff", + Effect.fn(function* (ctx) { + return { + data: yield* session.diff({ sessionID: ctx.params.sessionID, ...ctx.query }).pipe( + Effect.catchTag("Session.NotFoundError", missingSession), + Effect.catchTag("Session.MessageNotFoundError", missingMessage), + Effect.catchTag( + "Session.TurnRangeError", + (error) => new InvalidRequestError({ message: error.message, field: error.field }), + ), + Effect.catchTag("Snapshot.Error", failedSnapshot("diff session turn", ctx.params.sessionID)), + ), + } + }), + ) .handle( "session.inbox.list", Effect.fn(function* (ctx) { @@ -642,15 +621,7 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl Effect.fn(function* (ctx) { const message = yield* session.updateMessage({ ...ctx.params, content: ctx.payload.content }).pipe( Effect.catchTag("Session.NotFoundError", missingSession), - Effect.catchTag( - "Session.MessageNotFoundError", - (error) => - new MessageNotFoundError({ - sessionID: error.sessionID, - messageID: error.messageID, - message: `Message not found: ${error.messageID}`, - }), - ), + Effect.catchTag("Session.MessageNotFoundError", missingMessage), Effect.catchTag("Session.BusyError", busySession), Effect.catchTag( "Session.MessageNotAssistantError", diff --git a/packages/server/test/session-diff.test.ts b/packages/server/test/session-diff.test.ts new file mode 100644 index 000000000000..90f9efab2180 --- /dev/null +++ b/packages/server/test/session-diff.test.ts @@ -0,0 +1,98 @@ +import { expect, setDefaultTimeout } from "bun:test" +import { Agent } from "@opencode-ai/core/agent" +import { Bus } from "@opencode-ai/core/bus" +import { Model } from "@opencode-ai/core/model" +import { Provider } from "@opencode-ai/core/provider" +import { Session } from "@opencode-ai/core/session" +import { SessionEvent } from "@opencode-ai/core/session/event" +import { SessionExecution } from "@opencode-ai/core/session/execution" +import { SessionMessage } from "@opencode-ai/core/session/message" +import { Money } from "@opencode-ai/schema/money" +import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" +import { Effect, Layer } from "effect" +import { tmpdir } from "../../core/test/fixture/tmpdir" +import { it } from "../../core/test/lib/effect" +import { ServerFetch } from "../src/fetch" + +setDefaultTimeout(30_000) + +it.live("serves turn diffs by user message with range validation", () => + Effect.gen(function* () { + const tmp = yield* Effect.acquireDisposable(Effect.promise(() => tmpdir("opencode-session-diff-"))) + const ids = { user: SessionMessage.ID.create(), assistant: SessionMessage.ID.create() } + // Deliver the prompt and one step the way the runner would, without a model. + const execution = Layer.effect( + SessionExecution.Service, + Effect.gen(function* () { + const bus = yield* Bus.Service + return SessionExecution.Service.of({ + active: Effect.succeed(new Set()), + isActive: () => Effect.succeed(false), + resume: () => Effect.void, + wake: (sessionID) => + Effect.gen(function* () { + yield* bus.publish(SessionEvent.InboxDelivered, { sessionID, inboxID: ids.user }) + yield* bus.publish(SessionEvent.Step.Started, { + sessionID, + assistantMessageID: ids.assistant, + agent: Agent.defaultID, + model: { id: Model.ID.make("model"), providerID: Provider.ID.make("provider") }, + }) + yield* bus.publish(SessionEvent.Step.Ended, { + sessionID, + assistantMessageID: ids.assistant, + finish: "stop", + cost: Money.USD.zero, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + }) + }), + interrupt: () => Effect.succeed(false), + awaitIdle: () => Effect.void, + }) + }), + ) + const handler = yield* ServerFetch.make( + { + app: { version: "test-version" }, + database: { path: ":memory:" }, + fs: { filewatcher: false }, + models: { fetch: false }, + }, + { + overrides: [ + SessionExecution.node.replace( + makeGlobalNode({ service: SessionExecution.Service, layer: execution, deps: [Bus.node] }), + ), + ], + }, + ) + const request = (path: string, body?: unknown) => + Effect.promise(async () => { + const response = await handler( + new Request(`http://opencode.local${path}`, { + method: body === undefined ? "GET" : "POST", + headers: body === undefined ? undefined : { "content-type": "application/json" }, + body: body === undefined ? undefined : JSON.stringify(body), + }), + ) + return { status: response.status, body: (await response.json()) as Record } + }) + const created = yield* request("/api/session", { location: { directory: tmp.path } }) + const sessionID = Session.ID.make((created.body.data as { id: string }).id) + const diff = (query = "") => request(`/api/session/${sessionID}/diff${query}`) + + expect(yield* diff()).toEqual({ status: 200, body: { data: [] } }) + expect((yield* request(`/api/session/${sessionID}/prompt`, { id: ids.user, text: "prompt" })).status).toBe(200) + // Not a git repository, so steps record no snapshots and the turn has no diff. + expect(yield* diff(`?messageID=${ids.user}&context=3`)).toEqual({ status: 200, body: { data: [] } }) + expect(yield* diff(`?messageID=${ids.assistant}`)).toMatchObject({ + status: 400, + body: { _tag: "InvalidRequestError", field: "messageID" }, + }) + expect(yield* diff(`?messageID=${SessionMessage.ID.create()}`)).toMatchObject({ + status: 404, + body: { _tag: "MessageNotFoundError" }, + }) + expect((yield* request(`/api/session/${Session.ID.create()}/diff`)).status).toBe(404) + }), +) diff --git a/packages/session-ui/src/timeline/projection.ts b/packages/session-ui/src/timeline/projection.ts index 29ee7960158c..cc725f3048a1 100644 --- a/packages/session-ui/src/timeline/projection.ts +++ b/packages/session-ui/src/timeline/projection.ts @@ -16,7 +16,7 @@ export { TimelineRow, type PartGroup, type PartRef, type TimelineRowMap } export type ReasoningMode = "hidden" | "compact" | "full" -type Notice = Exclude +type Notice = Exclude type Entry = { type: "assistant"; message: SessionMessageAssistant } | { type: "notice"; message: Notice } type Content = SessionMessageAssistant["content"][number] type GroupRow = Extract @@ -765,7 +765,8 @@ function record(value: unknown): value is Record { } function isNotice(message: SessionMessageInfo): message is Notice { - if (message.type === "user" || message.type === "assistant" || message.type === "shell") return false + if (message.type === "user" || message.type === "assistant" || message.type === "shell" || message.type === "idle") + return false if (message.type !== "synthetic") return true return !!message.description?.trim() || timelineNoticeRequired(message) } diff --git a/packages/tui/src/routes/session/rows.ts b/packages/tui/src/routes/session/rows.ts index c1f61f849c15..cf55cea67512 100644 --- a/packages/tui/src/routes/session/rows.ts +++ b/packages/tui/src/routes/session/rows.ts @@ -305,6 +305,7 @@ export function reduceSessionRows(messages: SessionMessageInfo[], inputs = new S ...messages.filter(isInput), ].reduce((rows, message) => { if (message.type !== "assistant") { + if (message.type === "idle") return rows if (message.type === "synthetic" && !message.description?.trim()) return rows if (message.type === "compaction" && message.status === "completed" && usage) usage.previousTurnCache = undefined if (!pending.has(message.id)) completePrevious(rows) diff --git a/packages/www/openapi.json b/packages/www/openapi.json index 2b36ba43a68f..35b69f6e2bef 100644 --- a/packages/www/openapi.json +++ b/packages/www/openapi.json @@ -1621,14 +1621,7 @@ "content": { "application/json": { "schema": { - "anyOf": [ - { - "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" - }, - { - "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" - } - ] + "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" } } } @@ -3249,6 +3242,152 @@ "summary": "Get session context" } }, + "/api/session/{sessionID}/diff": { + "get": { + "tags": ["session"], + "operationId": "v2.session.diff", + "parameters": [ + { + "name": "sessionID", + "in": "path", + "schema": { + "type": "string", + "pattern": "^ses" + }, + "required": true + }, + { + "name": "messageID", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string", + "pattern": "^msg_" + }, + { + "type": "null" + } + ], + "description": "User message whose turn to diff. Defaults to the turn of the newest user message." + }, + "required": false + }, + { + "name": "to", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string", + "pattern": "^msg_" + }, + { + "type": "null" + } + ], + "description": "Later user message whose turn ends the range. Defaults to the turn of `messageID` alone." + }, + "required": false + }, + { + "name": "context", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "description": "Unchanged lines around each hunk. Omit for full-file patches." + }, + "required": false + } + ], + "security": [], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "data": { + "type": "array", + "items": { + "$ref": "#/components/schemas/FileDiff.Info" + } + } + }, + "required": ["data"], + "additionalProperties": false + } + } + } + }, + "400": { + "description": "InvalidRequestError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + }, + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + } + ] + } + } + } + }, + "401": { + "description": "UnauthorizedError", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UnauthorizedErrorEncoded" + } + } + } + }, + "404": { + "description": "MessageNotFoundError | SessionNotFoundError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/MessageNotFoundErrorEncoded" + }, + { + "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" + } + ] + } + } + } + }, + "500": { + "description": "UnknownError", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UnknownErrorEncoded" + } + } + } + } + }, + "description": "Structured per-file diffs of the files a turn changed. A turn runs from the first prompt after the session was last idle until its next idle marker, so prompts steered in while it was busy belong to the same turn; `to` extends the range through a later turn. Compares the range's first recorded snapshot with its last; a step still running in the active session compares against the working copy. Ranges that span a location change are rejected. In sessions without any idle marker, a prompt's turn spans until the next user message.", + "summary": "Diff session turns" + } + }, "/api/session/{sessionID}/inbox": { "get": { "tags": ["session"], @@ -18486,6 +18625,38 @@ "required": ["type", "id", "time", "status", "reason", "summary", "recent"], "additionalProperties": false }, + "Session.Message.Idle": { + "type": "object", + "properties": { + "id": { + "type": "string", + "pattern": "^msg_" + }, + "metadata": { + "type": "object" + }, + "time": { + "type": "object", + "properties": { + "created": { + "type": "number" + } + }, + "required": ["created"], + "additionalProperties": false + }, + "type": { + "type": "string", + "enum": ["idle"] + }, + "outcome": { + "type": "string", + "enum": ["succeeded", "failed", "interrupted"] + } + }, + "required": ["id", "time", "type", "outcome"], + "additionalProperties": false + }, "Session.Message.Info": { "anyOf": [ { @@ -18517,6 +18688,9 @@ }, { "$ref": "#/components/schemas/Session.Message.Compaction" + }, + { + "$ref": "#/components/schemas/Session.Message.Idle" } ] }, diff --git a/packages/www/public/openapi.json b/packages/www/public/openapi.json index 2b36ba43a68f..35b69f6e2bef 100644 --- a/packages/www/public/openapi.json +++ b/packages/www/public/openapi.json @@ -1621,14 +1621,7 @@ "content": { "application/json": { "schema": { - "anyOf": [ - { - "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" - }, - { - "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" - } - ] + "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" } } } @@ -3249,6 +3242,152 @@ "summary": "Get session context" } }, + "/api/session/{sessionID}/diff": { + "get": { + "tags": ["session"], + "operationId": "v2.session.diff", + "parameters": [ + { + "name": "sessionID", + "in": "path", + "schema": { + "type": "string", + "pattern": "^ses" + }, + "required": true + }, + { + "name": "messageID", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string", + "pattern": "^msg_" + }, + { + "type": "null" + } + ], + "description": "User message whose turn to diff. Defaults to the turn of the newest user message." + }, + "required": false + }, + { + "name": "to", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string", + "pattern": "^msg_" + }, + { + "type": "null" + } + ], + "description": "Later user message whose turn ends the range. Defaults to the turn of `messageID` alone." + }, + "required": false + }, + { + "name": "context", + "in": "query", + "schema": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "description": "Unchanged lines around each hunk. Omit for full-file patches." + }, + "required": false + } + ], + "security": [], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "data": { + "type": "array", + "items": { + "$ref": "#/components/schemas/FileDiff.Info" + } + } + }, + "required": ["data"], + "additionalProperties": false + } + } + } + }, + "400": { + "description": "InvalidRequestError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + }, + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + } + ] + } + } + } + }, + "401": { + "description": "UnauthorizedError", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UnauthorizedErrorEncoded" + } + } + } + }, + "404": { + "description": "MessageNotFoundError | SessionNotFoundError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/MessageNotFoundErrorEncoded" + }, + { + "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" + } + ] + } + } + } + }, + "500": { + "description": "UnknownError", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UnknownErrorEncoded" + } + } + } + } + }, + "description": "Structured per-file diffs of the files a turn changed. A turn runs from the first prompt after the session was last idle until its next idle marker, so prompts steered in while it was busy belong to the same turn; `to` extends the range through a later turn. Compares the range's first recorded snapshot with its last; a step still running in the active session compares against the working copy. Ranges that span a location change are rejected. In sessions without any idle marker, a prompt's turn spans until the next user message.", + "summary": "Diff session turns" + } + }, "/api/session/{sessionID}/inbox": { "get": { "tags": ["session"], @@ -18486,6 +18625,38 @@ "required": ["type", "id", "time", "status", "reason", "summary", "recent"], "additionalProperties": false }, + "Session.Message.Idle": { + "type": "object", + "properties": { + "id": { + "type": "string", + "pattern": "^msg_" + }, + "metadata": { + "type": "object" + }, + "time": { + "type": "object", + "properties": { + "created": { + "type": "number" + } + }, + "required": ["created"], + "additionalProperties": false + }, + "type": { + "type": "string", + "enum": ["idle"] + }, + "outcome": { + "type": "string", + "enum": ["succeeded", "failed", "interrupted"] + } + }, + "required": ["id", "time", "type", "outcome"], + "additionalProperties": false + }, "Session.Message.Info": { "anyOf": [ { @@ -18517,6 +18688,9 @@ }, { "$ref": "#/components/schemas/Session.Message.Compaction" + }, + { + "$ref": "#/components/schemas/Session.Message.Idle" } ] },