src/app/core/database/pouchdb/remote-pouch-database.ts

Description

An alternative implementation of PouchDatabase that directly makes HTTP requests to a remote CouchDB.

Extends

PouchDatabase

Relationships

Used by

No results matching.

Depends on

Index

Properties
Methods

Constructor

constructor(dbName: string, authService: KeycloakAuthService, globalSyncState?: SyncStateSubject, ngZone?: NgZone, alertService?: AlertService)
Parameters :
Name Type Optional
dbName string No
authService KeycloakAuthService No
globalSyncState SyncStateSubject Yes
ngZone NgZone Yes
alertService AlertService Yes

Properties

Optional trackLostPermissions
Type : boolean

Whether to track docs whose permissions were lost as reported by the server. Toggled by SyncedPouchDatabase to skip tracking on first sync.

adapter
Type : string
Default value : "indexeddb"
Inherited from PouchDatabase

The PouchDB adapter to use for local storage. Set by the factory/resolver before calling init(). Default: "indexeddb" (the newer adapter). Use "idb" for the legacy adapter.

Optional analytics
Type : function
Inherited from PouchDatabase

Optional accessor for usage analytics (see reportConflict).

Assigned by DatabaseFactoryService rather than injected: the database classes are constructed manually and have no access to Angular DI. It resolves lazily, and only once something is actually tracked, so that bootstrap never closes the cycle AnalyticsService -> ConfigService -> EntityMapperService -> DatabaseResolver -> DatabaseFactoryService.

Protected changesFeed
Type : Subject<any>
Inherited from PouchDatabase

An observable that emits a value whenever the PouchDB receives a new change. This change can come from the current user or remotely from the (live) synchronization

Protected databaseInitialized
Type : unknown
Default value : new Subject<void>()
Inherited from PouchDatabase
Protected Readonly destroy$
Type : unknown
Default value : new Subject<void>()
Inherited from PouchDatabase

trigger to unsubscribe any internal subscriptions

Protected indexPromises
Type : Promise<any>[]
Default value : []
Inherited from PouchDatabase

A list of promises that resolve once all the (until now saved) indexes are created

Protected pouchDB
Type : PouchDB.Database
Inherited from PouchDatabase

The reference to the PouchDB instance

Methods

collectAndClearLostPermissions
collectAndClearLostPermissions()

Returns all doc IDs whose permissions were lost since the last sync (as reported in _changes responses intercepted during that sync) and resets the internal list.

Returns : string[]
Protected Async findPage
findPage(findOptions: PouchDB.Find.FindRequest, page?: { limit?: number; bookmark?: string })
Inherited from PouchDatabase

Uses the PouchDB-find plugin https://github.com/apache/pouchdb/tree/master/packages/node_modules/pouchdb-find to query the remote CouchDB (via the replication-backend) using the Mango Query Language https://pouchdb.com/guides/mango-queries.html#query-language (see PouchDatabase.find).

Pagination uses CouchDB's real bookmark cursor: pass the bookmark returned by a previous call to continue right after those results. This is forward-only - there is no way to jump back to an earlier page - and (unlike skip) it also works correctly when the server applies permission filtering to the query (see https://github.com/Aam-Digital/replication-backend/pull/330).

Parameters :
Name Type Optional
findOptions PouchDB.Find.FindRequest<any> No
page { limit?: number; bookmark?: string } Yes
Returns : Promise<FindPage>
init
init(dbName?: string, config?: { unauthenticatedSession?: boolean; trackLostPermissions?: boolean })
Inherited from Database

Initializes the PouchDB with the http adapter to directly access a remote CouchDB without replication See https://pouchdb.com/adapters.html#pouchdb_over_http

Parameters :
Name Type Optional Description
dbName string Yes

(relative) path to the remote database

config { unauthenticatedSession?: boolean; trackLostPermissions?: boolean } Yes

optional configuration for the remote database session

Returns : void
Async put
put(object: any, forceOverwrite: unknown)
Inherited from Database

Unlike a synced database, which writes locally and hears about it on the local changes feed in the same tick, this database writes over HTTP and would only learn of its own write on the next poll - up to CHANGES_POLLING_INTERVAL later. Announcing the stored document here closes that gap, so the app reacts to its own writes immediately in both session types.

Parameters :
Name Type Optional Default value
object any No
forceOverwrite unknown No false
Returns : Promise<any>
Async putAll
putAll(objects: any[], forceOverwrite: unknown)
Inherited from Database
Parameters :
Name Type Optional Default value
objects any[] No
forceOverwrite unknown No false
Returns : Promise<any>
Async remove
remove(object: any)
Inherited from Database

