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);
}