Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
fix(api): use moving session history pages
  • Loading branch information
kitlangton committed Jun 26, 2026
commit e111f3d76556a903318aaf42f4ca81d9e8908b50
1 change: 0 additions & 1 deletion packages/client/src/effect.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,4 +10,3 @@ export { Session } from "@opencode-ai/schema/session"
export { SessionInput } from "@opencode-ai/schema/session-input"
export { SessionMessage } from "@opencode-ai/schema/session-message"
export { Prompt } from "@opencode-ai/schema/prompt"
export { SessionHistoryCursor } from "@opencode-ai/protocol/groups/session"
4 changes: 2 additions & 2 deletions packages/client/src/generated-effect/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -153,12 +153,12 @@ type Endpoint0_13Request = Parameters<RawClient["server.session"]["session.histo
type Endpoint0_13Input = {
readonly sessionID: Endpoint0_13Request["params"]["sessionID"]
readonly limit?: Endpoint0_13Request["query"]["limit"]
readonly cursor?: Endpoint0_13Request["query"]["cursor"]
readonly after?: Endpoint0_13Request["query"]["after"]
}
const Endpoint0_13 = (raw: RawClient["server.session"]) => (input: Endpoint0_13Input) =>
raw["session.history"]({
params: { sessionID: input.sessionID },
query: { limit: input.limit, cursor: input.cursor },
query: { limit: input.limit, after: input.after },
}).pipe(Effect.mapError(mapClientError))

type Endpoint0_14Request = Parameters<RawClient["server.session"]["session.events"]>[0]
Expand Down
2 changes: 1 addition & 1 deletion packages/client/src/generated/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -331,7 +331,7 @@ export function make(options: ClientOptions) {
{
method: "GET",
path: `/api/session/${encodeURIComponent(input.sessionID)}/history`,
query: { limit: input.limit, cursor: input.cursor },
query: { limit: input.limit, after: input.after },
successStatus: 200,
declaredStatuses: [404, 400, 401],
empty: false,
Expand Down
192 changes: 40 additions & 152 deletions packages/client/src/generated/types.ts

Large diffs are not rendered by default.

1 change: 0 additions & 1 deletion packages/client/src/index.ts
Original file line number Diff line number Diff line change
@@ -1,2 +1 @@
export * from "./generated/index"
export { SessionHistoryCursor } from "./session-history-cursor"
17 changes: 0 additions & 17 deletions packages/client/src/session-history-cursor.ts

This file was deleted.

31 changes: 16 additions & 15 deletions packages/client/test/effect.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@ import {
OpenCode,
Prompt,
Session,
SessionHistoryCursor,
SessionMessage,
} from "../src/effect"

Expand Down Expand Up @@ -48,8 +47,8 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
request,
Response.json(
historyPage === 1
? { data: [modelSwitchedEvent], cursor: { next: "opaque_history_cursor" } }
: { data: [], cursor: {} },
? { data: [modelSwitchedEvent], hasMore: true }
: { data: [], hasMore: false },
),
),
)
Expand Down Expand Up @@ -100,13 +99,13 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
const context = yield* client.sessions.context({ sessionID: Session.ID.make("ses_test") })
const history = yield* client.sessions.history({
sessionID: Session.ID.make("ses_test"),
cursor: SessionHistoryCursor.after(0),
after: 0,
limit: 1,
})
const historyNext = history.cursor.next
const historyNext = history.hasMore
? yield* client.sessions.history({
sessionID: Session.ID.make("ses_test"),
cursor: history.cursor.next,
after: history.data.at(-1)?.durable?.seq,
limit: 2,
})
: undefined
Expand All @@ -131,34 +130,36 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
expect(DateTime.toEpochMillis(result.admitted.timeCreated)).toBe(1_717_171_717_000)
expect(result.context).toEqual([])
expect(DateTime.toEpochMillis(result.history.data[0].data.timestamp)).toBe(1_717_171_717_000)
expect(result.history).toEqual(expect.objectContaining({ cursor: { next: "opaque_history_cursor" } }))
expect(result.historyNext).toEqual({ data: [], cursor: {} })
expect(historyQueries[0]).toEqual({ limit: "1", cursor: "eyJhZnRlciI6MH0" })
expect(historyQueries[1]).toEqual({ limit: "2", cursor: "opaque_history_cursor" })
expect(result.history).toEqual(expect.objectContaining({ hasMore: true }))
expect(result.historyNext).toEqual({ data: [], hasMore: false })
expect(historyQueries[0]).toEqual({ limit: "1", after: "0" })
expect(historyQueries[1]).toEqual({ limit: "2", after: "1" })
expect(DateTime.toEpochMillis(result.events[0].data.timestamp)).toBe(1_717_171_717_000)
expect(result.message).toEqual(expect.objectContaining({ id: "msg_model", type: "model-switched" }))
})