A deletion has the same delay as a write, so it is announced the same way. The emitted document is the tombstone the changes feed would deliver - id, revision and the deleted flag, without the data fields - which is what subscribers read to tell a removal apart from an update.

Parameters :
Name Type Optional
object any No
Returns : Promise<any>
Protected shouldSkipIndexUpdate
shouldSkipIndexUpdate(existingDesignDoc: any)
Inherited from PouchDatabase
Parameters :
Name Type Optional
existingDesignDoc any No
Returns : boolean
Protected Async subscribeChanges
subscribeChanges()
Inherited from PouchDatabase

Poll the _changes endpoint periodically to detect document changes. Emits individual documents that have changed since the last poll.

Overridden to use periodic polling instead of live long-polling. This avoids connection stability issues with remote-only (anonymous) sessions. Changes are fetched at regular intervals rather than maintaining a persistent connection.

Returns : any
supportsFind
supportsFind()
Inherited from Database
Returns : boolean
Protected Async withReadRetry
withReadRetry<T>(operation: () => void)
Inherited from PouchDatabase
Type parameters :
  • T

Retry idempotent reads that fail with a transient network error (e.g. an abort/timeout on a connection gone stale after the tab was suspended, or ERR_NETWORK_CHANGED).

Retrying lives here, at the operation level, rather than in the fetch wrapper: fetchWithTimeout can only abort while acquiring the response headers, but an abort during response-body streaming surfaces after the fetch has returned — outside the wrapper's reach. Re-issuing the whole read (fresh fetch and body) recovers both cases transparently. Only reads route through this hook; writes must run exactly once (see fetchWithTimeout). The one exception is creating a query index, which changes nothing when repeated.

Parameters :
Name Type Optional
operation function No
Returns : Promise<T>
Async allDocs
allDocs(options?: GetAllOptions)
Inherited from Database

Load all documents (matching the given PouchDB options) from the database. (see Database)

Normally you should rather use "getAll()" or another well typed method of this class instead of passing PouchDB specific options here because that will make your code tightly coupled with PouchDB rather than any other database provider.

Parameters :
Name Type Optional Description
options GetAllOptions Yes

PouchDB options object as in the normal PouchDB library

Returns : unknown
changes
changes()
Inherited from Database

Listen to changes to documents in the database. Use rxjs operators to filter for specific prefixes etc. if needed.

Returns : Observable<any>

observable which emits the filtered changes

Async destroy
destroy()
Inherited from Database

Destroy the database and all saved data

Returns : Promise<any>
Async find
find(prefix: string, query: PouchDB.Find.Selector, page?: { limit?: number; bookmark?: string }, sort?: { prop?: string; dir?: "asc" | "desc" })
Inherited from Database

Query a page of one entity type's documents with a Mango selector, optionally sorted (see Database.find).

Only available on databases that implement findPage, which differ only in how they fetch one page of results: RemotePouchDatabase uses CouchDB's real Mango bookmark cursor, MemoryPouchDatabase (for tests) the offset of the next document. A synced local database does not support this: PouchDB's local Mango query engine has no bookmark support at all - it always reports "nil" (see pouchdb-find/lib/index.js) - and paging by offset is currently not needed there.

When sorted, documents without a value for the sort property are included as well: last for "asc", first for "desc" (see sortedQueries).

Parameters :
Name Type Optional Default value
prefix string No ""
query PouchDB.Find.Selector No {}
page { limit?: number; bookmark?: string } Yes
sort { prop?: string; dir?: "asc" | "desc" } Yes
Returns : Promise<FindPage>
Async get
get(id: string, options: GetOptions, returnUndefined?: boolean)
Inherited from Database

Load a single document by id from the database. (see Database)

Parameters :
Name Type Optional Default value Description
id string No

The primary key of the document to be loaded

options GetOptions No {}

Optional PouchDB options for the request

returnUndefined boolean Yes

(Optional) return undefined instead of throwing error if doc is not found in database

Returns : Promise<any>
getPouchDB
getPouchDB()
Inherited from PouchDatabase

Get the actual instance of the PouchDB

Returns : PouchDB.Database
Async getPouchDBOnceReady
getPouchDBOnceReady()
Inherited from PouchDatabase
isEmpty
isEmpty()
Inherited from Database

Check if a database is new/empty. Returns true if there are no documents in the database

