import type { PublisherQuerySession } from "./publisher-query-session"; export class PublisherReadError extends Error { constructor( readonly status: number, readonly code: string, ) { super(code); } } export async function publisherJson( path: string, signal: AbortSignal, timeout: number | null = 15_000, ): Promise { signal.throwIfAborted(); if (typeof navigator !== "undefined" && !navigator.onLine) throw new Error("Publisher is offline"); const bounded = timeout ? AbortSignal.any([signal, AbortSignal.timeout(timeout)]) : signal; const response = await fetch(path, { credentials: "same-origin", cache: "no-store", signal: bounded, }); const data = await response.json(); bounded.throwIfAborted(); if (!response.ok) throw new PublisherReadError( response.status, typeof data.error === "string" ? data.error : "temporarily_unavailable", ); return data as T; } export async function publisherRead( session: PublisherQuerySession, path: string, signal: AbortSignal, owner: (data: T) => string, identityOnly = false, timeout: number | null = 15_000, ): Promise { const { epoch, ownerSubject } = session.getSnapshot(); try { const data = await publisherJson(path, signal, timeout); session.assertCurrent(epoch); const subject = owner(data); if (typeof subject !== "string" || !subject) throw new Error("Invalid owner response"); if (subject !== ownerSubject) { session.identify(subject); throw new Error("Publisher owner changed"); } return data; } catch (error) { signal.throwIfAborted(); session.assertCurrent(epoch); if ( identityOnly && error instanceof PublisherReadError && error.status === 401 ) session.identify(null); throw error; } }