1use crate::{
2 ExtractionBackfillRecord, ExtractionBackfillRun, ExtractionBackfillStatus, ExtractionPlan,
3 extraction_backfill_integrity_is_valid, extraction_input_digest,
4 extraction_plan_integrity_is_valid,
5};
6use schemars::JsonSchema;
7use serde::{Deserialize, Serialize};
8use std::collections::{BTreeMap, BTreeSet};
9use std::fmt;
10
11pub const EXTRACTION_RECONCILIATION_PROTOCOL: &str = "lenso.extraction-reconciliation.v1";
12
13#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
14#[serde(rename_all = "camelCase")]
15pub struct ExtractionRelationshipCount {
16 pub relationship: String,
17 pub count: u64,
18}
19
20impl ExtractionRelationshipCount {
21 #[must_use]
22 pub fn new(relationship: impl Into<String>, count: u64) -> Self {
23 Self {
24 relationship: relationship.into(),
25 count,
26 }
27 }
28}
29
30#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
31#[serde(rename_all = "camelCase")]
32pub struct ExtractionNormalizedField {
33 pub json_pointer: String,
34 pub reason: String,
35}
36
37#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
38#[serde(rename_all = "camelCase")]
39pub struct ExtractionBusinessInvariant {
40 pub invariant_id: String,
41 pub passed: bool,
42 pub evidence: String,
43}
44
45impl ExtractionBusinessInvariant {
46 #[must_use]
47 pub fn passed(invariant_id: impl Into<String>, evidence: impl Into<String>) -> Self {
48 Self {
49 invariant_id: invariant_id.into(),
50 passed: true,
51 evidence: evidence.into(),
52 }
53 }
54
55 #[must_use]
56 pub fn failed(invariant_id: impl Into<String>, evidence: impl Into<String>) -> Self {
57 Self {
58 invariant_id: invariant_id.into(),
59 passed: false,
60 evidence: evidence.into(),
61 }
62 }
63}
64
65#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
66#[serde(rename_all = "camelCase")]
67pub struct ExtractionSourceSnapshot {
68 pub source_high_water_mark: String,
69 pub records: Vec<ExtractionBackfillRecord>,
70 #[serde(default)]
71 pub relationship_counts: Vec<ExtractionRelationshipCount>,
72}
73
74#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
75#[serde(rename_all = "camelCase")]
76pub struct ExtractionReconciliationInputs {
77 pub backfill: ExtractionBackfillRun,
78 pub source: ExtractionSourceSnapshot,
79 #[serde(default, skip_serializing_if = "Option::is_none")]
80 pub destination_records: Option<Vec<ExtractionBackfillRecord>>,
81 #[serde(default)]
82 pub destination_relationship_counts: Vec<ExtractionRelationshipCount>,
83 #[serde(default)]
84 pub normalized_fields: Vec<ExtractionNormalizedField>,
85 #[serde(default)]
86 pub business_invariants: Vec<ExtractionBusinessInvariant>,
87}
88
89#[derive(Debug)]
90pub struct ExtractionReconciliationReadError {
91 pub message: String,
92}
93
94impl fmt::Display for ExtractionReconciliationReadError {
95 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
96 formatter.write_str(&self.message)
97 }
98}
99
100impl std::error::Error for ExtractionReconciliationReadError {}
101
102#[derive(
103 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
104)]
105#[serde(rename_all = "snake_case")]
106pub enum ExtractionReconciliationStatus {
107 Matched,
108 Blocked,
109}
110
111#[derive(
112 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
113)]
114#[serde(rename_all = "snake_case")]
115pub enum ExtractionReconciliationIssueCode {
116 BackfillIncomplete,
117 SourceStateChanged,
118 RecordCountMismatch,
119 StableIdentityMismatch,
120 FieldDigestMismatch,
121 RelationshipCountMismatch,
122 BusinessInvariantMismatch,
123 NormalizationReasonMissing,
124}
125
126#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
127#[serde(rename_all = "camelCase")]
128pub struct ExtractionReconciliationIssue {
129 pub code: ExtractionReconciliationIssueCode,
130 pub subject: String,
131 pub detail: String,
132 pub next_actions: Vec<String>,
133}
134
135#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
136#[serde(rename_all = "camelCase")]
137pub struct ExtractionReconciliationEvidence {
138 pub kind: String,
139 pub subject: String,
140 pub digest: String,
141 pub detail: String,
142}
143
144#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
145#[serde(rename_all = "camelCase")]
146pub struct ExtractionReconciliationEffects {
147 pub reads_source_snapshot: bool,
148 pub reads_candidate_snapshot: bool,
149 pub mutates_source: bool,
150 pub mutates_candidate: bool,
151 pub changes_authority: bool,
152}
153
154#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
155#[serde(rename_all = "camelCase")]
156pub struct ExtractionReconciliationResult {
157 pub protocol: String,
158 pub reconciliation_id: String,
159 pub reconciliation_digest: String,
160 pub status: ExtractionReconciliationStatus,
161 pub plan_id: String,
162 pub plan_digest: String,
163 pub source_high_water_mark: String,
164 pub destination_checkpoint: String,
165 pub source_record_count: u64,
166 pub destination_record_count: u64,
167 pub issues: Vec<ExtractionReconciliationIssue>,
168 pub evidence: Vec<ExtractionReconciliationEvidence>,
169 pub normalized_fields: Vec<ExtractionNormalizedField>,
170 pub linked_authority_remains_authoritative: bool,
171 pub candidate_writes_admitted: bool,
172 pub effects: ExtractionReconciliationEffects,
173}
174
175#[must_use]
176pub fn reconcile_extraction_data(
177 mut inputs: ExtractionReconciliationInputs,
178) -> ExtractionReconciliationResult {
179 inputs
180 .source
181 .records
182 .sort_by(|left, right| left.stable_id.cmp(&right.stable_id));
183 inputs
184 .source
185 .relationship_counts
186 .sort_by(|left, right| left.relationship.cmp(&right.relationship));
187 inputs
188 .destination_relationship_counts
189 .sort_by(|left, right| left.relationship.cmp(&right.relationship));
190 inputs
191 .normalized_fields
192 .sort_by(|left, right| left.json_pointer.cmp(&right.json_pointer));
193 inputs
194 .business_invariants
195 .sort_by(|left, right| left.invariant_id.cmp(&right.invariant_id));
196
197 let backfill = &inputs.backfill;
198 let destination_records = inputs
199 .destination_records
200 .as_ref()
201 .unwrap_or(&backfill.destination_records);
202 let checkpoint = backfill
203 .progress
204 .destination_checkpoint
205 .clone()
206 .unwrap_or_default();
207 let mut issues = Vec::new();
208 let mut evidence = Vec::new();
209 if !extraction_backfill_integrity_is_valid(backfill)
210 || backfill.status != ExtractionBackfillStatus::Succeeded
211 {
212 push_issue(
213 &mut issues,
214 ExtractionReconciliationIssueCode::BackfillIncomplete,
215 "backfill",
216 "Backfill is not complete and integrity-valid.",
217 "Resume the checkpointed backfill before reconciliation.",
218 );
219 }
220 if inputs.source.source_high_water_mark != backfill.progress.source_high_water_mark {
221 push_issue(
222 &mut issues,
223 ExtractionReconciliationIssueCode::SourceStateChanged,
224 "source-high-water-mark",
225 "The source snapshot no longer matches the backfill high-water mark.",
226 "Capture a new bounded source snapshot and backfill its delta.",
227 );
228 }
229 for field in &inputs.normalized_fields {
230 if field.json_pointer.trim().is_empty() || field.reason.trim().is_empty() {
231 push_issue(
232 &mut issues,
233 ExtractionReconciliationIssueCode::NormalizationReasonMissing,
234 &field.json_pointer,
235 "Every normalized or ignored field needs a specific reason.",
236 "Declare the exact field and reviewable normalization reason.",
237 );
238 }
239 }
240 let source_count = u64::try_from(inputs.source.records.len()).unwrap_or(u64::MAX);
241 let destination_count = u64::try_from(destination_records.len()).unwrap_or(u64::MAX);
242 if source_count != destination_count {
243 push_issue(
244 &mut issues,
245 ExtractionReconciliationIssueCode::RecordCountMismatch,
246 "records",
247 &format!("Source has {source_count} records; destination has {destination_count}."),
248 "Repair or resume backfill, then retry reconciliation.",
249 );
250 }
251 let source_ids = stable_ids(&inputs.source.records);
252 let destination_ids = stable_ids(destination_records);
253 if source_ids != destination_ids {
254 push_issue(
255 &mut issues,
256 ExtractionReconciliationIssueCode::StableIdentityMismatch,
257 "stable-identities",
258 "Source and destination stable record identities differ.",
259 "Compare missing and unexpected identities before retrying.",
260 );
261 }
262 let normalized_source = normalize_records(&inputs.source.records, &inputs.normalized_fields);
263 let normalized_destination = normalize_records(destination_records, &inputs.normalized_fields);
264 let source_digest = digest(&normalized_source);
265 let destination_digest = digest(&normalized_destination);
266 evidence.push(ExtractionReconciliationEvidence {
267 kind: "record_digest".to_owned(),
268 subject: "source".to_owned(),
269 digest: source_digest.clone(),
270 detail: format!(
271 "{source_count} records at {}",
272 inputs.source.source_high_water_mark
273 ),
274 });
275 evidence.push(ExtractionReconciliationEvidence {
276 kind: "record_digest".to_owned(),
277 subject: "destination".to_owned(),
278 digest: destination_digest.clone(),
279 detail: format!("{destination_count} records at {checkpoint}"),
280 });
281 if source_ids == destination_ids && source_digest != destination_digest {
282 push_issue(
283 &mut issues,
284 ExtractionReconciliationIssueCode::FieldDigestMismatch,
285 "declared-fields",
286 "Stable records have different declared field digests.",
287 "Inspect per-record differences without broadening normalization.",
288 );
289 }
290 let source_relationships = relationship_map(&inputs.source.relationship_counts);
291 let destination_relationships = relationship_map(&inputs.destination_relationship_counts);
292 if source_relationships != destination_relationships {
293 push_issue(
294 &mut issues,
295 ExtractionReconciliationIssueCode::RelationshipCountMismatch,
296 "relationships",
297 "Declared source and destination relationship counts differ.",
298 "Repair the missing relationship copy and retry reconciliation.",
299 );
300 }
301 for invariant in &inputs.business_invariants {
302 evidence.push(ExtractionReconciliationEvidence {
303 kind: "business_invariant".to_owned(),
304 subject: invariant.invariant_id.clone(),
305 digest: extraction_input_digest(invariant.evidence.as_bytes()),
306 detail: invariant.evidence.clone(),
307 });
308 if !invariant.passed {
309 push_issue(
310 &mut issues,
311 ExtractionReconciliationIssueCode::BusinessInvariantMismatch,
312 &invariant.invariant_id,
313 &invariant.evidence,
314 "Remediate the declared business invariant and retry reconciliation.",
315 );
316 }
317 }
318 issues.sort();
319 evidence.sort();
320 let status = if issues.is_empty() {
321 ExtractionReconciliationStatus::Matched
322 } else {
323 ExtractionReconciliationStatus::Blocked
324 };
325 let identity_digest = digest(&(
326 backfill.scope.plan_id.as_str(),
327 inputs.source.source_high_water_mark.as_str(),
328 checkpoint.as_str(),
329 ));
330 let mut result = ExtractionReconciliationResult {
331 protocol: EXTRACTION_RECONCILIATION_PROTOCOL.to_owned(),
332 reconciliation_id: format!("extraction-reconciliation:{identity_digest}"),
333 reconciliation_digest: String::new(),
334 status,
335 plan_id: backfill.scope.plan_id.clone(),
336 plan_digest: backfill.scope.plan_digest.clone(),
337 source_high_water_mark: inputs.source.source_high_water_mark,
338 destination_checkpoint: checkpoint,
339 source_record_count: source_count,
340 destination_record_count: destination_count,
341 issues,
342 evidence,
343 normalized_fields: inputs.normalized_fields,
344 linked_authority_remains_authoritative: true,
345 candidate_writes_admitted: false,
346 effects: ExtractionReconciliationEffects {
347 reads_source_snapshot: true,
348 reads_candidate_snapshot: true,
349 ..ExtractionReconciliationEffects::default()
350 },
351 };
352 result.reconciliation_digest = digest(&result_without_digest(&result));
353 result
354}
355
356pub async fn reconcile_postgres_extraction_service_data(
360 source_pool: &sqlx::PgPool,
361 destination_pool: &sqlx::PgPool,
362 plan: &ExtractionPlan,
363 backfill: ExtractionBackfillRun,
364 normalized_fields: Vec<ExtractionNormalizedField>,
365 business_invariants: Vec<ExtractionBusinessInvariant>,
366) -> Result<ExtractionReconciliationResult, ExtractionReconciliationReadError> {
367 if !extraction_plan_integrity_is_valid(plan)
368 || !extraction_backfill_integrity_is_valid(&backfill)
369 || backfill.scope.plan_id != plan.plan_id
370 || backfill.scope.plan_digest != plan.plan_digest
371 || plan.data_mapping.tables.len() != 1
372 {
373 return Err(read_error(
374 "PostgreSQL reconciliation requires one integrity-valid plan-scoped table mapping.",
375 ));
376 }
377 let mapping = &plan.data_mapping.tables[0];
378 let cursor = mapping
379 .cursors
380 .iter()
381 .find(|cursor| cursor.trustworthy)
382 .ok_or_else(|| read_error("PostgreSQL reconciliation requires a trustworthy cursor."))?;
383 let source_table = quoted_relation(&mapping.source_table)?;
384 let destination_table = quoted_relation(&mapping.destination_table)?;
385 let cursor_column = quoted_identifier(&cursor.column)?;
386 let cursor_name = cursor.column.as_str();
387 let high_water_cursor = format!(
388 "(jsonb_populate_record(null::{source_table}, jsonb_build_object('{cursor_name}', $1::text))).{cursor_column}"
389 );
390 let source_rows = sqlx::query_as::<_, (String, serde_json::Value)>(sqlx::AssertSqlSafe(
391 format!(
392 "select {cursor_column}::text, to_jsonb(source_row) from {source_table} source_row where {cursor_column} <= {high_water_cursor} order by {cursor_column}"
393 ),
394 ))
395 .bind(&backfill.progress.source_high_water_mark)
396 .fetch_all(source_pool)
397 .await
398 .map_err(|source| read_error(format!("Source snapshot failed: {source}")))?;
399 let destination_rows =
400 sqlx::query_as::<_, (String, serde_json::Value)>(sqlx::AssertSqlSafe(format!(
401 "select {cursor_column}::text, to_jsonb(destination_row) from {destination_table} destination_row order by {cursor_column}"
402 )))
403 .fetch_all(destination_pool)
404 .await
405 .map_err(|source| read_error(format!("Candidate snapshot failed: {source}")))?;
406 let records = |rows: Vec<(String, serde_json::Value)>| {
407 rows.into_iter()
408 .map(|(stable_id, value)| ExtractionBackfillRecord::new(stable_id, value))
409 .collect::<Vec<_>>()
410 };
411 Ok(reconcile_extraction_data(ExtractionReconciliationInputs {
412 source: ExtractionSourceSnapshot {
413 source_high_water_mark: backfill.progress.source_high_water_mark.clone(),
414 records: records(source_rows),
415 relationship_counts: Vec::new(),
416 },
417 destination_records: Some(records(destination_rows)),
418 backfill,
419 destination_relationship_counts: Vec::new(),
420 normalized_fields,
421 business_invariants,
422 }))
423}
424
425fn quoted_relation(value: &str) -> Result<String, ExtractionReconciliationReadError> {
426 value
427 .split('.')
428 .map(quoted_identifier)
429 .collect::<Result<Vec<_>, _>>()
430 .map(|parts| parts.join("."))
431}
432
433fn quoted_identifier(value: &str) -> Result<String, ExtractionReconciliationReadError> {
434 if value.is_empty()
435 || !value
436 .chars()
437 .all(|character| character.is_ascii_alphanumeric() || character == '_')
438 {
439 return Err(read_error(format!(
440 "Unsafe PostgreSQL identifier `{value}` in Extraction Plan."
441 )));
442 }
443 Ok(format!("\"{value}\""))
444}
445
446fn read_error(message: impl Into<String>) -> ExtractionReconciliationReadError {
447 ExtractionReconciliationReadError {
448 message: message.into(),
449 }
450}
451
452fn normalize_records(
453 records: &[ExtractionBackfillRecord],
454 normalized_fields: &[ExtractionNormalizedField],
455) -> Vec<serde_json::Value> {
456 records
457 .iter()
458 .map(|record| {
459 let mut value = record.value.clone();
460 for field in normalized_fields {
461 remove_pointer(&mut value, &field.json_pointer);
462 }
463 value
464 })
465 .collect()
466}
467
468fn remove_pointer(value: &mut serde_json::Value, pointer: &str) {
469 let Some((parent, key)) = pointer.rsplit_once('/') else {
470 return;
471 };
472 if let Some(serde_json::Value::Object(object)) = value.pointer_mut(parent) {
473 object.remove(key);
474 }
475}
476
477fn stable_ids(records: &[ExtractionBackfillRecord]) -> BTreeSet<&str> {
478 records
479 .iter()
480 .map(|record| record.stable_id.as_str())
481 .collect()
482}
483
484fn relationship_map(counts: &[ExtractionRelationshipCount]) -> BTreeMap<&str, u64> {
485 counts
486 .iter()
487 .map(|count| (count.relationship.as_str(), count.count))
488 .collect()
489}
490
491fn push_issue(
492 issues: &mut Vec<ExtractionReconciliationIssue>,
493 code: ExtractionReconciliationIssueCode,
494 subject: impl Into<String>,
495 detail: impl Into<String>,
496 next_action: impl Into<String>,
497) {
498 issues.push(ExtractionReconciliationIssue {
499 code,
500 subject: subject.into(),
501 detail: detail.into(),
502 next_actions: vec![next_action.into()],
503 });
504}
505
506fn digest(value: &impl Serialize) -> String {
507 extraction_input_digest(
508 &serde_json::to_vec(value).expect("Extraction reconciliation values must serialize"),
509 )
510}
511
512fn result_without_digest(
513 result: &ExtractionReconciliationResult,
514) -> ExtractionReconciliationResult {
515 let mut value = result.clone();
516 value.reconciliation_digest.clear();
517 value
518}
519
520#[must_use]
521pub fn extraction_reconciliation_integrity_is_valid(
522 result: &ExtractionReconciliationResult,
523) -> bool {
524 result.protocol == EXTRACTION_RECONCILIATION_PROTOCOL
525 && result.reconciliation_digest == digest(&result_without_digest(result))
526}