Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261import { createContext, useContext, useEffect, useState, useSyncExternalStore, type OctaneNode,} from "octane";import type { ConversationSummary } from "../../shared/conversations";import { createOwnerClient, OwnerSessionError } from "./owner-client";import type { CustomerQuerySession } from "./customer-query-session";import { useCustomerQuerySession } from "./customer-query-provider";import { conversationOptions, useConversationsQuery,} from "./queries/use-conversations-query";import { callOwner } from "./owner-rpc";
type Connection = ReturnType<typeof createOwnerClient>;type View = { status: "connecting" | "connected" | "reconnecting" | "offline" | "unauthorized"; creating: boolean; error: string;};const initial: View = { status: "connecting", creating: false, error: "" };
/** Native connection and conversation mutations. Query owns all summaries. */class ShellSession { private cache: CustomerQuerySession; constructor(cache: CustomerQuerySession) { this.cache = cache; } private view = initial; private listeners = new Set<() => void>(); private connection: Connection | null = null; private active = false; private retryTimer: ReturnType<typeof setTimeout> | undefined; getSnapshot = () => this.view; getServerSnapshot = () => initial; subscribe = (listener: () => void) => { this.listeners.add(listener); return () => { this.listeners.delete(listener); }; }; getConnection = () => this.connection; options = () => conversationOptions(this.cache, this.connection); private publish(patch: Partial<View>) { this.view = { ...this.view, ...patch }; this.listeners.forEach((listener) => listener()); } private retire() { clearTimeout(this.retryTimer); this.connection?.close(); this.connection = null; void this.cache.client.cancelQueries({ queryKey: this.options().queryKey }); } start = () => { this.active = true; this.reconnect(); window.addEventListener("online", this.reconnect); window.addEventListener("offline", this.offline); return () => { this.active = false; this.retire(); window.removeEventListener("online", this.reconnect); window.removeEventListener("offline", this.offline); }; }; private offline = () => { this.retire(); this.publish({ status: "offline", creating: false }); }; reconnect = () => { if (!this.active) return; if (!navigator.onLine) return this.offline(); this.retire(); this.publish({ status: this.view.status === "connecting" ? "connecting" : "reconnecting", creating: false, error: "", }); const connection = createOwnerClient((event) => { if (!this.active || this.connection !== connection) return; this.retire(); const unauthorized = event.code === 4001; this.publish({ status: unauthorized ? "unauthorized" : "reconnecting", creating: false, error: unauthorized ? new OwnerSessionError().message : "", }); if (!unauthorized) this.retryTimer = setTimeout(this.reconnect, 3000); }, this.cache); this.connection = connection; void connection.ready .then(async (client) => { if ( !this.active || this.connection !== connection || !connection.isCurrent() ) return; client.addEventListener("message", (event) => { if (!this.active || this.connection !== connection) return; try { if (JSON.parse(String(event.data)).type === "conversations-changed") this.refresh(); } catch { /* Other native protocol frames do not affect metadata. */ } }); const options = this.options(); await this.cache.client.invalidateQueries({ queryKey: options.queryKey, refetchType: "none", }); if (!connection.isCurrent()) return; // Publish native readiness for observers, but Connected still waits for a list. this.publish({}); }) .catch((error: unknown) => { if (!this.active || this.connection !== connection) return; this.retire(); const unauthorized = error instanceof OwnerSessionError; this.publish({ status: unauthorized ? "unauthorized" : "reconnecting", error: unauthorized ? error.message : "Could not reach Flarebot. Retrying…", }); if (!unauthorized) this.retryTimer = setTimeout(this.reconnect, 3000); }); }; confirmReady = () => { if ( !this.active || !this.connection?.isReady() || this.view.status === "connected" ) return; // Query errors belong to the read UI. Only native failures retire the socket. this.publish({ status: "connected" }); }; refresh = () => { if (!this.connection?.isReady()) return; const queryKey = this.options().queryKey; this.publish({ error: "" }); // Broadcasts must retire even an initial request that captured older metadata. void this.cache.client.cancelQueries({ queryKey }); void this.cache.client.invalidateQueries({ queryKey }); }; createConversation = async () => { const connection = this.connection; if ( !connection?.isReady() || this.view.status !== "connected" || this.view.creating ) return; const queryKey = this.options().queryKey; this.publish({ creating: true, error: "" }); try { const conversation = await callOwner<ConversationSummary>( connection, "createConversation", ); if ( !this.active || this.connection !== connection || !connection.isCurrent() ) return; await this.cache.client.cancelQueries({ queryKey }); if ( !this.active || this.connection !== connection || !connection.isCurrent() ) return; this.cache.client.setQueryData<ConversationSummary[]>( queryKey, (previous) => [ conversation, ...(previous ?? []).filter((item) => item.id !== conversation.id), ], ); void this.cache.client.invalidateQueries({ queryKey, refetchType: "none", }); return conversation; } catch { if (this.active && this.connection === connection) this.publish({ error: "Could not create a conversation. Your current conversation is still here. Refresh conversations before trying again.", }); } finally { if (this.active && this.connection === connection) this.publish({ creating: false }); } }; deleteConversation = async (id: string) => { const connection = this.connection; const queryKey = this.options().queryKey; await callOwner<void>(connection, "deleteConversation", [id]); await this.cache.client.cancelQueries({ queryKey }); if ( !this.active || this.connection !== connection || !connection?.isCurrent() ) throw new Error("Conversation connection retired"); this.cache.client.setQueryData<ConversationSummary[]>( queryKey, (previous) => (previous ?? []).filter((item) => item.id !== id), ); void this.cache.client.invalidateQueries({ queryKey }); };}const Context = createContext<ShellSession | null>(null);export function ShellSessionProvider({ children }: { children: OctaneNode }) { const cache = useCustomerQuerySession(); const [session] = useState(() => new ShellSession(cache)); const view = useSyncExternalStore( session.subscribe, session.getSnapshot, session.getServerSnapshot, ); const query = useConversationsQuery(session.getConnection()); useEffect(() => session.start(), [session]); useEffect(() => { if (query.isSuccess && !query.isFetching) session.confirmReady(); }, [session, view.status, query.isSuccess, query.isFetching]); return <Context.Provider value={session}>{children}</Context.Provider>;}export function useShellSession() { const session = useContext(Context); if (!session) throw new Error("Shell session provider is missing"); const view = useSyncExternalStore( session.subscribe, session.getSnapshot, session.getServerSnapshot, ); const query = useConversationsQuery(session.getConnection()); const listError = query.isError ? "Could not refresh conversations. Try again." : ""; return { ...view, conversations: query.data ?? [], loading: query.isPending && (view.status === "connecting" || view.status === "reconnecting"), error: view.error || listError, listError, session, };}