Returns : Promise<boolean>
isInitialized
isInitialized()
Inherited from Database
Returns : boolean
Protected isNotificationsDatabase
isNotificationsDatabase()
Inherited from PouchDatabase

Check if this is a notifications database based on the database name. (may be used for special handling of notification DBs)

Returns : boolean
Async purge
purge(id: string)
Inherited from Database

Permanently purge a document and all its revisions from the local database.

Unlike remove, which creates a deletion tombstone that is synced, purge completely removes local data without affecting the remote database. This also emits a synthetic deletion event via the changes feed so that in-memory caches (entity stores) drop the purged entity.

Adapter limitation: The underlying PouchDB purge() API is only available on the indexeddb adapter (PouchDB 8+). On unsupported adapters PouchDB will throw its own error, which callers should handle (e.g. via try/catch + logging).

Example :
     false if the document did not exist locally (benign — desired state already achieved)
Parameters :
Name Type Optional Description
id string No

The document ID to purge

Returns : Promise<boolean>

true if the document was purged, false if the document did not exist locally (benign — desired state already achieved)

query
query(fun: string | unknown, options: QueryOptions)
Inherited from Database

Query data from the database based on a more complex, indexed request. (see Database)

This is directly calling the PouchDB implementation of this function. Also see the documentation there: https://pouchdb.com/api.html#query_database

Parameters :
Name Type Optional Description
fun string | unknown No

The name of a previously saved database index

options QueryOptions No

Additional options for the query, like a key. See the PouchDB docs for details.

Returns : Promise<any>
Async reset
reset()
Inherited from Database

Reset the database state so a new one can be opened.

Returns : any
saveDatabaseIndex
saveDatabaseIndex(designDoc: any)
Inherited from Database

Create a database index to query() certain data more efficiently in the future. (see Database)

Also see the PouchDB documentation regarding indices and queries: https://pouchdb.com/api.html#query_database

Parameters :
Name Type Optional Description
designDoc any No

The PouchDB style design document for the map/reduce query

Returns : Promise<void>
getAll
getAll(prefix: string)
Inherited from Database

Load all documents (with the given prefix) from the database.

Parameters :
Name Type Optional Default value Description
prefix string No ""

The string prefix of document ids that should be retrieved

Returns : Promise<Array<any>>
Async removeAll
removeAll<T>(shouldDelete: (doc: T) => void)
Inherited from Database
Type parameters :
  • T

Delete all documents of the database matching the given filter.

Note that getAll() also returns the database's own "_design/..." index documents, so a filter that does not exclude them deletes the indices as well. They are only recreated when the app starts up again (see DatabaseIndexingService), so only delete them if the calling code reloads the app or restores the documents immediately afterwards.

Example :
   Pass a specialized document type if the filter reads fields other than `_id`.
Parameters :
Name Type Optional Description
shouldDelete function No

filter deciding for each document whether it is deleted. Pass a specialized document type if the filter reads fields other than _id.

Returns : Promise<number>

the number of deleted documents

import { DatabaseException, PouchDatabase } from "./pouch-database";
import { environment } from "../../../../environments/environment";
import PouchDB from "pouchdb-browser";
import { Logging } from "../../logging/logging.service";
import { HttpStatusCode } from "@angular/common/http";
import { KeycloakAuthService } from "../../session/auth/keycloak/keycloak-auth.service";
import { RemoteLoginNotAvailableError } from "../../session/auth/keycloak/remote-login-not-available.error";
import { SyncStateSubject } from "app/core/session/session-type";
import { SyncState } from "app/core/session/session-states/sync-state.enum";
import { NgZone } from "@angular/core";
import { timer } from "rxjs";
import { exhaustMap, takeUntil } from "rxjs/operators";
import { AlertService } from "../../alerts/alert.service";
import { isVersionNewer } from "./version-comparison.utils";
import { isConnectivityError } from "#src/app/utils/connectivity-error";
import { FindPage } from "./find-queries";
import {
  describeResponse,
  unexpectedResponseMessage,
} from "../../logging/http-response-logging";

/**
 * 4XX statuses that occur during normal operation
 * (and are handled by callers or the auth layer),
 * so they are not reported to remote logging.
 */
const EXPECTED_4XX_STATUSES: number[] = [
  HttpStatusCode.Unauthorized,
  HttpStatusCode.Forbidden,
  HttpStatusCode.NotFound,
];

/** The HTTP method of a fetch request, defaulting the way `fetch` itself does. */
function requestMethod(opts: RequestInit | undefined): string {
  return (opts?.method ?? "GET").toUpperCase();
}

