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

Description

Wrapper for a PouchDB instance to decouple the code from that external library.

Additional convenience functions on top of the PouchDB API should be implemented in the abstract Database.

Extends

Database

Relationships

Index

Properties
Methods

Constructor

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

Properties

adapter
Type : string
Default value : "indexeddb"

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

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>

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>()
Protected Readonly destroy$
Type : unknown
Default value : new Subject<void>()

trigger to unsubscribe any internal subscriptions

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

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

Protected pouchDB
Type : PouchDB.Database

The reference to the PouchDB instance

Methods

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>
Protected Optional findPage
findPage(findOptions: PouchDB.Find.FindRequest, page?: { limit?: number; bookmark?: string })

Run a single Mango query for one page of results. Implemented only by the databases that support find.

Parameters :
Name Type Optional
findOptions PouchDB.Find.FindRequest<any> No
page { limit?: number; bookmark?: string } 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()

Get the actual instance of the PouchDB

Returns : PouchDB.Database
Async getPouchDBOnceReady
getPouchDBOnceReady()
init
init(dbName?: string, options?: PouchDB.Configuration.DatabaseConfiguration | any, suppressSyncCompleted?: boolean)
Inherited from Database

Initialize the PouchDB with the IndexedDB/in-browser adapter (default). See {link https://github.com/pouchdb/pouchdb/tree/master/packages/node_modules/pouchdb-browser}

Parameters :
Name Type Optional Description
dbName string Yes

the name for the database under which the IndexedDB entries will be created

options PouchDB.Configuration.DatabaseConfiguration | any Yes

PouchDB options which are directly passed to the constructor

suppressSyncCompleted boolean Yes

whether to skip emitting a SyncState.COMPLETED event to the globalSyncState (because other logic for sync is building on top of this)

Returns : void
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()

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)

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

Save a document to the database. (see Database)

Parameters :
Name Type Optional Default value Description
object any No

The document to be saved

forceOverwrite unknown No false

(Optional) Whether conflicts should be ignored and an existing conflicting document forcefully overwritten.

Returns : Promise<any>
Async putAll
putAll(objects: any[], forceOverwrite: unknown)
Inherited from Database

Save an array of documents to the database The save can partially fail and return a mix of success and error states in the array (e.g. [{ ok: true, ... }, { error: true, ... }])

Parameters :
Name Type Optional Default value Description
objects any[] No

the documents to be saved

forceOverwrite unknown No false

whether conflicting versions should be overwritten

Returns : Promise<any>

array with the result for each object to be saved, if any item fails to be saved, this returns a rejected Promise. The save can partially fail and return a mix of success and error states in the array (e.g. [{ ok: true, ... }, { error: true, ... }])

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>
remove
remove(object: any)
Inherited from Database

Delete a document from the database (see Database)

Parameters :
Name Type Optional Description
object any No

The document to be deleted (usually this object must at least contain the _id and _rev)

Returns : 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>
Protected shouldSkipIndexUpdate
shouldSkipIndexUpdate(_existingDesignDoc: any)

Whether to skip updating a design doc that differs from the local version. Overridden in RemotePouchDatabase to prevent older clients from overwriting indexes created by a newer app version on a shared server.

Parameters :
Name Type Optional
_existingDesignDoc any No
Returns : boolean
Protected Async subscribeChanges
subscribeChanges()
Returns : any
Protected Async withReadRetry
withReadRetry<T>(operation: () => void)
Type parameters :
  • T

Hook wrapping an idempotent read (get/allDocs/query) so subclasses can transparently retry transient failures.

The base implementation runs the operation exactly once (no retry). RemotePouchDatabase overrides this to re-issue reads that fail with a transient network abort/timeout, which recovers e.g. a connection gone stale after the tab was suspended. Only reads use this hook — writes must run exactly once and are never routed through it, except for creating a query index, which changes nothing when repeated.

Parameters :
Name Type Optional
operation function No
Returns : Promise<T>
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

supportsFind
supportsFind()
Inherited from Database

Whether find (bookmark-based pagination) is supported by this implementation. Callers should fall back to loading all docs otherwise.

