File size: 4,785 Bytes
5710dd0 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 | /**
* REST client for the history endpoint of the message protocol:
* `GET {baseUrl}/api/v1/sessions/{sessionId}/history`.
*
* This is the ONLY source of persisted (completed) timeline state: the
* initial load fetches the newest page, "load earlier" pages further with a
* `before_turn` cursor, and a reconnect catch-up pages forward from an
* `after_step` cursor. The in-flight step's entities arrive over the WS
* recovery payload instead (idempotent replace-by-id at the seam).
*
* Pages are flat entity-message slices (`{ messages, in_flight? }`,
* time-ordered, same schemas as the WS stream). There is deliberately no
* has-more flag: a page shorter than `page_size` is the end in that
* direction, an empty page is definitive.
*/
import { historyResponseSchema, type HistoryMessage } from '@moonshot-ai/kap-server/protocol';
export const HISTORY_PAGE_SIZE = 500;
export interface HistoryPage {
readonly messages: readonly HistoryMessage[];
/** Current streaming position of a live session; absent for idle/cold ones. */
readonly inFlight?: { turn_id: string; step_id: string };
}
export interface FetchHistoryPageOptions {
readonly baseUrl: string;
readonly token?: string;
readonly sessionId: string;
readonly agentId: string;
/** Turn-id cursor; fetches up to `pageSize` messages strictly older than that turn. */
readonly beforeTurn?: string;
/** Step-id cursor; fetches up to `pageSize` messages strictly newer than that step. */
readonly afterStep?: string;
readonly pageSize?: number;
/** Injectable for tests. */
readonly fetchImpl?: typeof fetch;
}
export async function fetchHistoryPage(opts: FetchHistoryPageOptions): Promise<HistoryPage> {
const params = new URLSearchParams({
agent_id: opts.agentId,
page_size: String(opts.pageSize ?? HISTORY_PAGE_SIZE),
});
if (opts.beforeTurn !== undefined) params.set('before_turn', opts.beforeTurn);
if (opts.afterStep !== undefined) params.set('after_step', opts.afterStep);
const headers: Record<string, string> = {};
if (opts.token !== undefined && opts.token !== '') {
headers['authorization'] = `Bearer ${opts.token}`;
}
const doFetch = opts.fetchImpl ?? fetch;
const res = await doFetch(
`${opts.baseUrl}/api/v1/sessions/${encodeURIComponent(opts.sessionId)}/history?${params.toString()}`,
{ headers },
);
const envelope = (await res.json()) as { code: number; msg: string; data: unknown };
if (envelope.code !== 0) {
throw new Error(`history page failed (${envelope.code}): ${envelope.msg}`);
}
const parsed = historyResponseSchema.safeParse(envelope.data);
if (!parsed.success) {
throw new Error('history page: unexpected response shape');
}
return { messages: parsed.data.messages, inFlight: parsed.data.in_flight };
}
/**
* Read the agent's WHOLE history (newest page + `before_turn` paging to the
* beginning) in timeline order. On-demand debug reads only (plan lookup) —
* the chat channel pages lazily instead.
*/
export async function fetchFullHistory(opts: {
readonly baseUrl: string;
readonly token?: string;
readonly sessionId: string;
readonly agentId: string;
readonly pageSize?: number;
readonly fetchImpl?: typeof fetch;
}): Promise<readonly HistoryMessage[]> {
const pageSize = opts.pageSize ?? HISTORY_PAGE_SIZE;
const messages: HistoryMessage[] = [];
const seen = new Set<string>();
let beforeTurn: string | undefined;
for (;;) {
const page = await fetchHistoryPage({ ...opts, beforeTurn, pageSize });
if (page.messages.length === 0) break;
const fresh: HistoryMessage[] = [];
for (const message of page.messages) {
const key = historyEntityKey(message);
if (seen.has(key)) continue;
seen.add(key);
fresh.push(message);
}
messages.unshift(...fresh);
if (page.messages.length < pageSize) break;
const oldest = page.messages
.map((message) => ('turn_id' in message ? message.turn_id : undefined))
.find((turnId) => turnId !== undefined);
if (oldest === undefined || oldest === beforeTurn) break;
beforeTurn = oldest;
}
return messages;
}
function historyEntityKey(message: HistoryMessage): string {
switch (message.type) {
case 'turn':
return `turn:${message.turn_id}`;
case 'step':
return `step:${message.step_id}`;
case 'user':
case 'assistant':
case 'thinking':
return `${message.type}:${message.message_id}`;
case 'tool_call':
return `tool_call:${message.tool_call_id}`;
case 'system':
return `system:${message.system_id}`;
case 'interaction':
return `interaction:${message.interaction_id}`;
case 'task':
return `task:${message.task_id}`;
case 'todo':
return `todo:${message.todo_id}`;
}
}
|