/** * Per-spreadsheet provider coordinator that serializes remote mutations. * * The Google Sheets API is not transactional across requests, and concurrent * in-flight requests to one spreadsheet can race (duplicate-row appends, * interleaved update/delete batches) and burn quota with lock-style retries. * Before this coordinator, the outbound effect worker, the User_Input polling * supervisor, provisioning, and test controls each called the provider * directly, so several in-flight requests competed for the same Sheet context * at once. Under load that produced write races and tail latency dominated by * retries. * * The coordinator is a transparent decorator: it implements the same provider * boundary the worker and polling already use, so no caller changes are * required. It routes every mutation and recovery barrier through one * per-spreadsheet FIFO lane (default total concurrency 1) while leaving * lock-free value reads untouched. SQLite outbox, leases, receipts, CAS, and * fencing remain the durable authority; the coordinator only prevents the Node * side from issuing competing mutations. * * Invariants preserved: * - Total mutation concurrency is 1 per lane by default, giving one in-flight * write per spreadsheet at a time. * - Lock-free reads (`readSnapshot`, `readRows`, `readRowsBatch`) never wait on * the mutation lane. * - Recovery barriers (`readEffectPostcondition(s)`) use the mutation lane so a * postcondition probe proves an earlier timed-out critical section ended. * - Batch mutation (`observeSnapshots`) acquires every involved lane in a * stable sorted order, so it cannot deadlock with single-sheet mutations. * - Provisioning uses a dedicated lane so setup never races an append/write. * - A failing or throwing mutation releases the lane, so it never deadlocks. * - `runSerializedInner` runs an optional in-lane precondition (the effect * lease renewal) after the lane is acquired and before the inner provider * call, so queue time plus limiter waits cannot outlive the lease. */ import type { ApplySyncEffectsRequest, ApplySyncEffectsResult, EnsureSyncRowAnchorsRequest, EnsureSyncRowAnchorsResult, FastAppendRowsRequest, FastAppendRowsResult, PreparedApplyEffects, ReadSyncEffectPostconditionsRequest, ReadSyncRowChecksRequest, ReadSyncSnapshotRequest, ReadSyncTableRowsRequest, SyncEffectPostcondition, SyncProjectionEffect, SyncEffectPostconditionResult, SyncRowChecksResult, SyncSheetsSnapshot, SyncObservedSnapshot, SyncSheetsProvider, SyncSheetsObservationBatchProvider, SyncSheetsTableReader, SyncTableRowsResult, SyncEffectWorkerProvider } from "../syncSheets.js"; import { AsyncMutex } from "./asyncMutex.js"; import { type CoordinatorLaneEvent } from "./laneTelemetry.js"; /** Provider boundary the coordinator wraps (effect + observation + table read). */ export type CoordinatedSheetsInner = SyncSheetsProvider & SyncSheetsTableReader; /** * Thrown when an in-lane precondition callback rejects before the remote * call. Carries only the static redacted message above. */ export declare class CoordinatedLanePreconditionError extends Error { constructor(); } /** * Raised when a coordinated provider wraps an inner provider that does not * implement the split apply stages. The dispatcher catches this and falls * back to the inner's single legacy `applyEffects` path. */ export declare class CoordinatedSplitApplyUnsupportedError extends Error { constructor(); } /** Options for constructing a per-spreadsheet coordinator. */ export interface CoordinatedSheetsProviderOptions { readonly inner: TInner; /** * Derives the mutation-lane key for one physical sheet. Defaults to a * constant key, giving total mutation concurrency 1 across the process. Pass * a per-spreadsheet resolver when one coordinator serves multiple * spreadsheets so unrelated sheets do not serialize against each other. */ readonly mutationKeyForPhysicalSheet?: (physicalSheetId: string) => string; /** Diagnostic lane observer; failures are swallowed and never alter results. */ readonly onLaneEvent?: (event: CoordinatorLaneEvent) => void; /** Injectable clock used by deterministic telemetry tests. */ readonly clock?: () => number; } /** * Serializes provider mutations per lane while keeping value reads lock-free. * * Implements the full provider, table-reader, observation-batch, and provisioner * boundary, so it can wrap both the real Google Sheets API provider and fakes. * When the inner provider lacks the one-request observation path, observation is * routed through `ensureRowAnchors` + `readSnapshot` inside the same held lane. */ export declare class CoordinatedSheetsProvider implements SyncSheetsProvider, SyncSheetsTableReader, SyncSheetsObservationBatchProvider { private readonly inner; private readonly resolveKey; private readonly onLaneEvent; private readonly now; private readonly lanes; /** * Narrow row-check read (the formula check-column polling gate), exposed * ONLY when the inner provider supports it. Lock-free like every value * read (never waits on a mutation lane). `isSyncSheetsRowChecksReader` * against the coordinator therefore reports the INNER's capability, which * is what polling's feature detection needs. */ readonly readRowChecksBatch?: (requests: readonly ReadSyncRowChecksRequest[]) => Promise; constructor(options: CoordinatedSheetsProviderOptions); /** Returns a snapshot of all known lane metrics; diagnostic only. */ laneMetrics(): ReadonlyMap>; /** * Runs an internal control-plane operation on the same mutation lane as * effects and anchor writes. Test/admin controls use this to avoid issuing a * competing provider request outside the coordinator. */ runSerializedControl(physicalSheetId: string, operation: string, task: () => Promise): Promise; /** * Runs one internal remote task on the mutation lane with an optional * in-lane precondition hook. * * The effect dispatcher uses this so the worker's lease-renewal callback * runs AFTER the lane is acquired (covering queue wait time) but BEFORE the * inner provider's remote call (covering shared limiter waits and the call * itself). The `remote` closure receives the INNER provider and must call * it directly, so the call never re-enters this coordinator's lane and * cannot deadlock. When `beforeRemote` resolves false (or throws), a * `CoordinatedLanePreconditionError` is raised before any remote request * and the lane is released normally; lane telemetry still records the * aborted dispatch with a delivery-uncertain outcome. */ runSerializedInner(physicalSheetId: string, operation: string, remote: (inner: TInner) => Promise, beforeRemote?: () => Promise): Promise; /** * Runs one internal remote task across ALL distinct mutation lanes for a * multi-tab batch, with the same optional in-lane precondition hook as * `runSerializedInner`. Acquiring every involved lane (in stable sorted * order) prevents another writer from interleaving on any tab while the * combined preflight/write or recovery read runs. Single-route callers keep * the single-lane `runSerializedInner` path. */ runSerializedInnerForRoutes(physicalSheetIds: readonly string[], operation: string, remote: (inner: TInner) => Promise, beforeRemote?: () => Promise): Promise; fastAppendRows(request: FastAppendRowsRequest): Promise; applyEffects(request: ApplySyncEffectsRequest): Promise; /** * Lock-free read+plan stage of one apply request. * * Preflight is a read-only stage (sheet enumeration plus a ranged data * read) and must NOT hold the mutation lane: a read-ahead worker wants to * run another route's preflight concurrently with a write. The subsequent * write+verify stage is serialized under the lanes (by the dispatcher via * `runSerializedInner`, and by this coordinator's own `applyPreparedEffects` * for direct callers), and the CAS guards already make a stale read safe. */ preflightApplyEffects(request: ApplySyncEffectsRequest): Promise; /** * Write+verify stage of one apply request, serialized under the mutation * lanes derived from the prepared state's own effect routes. * * The write stage is a REMOTE MUTATION, so unlike the lock-free preflight * above it must never bypass the lanes: a direct consumer of this * coordinator that calls `applyPreparedEffects` itself (instead of going * through the effect dispatcher) previously reached the inner provider * unserialized and could interleave with a concurrent fast append or * regular write on the same spreadsheet. The call therefore acquires every * involved lane (stable sorted order, same as `observeSnapshots`) before * delegating. * * The dispatcher path cannot deadlock on this: it dispatches through * `runSerializedInner`/`runSerializedInnerForRoutes`, whose `remote` closure * receives the INNER provider and calls it directly, so it never re-enters * this method while its lanes are already held. No caller may invoke this * method from inside a `runSerializedControl`/`runSerializedInner` task (the * lanes are not reentrant); the only production writer path is the * dispatcher's out-of-lane preflight plus in-lane write, which this change * keeps byte-identical. */ applyPreparedEffects(prepared: PreparedApplyEffects): Promise; readEffectPostcondition(effect: SyncProjectionEffect): Promise; readEffectPostconditions(request: ReadSyncEffectPostconditionsRequest): Promise; ensureRowAnchors(request: EnsureSyncRowAnchorsRequest): Promise; observeSnapshot(request: ReadSyncSnapshotRequest): Promise; observeSnapshots(requests: readonly ReadSyncSnapshotRequest[]): Promise; /** Lock-free snapshot read; never waits on the mutation lane. */ readSnapshot(request: ReadSyncSnapshotRequest): Promise; /** Lock-free values-only read; never waits on the mutation lane. */ readRows(request: ReadSyncTableRowsRequest): Promise; /** Lock-free batched values-only read; never waits on the mutation lane. */ readRowsBatch(requests: readonly ReadSyncTableRowsRequest[]): Promise; /** * Observes one snapshot through the inner provider, using its single-request * capability when present and falling back to ensure+read otherwise. Both * calls run inside the caller's already-held lane, so observation stays * serialized against competing mutations. */ private observeOneInner; private laneFor; private runMutation; private runInLanes; private emitSuccess; private emitFailure; private emit; } /** The in-lane dispatch capability the effect dispatcher detects on a provider. */ export type CoordinatedSerializedInnerProvider = Pick, "runSerializedInner"> & { /** Multi-route variant; absent on legacy coordinators that predate it. */ runSerializedInnerForRoutes?: CoordinatedSheetsProvider["runSerializedInnerForRoutes"]; }; /** * Returns whether a provider exposes the coordinator's in-lane dispatch hook. * * The effect dispatcher uses this to route its `beforeRemoteDispatch` lease * renewal through the acquired mutation lane; providers without the hook * (bare fakes and test doubles) keep the direct call path and simply ignore * the hook. */ export declare function hasCoordinatedSerializedInner(provider: SyncEffectWorkerProvider): provider is SyncEffectWorkerProvider & CoordinatedSerializedInnerProvider; /** * Returns whether a provider exposes the coordinator's multi-route in-lane * dispatch hook (`runSerializedInnerForRoutes`). * * The effect dispatcher uses this to decide whether a multi-tab call can * acquire every involved mutation lane in one serialized pass. Legacy * coordinators that predate the multi-route hook expose only * `runSerializedInner`; the dispatcher falls back to that single-lane path * instead of crashing on a missing method. */ export declare function hasCoordinatedSerializedInnerForRoutes(provider: SyncEffectWorkerProvider): provider is SyncEffectWorkerProvider & Pick, "runSerializedInnerForRoutes">; /** * Persistent error codes for opaque prepared-token invariant guards. * * Classification (survey §c (I)): a prepared token crosses an opaque API * boundary; a malformed token is foreign-input failure, not a local * programmer assertion, but the thrown error is still a TypeError because the * caller fed structurally invalid data into an internal invariant slot. */ export declare const COORDINATED_PREPARED_STATE_ERROR_CODES: { readonly PREPARED_APPLY_ROUTES_REQUIRED: "coordinated_prepared_apply_routes_required"; readonly PREPARED_APPLY_EFFECT_SHEET_ID_REQUIRED: "coordinated_prepared_apply_effect_sheet_id_required"; }; /** Union type derived from the coordinated prepared-state error code table. */ export type CoordinatedPreparedStateErrorCode = (typeof COORDINATED_PREPARED_STATE_ERROR_CODES)[keyof typeof COORDINATED_PREPARED_STATE_ERROR_CODES]; /** Error thrown when a prepared apply token violates structural invariants. */ export declare class CoordinatedPreparedStateError extends TypeError { readonly code: CoordinatedPreparedStateErrorCode; constructor(code: CoordinatedPreparedStateErrorCode); } //# sourceMappingURL=CoordinatedSheetsProvider.d.ts.map