|
@@ -18,7 +18,7 @@ const configSchema = {
|
|
|
timestampField: {
|
|
timestampField: {
|
|
|
type: 'string',
|
|
type: 'string',
|
|
|
title: 'Timestamp Field',
|
|
title: 'Timestamp Field',
|
|
|
- description: 'Field holding the last-modified time. Defaults to the _updated_at the database maintains itself, in nanoseconds',
|
|
|
|
|
|
|
+ description: 'Field holding the last-modified time. Defaults to the _updated_at the database maintains itself, in nanoseconds. A document with no readable value in this field is never reported',
|
|
|
default: '_updated_at'
|
|
default: '_updated_at'
|
|
|
},
|
|
},
|
|
|
filter: {
|
|
filter: {
|
|
@@ -118,12 +118,24 @@ async function execute(config, input, context) {
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
let highWater = since;
|
|
let highWater = since;
|
|
|
|
|
+ let noStamp = 0;
|
|
|
for (const document of documents) {
|
|
for (const document of documents) {
|
|
|
const stamp = stampOf(document);
|
|
const stamp = stampOf(document);
|
|
|
|
|
+ if (stamp === 0) {
|
|
|
|
|
+ noStamp++;
|
|
|
|
|
+ }
|
|
|
if (stamp > highWater) {
|
|
if (stamp > highWater) {
|
|
|
highWater = stamp;
|
|
highWater = stamp;
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
+ if (noStamp > 0) {
|
|
|
|
|
+ // A document with no readable value in timestampField sorts as if it
|
|
|
|
|
+ // predates everything, and once the mark advances past 0 it can never
|
|
|
|
|
+ // satisfy the strict > filter below - it is never reported, silently,
|
|
|
|
|
+ // unless this is logged.
|
|
|
|
|
+ smartbotic.log.warn('Database Change: ' + noStamp + ' document(s) in ' + collection +
|
|
|
|
|
+ ' had no readable "' + timestampField + '" value and will never be reported');
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
let changed = [];
|
|
let changed = [];
|
|
|
if (isFirstRun && config.emitOnFirstRun !== true) {
|
|
if (isFirstRun && config.emitOnFirstRun !== true) {
|
|
@@ -134,9 +146,21 @@ async function execute(config, input, context) {
|
|
|
.filter(function (document) { return stampOf(document) > since; })
|
|
.filter(function (document) { return stampOf(document) > since; })
|
|
|
.sort(function (left, right) { return stampOf(left) - stampOf(right); });
|
|
.sort(function (left, right) { return stampOf(left) - stampOf(right); });
|
|
|
if (changed.length > maxDocuments) {
|
|
if (changed.length > maxDocuments) {
|
|
|
- smartbotic.log.warn('Database Change: ' + changed.length + ' documents changed, reporting the oldest ' +
|
|
|
|
|
- maxDocuments + '. The rest arrive on the next run');
|
|
|
|
|
- changed = changed.slice(0, maxDocuments);
|
|
|
|
|
|
|
+ let cut = maxDocuments;
|
|
|
|
|
+ const boundary = stampOf(changed[cut - 1]);
|
|
|
|
|
+ // Cutting through a group that shares one timestamp would strand
|
|
|
|
|
+ // the rest of that group: the mark advances to the shared stamp
|
|
|
|
|
+ // and the next run's strict > excludes them forever. The slice
|
|
|
|
|
+ // grows to take the whole tie, even though that overshoots the
|
|
|
|
|
+ // limit. If every changed document shares one stamp, cut grows to
|
|
|
|
|
+ // all of them, which is correct - nothing is stranded.
|
|
|
|
|
+ while (cut < changed.length && stampOf(changed[cut]) === boundary) {
|
|
|
|
|
+ cut++;
|
|
|
|
|
+ }
|
|
|
|
|
+ smartbotic.log.warn('Database Change: ' + changed.length + ' documents changed, reporting ' +
|
|
|
|
|
+ cut + (cut > maxDocuments ? ' (over the configured ' + maxDocuments + ', to avoid splitting a tied timestamp)' : '') +
|
|
|
|
|
+ '. The rest arrive on the next run');
|
|
|
|
|
+ changed = changed.slice(0, cut);
|
|
|
// The mark only advances as far as what was actually reported, or
|
|
// The mark only advances as far as what was actually reported, or
|
|
|
// the remainder would be skipped rather than deferred.
|
|
// the remainder would be skipped rather than deferred.
|
|
|
highWater = stampOf(changed[changed.length - 1]);
|
|
highWater = stampOf(changed[changed.length - 1]);
|