Returns : boolean
import { Database, GetAllOptions, GetOptions, QueryOptions } from "../database";
import { Logging } from "../../logging/logging.service";
import { DatabaseException } from "./database-exception";
import PouchDB from "pouchdb-browser";
import indexeddbAdapter from "pouchdb-adapter-indexeddb";
import pouchdbFind from "pouchdb-find";
import { NgZone } from "@angular/core";
import { PerformanceAnalysisLogging } from "../../../utils/performance-analysis-logging";
import { firstValueFrom, Observable, Subject } from "rxjs";
import { HttpStatusCode } from "@angular/common/http";
import { environment } from "environments/environment";
import { SyncState } from "app/core/session/session-states/sync-state.enum";
import { SyncStateSubject } from "app/core/session/session-type";
import { NotificationEvent } from "#src/app/features/notification/model/notification-event";
// type-only, so the database layer gains no runtime dependency on analytics
import type { AnalyticsService } from "../../analytics/analytics.service";
import {
  FindPage,
  IndexedQuery,
  sortedQueries,
  typeSelector,
} from "./find-queries";

// Register the newer "indexeddb" adapter alongside the default "idb" adapter
PouchDB.plugin(indexeddbAdapter);
PouchDB.plugin(pouchdbFind);

/** Marks a bookmark of the second query (see {@link PouchDatabase.find}). */
const SECOND_QUERY_BOOKMARK = "q2:";

function markAsSecondQuery(res: FindPage): FindPage {
  return {
    docs: res.docs,
    bookmark: SECOND_QUERY_BOOKMARK + (res.bookmark ?? ""),
  };
}

/**
 * What happened to a document whose update was rejected as a conflict
 * (see {@link PouchDatabase.resolveConflict}).
 */
export type ConflictOutcome = "merged" | "overwritten" | "unresolved";

/**
 * Wrapper for a PouchDB instance to decouple the code from
 * that external library.
 *
 * Additional convenience functions on top of the PouchDB API
 * should be implemented in the abstract {@link Database}.
 */
export class PouchDatabase extends Database {
  /**
   * The reference to the PouchDB instance
   * @private
   */
  protected pouchDB: PouchDB.Database;

  /**
   * A list of promises that resolve once all the (until now saved) indexes are created
   * @private
   */
  protected indexPromises: Promise<any>[] = [];

  /**
   * 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
   * @private
   */
  protected changesFeed: Subject<any>;

  protected databaseInitialized = new Subject<void>();

  /** trigger to unsubscribe any internal subscriptions */
  protected readonly destroy$ = new Subject<void>();

  /**
   * 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.
   */
  adapter: string = "indexeddb";

  /**
   * Optional accessor for usage analytics (see {@link 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.
   */
  analytics?: () => Promise<AnalyticsService | null>;

  constructor(
    dbName: string,
    protected globalSyncState?: SyncStateSubject,
    protected ngZone?: NgZone,
  ) {
    super(dbName);
  }

  /**
   * Initialize the PouchDB with the IndexedDB/in-browser adapter (default).
   * See {link https://github.com/pouchdb/pouchdb/tree/master/packages/node_modules/pouchdb-browser}
   * @param dbName the name for the database under which the IndexedDB entries will be created
   * @param options PouchDB options which are directly passed to the constructor
   * @param suppressSyncCompleted whether to skip emitting a SyncState.COMPLETED event to the globalSyncState (because other logic for sync is building on top of this)
   */
  init(
    dbName?: string,
    options?: PouchDB.Configuration.DatabaseConfiguration | any,
    suppressSyncCompleted?: boolean,
  ) {
    this.pouchDB = new PouchDB(dbName ?? this.dbName, {
      adapter: this.adapter,
      ...options,
    });
    this.databaseInitialized.complete();

    if (!suppressSyncCompleted) {
      this.globalSyncState?.next(SyncState.COMPLETED);
    }
  }

  override isInitialized(): boolean {
    return !!this.pouchDB;
  }

  async getPouchDBOnceReady(): Promise<PouchDB.Database> {
    await firstValueFrom(this.databaseInitialized, {
      defaultValue: this.pouchDB,
    });
    return this.pouchDB;
  }

  /**
   * Get the actual instance of the PouchDB
   */
  getPouchDB(): PouchDB.Database {
    return this.pouchDB;
  }