/**
 * The generation number of a CouchDB revision (the `3` of `3-abc...`), which increases
 * with every write to a document. NaN if it cannot be read, so callers fall back to
 * emitting rather than suppressing a document they cannot place in order.
 */
function revisionNumber(rev: string | undefined): number {
  return Number.parseInt(rev ?? "", 10);
}

/**
 * An alternative implementation of PouchDatabase that directly makes HTTP requests to a remote CouchDB.
 */
export class RemotePouchDatabase extends PouchDatabase {
  /**
   * Whether the session is not logging in any user (e.g. for public forms).
   * @private
   */
  private unauthenticatedSession?: boolean;

  /**
   * Whether to track docs whose permissions were lost as reported by the server.
   * Toggled by {@link SyncedPouchDatabase} to skip tracking on first sync.
   */
  trackLostPermissions?: boolean;

  /**
   * Doc IDs whose permissions were lost as reported by the server in `_changes` responses.
   * Accumulated across all `_changes` calls during a sync and consumed after sync completes.
   */
  private pendingLostPermissions: string[] = [];

  /**
   * Polling interval for changes in milliseconds (for remote-only databases).
   * Avoids long-polling connection issues by using periodic polling instead.
   * @private
   */
  private readonly CHANGES_POLLING_INTERVAL = 10000; // 10 seconds

  /**
   * Document id to the most recent revision this client announced for it through
   * {@link announceOwnWrite}, used by {@link isAlreadyAnnounced} to decide what the
   * changes feed can skip.
   *
   * This covers the poll echoing a write back, a conflict retry announcing the
   * revision its caller then announces again, and a poll response that read the
   * server before a local write and only arrives afterwards. The last one matters
   * most - emitting a superseded revision after a newer one would leave subscribers
   * on stale data until that document changes again.
   */
  private readonly announcedRevisions = new Map<string, string>();

  /**
   * Upper bound for {@link announcedRevisions}, which holds one entry per document this
   * client has written and so grows over a long session. Dropping it past this size keeps
   * memory bounded; the only cost is that a duplicate may slip through, which is the
   * behaviour without this map.
   */
  private readonly MAX_ANNOUNCED_REVISIONS = 1000;

  /** Cooldown (ms) between user-facing connection issue alerts. */
  private readonly CONNECTION_ALERT_COOLDOWN_MS = 60000;
  private lastConnectionAlertTime = 0;

  /** Whether the user was already told that their session could not be renewed. */
  private sessionRenewalAlertShown = false;

  constructor(
    dbName: string,
    private authService: KeycloakAuthService,
    globalSyncState?: SyncStateSubject,
    ngZone?: NgZone,
    private readonly alertService?: AlertService,
  ) {
    super(dbName, globalSyncState, ngZone);
  }

  /**
   * Initializes the PouchDB with the http adapter to directly access a remote CouchDB without replication
   * See {@link https://pouchdb.com/adapters.html#pouchdb_over_http}
   * @param dbName (relative) path to the remote database
   * @param config optional configuration for the remote database session
   */
  override init(
    dbName?: string,
    config?: {
      unauthenticatedSession?: boolean;
      trackLostPermissions?: boolean;
    },
  ) {
    this.unauthenticatedSession = config?.unauthenticatedSession;
    this.trackLostPermissions = config?.trackLostPermissions;

    if (dbName) {
      this.dbName = dbName;
    }
    this.pendingLostPermissions = [];

    const options = {
      adapter: "http",
      skip_setup: true,
      fetch: (url: string | Request, opts: RequestInit) =>
        this.defaultFetch(url, opts),
    };
    // add the proxy prefix to the database name so that we get a correct remote URL
    this.pouchDB = new PouchDB(
      `${environment.DB_PROXY_PREFIX}/${this.dbName}`,
      options,
    );
    this.databaseInitialized.complete();

    // No local sync needed — immediately signal that data is available
    this.globalSyncState?.next(SyncState.COMPLETED);
  }

  /**
   * Maximum number of retries for transient network errors (e.g. ERR_NETWORK_CHANGED).
   */
  private readonly TRANSIENT_ERROR_RETRIES = 2;
  private readonly TRANSIENT_ERROR_DELAY_MS = 2000;

  /**
   * Per-request timeout in ms for READ requests only. If the server sends no
   * response within this window, the read fetch is aborted (and retried at the
   * operation level by {@link withReadRetry}). This prevents long-idle
   * connections from being killed unpredictably by Chrome (ERR_NETWORK_CHANGED)
   * or proxies. Writes are never aborted or retried — see
   * {@link fetchWithTimeout}.
   */
  private readonly FETCH_TIMEOUT_MS = 15000;

