Skip to main content

lenso_service/
extraction_reconciliation.rs

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
356/// Reads fresh plan-scoped source and candidate snapshots from PostgreSQL and
357/// reconciles those durable rows instead of trusting the Backfill receipt as
358/// the candidate state.
359pub 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}