  /**
   * Hook wrapping an idempotent read (`get`/`allDocs`/`query`) so subclasses
   * can transparently retry transient failures.
   *
   * The base implementation runs the operation exactly once (no retry).
   * {@link RemotePouchDatabase} overrides this to re-issue reads that fail with
   * a transient network abort/timeout, which recovers e.g. a connection gone
   * stale after the tab was suspended. Only reads use this hook — writes must
   * run exactly once and are never routed through it, except for creating a
   * query index, which changes nothing when repeated.
   */
  protected async withReadRetry<T>(operation: () => Promise<T>): Promise<T> {
    return operation();
  }

  /**
   * Load a single document by id from the database.
   * (see {@link Database})
   * @param id The primary key of the document to be loaded
   * @param options Optional PouchDB options for the request
   * @param returnUndefined (Optional) return undefined instead of throwing error if doc is not found in database
   */
  async get(
    id: string,
    options: GetOptions = {},
    returnUndefined?: boolean,
  ): Promise<any> {
    try {
      return await this.withReadRetry(async () =>
        (await this.getPouchDBOnceReady()).get(id, options),
      );
    } catch (err) {
      if (err.status === 404) {
        Logging.debug("Doc not found in database: " + id);
        if (returnUndefined) {
          return undefined;
        }
      }

      throw new DatabaseException(err, id);
    }
  }

  /**
   * Load all documents (matching the given PouchDB options) from the database.
   * (see {@link 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.
   *
   * @param options PouchDB options object as in the normal PouchDB library
   */
  async allDocs(options?: GetAllOptions) {
    try {
      const result = await this.withReadRetry(async () =>
        (await this.getPouchDBOnceReady()).allDocs(options),
      );
      return result.rows.map((row) => row.doc);
    } catch (err) {
      throw new DatabaseException(
        err,
        "allDocs; startkey: " + options?.["startkey"],
      );
    }
  }

  /**
   * Save a document to the database.
   * (see {@link Database})
   *
   * @param object The document to be saved
   * @param forceOverwrite (Optional) Whether conflicts should be ignored and an existing conflicting document forcefully overwritten.
   */
  async put(object: any, forceOverwrite = false): Promise<any> {
    if (forceOverwrite) {
      object._rev = undefined;
    }

    try {
      return await (await this.getPouchDBOnceReady()).put(object);
    } catch (err) {
      if (err.status === 409) {
        return this.resolveConflict(object, forceOverwrite, err);
      } else {
        throw new DatabaseException(err, object._id);
      }
    }
  }

  /**
   * Save an array of documents to the database
   * @param objects the documents to be saved
   * @param forceOverwrite whether conflicting versions should be overwritten
   * @returns array with the result for each object to be saved, if any item fails to be saved, this returns a rejected Promise.
   *          The save can partially fail and return a mix of success and error states in the array (e.g. `[{ ok: true, ... }, { error: true, ... }]`)
   */
  async putAll(objects: any[], forceOverwrite = false): Promise<any> {
    if (forceOverwrite) {
      objects.forEach((obj) => (obj._rev = undefined));
    }

    const pouchDB = await this.getPouchDBOnceReady();
    const results = await pouchDB.bulkDocs(objects);

    for (let i = 0; i < results.length; i++) {
      // Check if document update conflicts happened in the request
      const result = results[i] as PouchDB.Core.Error;
      if (result.status === 409) {
        results[i] = await this.resolveConflict(
          objects.find((obj) => obj._id === result.id),
          forceOverwrite,
          result,
        ).catch((e) => {
          if (e?.status === HttpStatusCode.Conflict) {
            // counted by reportConflict and still returned to the caller below,
            // so alerting on it here would only repeat what the single-document
            // path deliberately stopped reporting
            Logging.debug("could not resolve conflict during putAll", e);
          } else {
            Logging.warn("error during putAll", e, {
              // entity types only - a list of document IDs has no place in
              // remote monitoring (see #4174)
              entityTypes: [
                ...new Set(objects.map((x) => x._id?.split(":")[0])),
              ],
            });
          }
          return new DatabaseException(e);
        });
      }
    }

    if (results.some((r) => r instanceof Error)) {
      return Promise.reject(results);
    }
    return results;
  }

  /**
   * Delete a document from the database
   * (see {@link Database})
   *
   * @param object The document to be deleted (usually this object must at least contain the _id and _rev)
   */
  remove(object: any) {
    return this.getPouchDBOnceReady()
      .then((pouchDB) => pouchDB.remove(object))
      .catch((err) => {
        throw new DatabaseException(err, object["_id"]);
      });
  }

