Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934import type { Attachment } from "../../shared/attachments";import type { CapabilityApprovalRequest } from "../../shared/capability-approval";import { AgentClient } from "agents/client";import { WebSocketChatTransport } from "agents/chat/transport";import { MessageType, broadcastTransition, type BroadcastStreamState, type OutgoingMessage,} from "agents/chat";import { AbstractChat, type ChatInit, type ChatState, type ChatStatus, type UIMessage,} from "ai";import type { ToolActivity, ToolActivityPage, ToolActivityState,} from "../../shared/tool-activity";import type { ModelConfiguration } from "../../shared/model-providers";import type { ConversationModelSettings } from "../../worker/model-settings";
type Connection = "loading" | "connected" | "offline" | "signed-out" | "missing";export interface ConversationView { messages: UIMessage[]; status: ChatStatus; error?: Error; connection: Connection; notice?: string; unsentText?: string; unsentAttachments?: Attachment[]; serverStreaming: boolean; recovering: boolean; stopping: boolean; reconciling: boolean; activities: ReadonlyMap<string, ToolActivity>; activityCursor: number | null | undefined; activitiesLoading: boolean; modelSettings?: ConversationModelSettings; modelSelection?: ModelConfiguration | null;}class NativeChat extends AbstractChat<UIMessage> { constructor(state: ChatState<UIMessage>, options: ChatInit<UIMessage>) { super({ ...options, state }); }}class HistoryError extends Error { constructor(readonly status: number) { super("Could not load conversation"); }}const initialView = (): ConversationView => ({ messages: [], status: "ready", connection: "loading", serverStreaming: false, recovering: false, stopping: false, reconciling: true, activities: new Map(), activityCursor: undefined, activitiesLoading: false,});
/** Selected conversation only. Think is the transcript authority; this store is a * disposable immutable browser view, never a client transcript persistence API. */export class ConversationSession { private view = initialView(); private readonly serverView = this.view; private listeners = new Set<() => void>(); private lifetime = new AbortController(); private connectionLifetime?: AbortController; private generation = 0; private revision = 0; private historyRead = 0; private client?: AgentClient; private transport?: WebSocketChatTransport; private chat?: NativeChat; private observed: BroadcastStreamState = { status: "idle" }; private active = new Set<string>(); private acked = new Set<string>(); private protectedAssistant?: { id: string; anchor?: string }; private reconnectTimer?: ReturnType<typeof setTimeout>; private reconcileTimer?: ReturnType<typeof setTimeout>; private resumeOperation?: Promise<void>; private connecting = false; private reconnectQueued = false; private idleConfirmed = false; private probeTimer?: ReturnType<typeof setTimeout>; private operation?: { generation: number; cancelled: boolean }; private streamEpoch = 0; private pendingInput?: { id: string; text: string; attachments: Attachment[]; }; constructor(readonly id: string) {} getSnapshot = () => this.view; getServerSnapshot = () => this.serverView; subscribe = (listener: () => void) => { this.listeners.add(listener); return () => { this.listeners.delete(listener); }; }; private publish(patch: Partial<ConversationView>) { if (this.lifetime.signal.aborted) return; this.view = { ...this.view, ...patch }; for (const listener of this.listeners) listener(); } get busy() { return ( !!this.operation || this.view.status === "streaming" || this.view.status === "submitted" || this.view.serverStreaming || this.view.recovering || this.view.stopping ); } private current(generation: number) { return generation === this.generation && !this.lifetime.signal.aborted; } start = () => { window.addEventListener("online", this.online); window.addEventListener("offline", this.offline); void this.connect(); }; reconnect = () => { clearTimeout(this.reconnectTimer); void this.connect(); }; private get terminal() { return ( this.view.connection === "signed-out" || this.view.connection === "missing" ); } // An established socket may outlive a browser network loss, so the browser's // own signal detaches immediately and the next connection re-reads history // (and any session expiry) instead of trusting the stale connection. private offline = () => { if (this.lifetime.signal.aborted || this.terminal) return; clearTimeout(this.reconnectTimer); this.detach(); this.publish({ connection: "offline", reconciling: true, status: "ready", serverStreaming: false, recovering: false, stopping: false, notice: "Connection lost. Reconnecting to the saved conversation…", }); }; private online = () => { if (this.terminal) return; // A connection aborted by `offline` may still be unwinding. if (this.connecting) this.reconnectQueued = true; else this.reconnect(); }; private setMessages(messages: UIMessage[]) { this.revision++; this.publish({ messages }); } private state( generation: number, epoch = this.streamEpoch, ): ChatState<UIMessage> { const session = this; const valid = () => session.current(generation) && epoch === session.streamEpoch; return { get messages() { return session.view.messages; }, set messages(messages) { if (valid()) session.setMessages(messages); }, get status() { return session.view.status; }, set status(status) { if (valid()) session.publish({ status }); }, get error() { return session.view.error; }, set error(error) { if (valid()) session.publish({ error }); }, snapshot: <T>(value: T): T => structuredClone(value), pushMessage(message) { if (valid()) session.setMessages([ ...session.view.messages, structuredClone(message), ]); }, popMessage() { if (valid()) session.setMessages(session.view.messages.slice(0, -1)); }, replaceMessage(index, message) { if (valid()) session.setMessages( session.view.messages.map((current, i) => i === index ? structuredClone(message) : current, ), ); }, }; } private authority(messages: UIMessage[]) { let next = structuredClone(messages); const protect = this.protectedAssistant; if (protect) { const live = this.view.messages.find( (message) => message.id === protect.id, ); if (live) next = [...next.filter((message) => message.id !== protect.id), live]; } if (this.observed.status === "observing") next = this.observed.accumulator.mergeInto(next); this.setMessages(next); } private async history(generation: number, seed = false) { const revision = this.revision; const read = ++this.historyRead; const response = await fetch( `/agents/personal-agent/personal/sub/conversation/${this.id}/get-messages`, { credentials: "same-origin", cache: "no-store", signal: AbortSignal.any([ this.lifetime.signal, this.connectionLifetime!.signal, AbortSignal.timeout(10_000), ]), }, ); if (!response.ok) throw new HistoryError(response.status); const messages: UIMessage[] = await response.json(); if (!Array.isArray(messages)) throw new Error("Invalid conversation history"); if (!this.current(generation) || read !== this.historyRead) return; if (revision === this.revision) this.authority(messages); else if (!this.busy) this.scheduleReconcile(generation); if (seed) this.publish({ reconciling: true }); } private fail(error: unknown, generation: number) { if (!this.current(generation)) return; if ( error instanceof HistoryError && [401, 403, 404].includes(error.status) ) { this.detach(); this.setMessages([]); this.pendingInput = undefined; this.publish({ unsentText: undefined, unsentAttachments: undefined, connection: error.status === 404 ? "missing" : "signed-out", notice: error.status === 404 ? "This conversation was deleted or could not be found." : "Sign in to this installation to view this conversation.", activities: new Map(), status: "ready", error: undefined, serverStreaming: false, recovering: false, stopping: false, reconciling: false, }); return; } this.publish({ connection: "offline", reconciling: true, notice: "Connection lost. Reconnecting to the saved conversation…", }); clearTimeout(this.reconnectTimer); this.reconnectTimer = setTimeout(() => { void this.connect(); }, 3000); } private detach() { this.generation++; this.idleConfirmed = false; clearTimeout(this.probeTimer); this.operation = undefined; this.connectionLifetime?.abort(); this.transport?.resetResumeState(); void this.chat?.stop(); // Local detach only; cancelOnClientAbort is false. this.client?.close(); this.client = undefined; this.chat = undefined; this.transport = undefined; this.resumeOperation = undefined; this.active = new Set(); this.acked = new Set(); this.observed = { status: "idle" }; this.protectedAssistant = undefined; clearTimeout(this.reconcileTimer); } private async connect() { if (this.lifetime.signal.aborted || this.connecting) return; this.connecting = true; this.detach(); const generation = this.generation; this.connectionLifetime = new AbortController(); this.publish({ reconciling: true, status: "ready", serverStreaming: false, recovering: false, }); try { await this.history(generation, true); if (!this.current(generation)) return; const client = new AgentClient<ToolActivityState>({ host: window.location.host, protocol: window.location.protocol === "https:" ? "wss" : "ws", agent: "PersonalAgent", basePath: `agents/personal-agent/personal/sub/conversation/${this.id}`, startClosed: true, defaultCallTimeout: 10_000, shouldReconnectOnClose: () => false, onStateUpdate: (state) => { if (this.current(generation) && state?.toolActivityVersion === 1) this.upsertActivities(state.toolActivities); }, }); this.client = client; const transport = new WebSocketChatTransport({ agent: client, activeRequestIds: this.active, cancelOnClientAbort: false, }); this.transport = transport; this.chat = this.createChat(generation, transport); // Register before reconnect and before transport-owned message listeners. client.addEventListener("message", (event) => { if (this.current(generation)) this.message(event, generation); }); client.addEventListener("close", () => { if (this.current(generation)) this.fail(new Error("Disconnected"), generation); }); client.reconnect(); const signal = AbortSignal.any([ this.connectionLifetime.signal, AbortSignal.timeout(10_000), ]); await new Promise<void>((resolve, reject) => { const abort = () => reject(signal.reason); signal.addEventListener("abort", abort, { once: true }); client.ready .then(resolve, reject) .finally(() => signal.removeEventListener("abort", abort)); }); if (!this.current(generation)) return; const modelSettings = await client.call<ConversationModelSettings>( "getConversationModelSettings", ); if (!this.current(generation)) return; this.publish({ modelSettings, modelSelection: this.view.modelSelection === undefined ? modelSettings.override : this.view.modelSelection, }); this.publish({ connection: "connected", notice: undefined }); await this.history(generation); if (!this.current(generation)) return; // An unsolicited onConnect offer may already be observing the turn. if (this.observed.status === "idle") this.resume(generation); else this.publish({ reconciling: false }); } catch (error) { this.fail(error, generation); } finally { this.connecting = false; if (this.reconnectQueued) { this.reconnectQueued = false; if (this.view.connection === "offline") this.reconnect(); } } } private createChat(generation: number, transport: WebSocketChatTransport) { const epoch = this.streamEpoch; return new NativeChat(this.state(generation), { id: this.id, transport, onFinish: ({ isError, isDisconnect }) => { if (!this.current(generation) || epoch !== this.streamEpoch) return; this.restoreAssistant(); this.publish({ reconciling: true }); if (isError || isDisconnect) this.publish({ notice: "The response was interrupted. Checking saved progress…", reconciling: true, }); this.scheduleReconcile(generation); }, }); } private resume(generation: number) { if (!this.current(generation) || !this.chat || this.resumeOperation) return; const operation = this.chat.resumeStream(); this.resumeOperation = operation; void operation .catch(() => {}) .finally(() => { if (!this.current(generation) || this.resumeOperation !== operation) return; this.resumeOperation = undefined; if (!this.idleConfirmed && this.observed.status === "idle") { // A missing/uncorrelated probe reply is not evidence that durable work stopped. this.publish({ reconciling: true }); clearTimeout(this.probeTimer); this.probeTimer = setTimeout(() => this.resume(generation), 3000); } else this.scheduleReconcile(generation); }); } private restoreAssistant() { const protectedRow = this.protectedAssistant; this.protectedAssistant = undefined; if (!protectedRow) return; const assistant = this.view.messages.find( (message) => message.id === protectedRow.id, ); if (!assistant) return; const messages = this.view.messages.filter( (message) => message.id !== protectedRow.id, ); const index = protectedRow.anchor ? messages.findIndex((message) => message.id === protectedRow.anchor) + 1 : 0; messages.splice(index, 0, assistant); this.setMessages(messages); } private scheduleReconcile(generation: number) { clearTimeout(this.reconcileTimer); this.reconcileTimer = setTimeout(() => { void this.reconcile(generation); }, 40); } private async reconcile(generation: number) { if ( !this.current(generation) || this.view.connection !== "connected" || !this.idleConfirmed ) return; if ( this.view.status === "streaming" || this.view.status === "submitted" || this.view.serverStreaming || this.view.recovering || this.view.stopping ) return; this.protectedAssistant = undefined; try { await this.history(generation); if (this.current(generation) && this.pendingInput) { const accepted = this.view.messages.some( (message) => message.id === this.pendingInput!.id, ); this.publish({ unsentText: accepted ? undefined : this.pendingInput.text, unsentAttachments: accepted ? undefined : this.pendingInput.attachments, }); this.pendingInput = undefined; } if (this.current(generation)) this.publish({ reconciling: false, stopping: false, notice: this.view.error ? "The response could not be completed. Your saved messages are here; you can retry the last turn." : undefined, }); } catch (error) { this.fail(error, generation); } } private message(event: MessageEvent, generation: number) { if (typeof event.data !== "string") return; let data: OutgoingMessage; try { data = JSON.parse(event.data); } catch { return; } const transport = this.transport!; switch (data.type) { case MessageType.CF_AGENT_CHAT_CLEAR: this.historyRead++; this.streamEpoch++; clearTimeout(this.probeTimer); this.resumeOperation = undefined; this.idleConfirmed = true; this.operation = undefined; this.pendingInput = undefined; this.active.clear(); this.acked.clear(); this.observed = { status: "idle" }; this.protectedAssistant = undefined; transport.resetResumeState(); void this.chat?.stop(); this.chat = this.createChat(generation, transport); this.setMessages([]); this.publish({ unsentText: undefined, unsentAttachments: undefined, status: "ready", activities: new Map(), activityCursor: undefined, activitiesLoading: false, serverStreaming: false, recovering: false, stopping: false, error: undefined, reconciling: false, }); break; case MessageType.CF_AGENT_CHAT_RECOVERING: this.publish({ recovering: Boolean(data.recovering) }); if (!data.recovering) this.scheduleReconcile(generation); break; case MessageType.CF_AGENT_CHAT_MESSAGES: if (Array.isArray(data.messages)) this.authority([...data.messages]); break; case MessageType.CF_AGENT_MESSAGE_UPDATED: { const updated = data.message; const calls = new Set( updated.parts .filter((part) => "toolCallId" in part) .map((part) => (part as { toolCallId: string }).toolCallId), ); const index = this.view.messages.findIndex( (message) => message.id === updated.id || message.parts.some( (part) => "toolCallId" in part && calls.has(part.toolCallId), ), ); if ( index >= 0 && this.view.messages[index].id !== this.protectedAssistant?.id ) this.setMessages( this.view.messages.map((message, i) => i === index ? { ...structuredClone(updated), id: message.id } : message, ), ); break; } case MessageType.CF_AGENT_STREAM_PENDING: this.idleConfirmed = false; this.publish({ reconciling: true }); transport.handleStreamPending(); if (data.id) transport.observeServerTurn(data.id); break; case MessageType.CF_AGENT_STREAM_RESUME_NONE: if ( transport.handleStreamResumeNone(data) && data.reason === "idle" && data.probeId ) { this.idleConfirmed = true; this.observed = { status: "idle" }; this.publish({ serverStreaming: false, recovering: false, stopping: false, }); this.scheduleReconcile(generation); } break; case MessageType.CF_AGENT_STREAM_RESUMING: this.idleConfirmed = false; if ( transport.handleStreamResuming(data) || this.active.has(data.id) || this.acked.has(data.id) ) break; this.observed = broadcastTransition(this.observed, { type: "resume-fallback", streamId: data.id, messageId: crypto.randomUUID(), }).state; transport.observeServerTurn(data.id); this.acked.add(data.id); this.publish({ serverStreaming: true, reconciling: false }); this.client!.send( JSON.stringify({ type: MessageType.CF_AGENT_STREAM_RESUME_ACK, id: data.id, }), ); break; case MessageType.CF_AGENT_USE_CHAT_RESPONSE: { if (!data.done && !data.error) this.idleConfirmed = false; let chunk: { type?: string; messageId?: string } | undefined; try { if (data.body?.trim()) chunk = JSON.parse(data.body); } catch { if (!data.error && !data.done) return; } const owned = this.active.has(data.id); if ( owned && chunk?.type === "start" && typeof chunk.messageId === "string" ) { const index = this.view.messages.findIndex( (message) => message.id === chunk!.messageId, ); this.protectedAssistant = { id: chunk.messageId, anchor: index >= 0 ? this.view.messages[index - 1]?.id : this.view.messages.at(-1)?.id, }; if (data.replay && !data.continuation) this.setMessages( this.view.messages.map((message) => message.id === chunk!.messageId ? { ...message, parts: [] } : message, ), ); } if (!owned) { if ( data.replay && this.observed.status === "idle" && !this.acked.has(data.id) ) break; transport.observeServerTurn(data.id); const result = broadcastTransition(this.observed, { type: "response", streamId: data.id, messageId: crypto.randomUUID(), chunkData: chunk, done: data.done, error: data.error, replay: data.replay, replayComplete: data.replayComplete, continuation: data.continuation, }); this.observed = result.state; if (result.messagesUpdate) this.setMessages( result .messagesUpdate(this.view.messages) .map((message) => message.id === (result.state.status === "observing" ? result.state.accumulator.messageId : "") ? structuredClone(message) : message, ), ); this.publish({ serverStreaming: result.isStreaming }); } if (data.done || data.error) { // A fallback observer can become transport-owned during resume. Its // accumulator must retire with that stream before final history is // applied, or it would overwrite saved partials with replayed text. if ( this.observed.status === "observing" && this.observed.streamId === data.id ) this.observed = { status: "idle" }; this.idleConfirmed = true; transport.handleServerTurnCompleted(data.id); this.acked.delete(data.id); this.publish({ serverStreaming: false, recovering: false, stopping: false, reconciling: true, ...(data.error ? { error: new Error("Response failed") } : {}), }); this.scheduleReconcile(generation); } break; } } } private upsertActivities(activities: ToolActivity[]) { const next = new Map(this.view.activities); for (const activity of activities) if ( !next.has(activity.toolCallId) || next.get(activity.toolCallId)!.updatedAt <= activity.updatedAt ) next.set(activity.toolCallId, structuredClone(activity)); this.publish({ activities: next }); } loadActivities = async () => { if ( !this.client || this.view.connection !== "connected" || this.view.activitiesLoading || this.view.activityCursor === null ) return; const generation = this.generation; const revision = this.streamEpoch; this.publish({ activitiesLoading: true }); try { const page = await this.client.call<ToolActivityPage>( "listToolActivities", this.view.activityCursor === undefined ? [] : [this.view.activityCursor], ); if (!this.current(generation) || revision !== this.streamEpoch) return; this.upsertActivities(page.activities); this.publish({ activityCursor: page.nextCursor }); } catch { if (this.current(generation)) this.publish({ notice: "Could not load earlier activity. Try again." }); } finally { if (this.current(generation)) this.publish({ activitiesLoading: false }); } }; // Native RPC has no per-call AbortSignal. Fence both sides of the call so // a retired socket cannot populate Query or complete a newer UI operation. private async approvalRpc<T>( method: "listCapabilityApprovals" | "decideCapabilityApproval", args: unknown[], signal?: AbortSignal, ): Promise<T> { const generation = this.generation; const client = this.client; const check = () => { signal?.throwIfAborted(); if ( !client || client !== this.client || !this.current(generation) || this.view.connection !== "connected" || !navigator.onLine ) throw new Error("Conversation connection unavailable"); }; check(); const result = await client!.call<T>(method, args, { timeout: method === "decideCapabilityApproval" ? 120_000 : 30_000, }); check(); return result; } listCapabilityApprovals = (signal?: AbortSignal) => this.approvalRpc<CapabilityApprovalRequest[]>( "listCapabilityApprovals", [], signal, ); decideCapabilityApproval = (decision: { executionId: string; decision: "approve" | "deny"; }) => this.approvalRpc<{ resolved: boolean }>("decideCapabilityApproval", [ decision, ]); selectModel = (configuration: ModelConfiguration | null) => { this.publish({ modelSelection: configuration }); }; private async modelBody(generation: number) { const selection = this.view.modelSelection; if (selection) return { modelConfiguration: selection, useDefaultModel: false }; // Refresh inherited defaults just before submission, not only when opening // the chat. The resulting snapshot travels with the native request. const settings = await this.client!.call<ConversationModelSettings>( "getConversationModelSettings", ); if (this.current(generation)) this.publish({ modelSettings: settings }); return { modelConfiguration: settings.configuration, useDefaultModel: true, }; } send = async (text: string, attachments: Attachment[] = []) => { if ( !this.chat || this.busy || this.view.reconciling || this.view.connection !== "connected" || (!text.trim() && !attachments.length) ) return false; const generation = this.generation; const operation = { generation, cancelled: false }; this.operation = operation; const id = crypto.randomUUID(); this.idleConfirmed = false; this.pendingInput = { id, text, attachments }; this.publish({ notice: undefined, error: undefined, unsentText: undefined, unsentAttachments: undefined, status: "submitted", }); try { const body = await this.modelBody(generation); if (!this.current(generation) || operation.cancelled) return false; await this.chat.sendMessage( { id, role: "user", parts: [ { type: "text", text: text || "Please inspect the attached files.", }, ], ...(attachments.length ? { metadata: { attachments } } : {}), }, { body }, ); } catch { if (this.current(generation)) this.publish({ error: new Error("Could not submit the message"), unsentText: text, unsentAttachments: attachments, }); } finally { if (this.operation === operation) this.operation = undefined; if (this.current(generation)) { if (this.view.error) { this.publish({ reconciling: true }); void this.connect(); } else this.scheduleReconcile(generation); } } return true; }; stop = async () => { if ( !this.chat || !this.transport || !this.busy || this.view.connection !== "connected" ) return; this.publish({ stopping: true }); if (this.operation) this.operation.cancelled = true; const cancelled = this.transport.cancelActiveServerTurn(); await this.chat.stop(); if (!cancelled) { this.publish({ stopping: false, notice: "Checking the active response before stopping. Try Stop again when connected.", reconciling: true, }); void this.connect(); } // The shared native listener (not the detached AI reader) confirms terminal status. }; retry = async () => { if ( !this.chat || this.busy || this.view.reconciling || this.view.connection !== "connected" ) return; const last = this.view.messages.findLast( (message) => message.role === "assistant" || message.role === "user", ); if (!last) return; const retained = this.view.messages; this.idleConfirmed = false; const generation = this.generation; const operation = { generation, cancelled: false }; this.operation = operation; this.publish({ error: undefined, notice: undefined, status: "submitted" }); try { const body = await this.modelBody(generation); if (!this.current(generation) || operation.cancelled) return; await this.chat.regenerate({ messageId: last.id, body, }); } catch { if (this.current(generation)) this.publish({ error: new Error("Could not retry the response") }); } finally { if (this.operation === operation) this.operation = undefined; if (this.current(generation)) { if (this.view.error) { this.setMessages(retained); this.publish({ reconciling: true }); void this.connect(); } else this.scheduleReconcile(generation); } } }; close = () => { window.removeEventListener("online", this.online); window.removeEventListener("offline", this.offline); this.detach(); this.lifetime.abort(); clearTimeout(this.reconnectTimer); clearTimeout(this.reconcileTimer); clearTimeout(this.probeTimer); this.listeners.clear(); };}