src/app/core/database/pouchdb/remote-pouch-database.ts
An alternative implementation of PouchDatabase that directly makes HTTP requests to a remote CouchDB.
No results matching.
Properties |
|
Methods |
|
constructor(dbName: string, authService: KeycloakAuthService, globalSyncState?: SyncStateSubject, ngZone?: NgZone, alertService?: AlertService)
|
||||||||||||||||||
|
Parameters :
|
| 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 |
| 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 |
| Protected Async findPage | |||||||||
findPage(findOptions: PouchDB.Find.FindRequest
|
|||||||||
|
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
Parameters :
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 :
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 :
Returns :
Promise<any>
|
| Async putAll | ||||||||||||
putAll(objects: any[], forceOverwrite: unknown)
|
||||||||||||
|
Inherited from
Database
|
||||||||||||
|
Parameters :
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 :
Returns :
Promise<any>
|
| Protected shouldSkipIndexUpdate | ||||||
shouldSkipIndexUpdate(existingDesignDoc: any)
|
||||||
|
Inherited from
PouchDatabase
|
||||||
|
Parameters :
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 :
|
||||||
|
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 :
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 :
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 When sorted, documents without a value for the sort property are included as well: last for "asc", first for "desc" (see sortedQueries).
Parameters :
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 :
Returns :
Promise<any>
|
| getPouchDB |
getPouchDB()
|
|
Inherited from
PouchDatabase
|
|
Get the actual instance of the PouchDB
Returns :
PouchDB.Database
|
| Async getPouchDBOnceReady |
getPouchDBOnceReady()
|
|
Inherited from
PouchDatabase
|
|
Returns :
Promise<PouchDB.Database>
|
| 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
Parameters :
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 :
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 Also see the PouchDB documentation regarding indices and queries: https://pouchdb.com/api.html#query_database
Parameters :
Returns :
Promise<void>
|
| getAll | ||||||||||
getAll(prefix: string)
|
||||||||||
|
Inherited from
Database
|
||||||||||
|
Load all documents (with the given prefix) from the database.
Parameters :
Returns :
Promise<Array<any>>
|
| Async removeAll | ||||||||
removeAll<T>(shouldDelete: (doc: T) => void)
|
||||||||
|
Inherited from
Database
|
||||||||
Type parameters :
|
||||||||
|
Delete all documents of the database matching the given filter. Note that
Parameters :
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)`,
);
}
}