  /**
   * Permanently purge a document and all its revisions from the local database.
   *
   * Unlike {@link 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).
   *
   * @param id The document ID to purge
   * @returns true if the document was purged,
   *          false if the document did not exist locally (benign — desired state already achieved)
   */
  override async purge(id: string): Promise<boolean> {
    const db = await this.getPouchDBOnceReady();
    let localDoc: PouchDB.Core.IdMeta & PouchDB.Core.GetMeta;
    try {
      localDoc = await db.get(id, { conflicts: true });
    } catch (err) {
      if (err.status === 404) {
        return false;
      }
      throw err;
    }

    // Purge the winning revision and any conflicting leaf revisions
    // so the document is fully removed from local storage.
    const revsToPurge = [
      localDoc._rev,
      ...((localDoc as any)._conflicts ?? []),
    ];
    for (const rev of revsToPurge) {
      await (db as any).purge(id, rev);
    }

    // PouchDB purge() does not emit change events, so manually notify
    // the changes feed so in-memory caches drop the purged entity.
    if (this.changesFeed) {
      const deletionEvent = { _id: id, _rev: localDoc._rev, _deleted: true };
      if (this.ngZone) {
        this.ngZone.run(() => this.changesFeed.next(deletionEvent));
      } else {
        this.changesFeed.next(deletionEvent);
      }
    }

    return true;
  }

  /**
   * Check if a database is new/empty.
   * Returns true if there are no documents in the database
   */
  isEmpty(): Promise<boolean> {
    return this.getPouchDBOnceReady()
      .then((pouchDB) => pouchDB.info())
      .then((res) => res.doc_count === 0);
  }

  /**
   * Listen to changes to documents in the database.
   * Use rxjs operators to filter for specific prefixes etc. if needed.
   * @returns observable which emits the filtered changes
   */
  changes(): Observable<any> {
    if (!this.changesFeed) {
      this.changesFeed = new Subject();

      // trigger subscription only once DB ready, to go to the right instance (e.g. remote only)
      this.getPouchDBOnceReady().then(() => this.subscribeChanges());
    }
    return this.changesFeed;
  }

  protected async subscribeChanges() {
    const runSubscription = async () => {
      const db = await this.getPouchDBOnceReady();
      db.changes({
        live: true,
        since: "now",
        include_docs: true,
      })
        .addListener("change", (change) => {
          // Emit changes inside Angular zone to trigger change detection
          if (this.ngZone) {
            this.ngZone.run(() => this.changesFeed.next(change.doc));
          } else {
            this.changesFeed.next(change.doc);
          }
        })
        .catch((err) => {
          if (
            err.statusCode === HttpStatusCode.Unauthorized ||
            err.statusCode === HttpStatusCode.GatewayTimeout
          ) {
            Logging.warn(err);
          } else {
            Logging.error(err);
          }

          // retry
          setTimeout(() => this.subscribeChanges(), 10000);
        });
    };

    // run PouchDB change listener outside Angular zone to avoid excessive change detection
    if (this.ngZone) {
      this.ngZone.runOutsideAngular(() => runSubscription());
    } else {
      runSubscription();
    }
  }

  /**
   * Destroy the database and all saved data
   */
  async destroy(): Promise<any> {
    this.destroy$.next();

    await Promise.all(this.indexPromises);
    if (this.pouchDB) {
      return this.pouchDB.destroy();
    }
  }

  /**
   * Reset the database state so a new one can be opened.
   */
  async reset() {
    this.destroy$.next();

    this.pouchDB = undefined;
    // keep this.changesFeed because some services are already subscribed to this reference
    this.databaseInitialized = new Subject();
  }

  /**
   * Query a page of one entity type's documents with a Mango selector,
   * optionally sorted (see {@link Database.find}).
   *
   * Only available on databases that implement {@link 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 {@link sortedQueries}).
   */
  async find(
    prefix = "",
    query: PouchDB.Find.Selector = {},
    page?: { limit?: number; bookmark?: string },
    sort?: { prop?: string; dir?: "asc" | "desc" },
  ): Promise<FindPage> {
    if (!this.findPage) {
      throw new Error(
        "find() is not supported by this database, as it cannot page through query results",
      );
    }

    if (!sort?.prop) {
      // without a sort, CouchDB's built-in `_id` index serves the type range,
      // which (unlike an index on the sort property) skips no docs - so neither
      // an index of our own nor a second query is needed
      return this.findPage(
        { selector: { ...query, ...typeSelector(prefix) } },
        page,
      );
    }

    const [first, second] = sortedQueries(prefix, query, {
      prop: sort.prop,
      dir: sort.dir,
    });
    return this.findInSequence(first, second, page);
  }