  private defaultFetch: Fetch = async (url: string | Request, opts: any) => {
    if (typeof url !== "string") {
      const err = new Error("PouchDatabase.fetch: url is not a string");
      err["details"] = url;
      throw err;
    }

    const remoteUrl =
      environment.DB_PROXY_PREFIX + url.split(environment.DB_PROXY_PREFIX)[1];
    this.authService.addAuthHeader(opts.headers);
    // bypass Angular service worker to avoid synthetic 504 errors on network blips
    if (opts.headers?.set && typeof opts.headers.set === "function") {
      opts.headers.set("ngsw-bypass", "true");
    } else if (opts.headers) {
      opts.headers["ngsw-bypass"] = "true";
    }

    let result: Response;
    try {
      result = await this.fetchWithTimeout(remoteUrl, opts);
    } catch (err) {
      Logging.debug("Failed initial fetch from DB", err);
      Logging.debug("navigator.onLine", navigator.onLine);
      this.showConnectionIssueAlert();
    }

    // Retry login if request failed with unauthorized.
    // This will redirect to Keycloak if the token is expired or missing,
    // which is intentional — it ensures users re-authenticate online
    // when connectivity is available (including after an offline login).
    if (
      result?.status === HttpStatusCode.Unauthorized &&
      !this.unauthenticatedSession
    ) {
      result = await this.retryWithRenewedSession(remoteUrl, opts, result);
    }
    if (result && result.status !== HttpStatusCode.Unauthorized) {
      // the session is valid (again) - whichever request or tab renewed it
      this.sessionRenewalAlertShown = false;
    }

    const method = requestMethod(opts);

    if (!result || result.status >= 500) {
      Logging.debug("Actual DB Fetch response", result);
      Logging.debug("navigator.onLine", navigator.onLine);
      throw new DatabaseException({
        message: "Failed to fetch from DB",
        requestedUrl: remoteUrl,
        actualResponse: JSON.stringify(describeResponse(result, method)),
        actualResponseBody: await result?.text(),
      });
    }

    // additional output for debugging
    if (result?.status >= 400) {
      const response = describeResponse(result, method);
      if (this.isNotificationsDatabase() && result.status === 404) {
        Logging.debug(
          "Notifications database not found (404) - may be expected",
        );
      } else if (EXPECTED_4XX_STATUSES.includes(result.status)) {
        // expired session (401), permission-filtered doc (403) and missing doc (404)
        // are part of normal operation and handled by callers
        Logging.debug("Expected 4XX response from DB", response);
      } else if (result.status === HttpStatusCode.Conflict) {
        // A 409 is CouchDB rejecting a write whose revision is out of date. This
        // layer cannot tell whether that cost anyone anything, but whoever made
        // the request always can:
        //  - remote-only session: the write came from PouchDatabase.put(), and
        //    PouchDatabase.resolveConflict() decides between merging,
        //    overwriting and failing - and counts the outcome it chose.
        //  - synced session: writes go to the local database, so the only 409s
        //    reaching here are replication's own `PUT /_local/<checkpoint>`
        //    races between concurrent tabs, which PouchDB retries internally.
        // Reporting it here would therefore be an alert nobody can act on, and
        // it is what buried the genuinely unexpected statuses in the same issue.
        Logging.debug("Document update conflict from DB", response);
      } else {
        Logging.warn(unexpectedResponseMessage(result.status), {
          ...response,
          requestedUrl: remoteUrl,
        });
      }
    }

    if (
      this.trackLostPermissions &&
      result?.status === HttpStatusCode.Ok &&
      remoteUrl.includes("_changes")
    ) {
      await this.extractLostPermissions(result.clone());
    }

    return result;
  };

