import * as interfaces from '../../ts_interfaces/index.js'; export interface IOpsRealtimeTicket { serverEpoch: string | null; revisions: Partial; } export type TOpsRealtimeListener = ( invalidationsArg: interfaces.requests.IOpsRealtimeInvalidation[], ) => void; const emptyRevisions = (): interfaces.requests.IOpsRealtimeRevisionMap => { return Object.fromEntries( interfaces.requests.opsRealtimeResources.map((resource) => [resource, 0]), ) as unknown as interfaces.requests.IOpsRealtimeRevisionMap; }; /** * Orders pushes and HTTP reads. A read may apply only while every captured * resource revision is still current. */ export class OpsRealtimeFreshness { private serverEpoch: string | null = null; private revisions = emptyRevisions(); private listeners = new Set(); public getKnownState(): { knownServerEpoch?: string; knownRevisions: interfaces.requests.IOpsRealtimeRevisionMap; } { return { knownServerEpoch: this.serverEpoch || undefined, knownRevisions: { ...this.revisions }, }; } public capture( resourcesArg: | interfaces.requests.TOpsRealtimeResource | interfaces.requests.TOpsRealtimeResource[], ): IOpsRealtimeTicket { const resources = Array.isArray(resourcesArg) ? resourcesArg : [resourcesArg]; return { serverEpoch: this.serverEpoch, revisions: Object.fromEntries( resources.map((resource) => [resource, this.revisions[resource]]), ), }; } public isCurrent(ticketArg: IOpsRealtimeTicket): boolean { if (ticketArg.serverEpoch !== this.serverEpoch) return false; return Object.entries(ticketArg.revisions).every(([resource, revision]) => { return this.revisions[resource as interfaces.requests.TOpsRealtimeResource] === revision; }); } public accept( serverEpochArg: string, invalidationsArg: interfaces.requests.IOpsRealtimeInvalidation[], ): interfaces.requests.IOpsRealtimeInvalidation[] { const epochChanged = this.serverEpoch !== serverEpochArg; if (epochChanged) { this.serverEpoch = serverEpochArg; this.revisions = emptyRevisions(); } const accepted: interfaces.requests.IOpsRealtimeInvalidation[] = []; for (const invalidation of invalidationsArg) { if (!interfaces.requests.opsRealtimeResources.includes(invalidation.resource)) { continue; } if (invalidation.revision <= this.revisions[invalidation.resource]) { continue; } this.revisions[invalidation.resource] = invalidation.revision; accepted.push(invalidation); } if (accepted.length > 0) { for (const listener of Array.from(this.listeners)) { listener(accepted); } } return accepted; } public subscribe(listenerArg: TOpsRealtimeListener): () => void { this.listeners.add(listenerArg); return () => this.listeners.delete(listenerArg); } public reset(): void { this.serverEpoch = null; this.revisions = emptyRevisions(); } } export const opsRealtimeFreshness = new OpsRealtimeFreshness(); export const captureRealtimeTicket = ( resourcesArg: | interfaces.requests.TOpsRealtimeResource | interfaces.requests.TOpsRealtimeResource[], ): IOpsRealtimeTicket => opsRealtimeFreshness.capture(resourcesArg); export const isRealtimeTicketCurrent = (ticketArg: IOpsRealtimeTicket): boolean => { return opsRealtimeFreshness.isCurrent(ticketArg); };