  /**
   * Run a single Mango query for one page of results.
   * Implemented only by the databases that support {@link find}.
   */
  protected findPage?(
    findOptions: PouchDB.Find.FindRequest<any>,
    page?: { limit?: number; bookmark?: string },
  ): Promise<FindPage>;

  /**
   * Page through two queries one after the other, as if they were a single one.
   *
   * The bookmark stays opaque to callers: a bookmark of the second query is
   * marked with a prefix, so that the next page continues there.
   */
  private async findInSequence(
    first: IndexedQuery,
    second: IndexedQuery,
    page?: { limit?: number; bookmark?: string },
  ): Promise<FindPage> {
    const bookmark = page?.bookmark ?? "";
    if (bookmark.startsWith(SECOND_QUERY_BOOKMARK)) {
      const secondPage = await this.findIndexed(second, {
        limit: page.limit,
        bookmark: bookmark.slice(SECOND_QUERY_BOOKMARK.length),
      });
      return markAsSecondQuery(secondPage);
    }

    const firstPage = await this.findIndexed(first, page);
    if (page?.limit !== undefined && firstPage.docs.length >= page.limit) {
      // the first query may have more: continue there on the next page
      return firstPage;
    }

    // the first query is exhausted: fill up the page from the second one
    const secondPage = await this.findIndexed(second, {
      limit:
        page?.limit === undefined
          ? undefined
          : page.limit - firstPage.docs.length,
    });
    return markAsSecondQuery({
      docs: [...firstPage.docs, ...secondPage.docs],
      bookmark: secondPage.bookmark,
    });
  }

  /** Create a query's index (unless it exists already) and run the query through it. */
  private async findIndexed(
    { index, findOptions }: IndexedQuery,
    page?: { limit?: number; bookmark?: string },
  ): Promise<FindPage> {
    const pouchDB = await this.getPouchDBOnceReady();
    // Unlike other writes this is safe to retry like a read: creating an index
    // that already exists changes nothing (CouchDB answers "exists").
    const created = await this.withReadRetry(() =>
      pouchDB.createIndex({ index }),
    ).catch((err) => {
      throw new DatabaseException(err);
    });
    // the installed @types/pouchdb-find does not declare `id`
    return this.findPage({ ...findOptions, use_index: created["id"] }, page);
  }

  /**
   * Query data from the database based on a more complex, indexed request.
   * (see {@link Database})
   *
   * This is directly calling the PouchDB implementation of this function.
   * Also see the documentation there: {@link https://pouchdb.com/api.html#query_database}
   *
   * @param fun The name of a previously saved database index
   * @param options Additional options for the query, like a `key`. See the PouchDB docs for details.
   */
  query(
    fun: string | ((doc: any, emit: any) => void),
    options: QueryOptions,
  ): Promise<any> {
    return this.withReadRetry(() =>
      this.getPouchDBOnceReady().then((pouchDB) => pouchDB.query(fun, options)),
    ).catch((err) => {
      throw new DatabaseException(
        err,
        typeof fun === "string" ? fun : undefined,
      );
    });
  }

  /**
   * Create a database index to `query()` certain data more efficiently in the future.
   * (see {@link Database})
   *
   * Also see the PouchDB documentation regarding indices and queries: {@link https://pouchdb.com/api.html#query_database}
   *
   * @param designDoc The PouchDB style design document for the map/reduce query
   */
  saveDatabaseIndex(designDoc: any): Promise<void> {
    const creationPromise = this.createOrUpdateDesignDoc(designDoc).catch(
      (err) => {
        Logging.debug(
          "Could not create/update database index (may be expected in online-only mode)",
          err,
        );
      },
    );
    this.indexPromises.push(creationPromise);
    return creationPromise;
  }