  /**
   * Retry idempotent reads that fail with a transient network error
   * (e.g. an abort/timeout on a connection gone stale after the tab was
   * suspended, or ERR_NETWORK_CHANGED).
   *
   * Retrying lives here, at the operation level, rather than in the fetch
   * wrapper: {@link fetchWithTimeout} can only abort while acquiring the
   * response headers, but an abort during response-body streaming surfaces
   * after the fetch has returned — outside the wrapper's reach. Re-issuing the
   * whole read (fresh fetch and body) recovers both cases transparently.
   * Only reads route through this hook; writes must run exactly once
   * (see {@link fetchWithTimeout}). The one exception is creating a query
   * index, which changes nothing when repeated.
   */
  protected override async withReadRetry<T>(
    operation: () => Promise<T>,
  ): Promise<T> {
    for (let attempt = 0; ; attempt++) {
      try {
        return await operation();
      } catch (err) {
        if (
          !isConnectivityError(err) ||
          attempt >= this.TRANSIENT_ERROR_RETRIES
        ) {
          throw err;
        }
        Logging.debug(
          `Transient DB read error (attempt ${attempt + 1}/${this.TRANSIENT_ERROR_RETRIES}), retrying...`,
          err,
        );
        await new Promise((resolve) =>
          setTimeout(resolve, this.TRANSIENT_ERROR_DELAY_MS),
        );
      }
    }
  }

  /**
   * Fetch a request with a per-request timeout so a hung/stale connection is
   * aborted cleanly instead of hanging indefinitely (e.g. after the tab was
   * suspended, or ERR_NETWORK_CHANGED). The abort surfaces as a transient
   * error that {@link withReadRetry} retries at the operation level.
   *
   * Only safe/idempotent read methods (GET, HEAD) get the abort timeout.
   * Non-idempotent writes (PUT, POST, DELETE) are run exactly once with no
   * client-side timeout: a write may have already committed on the server even
   * when the client never sees the response, so aborting it can turn a "create"
   * into a forbidden "update" (the public role is create-only) or create a
   * duplicate — surfacing as spurious "unauthorized" errors. Letting the write
   * run to completion ensures the client reliably learns the new `_rev`.
   */
  private async fetchWithTimeout(
    url: string,
    opts: RequestInit,
  ): Promise<Response> {
    const method = requestMethod(opts);
    const isSafeMethod = method === "GET" || method === "HEAD";

    if (!isSafeMethod) {
      // Write: run once, to completion, without abort timeout.
      return PouchDB.fetch(url, opts);
    }

    const controller = new AbortController();
    const timeoutId = setTimeout(
      () => controller.abort(),
      this.FETCH_TIMEOUT_MS,
    );
    try {
      return await PouchDB.fetch(url, { ...opts, signal: controller.signal });
    } finally {
      clearTimeout(timeoutId);
    }
  }

  /**
   * Renew the session and repeat a request that the server rejected as unauthorized.
   *
   * @returns the repeated request's response; the original `unauthorized`
   *   response if the session could not be renewed; or `undefined` if the
   *   repeated request did not get a response at all (like a failed initial fetch)
   */
  private async retryWithRenewedSession(
    remoteUrl: string,
    opts: RequestInit,
    unauthorized: Response,
  ): Promise<Response | undefined> {
    try {
      await this.authService.login();
    } catch (err) {
      if (err instanceof RemoteLoginNotAvailableError) {
        // transient - the user is told, and background sync/polling keeps retrying
        Logging.debug(
          "Could not renew session after 401 (login unavailable)",
          err,
        );
        this.showSessionRenewalAlert();
      } else {
        Logging.warn("Could not renew session after 401", err);
      }
      return unauthorized;
    }

    this.authService.addAuthHeader(opts.headers);
    try {
      return await PouchDB.fetch(remoteUrl, opts);
    } catch (err) {
      Logging.debug("Failed retried fetch from DB after 401", err);
      this.showConnectionIssueAlert();
      return undefined;
    }
  }

  private showConnectionIssueAlert(): void {
    const now = Date.now();
    if (
      now - this.lastConnectionAlertTime <
      this.CONNECTION_ALERT_COOLDOWN_MS
    ) {
      return;
    }
    this.lastConnectionAlertTime = now;
    this.alertService?.addWarning(
      $localize`We are observing connection issues while syncing your data. Sync continues and retries automatically but may take longer than usual.`,
    );
  }

  /**
   * Tell the user that the server rejects their requests because the online
   * session could not be renewed, which otherwise just looks like data that
   * silently fails to load or sync.
   *
   * Shown once until the server accepts a request again, rather than on every
   * rejected request, as the background polling / sync keeps running into it.
   */
  private showSessionRenewalAlert(): void {
    if (this.sessionRenewalAlertShown) {
      return;
    }
    this.sessionRenewalAlertShown = true;
    this.alertService?.addWarning(
      $localize`:Alert when the login session could not be renewed:Your online login session could not be renewed because the login server cannot be reached right now. Until then, data cannot be loaded from or saved to the server. We keep retrying automatically.`,
    );
  }

