Remove rxjs based caching modules for now

This commit is contained in:
guerler
2022-02-10 11:58:38 -05:00
parent 0bf29229a2
commit 205af02fb8
32 changed files with 10 additions and 2692 deletions
@@ -15,7 +15,6 @@
import DscUI from "./DscUI";
import { DatasetCollection } from "../../model";
import { deleteDatasetCollection, updateContentFields } from "../../model/queries";
import { cacheContent } from "components/providers/History/caching";
export default {
components: {
@@ -33,13 +32,7 @@ export default {
async onDelete(flags = {}) {
const { recursive = false, purge = false } = flags;
const collection = this.item;
const result = await deleteDatasetCollection(collection, recursive, purge);
if (result.deleted) {
const newFields = Object.assign(collection, {
isDeleted: result.deleted,
});
await cacheContent(newFields);
}
await deleteDatasetCollection(collection, recursive, purge);
},
async onUnhide() {
await this.onUpdate({ visible: true });
@@ -48,8 +41,7 @@ export default {
await this.onUpdate({ deleted: false });
},
async onUpdate(changes) {
const newContent = await updateContentFields(this.item, changes);
await cacheContent(newContent);
await updateContentFields(this.item, changes);
},
},
};
@@ -19,7 +19,6 @@
</template>
<script>
import { deleteDatasetCollection, updateContentFields } from "../../model/queries";
import { cacheContent } from "components/providers/History/caching";
import { DatasetCollection } from "../../model/DatasetCollection";
import DscUI from "components/History/ContentItem/DatasetCollection/DscUI";
import { DatasetCollectionContentProvider } from "components/providers";
@@ -57,41 +56,17 @@ export default {
},
async onDelete(collection, flags = {}) {
const { recursive = false, purge = false } = flags;
const result = await deleteDatasetCollection(collection, recursive, purge);
if (result.deleted) {
const newFields = Object.assign(collection, {
isDeleted: result.deleted,
});
await cacheContent(newFields);
}
await deleteDatasetCollection(collection, recursive, purge);
},
async onUndelete(collection) {
const result = await updateContentFields(collection, { deleted: false });
if (result.deleted === false) {
const newFields = Object.assign(collection, {
isDeleted: result.deleted,
});
await cacheContent(newFields);
}
await updateContentFields(collection, { deleted: false });
},
async onHide(collection) {
const result = await updateContentFields(collection, { visible: false });
if (result.visible === false) {
const newFields = Object.assign(collection, {
isVisible: result.visible,
});
await cacheContent(newFields);
}
await updateContentFields(collection, { visible: false });
},
async onUnhide(collection) {
const result = await updateContentFields(collection, { visible: true });
if (result.visible) {
const newFields = Object.assign(collection, {
isVisible: result.visible,
});
await cacheContent(newFields);
}
await updateContentFields(collection, { visible: true });
},
},
};
@@ -18,7 +18,6 @@
import DatasetUI from "components/History/ContentItem/Dataset/DatasetUI";
import { Dataset } from "../../model";
import { deleteContent, updateContentFields } from "../../model/queries";
import { cacheContent } from "components/providers/History/caching";
export default {
components: {
@@ -44,8 +43,7 @@ export default {
this.expanded = !this.expanded;
},
async onDelete(opts = {}) {
const ajaxResult = await deleteContent(this.dataset, opts);
await cacheContent(ajaxResult);
await deleteContent(this.dataset, opts);
},
async onUnhide() {
await this.onUpdate({ visible: true });
@@ -54,8 +52,7 @@ export default {
await this.onUpdate({ deleted: false });
},
async onUpdate(changes = {}) {
const newContent = await updateContentFields(this.dataset, changes);
await cacheContent(newContent);
await updateContentFields(this.dataset, changes);
},
},
};
@@ -15,7 +15,6 @@ cache first to see if we already have the data -->
<script>
import { Dataset } from "../model";
import { getContentByTypeId, cacheContent } from "components/providers/History/caching";
import { getContentDetails } from "../model/queries";
import DatasetUI from "./Dataset/DatasetUI";
@@ -45,15 +44,8 @@ export default {
this.$emit("update:expand", val);
},
async loadDetails() {
const { type_id, history_id } = this.item;
const existingDataset = await getContentByTypeId(history_id, type_id);
if (existingDataset) {
this.localItem = existingDataset;
} else {
const loadedDataset = await getContentDetails(this.item);
await cacheContent(loadedDataset);
this.localItem = loadedDataset;
}
await getContentDetails(this.item);
this.localItem = loadedDataset;
},
},
};
@@ -168,7 +168,6 @@ import {
purgeAllDeletedContent,
} from "./model";
import { createDatasetCollection } from "./model/queries";
import { cacheContent } from "components/providers/History/caching";
import { legacyNavigationMixin } from "components/plugins/legacyNavigation";
import { buildCollectionModal } from "./adapters/buildCollectionModal";
import ContentFilters from "./ContentFilters";
@@ -283,14 +282,8 @@ export default {
const modalResult = await buildCollectionModal(collectionTypeCode, this.history.id, this.contentSelection);
const newCollection = await createDatasetCollection(this.history, modalResult);
// cache the collection
await cacheContent(newCollection);
// have to hide the source items if that was requested
if (modalResult.hide_source_items) {
this.contentSelection.forEach(async (dataset) => {
await cacheContent({ ...dataset, visible: false }, true);
});
this.$emit("resetSelection");
}
},
@@ -1,29 +0,0 @@
/**
* Separated api from Worker definition so that we can use it without the worker
* if desired.
*
* Note that database events are only consistent within the same database
* instance, so a cache Event from outside the worker will not be heard by the
* watchers inside the worker because those are actually 2 separate db
* instances.
*/
import { content$, dscContent$ } from "./db/observables";
import { monitorQuery } from "./db/monitorQuery";
export * from "./db/promises";
export { loadDscContent } from "./loadDscContent";
export { loadHistoryContents, clearHistoryDateStore } from "./loadHistoryContents";
export { monitorHistoryContent, monitorCollectionContent } from "./monitorHistoryContent";
export { wipeDatabase } from "./db/wipeDatabase";
// generic content query monitor
export const monitorContentQuery = (cfg = {}) => {
return monitorQuery({ db$: content$, ...cfg });
};
// generic collection content monitor
export const monitorDscQuery = (cfg = {}) => {
return monitorQuery({ db$: dscContent$, ...cfg });
};
@@ -1,10 +0,0 @@
// For testing, bypass the worker convrsions
// and wrap each function in a jest mock wrapper
const cacheapi = require("../CacheApi");
const mockedModule = Object.entries(cacheapi).reduce((mod, [fnName, orig]) => {
return { ...mod, [fnName]: jest.fn(orig) };
}, {});
module.exports = mockedModule;
@@ -1,48 +0,0 @@
/**
* Worker client utility
*
* asObservable takes an observable transformation and persists its state while
* it is running inside the worker. Normally when threads.js invokes an
* observable, it runs once, its subscription ends and subsequent calls to that
* same exposed function will start a new subscription.
*
* With asObservable, a subject is created inside the worker and subsequent
* method calls put their new value on that Subject, which is connected to the
* wrapped observable. Therefore the observable state is preserved until it is
* explicitly unsubscribed from the outside, and distinct(), scan(), or other
* stateful observable operators will continue to work.
*/
import { Subject, of } from "rxjs";
import { startWith } from "rxjs/operators";
// prettier-ignore
export const asObservable = (operator) => {
const currentSubs = new Map();
// process notifications
// materialize exposed a "kind" variable for all observable messages,
// it's either N,C,E for next, complete, error
return ({ id, cfg = {}, value, kind }) => {
if (kind == "N") {
if (!currentSubs.has(id)) {
const input$ = new Subject();
const output$ = input$.pipe(startWith(value), operator(cfg));
currentSubs.set(id, { input$, output$ });
return output$;
}
const sub = currentSubs.get(id);
sub.input$.next(value);
}
if (kind == "C" || kind == "E") {
const sub = currentSubs.get(id);
if (sub) {
sub.input$.complete();
currentSubs.delete(id);
}
}
return of(null);
};
};
@@ -1,28 +0,0 @@
/**
* Indices are required by pouchdb to use the pouuchdb.find functionality. These
* are the indices we keep on the main history content database
*/
export const contentIndices = [
{
index: {
fields: [{ cached_at: "desc" }],
},
name: "by cache time",
ddoc: "idx-content-history-cached_at",
},
];
/**
* ...and the collection contents
*/
export const dscIndices = [
{
index: {
fields: [{ cached_at: "desc" }],
},
name: "by cache time",
ddoc: "idx-dsc-contents-cached_at",
},
];
@@ -1,45 +0,0 @@
import { pipe, Observable } from "rxjs";
import { switchMap, filter, share } from "rxjs/operators";
// feed observables, keyed by underlying database instance
export const feeds = new Map();
/**
* Returns an observable with all the change events from the indicated database
* @param {Observable} db$ Pouch database observable
*/
// prettier-ignore
export const changes = (cfg = {}) => {
return pipe(
switchMap((db) => {
if (!feeds.has(db)) {
feeds.set(db, buildFeed(db, cfg));
}
return feeds.get(db);
}),
// filter out index creation which can appear as a change
filter(({ id }) => !id.includes("_design")),
);
};
/**
* Creates an observable that emits changes to the indicated db instance
* @param {PouchDB} db
* @param {object} cfg
*/
const buildFeed = (db, cfg = {}) => {
const { live = true, returnDocs = true, include_docs = true, since = "now", timeout = false } = cfg;
const feed$ = new Observable((obs) => {
const changeOpts = { live, include_docs, returnDocs, since, timeout };
const feed = db.changes(changeOpts);
feed.on("change", (update) => obs.next(update));
feed.on("error", (err) => obs.error(err));
return () => {
feed.cancel();
feeds.delete(db);
};
});
return feed$.pipe(share());
};
@@ -1,172 +0,0 @@
import { timer } from "rxjs";
import { takeWhile, share, takeUntil } from "rxjs/operators";
import { content$, dscContent$ } from "./observables";
import { wipeDatabase } from "./wipeDatabase";
import {
bulkCacheContent,
cacheContent,
getCachedContent,
bulkCacheDscContent,
getCachedCollectionContent,
cacheCollectionContent,
} from "./promises";
import { changes, feeds } from "./changes";
import { wait } from "jest/helpers";
// test data
import historyContent from "components/providers/History/test/json/historyContent.json";
import collectionContent from "components/providers/History/test/json/collectionContent.json";
// https://github.com/hirezio/observer-spy/blob/master/README.md
import { ObserverSpy } from "@hirez_io/observer-spy";
beforeEach(wipeDatabase);
afterEach(wipeDatabase);
describe("changes operator", () => {
describe("history content (content$)", () => {
test("hears updates when we add stuff to the history content cache", async () => {
// listen to changes until dumb fake prop appears
const update$ = content$.pipe(
changes(),
takeWhile((change) => change.doc.floobar == undefined, true),
share()
);
// spy on emissions
const spy = new ObserverSpy();
update$.subscribe(spy);
// put some stuff in th cache
const cacheResults = await bulkCacheContent(historyContent);
expect(cacheResults.length).toEqual(historyContent.length);
// pull out one item, modify it
const testItem = await getCachedContent(cacheResults[0].id);
expect(testItem._id).toEqual(cacheResults[0].id);
testItem.floobar = Math.random();
const updateResult = await cacheContent(testItem);
expect(updateResult.updated).toBeTruthy();
// wait for updates to complete, which should happen since we added "floobar"
await spy.onComplete();
// last update should be the testItem we saved
const emits = spy.getValues();
const lastEmit = spy.getLastValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(lastEmit.id).toEqual(testItem._id);
expect(lastEmit.doc.floobar).toEqual(testItem.floobar);
expect(emits.length).toEqual(historyContent.length + 1);
// original insert + later update
const updates = emits.filter((update) => update.doc._id == testItem._id);
expect(updates.length).toEqual(2);
// only one should have the "floobar" field
const floobars = emits.filter((update) => update.doc.floobar !== undefined);
expect(floobars.length).toEqual(1);
});
});
describe("collection content (dscContent$)", () => {
test("hears updates when we add stuff to the collection content cache", async () => {
// listen to changes until dumb fake prop appears
const update$ = dscContent$.pipe(
changes(),
takeWhile((change) => change.doc.floobar == undefined, true),
share()
);
// spy on emissions
const spy = new ObserverSpy();
update$.subscribe(spy);
// preprocess collection content to add parent_url
const fakeParent = "/api/123/blah";
const processedCollection = collectionContent.map((doc) => {
doc.parent_url = fakeParent;
return doc;
});
// put some stuff in the cache
const cacheResults = await bulkCacheDscContent(processedCollection);
expect(cacheResults.length).toEqual(processedCollection.length);
// pull out one item, modify it
const testItem = await getCachedCollectionContent(cacheResults[0].id);
expect(testItem._id).toEqual(cacheResults[0].id);
testItem.floobar = Math.random();
const updateResult = await cacheCollectionContent(testItem);
expect(updateResult.updated).toBeTruthy();
// wait for updates to complete, which should happen since we added "floobar"
await spy.onComplete();
// last update should be the testItem we saved
const emits = spy.getValues();
const lastEmit = spy.getLastValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(lastEmit.id).toEqual(testItem._id);
expect(lastEmit.doc.floobar).toEqual(testItem.floobar);
expect(emits.length).toEqual(collectionContent.length + 1);
// original insert + later update
const updates = emits.filter((update) => update.doc._id == testItem._id);
expect(updates.length).toEqual(2);
// only one should have the "floobar" field
const floobars = emits.filter((update) => update.doc.floobar !== undefined);
expect(floobars.length).toEqual(1);
});
});
describe("change feed should be shared", () => {
const lifeTime = 300;
const spinUpTime = 100;
test("subscriging should show 1 feed, unsubscribing should show 0", async () => {
// subscribe to changes once
const feed$ = content$.pipe(changes(), takeUntil(timer(lifeTime)));
const spy = new ObserverSpy();
feed$.subscribe(spy);
// give it a little time to put itself together
await wait(spinUpTime);
expect(feeds.size).toEqual(1);
// complete, share() should now remove instance
await spy.onComplete();
expect(feeds.size).toEqual(0);
});
test("subscribing 2 times should result in one feed if it's the same DB", async () => {
// subscribe to changes onece
const feed$ = content$.pipe(changes(), takeUntil(timer(2 * lifeTime)));
const spy = new ObserverSpy();
feed$.subscribe(spy);
// give it a little time to put itself together
await wait(spinUpTime);
expect(feeds.size).toEqual(1);
// subscribe again
const feed2$ = content$.pipe(changes(), takeUntil(timer(lifeTime)));
const spy2 = new ObserverSpy();
feed2$.subscribe(spy2);
await wait(spinUpTime);
expect(feeds.size).toEqual(1);
// unsub from 2nd feed
await spy2.onComplete();
expect(feeds.size).toEqual(1);
// unsub from 1st feed
await spy.onComplete();
expect(feeds.size).toEqual(0);
});
});
});
@@ -1,251 +0,0 @@
import isPromise from "is-promise";
import { isObservable } from "rxjs";
import { wipeDatabase } from "./wipeDatabase";
import { content$, dscContent$, buildContentId, buildCollectionId, prepContent, prepDscContent } from "./observables";
import { firstValueFrom } from "utils/observable/firstValueFrom";
import {
cacheContent,
getCachedContent,
uncacheContent,
bulkCacheContent,
getContentByTypeId,
bulkCacheDscContent,
} from "./promises";
// test data
import historyContent from "components/providers/History/test/json/historyContent.json";
import collectionContent from "components/providers/History/test/json/collectionContent.json";
beforeEach(wipeDatabase);
afterEach(wipeDatabase);
describe("database observables", () => {
describe("content$", () => {
it("should be an observable", () => {
expect(content$).toBeDefined;
expect(isObservable(content$)).toBeTruthy;
});
it("should yield a pouchdb instance", async () => {
const dbPromise = firstValueFrom(content$);
expect(isPromise(dbPromise)).toBeTrue;
const db = await dbPromise;
expect(db).toBeDefined;
expect(db.constructor.name).toEqual("PouchDB");
});
it("should only build once", async () => {
const stuff = await firstValueFrom(content$);
const stuff2 = await firstValueFrom(content$);
const stuff3 = await firstValueFrom(content$);
expect(stuff).toStrictEqual(stuff2);
expect(stuff).toStrictEqual(stuff3);
});
});
describe("dscContent$", () => {
it("should be an observable", () => {
expect(dscContent$).toBeDefined;
expect(isObservable(dscContent$)).toBeTruthy;
});
it("should yield a pouchdb instance", async () => {
const dbPromise = firstValueFrom(dscContent$);
expect(isPromise(dbPromise)).toBeTrue;
const db = await dbPromise;
expect(db).toBeDefined;
expect(db.constructor.name).toEqual("PouchDB");
});
it("should only build once", async () => {
const stuff = await firstValueFrom(dscContent$);
const stuff2 = await firstValueFrom(dscContent$);
const stuff3 = await firstValueFrom(dscContent$);
expect(stuff).toStrictEqual(stuff2);
expect(stuff).toStrictEqual(stuff3);
});
});
});
describe("history content operators and functions", () => {
const testDoc = historyContent[0];
const testId = buildContentId(testDoc);
describe("buildContentId", () => {
it("should add the history and hid together", () => {
const id = buildContentId(testDoc);
expect(id).toContain(String(testDoc.history_id));
expect(id).toContain(String(testDoc.hid));
});
});
describe("prepContent", () => {
const doc = prepContent(testDoc);
it("should generate an _id for the content", () => {
expect(doc._id).toContain(String(testDoc.history_id));
expect(doc._id).toContain(String(testDoc.hid));
});
it("should rename the deleted field because pouchdb reserves that prop", () => {
expect(doc.deleted).toBeUndefined;
expect(doc.isDeleted).toBeDefined;
expect(doc.isDeleted).toEqual(testDoc.deleted);
});
});
describe("cacheContent", () => {
it("should cache a single document", async () => {
const cacheSummary = await cacheContent(testDoc);
expect(cacheSummary).toBeDefined;
expect(cacheSummary.id).toEqual(testId);
});
});
describe("getCachedContent", () => {
it("should lookup cached content", async () => {
const cacheSummary = await cacheContent(testDoc);
expect(cacheSummary).toBeDefined;
expect(cacheSummary.id).toEqual(testId);
const lookup = await getCachedContent(testId);
expect(lookup).toBeDefined;
expect(lookup._id).toEqual(testId);
});
});
describe("uncacheContent", () => {
it("should erase a previously cached item", async () => {
await cacheContent(testDoc);
const cachedDoc = await getCachedContent(testId);
const deleteResult = await uncacheContent(cachedDoc);
expect(deleteResult.ok).toBeTruthy;
const lookup = await getCachedContent(testId);
expect(lookup).toEqual(null);
});
});
describe("getContentByTypeId", () => {
it("should retrieve cached content using the type_id", async () => {
// cache some content
const { id } = await cacheContent(testDoc);
// look it up using _id
const { history_id, type_id } = await getCachedContent(id);
// lookup by type_id
const docByType = await getContentByTypeId(history_id, type_id);
expect(docByType.type_id).toBe(type_id);
});
});
describe("cacheContent and return doc", () => {
it("should cache raw props and return with an _id and a cached_at time", async () => {
const doc = await cacheContent(testDoc, true);
expect(doc._id).toEqual(buildContentId(testDoc));
expect(doc.cached_at).toExist;
expect(typeof doc.cached_at).toEqual("number");
expect(doc.cached_at).toBeGreaterThan(0);
});
});
describe("bulkCacheContent", () => {
it("should put a whole line of stuff into the cache", async () => {
const bulkResult = await bulkCacheContent(historyContent);
expect(bulkResult.length).toEqual(historyContent.length);
const cachedIds = bulkResult.map((item) => item.id);
const inputIds = historyContent.map((doc) => buildContentId(doc));
expect(cachedIds).toEqual(inputIds);
});
});
});
describe("collection content operators and functions", () => {
describe("prepDscContent", () => {
it("should process a raw doc (insert from ajax)", () => {
const raw = collectionContent[0];
raw.parent_url = "foo"; // required step
const dsc = prepDscContent(raw);
expect(dsc).toExist;
expect(dsc._id).toContain(raw.parent_url);
expect(dsc._id).toContain(String(raw.element_index));
});
// Warrants more testing because we have to untwist the api response format as
// well as working on existing (already converted) objects
test("should return valid results for pre-processed cache results (on update)", () => {
// simulate a collection that has already been cached
const updateMe = {
id: "fb85969571388350",
type_id: "dataset_collection-fb85969571388350",
model_class: "DatasetCollection",
history_content_type: "dataset_collection",
name: "ABC",
element_identifier: "ABC",
element_index: 10000,
element_type: "dataset_collection",
parent_url: "/abc/def/ghi",
contents_url: "/api/dataset_collections/5a1cff6882ddb5b2/contents/fb85969571388350",
collection_type: "paired",
cached_at: 1602265731626,
_id: "/abc/def/ghi-000000010000",
_rev: "6-377d7786587cd5e9e4524a736c228a90",
};
const processed = prepDscContent(updateMe);
expect(processed._id).toEqual(buildCollectionId(updateMe));
expect(processed.history_content_type).toBe("dataset_collection");
expect(processed.id).toEqual(updateMe.id);
expect(processed.type_id).toEqual(updateMe.type_id);
});
});
describe("bulkCacheDscContent", () => {
// we need to add parent_url to every item in the collection contents cache
// becahse the contents_url for the collection the contents came from is
// going to be the parent_url and part of the _id
const parent_url = "fakecollection";
const collectionData = collectionContent.map((doc) => {
doc.parent_url = parent_url;
return doc;
});
it("should put a whole line of stuff into the cache", async () => {
const bulkSummary = await bulkCacheDscContent(collectionData);
expect(bulkSummary.length).toEqual(collectionContent.length);
const cachedIds = bulkSummary.map((item) => item.id);
const inputIds = collectionData.map((doc) => buildCollectionId(doc));
expect(cachedIds).toEqual(inputIds);
});
it("should retrieve that stuff later", async () => {
const bulkSummary = await bulkCacheDscContent(collectionData);
expect(bulkSummary.length).toEqual(collectionContent.length);
});
});
});
describe("wipeDatabase", () => {
const testDoc = historyContent[0];
const testId = buildContentId(testDoc);
it("should clear out stuff so each test is clean", async () => {
const cacheSummary = await cacheContent(testDoc);
expect(cacheSummary.id).toEqual(testId);
// retrieve the cached doc to make sure it's in there
const retrievedDoc = await getCachedContent(testId);
expect(retrievedDoc).toExist;
expect(retrievedDoc._id).toEqual(testId);
// wipe
await wipeDatabase();
// try to retrieve again
const lookupAgain = await getCachedContent(testId);
expect(lookupAgain).toEqual(null);
});
});
@@ -1,42 +0,0 @@
import { pipe } from "rxjs";
import { mergeMap } from "rxjs/operators";
import { show, needs } from "utils/observable";
/**
* Pouchdb-find as an operator
* https://pouchdb.com/guides/mango-queries.html
*
* @param {Observable} db$ Observable pouchDb instance
*/
export const find = (db$, cfg = {}) => {
const { label = "find", debug = false } = cfg;
return pipe(
show(debug, (request) => console.log(`${label} -> request`, request)),
needs(db$),
mergeMap(async (inputs) => {
const [request, db] = inputs;
const { index } = request;
if (index !== undefined) {
const indexResponse = await db.createIndex({ index });
const { result: idxResult } = indexResponse;
if (idxResult !== "created" && idxResult !== "exists") {
throw new Error("Unknown index creation result", indexResponse);
}
}
let docs;
try {
const findResponse = await db.find(request);
docs = findResponse.docs || [];
} catch (err) {
console.warn("find() error", err, request, db.name);
throw err;
}
return docs;
}),
show(debug, (result) => console.log(`${label} -> result`, result))
);
};
@@ -1,368 +0,0 @@
import { Subject } from "rxjs";
import { take } from "rxjs/operators";
import { ObserverSpy } from "@hirez_io/observer-spy";
import { wipeDatabase } from "./wipeDatabase";
import { bulkCacheContent, bulkCacheDscContent } from "./promises";
import { content$, dscContent$, buildContentId } from "./observables";
import { find } from "./find";
// test data
import historyContent from "components/providers/History/test/json/historyContent.json";
import collectionContent from "components/providers/History/test/json/collectionContent.json";
beforeEach(wipeDatabase);
afterEach(wipeDatabase);
describe("find operator", () => {
// setup request subject, find observable, and spy to look at output
const request$ = new Subject();
// unsub, if requred
let obs$;
let spy;
describe("content database (content$) queries", () => {
beforeEach(async () => await bulkCacheContent(historyContent));
beforeEach(() => {
obs$ = request$.pipe(find(content$), take(1));
spy = new ObserverSpy();
obs$.subscribe(spy);
});
test("should select everything with a blank selector", async () => {
const request = { selector: {} };
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(results.length).toEqual(historyContent.length);
});
test("should return everything with a blank selector, while building a custom index", async () => {
const request = {
selector: {},
sort: [{ _id: "desc" }],
index: {
fields: ["_id"],
sort: [{ _id: "desc" }],
name: "foo index",
ddoc: "idx-foo",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(results.length).toEqual(historyContent.length);
});
test("should return with a limit", async () => {
const limit = 3;
const request = {
selector: {},
sort: [{ _id: "desc" }],
limit,
index: {
fields: ["_id"],
sort: [{ _id: "desc" }],
name: "foo index",
ddoc: "idx-foo",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(results.length).toEqual(limit);
});
test("descending history-hid query", async () => {
const { history_id, hid } = historyContent[0];
const request = {
selector: {
_id: { $lte: `${history_id}-${hid}` },
},
sort: [{ _id: "desc" }],
limit: 3,
index: {
fields: ["_id"],
sort: [{ _id: "desc" }],
name: "content history_id and hid descending",
ddoc: "idx-historyid-hid-desc",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(results.length).toEqual(3);
results.forEach((doc) => {
expect(doc.history_id).toEqual(history_id);
expect(doc.hid).toBeLessThanOrEqual(hid);
});
});
test("descending with visible flag", async () => {
const { history_id, hid } = historyContent[0];
const request = {
selector: {
_id: { $lte: `${history_id}-${hid}` },
visible: false,
},
sort: [{ _id: "desc" }],
limit: 3,
index: {
fields: ["_id"],
sort: [{ _id: "desc" }],
name: "content history_id and hid descending",
ddoc: "idx-historyid-hid-desc",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
// only put two visible items in the test data
const visibleRows = historyContent.filter((row) => row.visible == false);
expect(results.length).toEqual(visibleRows.length);
results.forEach((doc) => {
expect(doc.history_id).toEqual(history_id);
expect(doc.hid).toBeLessThanOrEqual(hid);
expect(doc.visible).toEqual(false);
});
});
test("descending with isDeleted flag", async () => {
const { history_id, hid } = historyContent[0];
const request = {
selector: {
_id: { $lte: `${history_id}-${hid}` },
isDeleted: true,
},
sort: [{ _id: "desc" }],
limit: 3,
index: {
fields: ["_id"],
sort: [{ _id: "desc" }],
name: "content history_id and hid descending",
ddoc: "idx-historyid-hid-desc",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
// console.log(results);
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
// only put two visible items in the test data
const undeletedRows = historyContent.filter((row) => {
return row.hid <= hid && row.deleted === true;
});
expect(results.length).toEqual(undeletedRows.length);
results.forEach((doc) => {
expect(doc.history_id).toEqual(history_id);
expect(doc.hid).toBeLessThanOrEqual(hid);
expect(doc.isDeleted).toEqual(true);
});
});
test("ascending history-hid query", async () => {
const doc = historyContent[5];
const { history_id, hid } = doc;
const targetHid = buildContentId(doc);
const request = {
selector: {
_id: { $gt: targetHid },
},
sort: [{ _id: "desc" }],
limit: 3,
index: {
fields: ["_id"],
sort: [{ _id: "desc" }],
name: "content history_id and hid descending",
ddoc: "idx-historyid-hid-desc",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
// no additional fitlers should show all undeleted and visible rows
const expectedRows = historyContent.filter((row) => {
return row.hid > hid && row.deleted === false && row.visible == true;
});
expect(results.length).toEqual(expectedRows.length);
results.forEach((doc) => {
expect(doc.history_id).toEqual(history_id);
expect(doc.hid).toBeGreaterThan(hid);
});
});
test("ascending with visible flag", async () => {
// start at the 5th row and look upwards
const doc = historyContent[5];
const { hid } = doc;
const targetId = buildContentId(doc);
const request = {
selector: {
_id: { $gt: targetId },
visible: false,
},
sort: [{ _id: "desc" }],
limit: 3,
index: {
fields: ["_id"],
sort: [{ _id: "desc" }],
name: "content history_id and hid descending",
ddoc: "idx-historyid-hid-desc",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
// count the invisible rows where the hid > targetId;
const expectedRows = historyContent.filter((row) => {
return row.hid > hid && row.visible === false;
});
expect(results.length).toEqual(expectedRows.length);
});
test("ascending with isDeleted flag", async () => {
const doc = historyContent[5];
const { history_id, hid } = doc;
const targetId = buildContentId(doc);
const request = {
selector: {
_id: { $gt: targetId },
isDeleted: true,
},
sort: [{ _id: "desc" }],
limit: 3,
index: {
fields: ["_id"],
sort: [{ _id: "desc" }],
name: "content history_id and hid descending",
ddoc: "idx-historyid-hid-desc",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
// only put one deleted near the top
const expectedRows = historyContent.filter((row) => {
return row.hid > hid && row.deleted === true;
});
expect(results.length).toEqual(expectedRows.length);
results.forEach((doc) => {
expect(doc.history_id).toEqual(history_id);
expect(doc.hid).toBeGreaterThan(hid);
expect(doc.isDeleted).toEqual(true);
});
});
});
describe("collection content database (dscContent$) queries", () => {
// put some stuff in the cache
// preprocess the collection content
const fakeParent = "/foo/bar";
const dscTestContent = collectionContent.map((props) => {
props.parent_url = fakeParent;
return props;
});
beforeEach(async () => await bulkCacheDscContent(dscTestContent));
beforeEach(() => {
obs$ = request$.pipe(find(dscContent$), take(1));
spy = new ObserverSpy();
obs$.subscribe(spy);
});
test("should select everything with a blank selector", async () => {
const request = { selector: {} };
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(results.length).toEqual(collectionContent.length);
});
test("ascending element_index", async () => {
const request = {
selector: {
_id: { $gte: `${fakeParent}-0` },
},
sort: [{ _id: "asc" }],
index: {
fields: ["_id"],
sort: [{ _id: "asc" }],
name: "ascending index",
ddoc: "idx-collection-ascending",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(results.length).toEqual(collectionContent.length);
});
test("ascending element_index with name regex", async () => {
const request = {
selector: {
_id: { $gte: `${fakeParent}-0` },
element_identifier: { $regex: /M117C1/i },
},
sort: [{ _id: "asc" }],
index: {
fields: ["_id"],
sort: [{ _id: "asc" }],
name: "ascending index",
ddoc: "idx-collection-ascending",
},
};
request$.next(request);
await spy.onComplete();
const results = spy.getFirstValue();
expect(spy.receivedNext()).toBe(true);
expect(spy.receivedComplete()).toBe(true);
expect(results.length).toBeLessThan(collectionContent.length);
});
});
});
@@ -1,16 +0,0 @@
export {
content$,
dscContent$,
getCachedContent,
cacheContent,
uncacheContent,
bulkCacheContent,
bulkCacheDscContent,
} from "./observables";
export { find } from "./find";
export { changes } from "./changes";
export { monitorQuery } from "./monitorQuery";
// utility function
export { wipeDatabase } from "./wipeDatabase";
@@ -1,56 +0,0 @@
import deepEqual from "deep-equal";
import { isObservable, concat, of } from "rxjs";
import { concatAll, map, switchMap, debounceTime, distinctUntilChanged, pluck, filter } from "rxjs/operators";
import { matchesSelector } from "pouchdb-selector-core";
import { find } from "./find";
import { changes } from "./changes";
import { show } from "utils/observable";
/**
* Turns a selector into a live cache result emitter. Emits objects with
* a document object and a match boolean flag ({ doc, match: true }) indicating
* whether or not this document matches the request$ selector
*
* @param {Observable} db$ Observable PouchDB instance
* @param {Observable} request$ Observable Pouchdb-find configuration
*/
// prettier-ignore
export const monitorQuery = (cfg = {}) => (request$) => {
const { db$, inputDebounce = 0, debug = false, label } = cfg;
if (!isObservable(db$)) {
throw new Error("Please pass a pouch database observable to monitorQuery");
}
const debouncedRequest$ = request$.pipe(
debounceTime(inputDebounce),
distinctUntilChanged(deepEqual),
);
return debouncedRequest$.pipe(
switchMap((request) => {
const { selector } = request;
const docMatch = doc => matchesSelector(doc, selector);
// do a search of the cache first
const initial$ = of(request).pipe(
find(db$, { debug, label }),
concatAll(),
map((doc) => ({ doc, match: true, initial: true })),
show(debug, (change) => console.log("monitorQuery: initial", change)),
);
// later changes
const updates$ = db$.pipe(
changes({ debug, label }),
pluck("doc"),
filter(Boolean),
map((doc) => ({ doc, match: docMatch(doc), update: true })),
show(debug, (change) => console.log("monitorQuery: update", change)),
);
return concat(initial$, updates$);
}),
show(debug, (change) => console.log("monitorQuery", change)),
);
};
@@ -1,398 +0,0 @@
import { of, timer } from "rxjs";
import { takeUntil } from "rxjs/operators";
import { ObserverSpy } from "@hirez_io/observer-spy";
import { wait } from "jest/helpers";
import { wipeDatabase } from "./wipeDatabase";
import { monitorQuery } from "./monitorQuery";
import { content$, dscContent$ } from "./observables";
import { cacheContent, cacheCollectionContent, bulkCacheContent, bulkCacheDscContent } from "./promises";
// test data
import historyContent from "components/providers/History/test/json/historyContent.json";
import collectionContent from "components/providers/History/test/json/collectionContent.json";
jest.mock("app");
jest.mock("../../caching");
const monitorSpinUp = 100;
const monitorSafetyTimeout = 800;
const pluckAll = (arr, prop) => arr.map((o) => o[prop]);
afterEach(wipeDatabase);
// prettier-ignore
describe("monitorQuery: history content", () => {
describe("initial results", () => {
test("first monitor event has initial results reflecting existing matches in db", async () => {
const cachedContent = await bulkCacheContent(historyContent, true);
const cachedContentIds = new Set(pluckAll(cachedContent, "_id"));
// create new monitor with query selector hid > 160
const cutoffHid = 120;
const selector = { hid: { $gt: cutoffHid } };
const monitor$ = of({ selector }).pipe(
monitorQuery({ db$: content$ }),
takeUntil(timer(monitorSafetyTimeout)),
);
// // listen to observable, wait for complete
const spy = new ObserverSpy();
monitor$.subscribe(spy);
await spy.onComplete();
// should get one emit out with just the intial matches
const emits = spy.getValues();
expect(spy.getValuesLength()).toBeGreaterThan(0);
expect(spy.getValuesLength()).toBeLessThanOrEqual(cachedContentIds.size);
emits.forEach(({ doc, match, initial }) => {
// we aren't changing anything, results should be initial matches
expect(match).toBeTruthy();
// initial results have an identifying flag
expect(initial).toBe(true);
// every doc should match the selector
expect(doc => doc.hid > cutoffHid).toBeTruthy();
// every doc should be in the initially cached items
expect(cachedContentIds.has(doc._id)).toBeTruthy();
});
});
});
describe("updates", () => {
test("INSERT: adding to cache after instantiation should emit a new doc", async () => {
const cutoffHid = 50;
const selector = { hid: { $gt: cutoffHid } };
const monitor$ = of({ selector }).pipe(
monitorQuery({ db$: content$ }),
takeUntil(timer(monitorSafetyTimeout)),
);
const spy = new ObserverSpy();
monitor$.subscribe(spy);
// wait for the monitor to spin-up so the insert doesn't look like the initial vals
await wait(monitorSpinUp);
// Insert a row, make sure it should match the selector
const insertedContent = await cacheContent(historyContent[0], true);
expect(insertedContent.hid > cutoffHid).toBeTruthy();
// wait for observable end
await spy.onComplete();
const emits = spy.getValues();
expect(spy.getValuesLength()).toEqual(1);
emits.forEach(({ doc, match, initial, update }) => {
expect(doc.hid > cutoffHid).toBeTruthy();
expect(match).toBe(true);
expect(initial).toBeUndefined();
expect(update).toBe(true);
})
});
test("UPDATE: updating a previously emitted doc should emit an update event", async () => {
const cutoffHid = 50;
// Insert a row, make sure it should match the selector
const insertedContent = await cacheContent(historyContent[0], true);
expect(insertedContent.hid > cutoffHid).toBeTruthy();
// subscribe to listener
const selector = { hid: { $gt: cutoffHid } };
const monitor$ = of({ selector }).pipe(
monitorQuery({ db$: content$ }),
takeUntil(timer(monitorSafetyTimeout)),
);
const spy = new ObserverSpy();
monitor$.subscribe(spy);
// wait for the monitor to spin-up so the insert doesn't look like the initial vals
await wait(monitorSpinUp);
// update the content
const fooVal = 123;
insertedContent.foo = fooVal;
const updatedContent = await cacheContent(insertedContent, true);
expect(updatedContent._id).toEqual(insertedContent._id);
expect(updatedContent.foo).toEqual(fooVal);
// wait for observable end
await spy.onComplete();
expect(spy.getValuesLength()).toEqual(2); // initial insert + later update
// check first event
const firstEvent = spy.getValueAt(0);
{
const { doc, match, initial, update } = firstEvent;
expect(doc.foo).toBeUndefined();
expect(match).toBe(true);
expect(initial).toBe(true);
expect(update).toBeUndefined();
}
// check update event
const secondEvent = spy.getValueAt(1);
{
const { doc, match, initial, update } = secondEvent;
expect(doc.foo).toBe(fooVal);
expect(match).toBe(true);
expect(initial).toBeUndefined();
expect(update).toBe(true);
}
});
test("DELETE: cache update causes element to no longer be in selection", async () => {
const cutoffHid = 50;
// Insert a row, make sure it should match the selector
const insertedContent = await cacheContent(historyContent[0], true);
expect(insertedContent.hid > cutoffHid).toBeTruthy();
// subscribe to listener
const selector = { hid: { $gt: cutoffHid } };
const monitor$ = of({ selector }).pipe(
monitorQuery({ db$: content$ }),
takeUntil(timer(monitorSafetyTimeout)),
);
const spy = new ObserverSpy();
monitor$.subscribe(spy);
// wait for the monitor to spin-up so the insert doesn't look like the initial vals
await wait(monitorSpinUp);
// update the content to be outside the selector
const badHid = -10000;
insertedContent.hid = badHid; // no longer > cutoffHid
const updatedContent = await cacheContent(insertedContent, true);
expect(updatedContent._id).toEqual(insertedContent._id);
expect(updatedContent.hid).toEqual(badHid);
// wait for observable end
await spy.onComplete();
expect(spy.getValuesLength()).toEqual(2); // initial insert + later update
// check first event
const firstEvent = spy.getValueAt(0);
{
const { doc, match, initial, update } = firstEvent;
expect(doc).toBeDefined();
expect(doc).toBeInstanceOf(Object);
expect(doc.hid).not.toEqual(badHid);
expect(match).toBe(true);
expect(initial).toBe(true);
expect(update).toBeUndefined();
}
// check update event
const secondEvent = spy.getValueAt(1);
{
const { doc, match, initial, update } = secondEvent;
expect(doc).toBeDefined();
expect(doc).toBeInstanceOf(Object);
expect(doc.hid).toEqual(badHid);
expect(match).toBe(false); // no longer matches selector
expect(initial).toBeUndefined();
expect(update).toBe(true);
}
});
});
});
describe("monitorQuery: collection content", () => {
// doctor and cache sample content
const fakeParentUrl = "/abc/def/ghi";
const testContent = collectionContent.map((doc) => {
doc.parent_url = fakeParentUrl;
return doc;
});
describe("initial results", () => {
test("first monitor event has initial results reflecting existing matches in db", async () => {
// insert some stuff
const cachedDscContent = await bulkCacheDscContent(testContent, true);
const cachedContentIds = new Set(pluckAll(cachedDscContent, "_id"));
// create new monitor
const limitIndex = 1;
const selector = {
parent_url: fakeParentUrl,
element_index: { $gt: limitIndex },
};
const monitor$ = of({ selector }).pipe(
monitorQuery({ db$: dscContent$ }),
takeUntil(timer(monitorSafetyTimeout))
);
// wait til done
const spy = new ObserverSpy();
monitor$.subscribe(spy);
await spy.onComplete();
// check emits
const emits = spy.getValues();
expect(spy.getValuesLength()).toBeGreaterThan(0);
expect(spy.getValuesLength()).toBeLessThanOrEqual(cachedContentIds.size);
emits.forEach(({ doc, match, initial, update }) => {
expect(match).toBe(true);
expect(initial).toBe(true);
expect(update).toBeUndefined();
expect(cachedContentIds.has(doc._id)).toBeTruthy();
expect(doc.element_index > limitIndex).toBeTruthy();
expect(doc.parent_url).toEqual(fakeParentUrl);
});
});
});
describe("updates", () => {
test("INSERT: adding to cache after instantiation should emit a new doc", async () => {
// createa monitor
const limitIndex = 0;
const selector = {
parent_url: fakeParentUrl,
element_index: { $gt: limitIndex },
};
const monitor$ = of({ selector }).pipe(
monitorQuery({ db$: dscContent$ }),
takeUntil(timer(monitorSafetyTimeout))
);
const spy = new ObserverSpy();
monitor$.subscribe(spy);
// wait for the monitor to spin-up so the insert doesn't look like the initial vals
await wait(monitorSpinUp);
// Insert a row, make sure it should match the selector
const insertedContent = await cacheCollectionContent(testContent[3], true);
expect(insertedContent.element_index > limitIndex).toBeTruthy();
// wait for observable end
await spy.onComplete();
const emits = spy.getValues();
expect(spy.getValuesLength()).toEqual(1);
emits.forEach(({ doc, match, initial, update }) => {
expect(doc.element_index > limitIndex).toBeTruthy();
expect(doc._id).toEqual(insertedContent._id);
expect(match).toBe(true);
expect(initial).toBeUndefined();
expect(update).toBe(true);
});
});
test("UPDATE: updating a previously emitted doc should emit an update event", async () => {
const limitIndex = 0;
// Insert a row, make sure it should match the selector
const insertedContent = await cacheCollectionContent(testContent[1], true);
expect(insertedContent.element_index > limitIndex).toBeTruthy();
// create monitor
const selector = {
parent_url: fakeParentUrl,
element_index: { $gt: limitIndex },
};
const monitor$ = of({ selector }).pipe(
monitorQuery({ db$: dscContent$ }),
takeUntil(timer(monitorSafetyTimeout))
);
const spy = new ObserverSpy();
monitor$.subscribe(spy);
// wait for the monitor to spin-up so the insert doesn't look like the initial vals
await wait(monitorSpinUp);
// update the content
const fooVal = 123;
insertedContent.foo = fooVal;
const updatedContent = await cacheCollectionContent(insertedContent, true);
expect(updatedContent._id).toEqual(insertedContent._id);
expect(updatedContent.foo).toEqual(fooVal);
// wait for observable end
await spy.onComplete();
expect(spy.getValuesLength()).toEqual(2); // initial insert + later update
// check first event
const firstEvent = spy.getValueAt(0);
{
const { doc, match, initial, update } = firstEvent;
expect(doc.foo).toBeUndefined();
expect(match).toBe(true);
expect(initial).toBe(true);
expect(update).toBeUndefined();
}
// check update event
const secondEvent = spy.getValueAt(1);
{
const { doc, match, initial, update } = secondEvent;
expect(doc.foo).toBe(fooVal);
expect(match).toBe(true);
expect(initial).toBeUndefined();
expect(update).toBe(true);
}
});
test("DELETE: cache update causes element to no longer be in selection", async () => {
const limitIndex = 0;
// Insert a row, make sure it should match the selector
const insertedContent = await cacheCollectionContent(testContent[1], true);
expect(insertedContent.element_index > limitIndex).toBeTruthy();
// create monitor
const selector = {
parent_url: fakeParentUrl,
element_index: { $gt: limitIndex },
};
const monitor$ = of({ selector }).pipe(
monitorQuery({ db$: dscContent$ }),
takeUntil(timer(monitorSafetyTimeout))
);
const spy = new ObserverSpy();
monitor$.subscribe(spy);
// update the content to be outside the selector
const badIndex = -10000;
insertedContent.element_index = badIndex; // no longer > cutoffHid
const updatedContent = await cacheCollectionContent(insertedContent, true);
expect(updatedContent._id).toEqual(insertedContent._id);
expect(updatedContent.element_index).toEqual(badIndex);
// wait for observable end
await spy.onComplete();
expect(spy.getValuesLength()).toEqual(2); // initial insert + later update
// check first event
const firstEvent = spy.getValueAt(0);
{
const { doc, match, initial, update } = firstEvent;
expect(doc).toBeDefined();
expect(doc).toBeInstanceOf(Object);
expect(doc.element_index).not.toEqual(badIndex);
expect(match).toBe(true);
expect(initial).toBe(true);
expect(update).toBeUndefined();
}
// check update event
const secondEvent = spy.getValueAt(1);
{
const { doc, match, initial, update } = secondEvent;
expect(doc).toBeDefined();
expect(doc).toBeInstanceOf(Object);
expect(doc.element_index).toEqual(badIndex);
expect(match).toBe(false); // no longer matches selector
expect(initial).toBeUndefined();
expect(update).toBe(true);
}
});
});
});
@@ -1,181 +0,0 @@
/**
* Observable api for caching the galaxy history contents and collection contents
*/
import config from "config";
import { pipe } from "rxjs";
import { map } from "rxjs/operators";
import { collection, cacheItem, uncacheItem, getItemByKey, bulkCache } from "./pouch";
import { contentIndices, dscIndices } from "./cacheIndices";
// helper operator that applies function to each element of sourced array
// prettier-ignore
const applyEach = (fn) => pipe(
map((list) => list.map(fn))
);
// pad zeros so ids sort properly
const zeroPad = (num, places) => String(num).padStart(places, "0");
/**
* History Content & associated operators
*/
export const content$ = collection(
{
name: "galaxy-content",
indexes: contentIndices,
},
config
);
// prettier-ignore
export const getCachedContent = (key = "_id") => pipe(
getItemByKey(content$, key)
);
// prettier-ignore
export const cacheContent = (returnDoc = false) => pipe(
map(prepContent),
cacheItem(content$, returnDoc),
map((result) => (returnDoc ? result.doc : result))
);
// prettier-ignore
export const uncacheContent = () => pipe(
uncacheItem(content$)
);
// prettier-ignore
export const bulkCacheContent = (returnDocs = false) => pipe(
applyEach(prepContent),
bulkCache(content$, returnDocs),
map((list) => (returnDocs ? list.map((result) => result.doc) : list))
);
export const buildContentId = ({ history_id, hid }) => {
const paddedHid = zeroPad(hid, 12);
return `${history_id}-${paddedHid}`;
};
export const prepContent = (props) => {
const { history_content_type, id, type_id: origTypeId, ...theRest } = props;
const type_id = origTypeId ? origTypeId : `${history_content_type}-${id}`;
const _id = buildContentId(props);
const doc = {
_id,
history_content_type,
id,
type_id,
...theRest,
};
return fixDeleted(doc);
};
/**
* Collection content (drill down into a collection). Since we don't really
* update dsc content, all I think we need is the bulk operator which gets run
* when the scroller loads new data
*/
export const dscContent$ = collection(
{
name: "galaxy-collection-content",
indexes: dscIndices,
},
config
);
// prettier-ignore
export const getCachedCollectionContent = (key = "_id") => pipe(
getItemByKey(dscContent$, key)
);
// prettier-ignore
export const cacheCollectionContent = (returnDoc) => pipe(
map(prepDscContent),
cacheItem(dscContent$, returnDoc),
map((result) => (returnDoc ? result.doc : result))
);
// prettier-ignore
export const bulkCacheDscContent = (returnDocs = false) => pipe(
applyEach(prepDscContent),
bulkCache(dscContent$, returnDocs),
map((list) => (returnDocs ? list.map((result) => result.doc) : list))
);
// need to deal with both the bent api format that comes in from the ajax call
// as well as re-processing an existing collection item on a subsequent update
// cache operation.
export const prepDscContent = (rawProps) => {
// console.log("prepDscContent", rawProps);
const { object = {}, ...rootProps } = rawProps;
const props = { ...rootProps, ...object };
const { id, model_class, element_identifier, ...otherObjectFields } = props;
// Validation
if (id === undefined) {
throw new Error("missing id");
}
if (model_class === undefined) {
throw new Error("missing model_class");
}
const history_content_type = model_class == "HistoryDatasetAssociation" ? "dataset" : "dataset_collection";
const type_id = `${history_content_type}-${id}`;
const newProps = {
// parent content url + counter as that is the most likely query
_id: buildCollectionId(props),
// make a type_id so we can re-use all our functions which depend on it
id,
type_id,
model_class,
history_content_type,
// need to rename for cache indexing
name: element_identifier,
element_identifier,
// move stuff out of "object" and into root of cached packet
...otherObjectFields,
};
const clean = fixDeleted(newProps);
return clean;
};
export const buildCollectionId = (props) => {
const { parent_url, element_index } = props;
// this is not part of the returned api data, we need to provide it after the
// ajax call from the contents_url we got this from
if (undefined == parent_url) {
throw new Error("missing required parent_url");
}
if (undefined == element_index) {
throw new Error("missing rquired element_index");
}
const paddedIndex = zeroPad(element_index, 12);
return `${parent_url}-${paddedIndex}`;
};
/**
* We need to rename the deleted property because pouchDB uses that field,
* which is unfortunate
* @param {Object} props Raw document props
*/
const fixDeleted = (props) => {
// console.log(props);
if (Object.prototype.hasOwnProperty.call(props, "deleted")) {
const { deleted: isDeleted, ...theRest } = props;
return { isDeleted, ...theRest };
}
return props;
};
@@ -1,233 +0,0 @@
/**
* Generic pouch db operators, These are low-level pouchdb access operators
* intended to be used to create more specific operators.
*/
import moment from "moment";
import deepEqual from "deep-equal";
import { defer, pipe, from } from "rxjs";
import { tap, filter, mergeMap, reduce, shareReplay } from "rxjs/operators";
import { needs } from "utils/observable";
import { dasherize } from "underscore.string";
import PouchDB from "pouchdb";
import PouchAdapterMemory from "pouchdb-adapter-memory";
import PouchUpsert from "pouchdb-upsert";
import PouchFind from "pouchdb-find";
import PouchErase from "pouchdb-erase";
// import PouchDebug from "pouchdb-debug";
PouchDB.plugin(PouchAdapterMemory);
PouchDB.plugin(PouchUpsert);
PouchDB.plugin(PouchFind);
PouchDB.plugin(PouchErase);
// PouchDB.plugin(PouchDebug);
// debugging stuff
// PouchDB.debug.enable('pouchdb:find');
// const show = (obj) => console.log(JSON.stringify(obj, null, 4));
// Instance storage map, keyed by database name
export const dbs = new Map();
/**
* Generate an observable that initializes and shares a pouchdb instance.
*
* @param {object} options pouchdb initialization configs
* @return {Observable} observable that emits the pouch instance
*/
export const collection = (opts, appConfig) => {
return defer(() => {
const coll = buildCollection(opts, appConfig);
const name = collectionName(opts, appConfig);
return from(coll).pipe(tap((db) => dbs.set(name, db)));
}).pipe(
filter((db) => db instanceof PouchDB),
shareReplay(1)
);
};
async function buildCollection(opts, appConfig) {
if (!opts) {
throw new Error("collection: Missing database config");
}
if (!appConfig) {
throw new Error("collection: Missing application config");
}
const name = collectionName(opts, appConfig);
// eslint-disable-next-line no-unused-vars
const { indexes = [], ...otherOpts } = opts;
const dbConfig = { ...appConfig.caching, ...otherOpts, name };
// make instance
const db = await new PouchDB(dbConfig);
// pre-install common indexes
await installCollectionIndexes(db, indexes);
return db;
}
function collectionName(opts, appConfig) {
const { name: dbName } = opts;
const { name: envName } = appConfig;
return dasherize(`${dbName} ${envName}`);
}
async function installCollectionIndexes(db, indexes = []) {
const promises = indexes.map((idx) => db.createIndex(idx));
if (promises.length) {
await Promise.all(promises);
}
return db;
}
/**
* Retrieves an object from the cache
*
* @param {Observable} db$ Observable of a pouchdb instance
* @param {string} keyName Field to lookup object by
* @returns {Function} Observable operator
*/
// prettier-ignore
export const getItemByKey = (db$, keyName = "_id") => pipe(
needs(db$),
mergeMap(async ([keyValue, db]) => {
let doc = null;
if (keyName == "_id") {
doc = await getDocByDocId(db, keyValue);
} else {
doc = await getDocByKeyValue(db, keyName, keyValue);
}
return doc;
})
);
async function getDocByDocId(db, docId) {
let doc = null;
try {
doc = await db.get(docId);
} catch (err) {
// pouch throws a 404 when a .get doesn't find anything,
// which in my mind is a perfectly legitimate result, so
// we're returning a null here, otherwise rethrow
const { status } = err;
if (status !== 404) {
throw err;
}
}
return doc;
}
async function getDocByKeyValue(db, keyName, keyValue) {
let doc = null;
// otherwise do a find by keyname
const searchConfig = {
selector: { [keyName]: keyValue },
limit: 1,
};
const response = await db.find(searchConfig);
if (response.warning) {
console.warn(response.warning, searchConfig);
}
const { docs } = response;
if (docs && docs.length > 0) {
doc = docs[0];
}
return doc;
}
/**
* Operator that caches all source documents in the indicated collection.
* Upserts new fields into the old document, so the props need not be complete,
* just needs to have the _id.
*
* Adds a cached_at timestamp and fixes pouchdb specific field names like
* "deleted"
*
* @source Observable stream of docs to cache
* @param {Observable} db$ Observable pouchdb instance
* @returns {Function} Observable operator
*/
// prettier-ignore
export const cacheItem = (db$, returnDoc = false) => pipe(
needs(db$),
mergeMap(async ([item, db]) => {
const result = await db.upsert(item._id, (existing) => {
// eslint-disable-next-line no-unused-vars
const { _rev, cached_at, ...existingFields } = existing;
// ignore if what we're caching is the same as what's in there
const same = deepEqual(item, existingFields);
if (same) {return false;}
return {
...existing,
...item,
cached_at: moment().valueOf(),
};
});
if (returnDoc) {
result.doc = await getDocByDocId(db, result.id);
}
return result;
})
);
/**
* Doing it the dumb way for now, will try a true bulkDocs call later.
*
* @source Observable of docs to cache
* @param {Observable} db$ Observable of a PouchDB instance
* @returns {Function} Observable operator
*/
// prettier-ignore
export const bulkCache = (db$, returnDocs = false) => pipe(
mergeMap((list) =>
from(list).pipe(
cacheItem(db$, returnDocs),
reduce((result, item) => {
result.push(item);
return result;
}, [])
)
)
);
/**
* Creates an operator that will delete the source document from the configured
* database observable. In pouchdb, deleting just means setting _deleted to
* true.
*
* @source Observable stream of documents to uncache
* @param {Observable} db$ Observable of the pouchdb instance
* @returns {Function} Observable operator
*/
// prettier-ignore
export const uncacheItem = (db$) => pipe(
needs(db$),
mergeMap(([doomedDoc, db]) => {
return db.remove(doomedDoc);
})
);
/**
* Delete existing indexes in a pouchdb database.
*
* @param {PouchDB} db pouch db instance
* @returns {Promise}
*/
export async function deleteIndexes(db) {
const response = await db.getIndexes();
const doomedIndexes = response.indexes.filter((idx) => idx.ddoc !== null);
const promises = doomedIndexes.map((idx) => db.deleteIndex(idx));
return await Promise.all(promises);
}
@@ -1,47 +0,0 @@
/**
* Same functionality as exposed in observables, but in promise form for one-off
* utilities. A little easier to use if all you need to do is one thing.
* However, for more complicated processing, I do recommend making the attempt
* work in streams as you will have a lot more control over the throughput.
*/
import { of } from "rxjs";
import { firstValueFrom } from "utils/observable/firstValueFrom";
import {
content$,
cacheContent as cacheContentOp,
getCachedContent as getCachedContentOp,
uncacheContent as uncacheContentOp,
bulkCacheContent as bulkCacheContentOp,
bulkCacheDscContent as bulkCacheDscContentOp,
getCachedCollectionContent as getCachedCollectionContentOp,
cacheCollectionContent as cacheCollectionContentOp,
} from "./observables";
import { find } from "./find";
// result is an object pouch returns with the ID a rev, and whether it was
// actually updated or not (won't be if identical)
export const cacheContent = (rawProps, returnDoc) => firstValueFrom(of(rawProps).pipe(cacheContentOp(returnDoc)));
export const getCachedContent = (docId) => firstValueFrom(of(docId).pipe(getCachedContentOp()));
export const uncacheContent = (doc) => firstValueFrom(of(doc).pipe(uncacheContentOp()));
export const bulkCacheContent = (list, returnDocs) => firstValueFrom(of(list).pipe(bulkCacheContentOp(returnDocs)));
// collection content, collections are supposedly immutable after creation so
// all we really have is a bulk-load
export const cacheCollectionContent = (rawProps, returnDoc) =>
firstValueFrom(of(rawProps).pipe(cacheCollectionContentOp(returnDoc)));
export const getCachedCollectionContent = (docId) => firstValueFrom(of(docId).pipe(getCachedCollectionContentOp()));
export const bulkCacheDscContent = (list, returnDocs) =>
firstValueFrom(of(list).pipe(bulkCacheDscContentOp(returnDocs)));
// examples of how to fine-tune these to get specialized functionality
export const getContentByTypeId = async (history_id, type_id) => {
const selector = { history_id, type_id };
const index = { fields: ["history_id", "type_id"] };
const request = { selector, index, limit: 1 };
const obs$ = of(request).pipe(find(content$));
const results = await firstValueFrom(obs$);
return results.length > 0 ? results[0] : null;
};
@@ -1,10 +0,0 @@
import { dbs } from "./pouch";
/**
* Erases all stored database instances
*/
export async function wipeDatabase() {
for (const db of dbs.values()) {
await db.erase();
}
}
@@ -1,9 +0,0 @@
// seek types
const seeks = { ASC: "asc", DESC: "desc" };
seeks.isValid = (val) => {
return Object.values(SEEK).some((v) => v == val);
};
export const SEEK = seeks;
@@ -1,18 +0,0 @@
// Exposes worker functions as promises or Observable operators
export { monitorContentQuery, monitorDscQuery, monitorHistoryContent, monitorCollectionContent } from "./CacheApi";
export { loadHistoryContents, loadDscContent } from "./CacheApi";
export {
cacheContent,
getCachedContent,
uncacheContent,
bulkCacheContent,
cacheCollectionContent,
getCachedCollectionContent,
bulkCacheDscContent,
getContentByTypeId,
} from "./CacheApi";
export { wipeDatabase, clearHistoryDateStore } from "./CacheApi";
@@ -1,50 +0,0 @@
/**
* Manual loading operators for content lists. These get run as the user scrolls
* or filters data in the listing for immediate lookups. All of them will run
* one or more ajax calls against the api and cache ther results.
*/
import { map, pluck, withLatestFrom, publish } from "rxjs/operators";
import { nth } from "utils/observable";
import { requestWithUpdateTime } from "./operators/requestWithUpdateTime";
import { bulkCacheDscContent } from "./db";
import { SearchParams } from "components/providers/History/SearchParams";
import { prependPath } from "utils/redirect";
import { summarizeCacheOperation, dateStore } from "./loadHistoryContents";
import { show } from "utils/observable";
/**
* Load collection content (drill down)
* Params: contents_url + search params + element_index
*/
export const loadDscContent = (cfg = {}) => {
const { debug = false } = cfg;
return publish((inputs$) => {
const url$ = inputs$.pipe(nth(0));
return inputs$.pipe(
map(buildDscContentUrl),
map(prependPath),
show(debug, (url) => console.log("Sending collection request:", url)),
requestWithUpdateTime({ dateStore }),
pluck("response"),
withLatestFrom(url$),
map(([rows, parent_url]) => {
return rows.map((row) => ({ ...row, parent_url }));
}),
bulkCacheDscContent(),
summarizeCacheOperation()
);
});
};
// Collection + params -> request url w/o update_time
// ignore params for now, we don't filter collection contents yet.
const buildDscContentUrl = (inputs) => {
const [base, , pagination] = inputs;
const { offset = 0, limit = SearchParams.pageSize } = pagination;
const skipClause = offset ? `offset=${offset}` : "";
const limitClause = limit ? `limit=${limit}` : "";
const qs = [skipClause, limitClause].filter((o) => o.length).join("&");
return `${base}?${qs}`;
};
@@ -1,141 +0,0 @@
import { zip } from "rxjs";
import { map, pluck, share, filter } from "rxjs/operators";
import { hydrate } from "utils/observable";
import { requestWithUpdateTime } from "./operators/requestWithUpdateTime";
import { prependPath } from "utils/redirect";
import { bulkCacheContent } from "./db";
import { SearchParams } from "components/providers/History/SearchParams";
import { createDateStore } from "components/History/model/DateStore";
// Shared datestore for the request opertor for both loading and polling
export const dateStore = createDateStore();
export const clearHistoryDateStore = async () => {
console.log("clearHistoryDateStore");
dateStore.clear();
};
/**
* Turn [historyId, params, pagination] into content requests. Cache results and send
* back statistics about what was returned for this query.
*/
// prettier-ignore
export const loadHistoryContents = (cfg = {}) => (rawInputs$) => {
const {
noInitial = false,
windowSize = SearchParams.pageSize
} = cfg;
const inputs$ = rawInputs$.pipe(
hydrate([undefined, SearchParams]),
);
const responseQualifier = ajaxResponse => ajaxResponse.status == 200 && ajaxResponse.response.length > 0;
const ajaxResponse$ = inputs$.pipe(
map(([id, params, hid]) => {
const baseUrl = `/api/histories/${id}/contents/near/${hid}/${windowSize}`;
return `${baseUrl}?${params.historyContentQueryString}`;
}),
map(prependPath),
requestWithUpdateTime({ dateStore, noInitial, responseQualifier, dateFieldName: "since" }),
);
const validResponses$ = ajaxResponse$.pipe(
filter(response => response.status == 200),
share(),
);
const cacheSummary$ = validResponses$.pipe(
pluck("response"),
bulkCacheContent(),
summarizeCacheOperation(),
);
return zip(validResponses$, cacheSummary$).pipe(
map(([ajaxResponse, summary]) => {
const { xhr, response = [] } = ajaxResponse;
const { max: maxContentHid, min: minContentHid } = getPropRange(response, "hid");
// return an int or undefined if the field does not exist in the headers
const headerInt = field => {
const raw = xhr.getResponseHeader(field);
return raw === null ? undefined : parseInt(raw);
};
const headerBool = field => {
const raw = xhr.getResponseHeader(field);
return (raw == "true" || raw == "1");
}
// header counts
const matches = headerInt("matches")
const matchesUp = headerInt("matches_up");
const matchesDown = headerInt("matches_down");
const totalMatches = headerInt("total_matches")
const totalMatchesUp = headerInt("total_matches_up");
const totalMatchesDown = headerInt("total_matches_down");
const minHid = headerInt("min_hid");
const maxHid = headerInt("max_hid");
const historySize = headerInt("history_size");
const historyEmpty = headerBool("history_empty");
return {
summary,
matches,
totalMatches,
minHid,
maxHid,
minContentHid, // minimum hid in the returned result
maxContentHid, // maximum hid in the returned result
matchesUp,
matchesDown,
totalMatchesUp,
totalMatchesDown,
// new history size
historySize,
historyEmpty
};
})
);
};
/**
* Once data was cached, there's no need to send everything back over into the
* main thread since the cache watcher will pick up those new values. This just
* summarizes what pouchdb did during the bulk cache and sends back some stats.
*/
// prettier-ignore
export const summarizeCacheOperation = () => {
return map((list) => {
const cached = list.filter((result) => result.updated || result.ok);
return {
updatedItems: cached.length,
totalReceived: list.length,
};
})
}
/**
* Gets min and max values from an array of objects
*
* @param {Array} list Array of objects, each with a propName
* @param {string} propName Name of prop to measure range
*/
export const getPropRange = (list, propName) => {
const narrowRange = (range, row) => {
const val = parseInt(row[propName], 10);
range.max = Math.max(range.max, val);
range.min = Math.min(range.min, val);
return range;
};
const everywhere = {
min: Infinity,
max: -Infinity,
};
return list.reduce(narrowRange, everywhere);
};
@@ -1,190 +0,0 @@
import { merge } from "rxjs";
import { map, distinctUntilChanged, publish, withLatestFrom, share } from "rxjs/operators";
import { content$, dscContent$, buildContentId, buildCollectionId } from "./db/observables";
import { monitorQuery } from "./db/monitorQuery";
import { SearchParams } from "components/providers/History/SearchParams";
import { deepEqual } from "deep-equal";
import { hydrate } from "utils/observable";
import { SEEK } from "./enums";
import { matchesSelector } from "pouchdb-selector-core";
// history contents monitor
export const monitorHistoryContent = (cfg = {}) => {
return twoWayMonitor({
...cfg,
db$: content$,
buildRequest: buildContentPouchRequest,
aggregationKeyField: "hid",
});
};
// collection contents monitor
export const monitorCollectionContent = (cfg = {}) => {
return twoWayMonitor({
...cfg,
db$: dscContent$,
buildRequest: buildCollectionPouchRequest,
aggregationKeyField: "element_index",
});
};
/**
* Search cache upward and downward from the scroll HID, filtering out results
* which do not match params, we don't want to return the entire result set for
* 2 reasons: it might be huge, and we might not yet have large sections of the
* contents if the user is rapidly dragging the scrollbar to different regions
*/
// prettier-ignore
export const twoWayMonitor = (cfg = {}) => (src$) => {
const {
db$,
aggregationKeyField,
buildRequest = buildContentPouchRequest,
pageSize = SearchParams.pageSize,
inputDebounce = 100,
debug = false,
} = cfg;
return src$.pipe(
hydrate([undefined, SearchParams]),
publish(input$ => {
const upRequest$ = input$.pipe(
map(buildRequest({ seek: SEEK.ASC, pageSize })),
distinctUntilChanged(deepEqual),
share(),
);
const downRequest$ = input$.pipe(
map(buildRequest({ seek: SEEK.DESC, pageSize })),
distinctUntilChanged(deepEqual),
share(),
);
// one page up
const up$ = upRequest$.pipe(
monitorQuery({ db$, inputDebounce, debug, label: "up" }),
alsoMatchWith(downRequest$),
);
// current page + next page = 2 pages down
const down$ = downRequest$.pipe(
monitorQuery({ db$, inputDebounce, debug, label: "down" }),
alsoMatchWith(upRequest$),
);
return merge(up$, down$)
}),
map((change) => {
// don't bother sending back non-matching results
// we're just going to delete them from the skiplist
// using the aggregation key, so no need to serialize them
const { match, doc, ...more} = change;
const key = doc[aggregationKeyField];
return match ? { key, match, doc, ...more } : { key, match, ...more };
})
);
};
// utility operator, checks to see if change matches with the other selector
const alsoMatchWith = (req$) => (change$) => {
// As the scroller moves, a given row may no longer match in one query
// but it may in the other, simply because it has moved from the "up" results
// to the "down" results
return change$.pipe(
withLatestFrom(req$),
map(([change, otherRequest]) => {
const result = { ...change };
result.match = change.match || matchesSelector(result.doc, otherRequest.selector);
return result;
})
);
};
/**
* Generate a PouchDB selector for a specified history, filters, and a rough
* guess of where to start looking.
*
* Can't use skip because there might be big un-cached regions of the history
* and we need to be able to select without loading everything
*/
export const buildContentPouchRequest =
(cfg = {}) =>
(inputs) => {
const { seek, pageSize = SearchParams.pageSize } = cfg;
const [history_id, params, hid] = inputs;
// look up or down from target hid
const targetId = buildContentId({ history_id, hid });
// SEEK.ASC means the top seek, HID > target
const comparator = seek == SEEK.ASC ? "$gt" : "$lte";
// one page above the target, then 2 after to get the current page and some buffer
const pages = seek == SEEK.ASC ? 1 : 2;
const limit = pages * pageSize;
// index, will build if not existent
const filterFields = Array.from(params.criteria.keys()).map((f) => params.getPouchFieldName(f));
const fields = ["_id", "history_id", ...filterFields];
const ddoc = "idx-" + fields.sort().join("-");
const request = {
selector: {
$and: [
// doc id + direction (up/down)
{ _id: { [comparator]: targetId } },
{ history_id: history_id },
...params.pouchFilters,
],
},
sort: [{ _id: seek }],
limit,
index: {
fields,
ddoc,
},
};
return request;
};
// filters currently unused, but they could be if we add some filter UI
export const buildCollectionPouchRequest =
(cfg = {}) =>
(inputs) => {
const { seek, pageSize = SearchParams.pageSize } = cfg;
// eslint-disable-next-line no-unused-vars
const [parent_url, filters, element_index] = inputs;
const targetId = buildCollectionId({ parent_url, element_index });
// SEEK.ASC means above the target means element_index < target
const comparator = seek == SEEK.ASC ? "$lt" : "$gte";
// 3 pages total, 1 before, the current one, 1 after;
const pages = seek == SEEK.ASC ? 1 : 2;
const limit = pages * pageSize;
// when searching upward we invert the sort because we want the
// results closest to the target
const idSort = seek == SEEK.ASC ? "desc" : "asc";
const request = {
selector: {
_id: { [comparator]: targetId },
parent_url: { $eq: parent_url },
},
sort: [{ _id: idSort }],
limit,
index: {
fields: ["_id"],
sort: [{ _id: "asc" }],
name: "collection contents by parent_url and element_index ascending",
ddoc: "idx-dsc-contents-by-url-and-element-asc",
},
};
return request;
};
@@ -1,122 +0,0 @@
import { timer, of } from "rxjs";
import { takeUntil } from "rxjs/operators";
import { firstValueFrom } from "utils/observable/firstValueFrom";
import { wipeDatabase } from "./db/wipeDatabase";
import { wait } from "jest/helpers";
import { ObserverSpy } from "@hirez_io/observer-spy";
import { SearchParams } from "components/providers/History/SearchParams";
import { content$ } from "./db/observables";
import { bulkCacheContent, cacheContent } from "./db/promises";
import { find } from "./db/find";
import { buildContentPouchRequest, monitorHistoryContent } from "./monitorHistoryContent";
// test data
import historyContent from "components/providers/History/test/json/historyContent.json";
jest.mock("app");
jest.mock("../caching");
beforeEach(wipeDatabase);
afterEach(wipeDatabase);
const monitorSpinUp = 400;
const monitorSafetyTimeout = 1000;
const selectorHasField = (selector, field) => selector.$and.some((row) => row[field] !== undefined);
describe("buildContentPouchRequest", () => {
test("should turn inputs into a pouch request suitable for find", () => {
const fn = buildContentPouchRequest();
expect(fn).toBeInstanceOf(Function);
const { history_id, hid } = historyContent[0];
const inputs = [history_id, new SearchParams(), hid];
const request = fn(inputs);
expect(request).toBeDefined();
expect(request.selector).toBeDefined();
expect(request.selector.$and).toBeDefined();
expect(selectorHasField(request.selector, "_id")).toBeTruthy();
expect(selectorHasField(request.selector, "history_id")).toBeTruthy();
});
});
describe("monitorHistoryContent", () => {
let cachedContentMap;
// cache sample content, store in a map for later, keyed by HID
beforeEach(async () => {
const cachedContent = await bulkCacheContent(historyContent, true);
const cachedContentEntries = cachedContent.map((o) => [o.hid, o]);
cachedContentMap = new Map(cachedContentEntries);
});
afterEach(() => (cachedContentMap = null));
test("sanity check all the data is in there", async () => {
const { history_id } = historyContent[0];
const request = { selector: { history_id } };
const obs$ = of(request).pipe(find(content$));
const result = await firstValueFrom(obs$);
expect(result.length).toEqual(historyContent.length);
});
test("should see the pre-inserted content if params match", async () => {
const { history_id, hid } = historyContent[0];
const monitor$ = of([history_id, new SearchParams(), hid]).pipe(
monitorHistoryContent({
inputDebounce: 50,
pageSize: 10, // get everything without pagination for now
}),
takeUntil(timer(monitorSafetyTimeout)) // safety
);
const spy = new ObserverSpy();
monitor$.subscribe(spy);
await spy.onComplete();
expect(spy.getValuesLength()).toBeGreaterThan(0);
spy.getValues().forEach((evt) => {
const { key, match, doc } = evt;
expect(match).toBe(true);
expect(doc.hid).toEqual(key);
expect(cachedContentMap.has(key)).toBe(true);
// we had default params
expect(doc.isDeleted).toBe(false);
expect(doc.visible).toBe(true);
});
});
test("should see subsequent updates", async () => {
const firstDoc = historyContent[0];
const { history_id, hid } = firstDoc;
const inputs = [history_id, new SearchParams(), hid];
const monitor$ = of(inputs).pipe(
monitorHistoryContent({ inputDebounce: 50 }),
takeUntil(timer(monitorSafetyTimeout)) // safety
);
// monitor observable
const spy = new ObserverSpy();
monitor$.subscribe(spy);
// wait for first emission before updating, this lets all
// the initial matches emit
await wait(monitorSpinUp);
// update one doc
const firstHid = historyContent[0].hid;
const updateMe = cachedContentMap.get(firstHid);
const fakeVal = 123;
updateMe.foobar = fakeVal;
const updatedContent = await cacheContent(updateMe, true);
expect(updatedContent.foobar).toEqual(fakeVal);
// wait until monitor completes
await spy.onComplete();
const lastEmit = spy.getLastValue();
expect(lastEmit.doc.foobar).toEqual(fakeVal);
});
});
@@ -1,16 +0,0 @@
import { pipe } from "rxjs";
import { map } from "rxjs/operators";
/**
* passing objects into the worker removes class information, Pass in an array
* of constructors to be assigned, in order of the input array.
*/
// prettier-ignore
export const hydrate = (constructors = []) => pipe(
map((inputs) => {
return inputs.map((rawInput, i) => {
const C = constructors[i];
return C ? new C(rawInput) : rawInput;
});
})
)
@@ -1,36 +0,0 @@
import moment from "moment";
import { pipe } from "rxjs";
import { map } from "rxjs/operators";
import { find } from "../db/find";
/**
* Returns latest cache date from the cache.
*
* This seemingly simple calculation is complicated because the galaxy API
* currently returns update_times that are more precise than javascript dates.
*
* For example. The server will return an update_time as:
* 2020-07-02T17:25:09.385026
*
* When parsing a javascript date, however, we don't have that many decimals.
* 2020-07-02T17:25:09Z
*
* This means that trying to perform an inequality filter:
* update_time-gt=2020-07-02T17:25:09Z
*
* ...will always fail because of those extra fractions of a millisecond. And
* given the granularity of the date storage it is not safe to "just add one
* more" to the outgoing value, given that there may indeed be lost records if we do that.
*
* source stream: pouchdb-find query config
*/
// prettier-ignore
export const lastCachedDate = (db$) => pipe(
find(db$),
map((docs) => {
if (!docs.length) {return null;}
const dates = docs.map((d) => d.cached_at);
const maxDate = Math.max(...dates);
return moment.utc(maxDate).toISOString();
})
)
@@ -1,12 +0,0 @@
// I'm actually unclear on why I needed this operator to make pouchdb
// observables work. This bears further investigation. To my knowledge
// withLatestFrom() does the same thing as this operator, yet this works and
// withLatestFrom() doesn't when used with pouchdb observable creation.
import { zip, of, pipe } from "rxjs";
import { concatMap } from "rxjs/operators";
// prettier-ignore
export const needs = (dep$) => pipe(
concatMap((src) => zip(of(src), dep$))
);
@@ -1,71 +0,0 @@
import moment from "moment";
import { of, pipe } from "rxjs";
import { tap, map, mergeMap } from "rxjs/operators";
import { ajax } from "rxjs/ajax";
import { createDateStore } from "components/History/model/DateStore";
/**
* Global url date-store, keeps track of the last time a specific url was
* requested Can reset the update_time tracking by passing in a new datestore to
* the operator configs
*/
export const requestDateStore = createDateStore("requestWithUpdateTime default");
/***
* Check datestore to get the last time we did this request. Append
* an update_time clause at the end of the GET url so we only take
* the fresh updates. Mark the time after the request for future requests
*/
// prettier-ignore
export const requestWithUpdateTime = (config = {}) => {
const {
dateStore = requestDateStore,
bufferSeconds = 0,
dateFieldName = "update_time-gt",
requestTime = moment.utc(),
// indicates we don't want initial results
noInitial = false,
responseQualifier = (response) => response.status == 200
} = config;
// mark and flag this update-time, append to next request with same base
return pipe(
tap((url) => {
if (noInitial && !dateStore.has(url)) {
dateStore.set(url, moment.utc());
}
}),
mergeMap((baseUrl) => of(baseUrl).pipe(
appendUpdateTime({ dateStore, bufferSeconds, dateFieldName }),
mergeMap(ajax),
tap((response) => {
if (responseQualifier(response)) {
dateStore.set(baseUrl, requestTime)
}
})
))
);
};
/**
* Takes a base URL appends an update_time-gt restriction based on the lst
* time this URL was requestd as indicated by the dateStore.
* (Async in case we choose to store the date in indexDb instead of localStorage)
*/
// prettier-ignore
const appendUpdateTime = (cfg = {}) => {
const {
dateStore = requestDateStore,
dateFieldName = "update_time-gt",
} = cfg;
return pipe(
map((baseUrl) => {
if (!dateStore.has(baseUrl)) {return baseUrl;}
const lastRequest = dateStore.get(baseUrl);
const parts = [baseUrl, `${dateFieldName}=${lastRequest.toISOString()}`];
const separator = baseUrl.includes("?") ? "&" : "?";
return parts.join(separator);
})
);
};
@@ -1,32 +0,0 @@
// same behavior as as utils/observable/throttleDistinct,
// but uses external state storage so we can put it in nested switchMaps, etc.
import moment from "moment";
import { pipe } from "rxjs";
import { filter } from "rxjs/operators";
import { createDateStore } from "../../model/DateStore";
export const throttleDistinctDateStore = createDateStore("throttleDistinct default");
// prettier-ignore
export const throttleDistinct = (config = {}) => {
const {
timeout = 1000,
dateStore = throttleDistinctDateStore
} = config;
return pipe(
filter((val) => {
const now = moment();
let ok = true;
if (dateStore.has(val)) {
const lastRequest = dateStore.get(val);
ok = now - lastRequest > timeout;
}
if (ok) {
dateStore.set(val, now);
}
return ok;
})
);
};