Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
implement live query on @orbit/record-cache
- Loading branch information
Showing
8 changed files
with
839 additions
and
0 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
32 changes: 32 additions & 0 deletions
32
packages/@orbit/record-cache/src/live-query/async-live-query.ts
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,32 @@ | ||
import { RecordException } from '@orbit/data'; | ||
import { QueryResult } from '../query-result'; | ||
import { AsyncRecordCache } from '../async-record-cache'; | ||
import { LiveQuery, LiveQuerySettings } from './live-query'; | ||
|
||
export interface AsyncLiveQuerySettings extends LiveQuerySettings { | ||
cache: AsyncRecordCache; | ||
} | ||
|
||
export class AsyncLiveQuery extends LiveQuery { | ||
cache: AsyncRecordCache; | ||
|
||
constructor(settings: AsyncLiveQuerySettings) { | ||
super(settings); | ||
this.cache = settings.cache; | ||
} | ||
|
||
get schema() { | ||
return this.cache.schema; | ||
} | ||
|
||
executeQuery( | ||
onNext: (result: QueryResult) => void, | ||
onError?: (error: RecordException) => void | ||
): void { | ||
this.cache.query(this.query).then(onNext, onError || defaultOnError); | ||
} | ||
} | ||
|
||
function defaultOnError(error: RecordException): void { | ||
throw error; | ||
} |
187 changes: 187 additions & 0 deletions
187
packages/@orbit/record-cache/src/live-query/live-query.ts
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,187 @@ | ||
import { Evented } from '@orbit/core'; | ||
import { | ||
QueryExpression, | ||
FindRecord, | ||
FindRecords, | ||
FindRelatedRecord, | ||
FindRelatedRecords, | ||
equalRecordIdentities, | ||
Query, | ||
Schema, | ||
RecordOperation, | ||
RecordException | ||
} from '@orbit/data'; | ||
|
||
import { QueryResult } from '../query-result'; | ||
import { recordOperationChange, RecordChange } from './utils'; | ||
|
||
export interface LiveQuerySettings { | ||
query: Query; | ||
} | ||
|
||
export class LiveQuery { | ||
cache: Evented; | ||
schema: Schema; | ||
query: Query; | ||
|
||
executeQuery( | ||
onNext: (result: QueryResult) => void, | ||
onError?: (error: RecordException) => void | ||
): void { | ||
throw new TypeError('executeQuery: Not Implemented.'); | ||
} | ||
|
||
on( | ||
onNext: (result: QueryResult) => void, | ||
onError?: (error: RecordException) => void | ||
): () => void { | ||
const executeQuery = onceTick(() => this.executeQuery(onNext, onError)); | ||
const unsubscribePatch = this.cache.on( | ||
'patch', | ||
(operation: RecordOperation) => { | ||
if (this.match(operation)) { | ||
executeQuery(); | ||
} | ||
} | ||
); | ||
|
||
const unsubscribeReset = this.cache.on('reset', () => { | ||
executeQuery(); | ||
}); | ||
|
||
executeQuery(); | ||
|
||
return function unsubscribe() { | ||
cancelTick(executeQuery); | ||
unsubscribePatch(); | ||
unsubscribeReset(); | ||
}; | ||
} | ||
|
||
constructor(settings: LiveQuerySettings) { | ||
this.query = settings.query; | ||
} | ||
|
||
match(operation: RecordOperation): boolean { | ||
const change = recordOperationChange(operation); | ||
return !!this.query.expressions.find(expression => | ||
this._queryExpressionMatchChange(expression, change) | ||
); | ||
} | ||
|
||
protected _queryExpressionMatchChange( | ||
expression: QueryExpression, | ||
change: RecordChange | ||
): boolean { | ||
switch (expression.op) { | ||
case 'findRecord': | ||
return this._findRecordQueryExpressionMatchChange( | ||
expression as FindRecord, | ||
change | ||
); | ||
case 'findRecords': | ||
return this._findRecordsQueryExpressionMatchChange( | ||
expression as FindRecords, | ||
change | ||
); | ||
case 'findRelatedRecord': | ||
return this._findRelatedRecordQueryExpressionMatchChange( | ||
expression as FindRelatedRecord, | ||
change | ||
); | ||
case 'findRelatedRecords': | ||
return this._findRelatedRecordsQueryExpressionMatchChange( | ||
expression as FindRelatedRecords, | ||
change | ||
); | ||
default: | ||
return true; | ||
} | ||
} | ||
|
||
protected _findRecordQueryExpressionMatchChange( | ||
expression: FindRecord, | ||
change: RecordChange | ||
): boolean { | ||
return equalRecordIdentities(expression.record, change); | ||
} | ||
|
||
protected _findRecordsQueryExpressionMatchChange( | ||
expression: FindRecords, | ||
change: RecordChange | ||
): boolean { | ||
if (expression.type) { | ||
return expression.type === change.type; | ||
} else if (expression.records) { | ||
for (let record of expression.records) { | ||
if (record.type === change.type) { | ||
return true; | ||
} | ||
} | ||
return false; | ||
} | ||
return true; | ||
} | ||
|
||
protected _findRelatedRecordQueryExpressionMatchChange( | ||
expression: FindRelatedRecord, | ||
change: RecordChange | ||
): boolean { | ||
return ( | ||
equalRecordIdentities(expression.record, change) && | ||
(change.relationships.includes(expression.relationship) || change.remove) | ||
); | ||
} | ||
|
||
protected _findRelatedRecordsQueryExpressionMatchChange( | ||
expression: FindRelatedRecords, | ||
change: RecordChange | ||
): boolean { | ||
const { type } = this.schema.getRelationship( | ||
expression.record.type, | ||
expression.relationship | ||
); | ||
|
||
if (Array.isArray(type) && type.find(type => type === change.type)) { | ||
return true; | ||
} else if (type === change.type) { | ||
return true; | ||
} | ||
|
||
return ( | ||
equalRecordIdentities(expression.record, change) && | ||
(change.relationships.includes(expression.relationship) || change.remove) | ||
); | ||
} | ||
} | ||
|
||
let resolvedPromise: Promise<void>; | ||
const nextTick = | ||
typeof process === 'object' && typeof process.nextTick === 'function' | ||
? function(fn: () => void) { | ||
if (!resolvedPromise) { | ||
resolvedPromise = Promise.resolve(); | ||
} | ||
resolvedPromise.then(() => { | ||
process.nextTick(fn); | ||
}); | ||
} | ||
: window.setImmediate || setTimeout; | ||
|
||
function onceTick(fn: () => void) { | ||
return function tick() { | ||
if (!ticks.has(tick)) { | ||
ticks.add(tick); | ||
nextTick(() => { | ||
fn(); | ||
cancelTick(tick); | ||
}); | ||
} | ||
}; | ||
} | ||
|
||
function cancelTick(tick: () => void) { | ||
ticks.delete(tick); | ||
} | ||
|
||
const ticks = new WeakSet(); |
36 changes: 36 additions & 0 deletions
36
packages/@orbit/record-cache/src/live-query/sync-live-query.ts
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,36 @@ | ||
import { RecordException } from '@orbit/data'; | ||
import { QueryResult } from '../query-result'; | ||
import { SyncRecordCache } from '../sync-record-cache'; | ||
import { LiveQuery, LiveQuerySettings } from './live-query'; | ||
|
||
export interface SyncLiveQuerySettings extends LiveQuerySettings { | ||
cache: SyncRecordCache; | ||
} | ||
|
||
export class SyncLiveQuery extends LiveQuery { | ||
cache: SyncRecordCache; | ||
|
||
constructor(settings: SyncLiveQuerySettings) { | ||
super(settings); | ||
this.cache = settings.cache; | ||
} | ||
|
||
get schema() { | ||
return this.cache.schema; | ||
} | ||
|
||
executeQuery( | ||
onNext: (result: QueryResult) => void, | ||
onError?: (error: RecordException) => void | ||
): void { | ||
try { | ||
onNext(this.cache.query(this.query)); | ||
} catch (error) { | ||
(onError || defaultOnError)(error); | ||
} | ||
} | ||
} | ||
|
||
function defaultOnError(error: RecordException) { | ||
throw error; | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,58 @@ | ||
import { | ||
Record, | ||
cloneRecordIdentity, | ||
RecordIdentity, | ||
RecordOperation | ||
} from '@orbit/data'; | ||
|
||
export interface RecordChange extends RecordIdentity { | ||
keys: string[]; | ||
attributes: string[]; | ||
relationships: string[]; | ||
remove: boolean; | ||
} | ||
|
||
export function recordOperationChange( | ||
operation: RecordOperation | ||
): RecordChange { | ||
const record = operation.record as Record; | ||
const change: RecordChange = { | ||
...cloneRecordIdentity(record), | ||
remove: false, | ||
keys: [], | ||
attributes: [], | ||
relationships: [] | ||
}; | ||
|
||
switch (operation.op) { | ||
case 'addRecord': | ||
case 'updateRecord': | ||
if (record.keys) { | ||
change.keys = Object.keys(record.keys); | ||
} | ||
if (record.attributes) { | ||
change.attributes = Object.keys(record.attributes); | ||
} | ||
if (record.relationships) { | ||
change.relationships = Object.keys(record.relationships); | ||
} | ||
break; | ||
case 'replaceAttribute': | ||
change.attributes = [operation.attribute]; | ||
break; | ||
case 'replaceKey': | ||
change.keys = [operation.key]; | ||
break; | ||
case 'replaceRelatedRecord': | ||
case 'replaceRelatedRecords': | ||
case 'addToRelatedRecords': | ||
case 'removeFromRelatedRecords': | ||
change.relationships = [operation.relationship]; | ||
break; | ||
case 'removeRecord': | ||
change.remove = true; | ||
break; | ||
} | ||
|
||
return change; | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.