  /**
   * Parse `lostPermissions` from a `_changes` response and accumulate them
   * for later retrieval via {@link collectAndClearLostPermissions}.
   *
   * Awaited inside `defaultFetch` to ensure items are collected before the
   * fetch resolves — preventing a race with `collectAndClearLostPermissions`.
   */
  private async extractLostPermissions(response: Response): Promise<void> {
    try {
      const body = await response.json();
      if (body.lostPermissions?.length) {
        this.pendingLostPermissions.push(...body.lostPermissions);
      }
    } catch (err) {
      Logging.debug(
        "Could not parse lostPermissions from _changes response",
        err,
      );
    }
  }

  /**
   * Returns all doc IDs whose permissions were lost since the last sync
   * (as reported in `_changes` responses intercepted during that sync)
   * and resets the internal list.
   */
  collectAndClearLostPermissions(): string[] {
    const collected = this.pendingLostPermissions;
    this.pendingLostPermissions = [];
    return collected;
  }

  override supportsFind(): boolean {
    return true;
  }

  /**
   * Uses the PouchDB-find plugin {@link https://github.com/apache/pouchdb/tree/master/packages/node_modules/pouchdb-find}
   * to query the remote CouchDB (via the replication-backend) using the Mango
   * Query Language {@link https://pouchdb.com/guides/mango-queries.html#query-language}
   * (see {@link PouchDatabase.find}).
   *
   * Pagination uses CouchDB's real `bookmark` cursor: pass the `bookmark`
   * returned by a previous call to continue right after those results. This
   * is forward-only - there is no way to jump back to an earlier page - and
   * (unlike `skip`) it also works correctly when the server applies
   * permission filtering to the query (see
   * {@link https://github.com/Aam-Digital/replication-backend/pull/330}).
   */
  protected override async findPage(
    findOptions: PouchDB.Find.FindRequest<any>,
    page?: { limit?: number; bookmark?: string },
  ): Promise<FindPage> {
    // the installed @types/pouchdb-find does not declare `bookmark`, although
    // both CouchDB and pouchdb-find's own request/response objects support it
    const request: PouchDB.Find.FindRequest<any> & { bookmark?: string } = {
      ...findOptions,
    };
    if (Number.isInteger(page?.limit)) {
      request.limit = page.limit;
    }
    if (page?.bookmark) {
      request.bookmark = page.bookmark;
    }

    const pouchDB = await this.getPouchDBOnceReady();
    const res = await this.withReadRetry(
      () =>
        pouchDB.find(request) as Promise<
          PouchDB.Find.FindResponse<any> & { bookmark?: string }
        >,
    ).catch((err) => {
      throw new DatabaseException(err);
    });
    return { docs: res.docs, bookmark: res.bookmark };
  }

  protected override shouldSkipIndexUpdate(existingDesignDoc: any): boolean {
    if (
      existingDesignDoc.aam_version &&
      isVersionNewer(existingDesignDoc.aam_version, environment.appVersion)
    ) {
      Logging.debug(
        `skipping index update for ${existingDesignDoc._id}: server has version ${existingDesignDoc.aam_version}, we are ${environment.appVersion}`,
      );
      return true;
    }
    return false;
  }

  /**
   * Unlike a synced database, which writes locally and hears about it on the local
   * changes feed in the same tick, this database writes over HTTP and would only learn
   * of its own write on the next poll - up to {@link CHANGES_POLLING_INTERVAL} later.
   * Announcing the stored document here closes that gap, so the app reacts to its own
   * writes immediately in both session types.
   */
  override async put(object: any, forceOverwrite = false): Promise<any> {
    const result = await super.put(object, forceOverwrite);
    this.announceOwnWrite(object, result);
    return result;
  }

  override async putAll(objects: any[], forceOverwrite = false): Promise<any> {
    try {
      const results = await super.putAll(objects, forceOverwrite);
      this.announceStoredDocs(objects, results);
      return results;
    } catch (results) {
      // putAll rejects *with* its results array when any document failed; the ones
      // that did store are still stored and must be announced like any other write.
      if (Array.isArray(results)) {
        this.announceStoredDocs(objects, results);
      }
      throw results;
    }
  }

  /**
   * A deletion has the same delay as a write, so it is announced the same way.
   * The emitted document is the tombstone the changes feed would deliver - id,
   * revision and the deleted flag, without the data fields - which is what
   * subscribers read to tell a removal apart from an update.
   */
  override async remove(object: any): Promise<any> {
    const result = await super.remove(object);
    this.announceOwnWrite({ _id: object?._id, _deleted: true }, result);
    return result;
  }