  private async createOrUpdateDesignDoc(designDoc): Promise<void> {
    designDoc.aam_version = environment.appVersion;

    const existingDesignDoc = await this.get(designDoc._id, {}, true);
    if (!existingDesignDoc) {
      Logging.debug("creating new database index");
    } else if (
      JSON.stringify(existingDesignDoc.views) ===
      JSON.stringify(designDoc.views)
    ) {
      // already up to date, nothing more to do
      return;
    } else if (this.shouldSkipIndexUpdate(existingDesignDoc)) {
      return;
    } else {
      Logging.debug("replacing existing database index");
      designDoc._rev = existingDesignDoc._rev;
    }

    await this.put(designDoc);

    // for faster initial loading we disable prebuilding views in development
    // TODO: check if this should be completely removed, also for production systems
    if (environment.production) {
      await this.prebuildViewsOfDesignDoc(designDoc);
    }
  }

  /**
   * Whether to skip updating a design doc that differs from the local version.
   * Overridden in RemotePouchDatabase to prevent older clients from
   * overwriting indexes created by a newer app version on a shared server.
   */
  protected shouldSkipIndexUpdate(_existingDesignDoc: any): boolean {
    return false;
  }

  @PerformanceAnalysisLogging
  private async prebuildViewsOfDesignDoc(designDoc: any): Promise<void> {
    for (const viewName of Object.keys(designDoc.views)) {
      const queryName = designDoc._id.replace(/_design\//, "") + "/" + viewName;
      await this.query(queryName, { key: "1" });
    }
  }

  /**
   * Attempt to intelligently resolve conflicting document versions automatically.
   * @param newObject
   * @param overwriteChanges
   * @param existingError
   */
  private async resolveConflict(
    newObject: any,
    overwriteChanges = false,
    existingError: any = {},
  ): Promise<any> {
    const existingObject = await this.get(newObject._id);
    const resolvedObject = this.mergeObjects(existingObject, newObject);
    if (resolvedObject) {
      Logging.debug(
        "resolved document conflict automatically (" + resolvedObject._id + ")",
      );
      // counted only once the write actually went through - see reportConflict
      const result = await this.put(resolvedObject);
      this.reportConflict("merged", newObject._id);
      return result;
    } else if (overwriteChanges) {
      Logging.debug(
        "overwriting conflicting document version (" + newObject._id + ")",
      );
      newObject._rev = existingObject._rev;
      const result = await this.put(newObject);
      this.reportConflict("overwritten", newObject._id);
      return result;
    } else {
      this.reportConflict("unresolved", newObject._id);
      // the document's ID is passed as entityId rather than appended to the
      // message: remote monitoring groups by message, so an ID in there would
      // fragment one recurring problem into a separate issue per document
      existingError.message = `${existingError.message} (unable to resolve)`;
      throw new DatabaseException(existingError, newObject._id);
    }
  }

  /**
   * Count a document update conflict in usage statistics.
   *
   * Two users editing the same record is normal operation for an offline-first
   * app, not a fault, so conflicts are counted rather than reported as errors:
   * how often they happen (and for which record types) is the signal worth
   * having, while an alert per occurrence only buries the failures that do need
   * attention - which is what happened in AAM-DIGITAL-77H. A user whose save was
   * rejected is told so by the form that attempted it, so nothing is silently
   * swallowed here.
   *
   * Only the entity type is reported, never the document ID: the point is to
   * see which record types collide, and record identifiers have no place in
   * monitoring (see #4174).
   *
   * Deliberately not awaited by its callers, and total: counting a conflict must
   * neither delay a save nor be the reason one fails.
   *
   * `merged` and `overwritten` are reported only after the write that resolved
   * the conflict succeeded. That write can conflict again (another writer got in
   * between), in which case the retry reports its own outcome - so reporting up
   * front would count a save that never happened, and count one conflict twice.
   */
  private reportConflict(outcome: ConflictOutcome, docId: string) {
    this.analytics?.()
      .then((analytics) =>
        analytics?.eventTrack(outcome, {
          category: "document_update_conflict",
          label: docId?.split(":")[0],
        }),
      )
      .catch((err) =>
        Logging.debug("could not report document update conflict", err),
      );
  }

  private mergeObjects(_existingObject: any, _newObject: any) {
    // TODO: implement automatic merging of conflicting entity versions
    return undefined;
  }

  /**
   * Check if this is a notifications database based on the database name.
   * (may be used for special handling of notification DBs)
   */
  protected isNotificationsDatabase(): boolean {
    return this.dbName?.startsWith(NotificationEvent.DATABASE) ?? false;
  }
}

export { DatabaseException } from "./database-exception";

results matching ""

    No results matching ""