/** * PushTransport — outbound shipping seam for the watcher daemon. * * The daemon wires file-watcher events into a transport that ships each * PushEvent to the cloud. Defining the interface here lets the daemon ship * with a swappable boundary — tests inject a fake, a later story swaps in a * concrete WebSocket/HTTP implementation, and the daemon entry point never * changes. * * Lifecycle * ───────── * - `start()` is awaited BEFORE the watcher is started. It's the place * to open sockets, refresh tokens, etc. * - `publish(event)` is called for every coalesced PushEvent the watcher * emits. Implementations decide whether to buffer, batch, or send * inline; the daemon awaits the returned promise so back-pressure can * be honored. * - `dispose()` is awaited DURING shutdown, AFTER the watcher has been * torn down. Implementations should drain in-flight publishes (with * their own internal timeout) and close any sockets. * - `connected` is a passive boolean used by the health endpoint. It MAY * flap during reconnect attempts — that's fine; consumers treat it as * advisory, not a contract. * * The `NoopPushTransport` shipped here is the default when no transport * is wired in: it counts publishes (so unit tests can assert delivery) * and logs nothing — observers should rely on the watcher's `onEvent` * counter on the daemon side, not the transport's internals. * * Ported from indigoai-us/hq-pro PR #112 (src/sync/push-transport.ts) into * @indigoai-us/hq-cloud (Path B) per project event-driven-sync-menubar US-007. */ import { type PushEvent } from "./push-event.js"; /** Safe, recursive representation of an error retained for publish diagnostics. */ export interface SerializedPublishError { name: string; message: string; stack?: string; cause?: SerializedPublishError; } /** Process-local identity of a global used at a publish boundary. */ export interface PublishGlobalIdentity { type: string; name?: string; length?: number; matchesModuleLoad: boolean; } /** * One publish attempt's diagnostic boundary record. It contains no request * body, token, or headers. */ export interface PublishAttemptOutcome { elapsedMs: number; abortedAtIssue: boolean; abort?: { observedAt: string; reason: SerializedPublishError; }; responseStatus?: number; error?: SerializedPublishError; globals: { fetch: PublishGlobalIdentity; abortController: PublishGlobalIdentity; }; } /** Server code for accounts that are deliberately denied realtime sync. */ export declare const REALTIME_UNAVAILABLE_CODE = "REALTIME_UNAVAILABLE"; /** Server code for a company push denied by the recipient membership guard. */ export declare const PUSH_SCOPE_FORBIDDEN_CODE = "cross-tenant-push-rejected"; /** * A server-directed, non-retryable realtime denial. The caller should remain * poll-only rather than treating this as a transient transport failure. */ export declare class RealtimeUnavailableError extends Error { readonly code = "REALTIME_UNAVAILABLE"; readonly retryable = false; readonly status: number; constructor(opts: { status: number; operation: string; }); } /** A server-directed, non-retryable denial for one company push scope. */ export declare class PushScopeForbiddenError extends Error { readonly code = "cross-tenant-push-rejected"; readonly retryable = false; readonly status: number; readonly companyScope: string; constructor(opts: { status: number; operation: string; companyScope: string; }); } /** Best-effort classifier for the server's coded realtime denial response. */ export declare function realtimeUnavailableErrorForResponse(status: number, body: string, operation: string): RealtimeUnavailableError | undefined; /** Best-effort classifier for the server's coded company-membership denial. */ export declare function pushScopeForbiddenErrorForResponse(status: number, body: string, operation: string, relativePath: string): PushScopeForbiddenError | undefined; export declare function isRealtimeUnavailableError(error: unknown): error is RealtimeUnavailableError; export declare function isPushScopeForbiddenError(error: unknown): error is PushScopeForbiddenError; /** Serialize an Error (including nested causes) without traversing arbitrary objects. */ export declare function serializePublishError(value: unknown, seen?: WeakSet): SerializedPublishError; /** Read an outcome attached to the original rejected value by HttpPushTransport. */ export declare function getPublishFailureOutcome(value: unknown): PublishAttemptOutcome | undefined; /** Read a successful outcome retained for the exact PushEvent instance. */ export declare function getPublishSuccessOutcome(event: PushEvent): PublishAttemptOutcome | undefined; export interface PushTransport { /** Open sockets, refresh tokens. Awaited before the watcher starts. */ start(): Promise; /** Ship one coalesced PushEvent. Awaited per event for back-pressure. */ publish(event: PushEvent): Promise; /** Drain + close. Awaited during daemon shutdown after watcher.dispose. */ dispose(): Promise; /** Advisory: is the transport currently believed to be connected? */ readonly connected: boolean; } /** * Default `PushTransport` used until a real implementation lands. * * Behavior: * - `start()` flips `connected` to true. * - `publish()` increments a counter (visible via `publishedCount`). * - `dispose()` flips `connected` back to false. * * Deliberately silent: when the daemon runs with this default, the * watcher's per-event log line (emitted by the daemon itself, not the * transport) is the only observability. That keeps the noop from drowning * the log when high-rate writes hit a dev machine. */ export declare class NoopPushTransport implements PushTransport { private _connected; private _count; get connected(): boolean; /** Test/observability hook: how many events have been published. */ get publishedCount(): number; start(): Promise; publish(_event: PushEvent): Promise; dispose(): Promise; } /** * Minimal fetch surface so tests can inject a mocked `fetch` without pulling * in DOM/undici types. Matches the subset of the global `fetch` we use. */ export type FetchLike = (input: string, init?: { method?: string; headers?: Record; body?: string; signal?: AbortSignal; }) => Promise<{ ok: boolean; status: number; text(): Promise; }>; /** * Async getter for the current Cognito access token. Long-running daemons MUST * pass a getter (not a captured string) so each publish resolves the freshest * token — mirrors {@link VaultServiceConfig.authToken} semantics in types.ts. * A static string is also accepted for short-lived tools and tests. */ export type AuthTokenSource = string | (() => string | Promise); export interface HttpPushTransportOptions { /** * Vault API base URL (e.g. `https://vault-api.example.com`). The same * `apiUrl` the runner already resolves for VaultClient. Trailing slashes * are stripped. Required — config/env-driven, no hard-coded default. */ apiUrl: string; /** * Endpoint path the PushEvent is POSTed to. Defaults to `/sync/push` * (matches the deployed hq-pro endpoint from US-006). Override for testing * or future endpoint moves. */ pushPath?: string; /** Cognito JWT — static string OR async getter. See {@link AuthTokenSource}. */ authToken: AuthTokenSource; /** * Optional extra headers (e.g. client identification) merged onto every * request. Authorization + Content-Type are always set by the transport. */ headers?: Record; /** * Per-request timeout in milliseconds. Default 60_000. On timeout the * publish rejects — the daemon treats a rejected publish as a transient * miss (the cadence poll still covers it) and MUST NOT crash. */ timeoutMs?: number; /** Injectable fetch (tests). Defaults to the global `fetch`. */ fetchImpl?: FetchLike; } /** * Real client `PushTransport` that POSTs encoded PushEvents to the deployed * `/sync/push` endpoint, authenticating with the menubar's existing Cognito * bearer token — the same auth path VaultClient uses (Authorization: Bearer * , token resolved per-request via the supplied getter so it * self-heals across refreshes). * * Failure posture * ─────────────── * `publish()` rejects on a network error, a non-2xx response, or a timeout. * The DAEMON is responsible for not letting that rejection crash it — the * watcher's emit path catches publish errors and logs them; the periodic * cadence poll remains the safety net that eventually ships the change. This * transport never swallows errors itself, so callers retain full visibility. * * `connected` flips true on `start()` and false on `dispose()`. It is purely * advisory (HTTP is connectionless); the health endpoint may read it but it * carries no delivery guarantee. */ export declare class HttpPushTransport implements PushTransport { private readonly apiUrl; private readonly pushPath; private readonly getAuthToken; private readonly extraHeaders; private readonly timeoutMs; private readonly fetchImpl; private _connected; constructor(opts: HttpPushTransportOptions); get connected(): boolean; start(): Promise; publish(event: PushEvent): Promise; dispose(): Promise; } //# sourceMappingURL=push-transport.d.ts.map