  private announceStoredDocs(objects: any[], results: any[]) {
    for (const result of results ?? []) {
      const stored = objects.find((obj) => obj._id === result?.id);
      if (stored) {
        this.announceOwnWrite(stored, result);
      }
    }
  }

  /**
   * Whether subscribers already have this document at this revision or a newer one,
   * because this client announced it. An unreadable revision is treated as new, so an
   * unexpected format delays nothing - it only forgoes the deduplication.
   */
  private isAlreadyAnnounced(doc: { _id?: string; _rev?: string }): boolean {
    const announced = this.announcedRevisions.get(doc?._id);
    if (announced === undefined) {
      return false;
    }
    if (announced === doc?._rev) {
      // the server echoing back the exact revision this client wrote
      return true;
    }

    // A different revision of the same generation is a conflicting sibling, not an
    // echo - two clients edited the same parent and the server picked a winner
    // between them. Suppressing it would hide that winner from every subscriber,
    // so only a strictly older generation counts as superseded.
    const incoming = revisionNumber(doc?._rev);
    const mine = revisionNumber(announced);
    return !isNaN(incoming) && !isNaN(mine) && incoming < mine;
  }

  /**
   * Emit a document this client just stored, tagged with the revision the server
   * assigned, so subscribers see the same shape the changes feed would deliver.
   */
  private announceOwnWrite(object: any, result: any) {
    if (!result?.ok || !this.changesFeed) {
      return;
    }
    if (typeof object?._id === "string" && object._id.startsWith("_design/")) {
      // index definitions, not entities: no subscriber reacts to them, and the
      // server's _changes does not echo them, so tracking them would only fill
      // announcedRevisions with entries that never match.
      return;
    }

    const doc = { ...object, _rev: result.rev };
    if (this.isAlreadyAnnounced(doc)) {
      // a conflict retry announced this same revision through the nested put()
      return;
    }

    if (this.announcedRevisions.size >= this.MAX_ANNOUNCED_REVISIONS) {
      this.announcedRevisions.clear();
    }
    this.announcedRevisions.set(doc._id, doc._rev);

    if (this.ngZone) {
      this.ngZone.run(() => this.changesFeed.next(doc));
    } else {
      this.changesFeed.next(doc);
    }
  }

  /**
   * Poll the _changes endpoint periodically to detect document changes.
   * Emits individual documents that have changed since the last poll.
   *
   * Overridden to use periodic polling instead of live long-polling.
   * This avoids connection stability issues with remote-only (anonymous) sessions.
   * Changes are fetched at regular intervals rather than maintaining a persistent connection.
   *
   * @private
   */
  protected override async subscribeChanges() {
    const db = await this.getPouchDBOnceReady();
    let lastSequence: string | number = "now";

    // Run the polling loop outside Angular's zone to avoid:
    //  - a full app-wide change-detection cycle for every poll fetch
    //  - PouchDB-internal promise rejections being routed to Angular's
    //    ErrorHandler / Sentry. We re-enter the zone explicitly only when
    //    emitting to changesFeed so subscribers still trigger CD.
    const startPolling = () =>
      timer(0, this.CHANGES_POLLING_INTERVAL)
        .pipe(
          exhaustMap(async () => {
            try {
              const result = await db.changes({
                since: lastSequence,
                include_docs: true,
              });

              if (result?.results) {
                result.results.forEach(
                  (change: PouchDB.Core.ChangesResponseChange<{}>) => {
                    if (this.isAlreadyAnnounced(change.doc)) {
                      // this client's own write, already emitted by announceOwnWrite,
                      // or a response that read the server before that write
                      return;
                    }
                    if (this.ngZone) {
                      this.ngZone.run(() => this.changesFeed.next(change.doc));
                    } else {
                      this.changesFeed.next(change.doc);
                    }
                  },
                );
                lastSequence = result.last_seq;
              }

              return result;
            } catch (err) {
              Logging.debug("Error polling changes from remote database", err);
              // Continue polling despite errors
              return null;
            }
          }),
          takeUntil(this.destroy$),
        )
        .subscribe();

    if (this.ngZone) {
      this.ngZone.runOutsideAngular(startPolling);
    } else {
      startPolling();
    }

    Logging.debug(
      `Started periodic changes polling for ${this.dbName} (interval: ${this.CHANGES_POLLING_INTERVAL}ms)`,
    );
  }
}

results matching ""

    No results matching ""