Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,15 @@ export const DiffQuery = Schema.Struct({
export const MessagesQuery = Schema.Struct({
...WorkspaceRoutingQueryFields,
limit: Schema.optional(Schema.NumberFromString.check(Schema.isInt(), Schema.isGreaterThanOrEqualTo(0))),
before: Schema.optional(Schema.String),
before: Schema.optional(Schema.String).annotate({
description: "Return messages older than this cursor (X-Next-Cursor of a previous page) or message ID",
}),
after: Schema.optional(Schema.String).annotate({
description: "Return messages newer than this cursor (X-Next-Cursor of a previous page) or message ID",
}),
order: Schema.optional(Schema.Literals(["asc", "desc"])).annotate({
description: "Without before or after: 'desc' (default) pages from the newest message, 'asc' from the oldest",
}),
})
export const StatusMap = Schema.Record(Schema.String, SessionStatus.Info)
export const UpdatePayload = Schema.Struct({
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -107,11 +107,14 @@ export const sessionHandlers = HttpApiBuilder.group(InstanceHttpApi, "session",
params: { sessionID: SessionID }
query: typeof MessagesQuery.Type
}) {
if (ctx.query.before && ctx.query.limit === undefined) return yield* new HttpApiError.BadRequest({})
if (ctx.query.before) {
const before = ctx.query.before
const anchor = ctx.query.after ?? ctx.query.before
if (ctx.query.before && ctx.query.after) return yield* new HttpApiError.BadRequest({})
if (ctx.query.after && ctx.query.order === "desc") return yield* new HttpApiError.BadRequest({})
if (ctx.query.before && ctx.query.order === "asc") return yield* new HttpApiError.BadRequest({})
if (anchor && ctx.query.limit === undefined) return yield* new HttpApiError.BadRequest({})
if (anchor && !anchor.startsWith("msg")) {
yield* Effect.try({
try: () => MessageV2.cursor.decode(before),
try: () => MessageV2.cursor.decode(anchor),
catch: () => new HttpApiError.BadRequest({}),
})
}
Expand All @@ -125,21 +128,35 @@ export const sessionHandlers = HttpApiBuilder.group(InstanceHttpApi, "session",
sessionID: ctx.params.sessionID,
limit: ctx.query.limit,
before: ctx.query.before,
after: ctx.query.after,
order: ctx.query.order,
}),
)
if (!page.cursor) return page.items
const total = yield* MessageV2.total(ctx.params.sessionID)
if (!page.cursor) {
return HttpServerResponse.jsonUnsafe(page.items, {
headers: {
"Access-Control-Expose-Headers": "X-Total-Count",
"X-Total-Count": total.toString(),
},
})
}

const request = yield* HttpServerRequest.HttpServerRequest
// toURL() honors the Host + x-forwarded-proto headers, so the Link
// header echoes the real origin instead of a hard-coded localhost.
const url = Option.getOrElse(HttpServerRequest.toURL(request), () => new URL(request.url, "http://localhost"))
const direction = ctx.query.after !== undefined || ctx.query.order === "asc" ? "after" : "before"
url.searchParams.set("limit", ctx.query.limit.toString())
url.searchParams.set("before", page.cursor)
url.searchParams.delete("order")
url.searchParams.delete(direction === "after" ? "before" : "after")
url.searchParams.set(direction, page.cursor)
return HttpServerResponse.jsonUnsafe(page.items, {
headers: {
"Access-Control-Expose-Headers": "Link, X-Next-Cursor",
"Access-Control-Expose-Headers": "Link, X-Next-Cursor, X-Total-Count",
Link: `<${url.toString()}>; rel="next"`,
"X-Next-Cursor": page.cursor,
"X-Total-Count": total.toString(),
},
})
})
Expand Down
62 changes: 51 additions & 11 deletions packages/opencode/src/session/message-v2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,9 @@ import { desc } from "drizzle-orm"
import { eq } from "drizzle-orm"
import { inArray } from "drizzle-orm"
import { lt } from "drizzle-orm"
import { gt } from "drizzle-orm"
import { asc } from "drizzle-orm"
import { count } from "drizzle-orm"
import { or } from "drizzle-orm"
import { MessageTable, PartTable, SessionTable } from "@opencode-ai/core/session/sql"
import { ProviderError } from "@/provider/error"
Expand Down Expand Up @@ -67,6 +70,7 @@ const Cursor = Schema.Struct({
type Cursor = typeof Cursor.Type

const decodeCursor = Schema.decodeUnknownSync(Cursor)
const decodeMessageID = Schema.decodeUnknownSync(MessageID)

export const cursor = {
encode(input: Cursor) {
Expand Down Expand Up @@ -95,6 +99,9 @@ const part = (row: typeof PartTable.$inferSelect) =>
const older = (row: Cursor) =>
or(lt(MessageTable.time_created, row.time), and(eq(MessageTable.time_created, row.time), lt(MessageTable.id, row.id)))

const newer = (row: Cursor) =>
or(gt(MessageTable.time_created, row.time), and(eq(MessageTable.time_created, row.time), gt(MessageTable.id, row.id)))

function hydrate(db: Database.Interface["db"], rows: (typeof MessageTable.$inferSelect)[]) {
const ids = rows.map((row) => row.id)
const partByMessage = new Map<string, Part[]>()
Expand Down Expand Up @@ -435,21 +442,32 @@ export function toModelMessages(
return Effect.runPromise(toModelMessagesEffect(input, model, options))
}

// Pages through a session's messages. Without an anchor the page is the newest
// `limit` messages, or the oldest with `order: "asc"`. `before` and `after` take
// a cursor from a previous page or a message ID, and return the messages strictly
// older or newer than it. Items are always in chronological order; `cursor`
// continues in the same direction.
export const page = Effect.fn("MessageV2.page")(function* (input: {
sessionID: SessionID
limit: number
before?: string
after?: string
order?: "asc" | "desc"
}) {
const { db } = yield* Database.Service
const before = input.before ? cursor.decode(input.before) : undefined
const where = before
? and(eq(MessageTable.session_id, input.sessionID), older(before))
: eq(MessageTable.session_id, input.sessionID)
const ascending = input.after !== undefined || input.order === "asc"
const anchorValue = input.after ?? input.before
const anchor = anchorValue ? yield* resolveAnchor(db, input.sessionID, anchorValue) : undefined
const session = eq(MessageTable.session_id, input.sessionID)
const rows = yield* db
.select()
.from(MessageTable)
.where(where)
.orderBy(desc(MessageTable.time_created), desc(MessageTable.id))
.where(anchor ? and(session, ascending ? newer(anchor) : older(anchor)) : session)
.orderBy(
...(ascending
? [asc(MessageTable.time_created), asc(MessageTable.id)]
: [desc(MessageTable.time_created), desc(MessageTable.id)]),
)
.limit(input.limit + 1)
.all()
.pipe(Effect.orDie)
Expand All @@ -461,16 +479,12 @@ export const page = Effect.fn("MessageV2.page")(function* (input: {
.get()
.pipe(Effect.orDie)
if (!row) return yield* new NotFoundError({ message: `Session not found: ${input.sessionID}` })
return {
items: [] as WithParts[],
more: false,
}
}

const more = rows.length > input.limit
const slice = more ? rows.slice(0, input.limit) : rows
const items = yield* hydrate(db, slice)
items.reverse()
if (!ascending) items.reverse()
const tail = slice.at(-1)
return {
items,
Expand All @@ -479,6 +493,32 @@ export const page = Effect.fn("MessageV2.page")(function* (input: {
}
})

export const total = Effect.fn("MessageV2.total")(function* (sessionID: SessionID) {
const { db } = yield* Database.Service
const row = yield* db
.select({ value: count() })
.from(MessageTable)
.where(eq(MessageTable.session_id, sessionID))
.get()
.pipe(Effect.orDie)
return row?.value ?? 0
})

function resolveAnchor(db: Database.Interface["db"], sessionID: SessionID, value: string) {
return Effect.gen(function* () {
if (!value.startsWith("msg")) return cursor.decode(value)
const id = decodeMessageID(value)
const row = yield* db
.select({ time: MessageTable.time_created })
.from(MessageTable)
.where(and(eq(MessageTable.session_id, sessionID), eq(MessageTable.id, id)))
.get()
.pipe(Effect.orDie)
if (!row) return yield* new NotFoundError({ message: `Message not found: ${value}` })
return { id, time: row.time }
})
}

export function stream(sessionID: SessionID) {
const size = 50
return Effect.gen(function* () {
Expand Down
35 changes: 35 additions & 0 deletions packages/opencode/test/server/session-messages.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,41 @@ describe("session messages endpoint", () => {
{ git: true },
)

it.instance(
"pages forward from the oldest message and reports the total",
withoutWatcher(
Effect.gen(function* () {
const session = yield* sessionScoped
const ids = yield* fill(session.id, 5)

const a = yield* request(`/session/${session.id}/message?limit=2&order=asc`)
expect(a.status).toBe(200)
expect((yield* json<SessionV1.WithParts[]>(a)).map((item) => item.info.id)).toEqual(ids.slice(0, 2))
expect(a.headers["x-total-count"]).toBe("5")
expect(a.headers["link"]).toContain("after=")

const b = yield* request(
`/session/${session.id}/message?limit=2&after=${encodeURIComponent(a.headers["x-next-cursor"]!)}`,
)
expect((yield* json<SessionV1.WithParts[]>(b)).map((item) => item.info.id)).toEqual(ids.slice(2, 4))

const last = yield* request(`/session/${session.id}/message?limit=2&after=${ids[3]}`)
expect((yield* json<SessionV1.WithParts[]>(last)).map((item) => item.info.id)).toEqual(ids.slice(4))
expect(last.headers["x-next-cursor"]).toBeUndefined()
expect(last.headers["x-total-count"]).toBe("5")

const older = yield* request(`/session/${session.id}/message?limit=2&before=${ids[3]}`)
expect((yield* json<SessionV1.WithParts[]>(older)).map((item) => item.info.id)).toEqual(ids.slice(1, 3))

const conflict = yield* request(`/session/${session.id}/message?limit=2&before=${ids[3]}&after=${ids[1]}`)
expect(conflict.status).toBe(400)
const mismatch = yield* request(`/session/${session.id}/message?limit=2&order=asc&before=${ids[3]}`)
expect(mismatch.status).toBe(400)
}),
),
{ git: true },
)

it.instance(
"keeps full-history responses when limit is omitted",
withoutWatcher(
Expand Down
48 changes: 48 additions & 0 deletions packages/opencode/test/session/messages-pagination.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,54 @@ describe("MessageV2.page", () => {
),
)

it.instance("pages forward from the oldest message", () =>
withSession(({ sessionID }) =>
Effect.gen(function* () {
const ids = yield* fill(sessionID, 5)

const a = yield* MessageV2.page({ sessionID, limit: 2, order: "asc" })
expect(a.items.map((item) => item.info.id)).toEqual(ids.slice(0, 2))
expect(a.more).toBe(true)

const b = yield* MessageV2.page({ sessionID, limit: 2, after: a.cursor! })
expect(b.items.map((item) => item.info.id)).toEqual(ids.slice(2, 4))

const c = yield* MessageV2.page({ sessionID, limit: 2, after: b.cursor! })
expect(c.items.map((item) => item.info.id)).toEqual(ids.slice(4))
expect(c.more).toBe(false)
expect(c.cursor).toBeUndefined()
}),
),
)

it.instance("anchors before and after a message ID", () =>
withSession(({ sessionID }) =>
Effect.gen(function* () {
const ids = yield* fill(sessionID, 6)

const older = yield* MessageV2.page({ sessionID, limit: 2, before: ids[3] })
expect(older.items.map((item) => item.info.id)).toEqual(ids.slice(1, 3))
expect(older.more).toBe(true)

const newer = yield* MessageV2.page({ sessionID, limit: 2, after: ids[3] })
expect(newer.items.map((item) => item.info.id)).toEqual(ids.slice(4, 6))
expect(newer.more).toBe(false)

expect(yield* MessageV2.total(sessionID)).toBe(6)
}),
),
)

it.instance("fails with NotFoundError for an anchor message outside the session", () =>
withSession(({ sessionID }) =>
Effect.gen(function* () {
yield* fill(sessionID, 1)
const error = yield* Effect.flip(MessageV2.page({ sessionID, limit: 2, before: "msg_missing" }))
expect(error).toBeInstanceOf(NotFoundError)
}),
),
)

it.instance("returns items in chronological order within a page", () =>
withSession(({ sessionID }) =>
Effect.gen(function* () {
Expand Down
4 changes: 4 additions & 0 deletions packages/sdk/js/src/v2/gen/sdk.gen.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3710,6 +3710,8 @@ export class Session2 extends HeyApiClient {
workspace?: string
limit?: number
before?: string
after?: string
order?: "asc" | "desc"
},
options?: Options<never, ThrowOnError>,
) {
Expand All @@ -3723,6 +3725,8 @@ export class Session2 extends HeyApiClient {
{ in: "query", key: "workspace" },
{ in: "query", key: "limit" },
{ in: "query", key: "before" },
{ in: "query", key: "after" },
{ in: "query", key: "order" },
],
},
],
Expand Down
2 changes: 2 additions & 0 deletions packages/sdk/js/src/v2/gen/types.gen.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9765,6 +9765,8 @@ export type SessionMessagesData = {
workspace?: string
limit?: number
before?: string
after?: string
order?: "asc" | "desc"
}
url: "/session/{sessionID}/message"
}
Expand Down
12 changes: 12 additions & 0 deletions packages/tui/src/config/index.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,17 @@ export const Cursor = Schema.Struct({
}),
}).annotate({ description: "Terminal cursor settings" })

export const TranscriptMaxMessagesDefault = 100
export const Transcript = Schema.Struct({
max_messages: Schema.optional(Schema.Int.check(Schema.isGreaterThan(0))).annotate({
description:
"Most recent messages kept loaded per session; older ones are hidden behind a divider that loads them back (default: 100)",
}),
keep_first_prompt: Schema.optional(Schema.Boolean).annotate({
description: "Keep the session's first prompt at the top when older messages are hidden (default: true)",
}),
}).annotate({ description: "Session transcript loading" })

export const AttentionSounds = Schema.Record(AttentionSoundName, Schema.optionalKey(Schema.String))
export type AttentionSoundPaths = Schema.Schema.Type<typeof AttentionSounds>
export const Attention = Schema.Struct({
Expand Down Expand Up @@ -71,6 +82,7 @@ export const Info = Schema.Struct({
scroll_acceleration: Schema.optional(ScrollAcceleration),
diff_style: Schema.optional(DiffStyle),
cursor: Schema.optional(Cursor),
transcript: Schema.optional(Transcript),
mouse: Schema.optional(Schema.Boolean).annotate({ description: "Enable or disable mouse capture (default: true)" }),
})
export type Info = Schema.Schema.Type<typeof Info>
Expand Down
6 changes: 6 additions & 0 deletions packages/tui/src/config/keybind.ts
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,9 @@ export const Definitions = {
messages_next: keybind("none", "Navigate to next message"),
messages_previous: keybind("none", "Navigate to previous message"),
messages_last_user: keybind("none", "Navigate to last user message"),
messages_hidden_above: keybind("none", "Load hidden messages above the divider"),
messages_hidden_below: keybind("none", "Load hidden messages below the divider"),
messages_hidden_all: keybind("none", "Load all hidden messages"),
messages_copy: keybind("<leader>y", "Copy message"),
messages_undo: keybind("<leader>u", "Undo message"),
messages_redo: keybind("<leader>r", "Redo message"),
Expand Down Expand Up @@ -348,6 +351,9 @@ export const CommandMap = {
messages_next: "session.message.next",
messages_previous: "session.message.previous",
messages_last_user: "session.messages_last_user",
messages_hidden_above: "session.hidden.above",
messages_hidden_below: "session.hidden.below",
messages_hidden_all: "session.hidden.all",
messages_copy: "messages.copy",
messages_undo: "session.undo",
messages_redo: "session.redo",
Expand Down
Loading
Loading