test("sessions.history retains the typed InvalidCursorError", async () => {
test("sessions.history retains the typed SessionNotFoundError", async () => {
const httpClient = HttpClient.make((request) =>
Effect.succeed(
HttpClientResponse.fromWeb(
request,
Response.json({ _tag: "InvalidCursorError", message: "Invalid cutoff" }, { status: 400 }),
Response.json(
{ _tag: "SessionNotFoundError", sessionID: "ses_missing", message: "Session not found" },
{ status: 404 },
),
),
),
)
const error = await Effect.gen(function* () {
const client = yield* OpenCode.make({ baseUrl: "http://localhost:3000" })
return yield* client.sessions
.history({
sessionID: Session.ID.make("ses_test"),
cursor: SessionHistoryCursor.after(0),
sessionID: Session.ID.make("ses_missing"),
})
.pipe(Effect.flip)
}).pipe(Effect.provideService(HttpClient.HttpClient, httpClient), Effect.runPromise)

expect(error._tag).toBe("InvalidCursorError")
expect(error._tag).toBe("SessionNotFoundError")
})

const session = {
Expand Down
44 changes: 21 additions & 23 deletions packages/client/test/promise.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { expect, test } from "bun:test"
import { isInvalidCursorError, isUnauthorizedError, OpenCode, SessionHistoryCursor } from "../src"
import { isSessionNotFoundError, isUnauthorizedError, OpenCode } from "../src"

test("sessions.get returns the wire projection", async () => {
const client = OpenCode.make({
Expand Down Expand Up @@ -34,8 +34,8 @@ test("session methods use the public HTTP contract", async () => {
historyPage++
return Response.json(
historyPage === 1
? { data: [modelSwitchedEvent], cursor: { next: "opaque_history_cursor" } }
: { data: [], cursor: {} },
? { data: [modelSwitchedEvent], hasMore: true }
: { data: [], hasMore: false },
)
}
if (url.includes("/prompt")) return Response.json(admission)
Expand All @@ -48,7 +48,7 @@ test("session methods use the public HTTP contract", async () => {
},
})

const page = await client.sessions.list({ limit: "10", order: "desc" })
const page = await client.sessions.list({ limit: 10, order: "desc" })
const active = await client.sessions.active()
const created = await client.sessions.create({ location: { directory: "/tmp/project" } })
await client.sessions.switchAgent({ sessionID: "ses_test", agent: "build" })
Expand All @@ -64,13 +64,13 @@ test("session methods use the public HTTP contract", async () => {
await client.sessions.compact({ sessionID: "ses_test" })
await client.sessions.wait({ sessionID: "ses_test" })
const context = await client.sessions.context({ sessionID: "ses_test" })
const initialCursor = SessionHistoryCursor.after(0)
const history = await client.sessions.history({ sessionID: "ses_test", cursor: initialCursor, limit: "1" })
const historyNext = history.cursor.next
? await client.sessions.history({ sessionID: "ses_test", cursor: history.cursor.next, limit: "2" })
const history = await client.sessions.history({ sessionID: "ses_test", after: 0, limit: 1 })
const historyAfter = history.data.at(-1)?.durable?.seq
const historyNext = history.hasMore
? await client.sessions.history({ sessionID: "ses_test", after: historyAfter, limit: 2 })
: undefined
const events = []
for await (const event of client.sessions.events({ sessionID: "ses_test", after: "0" })) events.push(event)
for await (const event of client.sessions.events({ sessionID: "ses_test", after: 0 })) events.push(event)
await client.sessions.interrupt({ sessionID: "ses_test" })
const message = await client.sessions.message({ sessionID: "ses_test", messageID: "msg_model" })

Expand All @@ -79,9 +79,8 @@ test("session methods use the public HTTP contract", async () => {
expect(created.id).toBe("ses_test")
expect(admitted.id).toBe("msg_test")
expect(context).toEqual([])
expect(initialCursor).toBe("eyJhZnRlciI6MH0")
expect(history).toEqual({ data: [modelSwitchedEvent], cursor: { next: "opaque_history_cursor" } })
expect(historyNext).toEqual({ data: [], cursor: {} })
expect(history).toEqual({ data: [modelSwitchedEvent], hasMore: true })
expect(historyNext).toEqual({ data: [], hasMore: false })
expect(events).toEqual([modelSwitchedEvent])
expect(message).toEqual(modelSwitchedMessage)
expect(requests.map((request) => [request.init?.method, request.url])).toEqual([
Expand All @@ -94,8 +93,8 @@ test("session methods use the public HTTP contract", async () => {
["POST", "http://localhost:3000/api/session/ses_test/compact"],
["POST", "http://localhost:3000/api/session/ses_test/wait"],
["GET", "http://localhost:3000/api/session/ses_test/context"],
["GET", `http://localhost:3000/api/session/ses_test/history?limit=1&cursor=${initialCursor}`],
["GET", "http://localhost:3000/api/session/ses_test/history?limit=2&cursor=opaque_history_cursor"],
["GET", "http://localhost:3000/api/session/ses_test/history?limit=1&after=0"],
["GET", "http://localhost:3000/api/session/ses_test/history?limit=2&after=1"],
["GET", "http://localhost:3000/api/session/ses_test/event?after=0"],
["POST", "http://localhost:3000/api/session/ses_test/interrupt"],
["GET", "http://localhost:3000/api/session/ses_test/message/msg_model"],
Expand Down Expand Up @@ -123,25 +122,24 @@ test("middleware errors remain declared client errors", async () => {
}
})

test("sessions.history decodes InvalidCursorError", async () => {
test("sessions.history decodes SessionNotFoundError", async () => {
const client = OpenCode.make({
baseUrl: "http://localhost:3000",
fetch: async () => Response.json({ _tag: "InvalidCursorError", message: "Invalid cutoff" }, { status: 400 }),
fetch: async () =>
Response.json(
{ _tag: "SessionNotFoundError", sessionID: "ses_missing", message: "Session not found" },
{ status: 404 },
),
})

try {
await client.sessions.history({ sessionID: "ses_test", cursor: "malformed" })
await client.sessions.history({ sessionID: "ses_missing" })
throw new Error("Expected request to fail")
} catch (error) {
expect(isInvalidCursorError(error)).toBe(true)
expect(isSessionNotFoundError(error)).toBe(true)
}
})

test("SessionHistoryCursor rejects invalid durable checkpoints", () => {
expect(() => SessionHistoryCursor.after(-1)).toThrow(RangeError)
expect(() => SessionHistoryCursor.after(1.5)).toThrow(RangeError)
})

const session = {
data: {
id: "ses_test",
Expand Down
60 changes: 16 additions & 44 deletions packages/core/src/event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ export * as EventV2 from "./event"
import { Cause, Context, Effect, Layer, Option, PubSub, Schema, Stream } from "effect"
import { Event } from "@opencode-ai/schema/event"
import type { Data, Definition, Payload } from "@opencode-ai/schema/event"
import { and, asc, eq, gt, inArray, lte } from "drizzle-orm"
import { and, asc, eq, gt, inArray } from "drizzle-orm"
import { Database } from "./database/database"
import { EventSequenceTable, EventTable } from "./event/sql"
import { Location } from "./location"
Expand Down Expand Up @@ -47,10 +47,6 @@ export class InvalidDurableEventError extends Schema.TaggedErrorClass<InvalidDur
},
) {}

export class InvalidCursorError extends Schema.TaggedErrorClass<InvalidCursorError>()("EventV2.InvalidCursorError", {
message: Schema.String,
}) {}

const decodeSerializedEvent = (event: SerializedEvent): Payload => {
const definition = Durable.get(event.type)
if (!definition?.durable) {
Expand All @@ -69,7 +65,6 @@ export const readAggregate = Effect.fn("EventV2.readAggregate")(function* <A>(
input: {
readonly aggregateID: string
readonly after?: number
readonly through?: number
readonly limit: number
readonly manifest: {
readonly definitions: ReadonlyMap<string, Definition>
Expand All @@ -78,43 +73,21 @@ export const readAggregate = Effect.fn("EventV2.readAggregate")(function* <A>(
},
) {
const after = input.after ?? -1
const result = yield* db
.transaction((tx) =>
Effect.gen(function* () {
const sequence = yield* tx
.select({ seq: EventSequenceTable.seq })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, input.aggregateID))
.get()
.pipe(Effect.orDie)
const head = sequence?.seq ?? -1
if (input.through !== undefined && input.through > head) {
return yield* new InvalidCursorError({ message: "History cutoff is above the current aggregate head" })
}
const through = input.through ?? head
if (through < after) {
return yield* new InvalidCursorError({ message: "History cutoff must not be less than the cursor" })
}
const rows = yield* tx
.select()
.from(EventTable)
.where(
and(
eq(EventTable.aggregate_id, input.aggregateID),
gt(EventTable.seq, after),
lte(EventTable.seq, through),
inArray(EventTable.type, Array.from(input.manifest.definitions.keys())),
),
)
.orderBy(asc(EventTable.seq))
.limit(input.limit + 1)
.all()
.pipe(Effect.orDie)
return { rows, through }
}),
const rows = yield* db
.select()
.from(EventTable)
.where(
and(
eq(EventTable.aggregate_id, input.aggregateID),
gt(EventTable.seq, after),
inArray(EventTable.type, Array.from(input.manifest.definitions.keys())),
),
)
.pipe(Effect.catchTag("SqlError", Effect.die))
const page = result.rows.slice(0, input.limit)
.orderBy(asc(EventTable.seq))
.limit(input.limit + 1)
.all()
.pipe(Effect.orDie)
const page = rows.slice(0, input.limit)
const decode = Schema.decodeUnknownSync(input.manifest.schema)
const events = page.map((event) =>
decode({
Expand All @@ -130,8 +103,7 @@ export const readAggregate = Effect.fn("EventV2.readAggregate")(function* <A>(
)
return {
events,
through: result.through,
nextAfter: result.rows.length > input.limit ? page.at(-1)?.seq : undefined,
hasMore: rows.length > input.limit,
}
})

Expand Down
21 changes: 4 additions & 17 deletions packages/core/src/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -103,18 +103,10 @@ export class PromptConflictError extends Schema.TaggedErrorClass<PromptConflictE
sessionID: SessionSchema.ID,
messageID: SessionMessage.ID,
}) {}
export class InvalidCursorError extends Schema.TaggedErrorClass<InvalidCursorError>()("Session.InvalidCursorError", {
message: Schema.String,
}) {}
export const MessageNotFoundError = SessionRevert.MessageNotFoundError
export type MessageNotFoundError = SessionRevert.MessageNotFoundError

export type Error =
| NotFoundError
| MessageDecodeError
| OperationUnavailableError
| PromptConflictError
| InvalidCursorError
export type Error = NotFoundError | MessageDecodeError | OperationUnavailableError | PromptConflictError

export interface Interface {
readonly list: (input?: ListInput) => Effect.Effect<SessionSchema.Info[]>
Expand Down Expand Up @@ -143,15 +135,10 @@ export interface Interface {
readonly history: (input: {
sessionID: SessionSchema.ID
after?: number
through?: number
limit: number
}) => Effect.Effect<
{
events: ReadonlyArray<SessionEvent.DurableEvent>
through: number
nextAfter?: number
},
NotFoundError | InvalidCursorError
{ events: ReadonlyArray<SessionEvent.DurableEvent>; hasMore: boolean },
NotFoundError
>
readonly switchAgent: (input: { sessionID: SessionSchema.ID; agent: string }) => Effect.Effect<void, NotFoundError>
readonly switchModel: (input: {
Expand Down Expand Up @@ -375,7 +362,7 @@ export const layer = Layer.unwrap(
...input,
aggregateID: input.sessionID,
manifest: SessionDurable,
}).pipe(Effect.mapError((error) => new InvalidCursorError({ message: error.message })))
})
}),
prompt: Effect.fn("V2Session.prompt")((input) =>
Effect.uninterruptible(
Expand Down
Loading
Loading