|
|
@@ -15,11 +15,13 @@ const session = (id: string, parentID?: string): Session => ({
|
|
|
})
|
|
|
|
|
|
type UserMessage = Extract<Message, { role: "user" }>
|
|
|
+type AssistantMessage = Extract<Message, { role: "assistant" }>
|
|
|
type TextPart = Extract<Part, { type: "text" }>
|
|
|
type MessageResponse = {
|
|
|
data: { info: Message; parts: Part[] }[]
|
|
|
response: { headers: Headers }
|
|
|
}
|
|
|
+type SingleMessageResponse = { data: MessageResponse["data"][number] }
|
|
|
|
|
|
const userMessage = (id: string, input: Partial<UserMessage> = {}): UserMessage => ({
|
|
|
id,
|
|
|
@@ -31,6 +33,22 @@ const userMessage = (id: string, input: Partial<UserMessage> = {}): UserMessage
|
|
|
...input,
|
|
|
})
|
|
|
|
|
|
+const assistantMessage = (id: string, parentID: string, input: Partial<AssistantMessage> = {}): AssistantMessage => ({
|
|
|
+ id,
|
|
|
+ sessionID: "child",
|
|
|
+ role: "assistant",
|
|
|
+ time: { created: Number(id.at(-1)), completed: Number(id.at(-1)) },
|
|
|
+ parentID,
|
|
|
+ modelID: "model",
|
|
|
+ providerID: "provider",
|
|
|
+ mode: "build",
|
|
|
+ agent: "build",
|
|
|
+ path: { cwd: "/repo", root: "/repo" },
|
|
|
+ cost: 0,
|
|
|
+ tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
|
|
|
+ ...input,
|
|
|
+})
|
|
|
+
|
|
|
const textPart = (messageID: string, input: Partial<TextPart> = {}): TextPart => ({
|
|
|
id: "part",
|
|
|
sessionID: "child",
|
|
|
@@ -45,6 +63,8 @@ const response = (data: MessageResponse["data"] = [], cursor?: string): MessageR
|
|
|
response: { headers: new Headers(cursor ? { "x-next-cursor": cursor } : undefined) },
|
|
|
})
|
|
|
|
|
|
+const singleResponse = (info: Message, parts: Part[] = []): SingleMessageResponse => ({ data: { info, parts } })
|
|
|
+
|
|
|
const deferredResponse = () => Promise.withResolvers<MessageResponse>()
|
|
|
|
|
|
function messageClient(...responses: Array<MessageResponse | Promise<MessageResponse>>) {
|
|
|
@@ -71,6 +91,40 @@ function messageClient(...responses: Array<MessageResponse | Promise<MessageResp
|
|
|
})
|
|
|
}
|
|
|
|
|
|
+function rootMessageClient(
|
|
|
+ pages: Array<MessageResponse | Promise<MessageResponse>>,
|
|
|
+ roots: Array<SingleMessageResponse | Promise<SingleMessageResponse>>,
|
|
|
+) {
|
|
|
+ let pageIndex = 0
|
|
|
+ let rootIndex = 0
|
|
|
+ const requests: unknown[] = []
|
|
|
+ const rootRequests: unknown[] = []
|
|
|
+ const rootWaiting = new Map<number, () => void>()
|
|
|
+ const client = {
|
|
|
+ session: {
|
|
|
+ get: async () => ({ data: session("child", "root") }),
|
|
|
+ messages: (input: unknown) => {
|
|
|
+ requests.push(input)
|
|
|
+ return pages[pageIndex++]
|
|
|
+ },
|
|
|
+ message: (input: unknown) => {
|
|
|
+ rootRequests.push(input)
|
|
|
+ rootWaiting.get(rootRequests.length)?.()
|
|
|
+ rootWaiting.delete(rootRequests.length)
|
|
|
+ return roots[rootIndex++]
|
|
|
+ },
|
|
|
+ },
|
|
|
+ } as unknown as OpencodeClient
|
|
|
+ return Object.assign(client, {
|
|
|
+ requests,
|
|
|
+ rootRequests,
|
|
|
+ rootRequested(count: number) {
|
|
|
+ if (rootRequests.length >= count) return Promise.resolve()
|
|
|
+ return new Promise<void>((resolve) => rootWaiting.set(count, resolve))
|
|
|
+ },
|
|
|
+ })
|
|
|
+}
|
|
|
+
|
|
|
const retryImmediately: typeof retry = async (task, options = {}) => {
|
|
|
const attempts = options.attempts ?? 3
|
|
|
for (let attempt = 0; ; attempt++) {
|
|
|
@@ -124,6 +178,240 @@ describe("server session", () => {
|
|
|
expect(ctx.store.data.message.root).toEqual([])
|
|
|
})
|
|
|
|
|
|
+ test("backfills an assistant-only initial page through its user root", async () => {
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const assistants = [assistantMessage("message-2", user.id), assistantMessage("message-3", user.id)]
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [
|
|
|
+ response(
|
|
|
+ assistants.map((info) => ({ info, parts: [] })),
|
|
|
+ "older",
|
|
|
+ ),
|
|
|
+ ],
|
|
|
+ [singleResponse(user)],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client)
|
|
|
+
|
|
|
+ await store.sync("child")
|
|
|
+
|
|
|
+ expect(client.requests).toEqual([
|
|
|
+ { sessionID: "child", limit: 2, before: undefined },
|
|
|
+ ])
|
|
|
+ expect(client.rootRequests).toEqual([{ sessionID: "child", messageID: user.id }])
|
|
|
+ expect(store.data.message.child).toEqual([user, ...assistants])
|
|
|
+ expect(store.history.more("child")).toBe(true)
|
|
|
+ })
|
|
|
+
|
|
|
+ test("does not let an optimistic user suppress initial root backfill", async () => {
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const part = textPart(user.id)
|
|
|
+ const assistants = [assistantMessage("message-2", user.id), assistantMessage("message-3", user.id)]
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [
|
|
|
+ response(
|
|
|
+ assistants.map((info) => ({ info, parts: [] })),
|
|
|
+ "older",
|
|
|
+ ),
|
|
|
+ ],
|
|
|
+ [singleResponse(user)],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client)
|
|
|
+ store.optimistic.add({ sessionID: "child", message: user, parts: [part] })
|
|
|
+
|
|
|
+ await store.sync("child")
|
|
|
+ store.optimistic.remove({ sessionID: "child", messageID: user.id })
|
|
|
+
|
|
|
+ expect(client.requests).toHaveLength(1)
|
|
|
+ expect(client.rootRequests).toHaveLength(1)
|
|
|
+ expect(store.data.message.child).toEqual([user, ...assistants])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("backfills the parent of fetched assistants when another user is cached", async () => {
|
|
|
+ const unrelated = userMessage("message-0", { time: { created: 0 } })
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const assistants = [assistantMessage("message-2", user.id), assistantMessage("message-3", user.id)]
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [
|
|
|
+ response([{ info: unrelated, parts: [] }]),
|
|
|
+ response(
|
|
|
+ assistants.map((info) => ({ info, parts: [] })),
|
|
|
+ "older",
|
|
|
+ ),
|
|
|
+ ],
|
|
|
+ [singleResponse(user)],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client)
|
|
|
+ await store.sync("child")
|
|
|
+
|
|
|
+ await store.sync("child", { force: true })
|
|
|
+
|
|
|
+ expect(client.requests).toHaveLength(2)
|
|
|
+ expect(client.rootRequests).toHaveLength(1)
|
|
|
+ expect(store.data.message.child).toEqual([unrelated, user, ...assistants])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("preserves cached history between an injected parent and the page boundary", async () => {
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const cached = userMessage("message-3", { time: { created: 3 } })
|
|
|
+ const assistant = assistantMessage("message-4", user.id)
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [response([{ info: cached, parts: [] }]), response([{ info: assistant, parts: [] }], "older")],
|
|
|
+ [singleResponse(user)],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client)
|
|
|
+ await store.sync("child")
|
|
|
+
|
|
|
+ await store.sync("child", { force: true })
|
|
|
+
|
|
|
+ expect(store.data.message.child).toEqual([user, cached, assistant])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("refreshes a cached parent omitted by an assistant-only replacement page", async () => {
|
|
|
+ const stale = userMessage("message-1", { summary: { title: "stale", diffs: [] } })
|
|
|
+ const fresh = { ...stale, summary: { title: "fresh", diffs: [] } }
|
|
|
+ const stalePart = textPart(stale.id, { text: "stale" })
|
|
|
+ const freshPart = { ...stalePart, text: "fresh" }
|
|
|
+ const assistant = assistantMessage("message-2", stale.id)
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [response([{ info: stale, parts: [stalePart] }]), response([{ info: assistant, parts: [] }], "older")],
|
|
|
+ [singleResponse(fresh, [freshPart])],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client)
|
|
|
+ await store.sync("child")
|
|
|
+
|
|
|
+ await store.sync("child", { force: true })
|
|
|
+
|
|
|
+ expect(client.rootRequests).toEqual([{ sessionID: "child", messageID: stale.id }])
|
|
|
+ expect(store.data.message.child).toEqual([fresh, assistant])
|
|
|
+ expect(store.data.part[stale.id]).toEqual([freshPart])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("refreshes a confirmed optimistic parent while preserving pending parts", async () => {
|
|
|
+ const stale = userMessage("message-1", { summary: { title: "stale", diffs: [] } })
|
|
|
+ const fresh = { ...stale, summary: { title: "fresh", diffs: [] } }
|
|
|
+ const confirmed = textPart(stale.id, { id: "confirmed", text: "stale" })
|
|
|
+ const refreshed = { ...confirmed, text: "fresh" }
|
|
|
+ const pending = textPart(stale.id, { id: "pending", text: "pending" })
|
|
|
+ const assistant = assistantMessage("message-2", stale.id)
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [response([{ info: stale, parts: [confirmed] }]), response([{ info: assistant, parts: [] }], "older")],
|
|
|
+ [singleResponse(fresh, [refreshed])],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client)
|
|
|
+ store.optimistic.add({ sessionID: "child", message: stale, parts: [confirmed, pending] })
|
|
|
+ await store.sync("child")
|
|
|
+
|
|
|
+ await store.sync("child", { force: true })
|
|
|
+
|
|
|
+ expect(client.rootRequests).toEqual([{ sessionID: "child", messageID: stale.id }])
|
|
|
+ expect(store.data.message.child).toEqual([fresh, assistant])
|
|
|
+ expect(store.data.part[stale.id]).toEqual([refreshed, pending])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("uses a parent received by SSE during the replacement load", async () => {
|
|
|
+ const pending = deferredResponse()
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const assistant = assistantMessage("message-2", user.id)
|
|
|
+ const client = rootMessageClient([pending.promise], [])
|
|
|
+ const store = createServerSession(client)
|
|
|
+ const loading = store.sync("child")
|
|
|
+
|
|
|
+ store.apply({ type: "message.updated", properties: { info: user } })
|
|
|
+ pending.resolve(response([{ info: assistant, parts: [] }], "older"))
|
|
|
+ await loading
|
|
|
+
|
|
|
+ expect(client.rootRequests).toEqual([])
|
|
|
+ expect(store.data.message.child).toEqual([user, assistant])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("uses a successful retry over events received by a failed backfill attempt", async () => {
|
|
|
+ const failed = deferredResponse()
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const live = { ...user, agent: "stale" }
|
|
|
+ const assistants = [assistantMessage("message-2", user.id), assistantMessage("message-3", user.id)]
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [
|
|
|
+ response(
|
|
|
+ assistants.map((info) => ({ info, parts: [] })),
|
|
|
+ "older",
|
|
|
+ ),
|
|
|
+ ],
|
|
|
+ [failed.promise.then((result) => ({ data: result.data[0]! })), singleResponse(user)],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client, { retry: retryImmediately })
|
|
|
+ const loading = store.sync("child")
|
|
|
+ await client.rootRequested(1)
|
|
|
+
|
|
|
+ store.apply({ type: "message.updated", properties: { info: live } })
|
|
|
+ failed.reject(new Error("retry"))
|
|
|
+ await loading
|
|
|
+
|
|
|
+ expect(client.requests).toHaveLength(1)
|
|
|
+ expect(client.rootRequests).toHaveLength(2)
|
|
|
+ expect(store.data.message.child).toEqual([user, ...assistants])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("preserves newer-page events across a failed parent retry", async () => {
|
|
|
+ const failed = deferredResponse()
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const assistant = assistantMessage("message-2", user.id)
|
|
|
+ const live = { ...assistant, cost: 1 }
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [response([{ info: assistant, parts: [] }], "older")],
|
|
|
+ [failed.promise.then((result) => ({ data: result.data[0]! })), singleResponse(user)],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client, { retry: retryImmediately })
|
|
|
+ const loading = store.sync("child")
|
|
|
+ await client.rootRequested(1)
|
|
|
+
|
|
|
+ store.apply({ type: "message.updated", properties: { info: live } })
|
|
|
+ failed.reject(new Error("retry"))
|
|
|
+ await loading
|
|
|
+
|
|
|
+ expect(store.data.message.child).toEqual([user, live])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("preserves unrelated message events across a failed parent retry", async () => {
|
|
|
+ const failed = deferredResponse()
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const assistant = assistantMessage("message-2", user.id)
|
|
|
+ const live = userMessage("message-4", { time: { created: 4 } })
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [response([{ info: assistant, parts: [] }], "older")],
|
|
|
+ [failed.promise.then((result) => ({ data: result.data[0]! })), singleResponse(user)],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client, { retry: retryImmediately })
|
|
|
+ const loading = store.sync("child")
|
|
|
+ await client.rootRequested(1)
|
|
|
+
|
|
|
+ store.apply({ type: "message.updated", properties: { info: live } })
|
|
|
+ failed.reject(new Error("retry"))
|
|
|
+ await loading
|
|
|
+
|
|
|
+ expect(store.data.message.child).toEqual([user, assistant, live])
|
|
|
+ })
|
|
|
+
|
|
|
+ test("preserves newer-page part events across a failed parent retry", async () => {
|
|
|
+ const failed = deferredResponse()
|
|
|
+ const user = userMessage("message-1")
|
|
|
+ const assistant = assistantMessage("message-2", user.id)
|
|
|
+ const stale = textPart(assistant.id, { text: "stale" })
|
|
|
+ const live = { ...stale, text: "live" }
|
|
|
+ const client = rootMessageClient(
|
|
|
+ [response([{ info: assistant, parts: [stale] }], "older")],
|
|
|
+ [failed.promise.then((result) => ({ data: result.data[0]! })), singleResponse(user)],
|
|
|
+ )
|
|
|
+ const store = createServerSession(client, { retry: retryImmediately })
|
|
|
+ const loading = store.sync("child")
|
|
|
+ await client.rootRequested(1)
|
|
|
+
|
|
|
+ store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: live, time: 2 } })
|
|
|
+ failed.reject(new Error("retry"))
|
|
|
+ await loading
|
|
|
+
|
|
|
+ expect(store.data.part[assistant.id]).toEqual([live])
|
|
|
+ })
|
|
|
+
|
|
|
test("merges live events into the initial page", async () => {
|
|
|
const pending = deferredResponse()
|
|
|
const user = userMessage("message-1")
|
|
|
@@ -905,6 +1193,26 @@ describe("server session", () => {
|
|
|
expect(store.data.message.child).toEqual([latest])
|
|
|
})
|
|
|
|
|
|
+ test("does not scan cached messages for user roots during history prepend", async () => {
|
|
|
+ const guard = { active: false }
|
|
|
+ const latest = new Proxy(userMessage("message-2", { time: { created: 2 } }), {
|
|
|
+ get(target, property, receiver) {
|
|
|
+ if (guard.active && property === "role") throw new Error("cached role accessed")
|
|
|
+ return Reflect.get(target, property, receiver)
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const older = userMessage("message-1")
|
|
|
+ const store = createServerSession(
|
|
|
+ messageClient(response([{ info: latest, parts: [] }], "older"), response([{ info: older, parts: [] }])),
|
|
|
+ )
|
|
|
+ await store.sync("child")
|
|
|
+ guard.active = true
|
|
|
+
|
|
|
+ await store.history.loadMore("child")
|
|
|
+
|
|
|
+ expect(store.data.message.child).toEqual([older, latest])
|
|
|
+ })
|
|
|
+
|
|
|
test("preserves loaded history during an incomplete refresh", async () => {
|
|
|
const older = userMessage("message-1")
|
|
|
const latest = userMessage("message-2", { time: { created: 2 } })
|