File size: 5,154 Bytes
3201ca6
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
import type { OpenHandsEvent } from "#/types/agent-server/core";

export const TRANSCRIPT_HISTORY_PAGE_SIZE = 100;

interface TranscriptEventSearchOptions {
  limit: number;
  sortOrder: "TIMESTAMP_DESC";
  pageId?: string;
  timestampLt?: string;
  strictPagination: true;
}

interface TranscriptEventPage {
  items: OpenHandsEvent[];
  next_page_id?: string | null;
}

type SearchTranscriptEvents = (
  options: TranscriptEventSearchOptions,
) => Promise<TranscriptEventPage>;

const compareEventTimestamps = (
  first: OpenHandsEvent,
  second: OpenHandsEvent,
): number => (first.timestamp ?? "").localeCompare(second.timestamp ?? "");

/**
 * Loads the persisted history from the newest page back to the beginning,
 * then merges any live store events that have not persisted yet. The timestamp
 * anchor matches the chat's existing history pagination, while id-based
 * de-duplication keeps the result stable if pages overlap.
 */
export const loadCompleteTranscriptEvents = async (
  loadedEvents: OpenHandsEvent[],
  searchEvents: SearchTranscriptEvents,
  expectedEventCount?: number,
): Promise<OpenHandsEvent[]> => {
  const persistedDescending: OpenHandsEvent[] = [];
  const fetchedEventIds = new Set<string>();
  // Id-less events can't be de-duplicated by id, so keep a separate ordered
  // list and a count so they survive the merge and count toward completeness.
  const idLessEvents: OpenHandsEvent[] = [];
  let idLessPersisted = 0;
  const seenPageIds = new Set<string>();
  let oldestTimestamp: string | undefined;
  let pageId: string | undefined;
  let usedCursor = false;
  let usingTimestampFallback = false;

  while (true) {
    const page = await searchEvents({
      limit: TRANSCRIPT_HISTORY_PAGE_SIZE,
      sortOrder: "TIMESTAMP_DESC",
      strictPagination: true,
      ...(pageId ? { pageId } : {}),
      ...(oldestTimestamp ? { timestampLt: oldestTimestamp } : {}),
    });

    if (!Array.isArray(page.items)) {
      throw new Error(
        "Invalid transcript history response: expected page.items to be an array.",
      );
    }

    persistedDescending.push(...page.items);
    let pageOldestTimestamp: string | undefined;
    let addedEvent = false;
    page.items.forEach((event) => {
      if (!fetchedEventIds.has(event.id ?? "")) {
        if (event.id) {
          fetchedEventIds.add(event.id);
        } else {
          idLessPersisted += 1;
        }
        addedEvent = true;
      }
      if (
        !pageOldestTimestamp ||
        (event.timestamp ?? "") < pageOldestTimestamp
      ) {
        pageOldestTimestamp = event.timestamp ?? "";
      }
    });

    if (page.next_page_id) {
      if (seenPageIds.has(page.next_page_id)) {
        throw new Error(
          "Transcript history pagination repeated a page cursor.",
        );
      }
      seenPageIds.add(page.next_page_id);
      pageId = page.next_page_id;
      oldestTimestamp = undefined;
      usedCursor = true;
      continue;
    }

    // Once a server supplies a cursor, a page without a next cursor is an
    // explicit exhaustion signal, even when the final page is exactly full.
    if (usedCursor && !usingTimestampFallback) break;
    if (page.items.length < TRANSCRIPT_HISTORY_PAGE_SIZE) break;

    // Some older servers omit cursors for filtered searches. A timestamp
    // fallback is safe to attempt only when an independent event count can
    // prove completeness; otherwise fail instead of exporting a partial tail.
    if (expectedEventCount === undefined) {
      throw new Error(
        "Transcript history pagination cannot prove that all events were loaded.",
      );
    }
    if (!pageOldestTimestamp) {
      throw new Error("Transcript history pagination did not advance.");
    }
    if (
      oldestTimestamp &&
      (!addedEvent || pageOldestTimestamp >= oldestTimestamp)
    ) {
      throw new Error("Transcript history pagination did not advance.");
    }

    pageId = undefined;
    oldestTimestamp = pageOldestTimestamp;
    usingTimestampFallback = true;
  }

  const eventsById = new Map<string, OpenHandsEvent>();
  persistedDescending
    .slice()
    .reverse()
    .forEach((event) => {
      if (event.id !== undefined) {
        if (!eventsById.has(event.id)) {
          eventsById.set(event.id, event);
        }
        return;
      }
      idLessEvents.push(event);
    });
  loadedEvents.forEach((event) => {
    if (event.id !== undefined) {
      if (!eventsById.has(event.id)) {
        eventsById.set(event.id, event);
      }
      return;
    }
    idLessEvents.push(event);
  });
  // Array.prototype.sort is stable, so equal-timestamp events keep the causal
  // order returned by the server/store rather than being reordered by id.
  const completeEvents = [...idLessEvents, ...eventsById.values()].sort(
    compareEventTimestamps,
  );
  if (
    expectedEventCount !== undefined &&
    fetchedEventIds.size + idLessPersisted < expectedEventCount
  ) {
    throw new Error(
      `Transcript history is incomplete: expected ${expectedEventCount} persisted events, received ${fetchedEventIds.size + idLessPersisted}.`,
    );
  }
  return completeEvents;
};