configure static method
void
configure({
- required DVCaptureConfig config,
- required DVDatabaseAdapter database,
- required List<
DVStudioModelSpec> models, - String? secret(
- String name
- DVDatabaseAdapter open(
- String connection
- DateTime clock()?,
Configures change capture from config over database, the
application's own.
models are the data models' specs as the generator wrote them; the
captured ones are recorded to the log, and a destination takes those it
names. secret reads a destination's connection by name -- DV.Secrets
when null -- and open opens it, by default as a database connection.
A destination whose connection is not set is skipped, and said with
DV-CDC-009: its changes wait in the log.
Throws DVCaptureConfigError (DV-CDC-008) for a destination naming a
data model that is not captured.
Implementation
static void configure({
required DVCaptureConfig config,
required DVDatabaseAdapter database,
required List<DVStudioModelSpec> models,
String? Function(String name)? secret,
DVDatabaseAdapter Function(String connection)? open,
DateTime Function()? clock,
}) {
final Map<String, DVStudioModelSpec> captured = <String, DVStudioModelSpec>{
for (final DVStudioModelSpec spec in models)
if (spec.capture) spec.id: spec,
};
config.checkModels(captured.keys.toSet());
final DVCapture log = DVCapture(
database: database,
retention: config.retention,
clock: clock,
);
final Map<String, DVRecordTable> tables = <String, DVRecordTable>{
for (final MapEntry<String, DVStudioModelSpec> e in captured.entries)
e.key: _table(e.value, log, database),
};
for (final DVRecordTable table in tables.values) {
log.track(table);
}
final String? Function(String) read = secret ?? const DVSecrets().maybeGet;
final DVDatabaseAdapter Function(String) opening =
open ?? (String url) => DVDatabaseConnection.parse(url).open();
final Map<String, String> skipped = <String, String>{};
final List<DVCaptureConsumer> consumers = <DVCaptureConsumer>[];
final Map<String, List<DVRecordTable>> takes =
<String, List<DVRecordTable>>{};
for (final DVCaptureDestination destination in config.destinations) {
final String? connection = read(destination.connection)?.trim();
DVDatabaseAdapter? store;
String? why;
if (connection == null || connection.isEmpty) {
why = '${destination.connection} is not set';
} else {
try {
store = opening(connection);
} on FormatException catch (error) {
// Not the value: it carries credentials.
why = '${destination.connection} cannot be read: ${error.message}';
}
}
if (store == null) {
skipped[destination.name] = 'DV-CDC-009';
DVObservability.log(
'Change capture destination ${destination.name} is not delivered '
'to: $why. Its changes wait in the log for '
'${config.retention.inHours} hours; set it before then, or the '
'destination is backfilled when it comes back.',
level: DVLogLevel.warn,
code: 'DV-CDC-009',
);
continue;
}
final List<DVRecordTable> chosen = <DVRecordTable>[
for (final MapEntry<String, DVRecordTable> e in tables.entries)
if (destination.models == null ||
destination.models!.contains(e.key))
e.value,
];
takes[destination.name] = chosen;
consumers.add(log.consumer(
destination.name,
sink: DVWarehouseSink(database: store, name: destination.name),
models: <String>{for (final DVRecordTable t in chosen) t.table},
lagThreshold: destination.lagThreshold,
));
}
DVCapture.configure(log);
// An erasure takes the erased records out of every destination now,
// rather than at the next delivery, and is incomplete when one of them
// cannot be reached.
DVPrivacyRuntime.installAdapters(<DVCapturePrivacyAdapter>[
DVCapturePrivacyAdapter(
capture: log,
sinks: <DVCaptureSink>[
for (final DVCaptureConsumer c in consumers) c.sink,
],
),
]);
_runtime = _Runtime(
log: log,
consumers: consumers,
takes: takes,
skipped: skipped,
);
}