lag method
How far behind this consumer is, recorded as the gauges
dv_capture_lag_seconds and dv_capture_lag_changes.
Measured rather than assumed: the failure here is not an error but a destination quietly hours old while everyone reads it as current.
Implementation
Future<DVCaptureLag> lag() async {
final int at = await checkpoint();
final int head = await capture.head();
final List<Map<String, Object?>> pending = await capture._published(
after: at,
through: head,
recordsOnly: true,
fields: const <String>['occurred_at'],
);
final Duration age = pending.isEmpty
? Duration.zero
: capture._clock().difference(
DateTime.parse('${pending.first['occurred_at']}'));
final Duration measured = age.isNegative ? Duration.zero : age;
final Map<String, String> labels = <String, String>{'consumer': name};
DVObservability.metrics
.gauge('dv_capture_lag_seconds', labels,
'Age of the oldest captured change not yet delivered')
.set(measured.inSeconds.toDouble());
DVObservability.metrics
.gauge('dv_capture_lag_changes', labels,
'Captured changes not yet delivered')
.set(pending.length.toDouble());
final List<String> codes = <String>[];
final Duration? threshold = lagThreshold;
if (threshold != null && measured > threshold) {
codes.add('DV-CDC-003');
DVObservability.log(
'$name is ${measured.inMinutes} minutes behind the capture log, past '
'its ${threshold.inMinutes}-minute threshold.',
level: DVLogLevel.warn,
code: 'DV-CDC-003',
context: <String, Object?>{'changes': pending.length},
);
}
return DVCaptureLag(changes: pending.length, age: measured, codes: codes);
}