Skip to main content

type_bridge_migration/
backfill.rs

1//! Backfill count derivation for [`StepKind::Backfill`] steps.
2//!
3//! TypeDB's `insert` answer is [`QueryResult::Ok`] — it carries no affected-row
4//! count.  This module derives matched/inserted/skipped counts via bracketing
5//! `reduce $c = count;` read queries around the write (D2).
6//!
7//! # Cost
8//!
9//! Two extra read-transaction count queries are issued per backfill step
10//! (guarded count before insert, total-source count before insert).  This is
11//! acceptable for one-shot migrations and is the only way to surface counts given
12//! TypeDB's write answer surface.
13//!
14//! # Invariant 2 compliance
15//!
16//! The count queries are composed directly from the carried `step.forward` match
17//! clause — no op-semantic re-derivation occurs here.  The backfill step's
18//! `forward` text is the single source of truth for the query shape.
19
20use type_bridge_orm::Database;
21use type_bridge_orm::session::TransactionContext;
22use type_bridge_orm::session::backend::QueryResult;
23
24use crate::error::MigrationError;
25use crate::plan::ExecutionStep;
26use crate::state::require_legacy_writer_open_in_transaction;
27
28use serde::{Deserialize, Serialize};
29use type_bridge_orm::TxType;
30
31/// Per-step backfill count result.
32#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
33pub struct BackfillResult {
34    /// Zero-based index of the step within the migration's step list.
35    pub step_index: usize,
36    /// Number of owner instances that matched the source predicate (the total
37    /// candidate set, including those already having the destination).
38    pub matched: u64,
39    /// Number of instances that were actually backfilled (destination inserted).
40    ///
41    /// Computed as the count of the guarded match (insert-if-absent inserts
42    /// exactly the guarded set).
43    pub inserted: u64,
44    /// Number of instances skipped because they already held the destination
45    /// attribute (`matched - inserted`).
46    pub skipped: u64,
47    /// Number of instances with conflicting source values.  Always `0` for the
48    /// single-valued v1 `CopyAttribute` op; reserved for future multi-valued ops.
49    pub conflicts: u64,
50}
51
52/// Prepared backfill whose write query has succeeded but is not yet committed.
53///
54/// Kept crate-private so the recovery executor can durably emit its
55/// before-commit event in the only safe gap between query success and commit.
56pub(crate) struct PreparedBackfill {
57    pub(crate) transaction: TransactionContext,
58    pub(crate) result: BackfillResult,
59}
60
61/// Execute a backfill [`ExecutionStep`] against `db`, deriving and returning
62/// matched/inserted/skipped counts.
63///
64/// # Execution sequence
65///
66/// 1. Run a `Read` count query over the **guarded** match clause (the
67///    insert-if-absent predicate — rows that will be inserted).
68/// 2. Run a `Read` count query over the **unguarded** match clause (all
69///    candidate rows, including those already having the destination).
70/// 3. Run the forward TypeQL under a `Write` transaction and commit.
71///
72/// Steps 1 and 2 run before the write so the counts reflect the pre-write state.
73///
74/// # Invariant 2
75///
76/// Both count queries are derived from `step.forward` by string manipulation of
77/// the carried match clause.  No `OperationSpec` is inspected here.
78pub async fn execute_backfill(
79    db: &Database,
80    step: &ExecutionStep,
81    step_index: usize,
82) -> Result<BackfillResult, MigrationError> {
83    let prepared = prepare_backfill(db, step, step_index).await?;
84    prepared
85        .transaction
86        .commit()
87        .await
88        .map_err(|e| MigrationError::BackfillQuery {
89            message: format!("backfill step {step_index}: backfill write commit failed: {e}"),
90        })?;
91    Ok(prepared.result)
92}
93
94/// Execute the read counts and write query for a backfill without committing.
95pub(crate) async fn prepare_backfill(
96    db: &Database,
97    step: &ExecutionStep,
98    step_index: usize,
99) -> Result<PreparedBackfill, MigrationError> {
100    // ── Decompose the carried forward text ──────────────────────────────────
101    //
102    // The planner writes CopyAttribute steps in the form:
103    //   "match\n  ...\n  not { ... };\n...\ninsert\n  ..."
104    //
105    // We split on "\ninsert\n" to extract the match section.
106    let (match_section, _insert_section) = step
107        .forward
108        .split_once("\ninsert\n")
109        .ok_or_else(|| MigrationError::BackfillQuery {
110            message: format!(
111                "backfill step {step_index}: forward TypeQL does not contain '\\ninsert\\n' separator; \
112                 cannot compose count queries without re-deriving semantics (invariant 2)"
113            ),
114        })?;
115
116    // ── Guarded count: rows the insert will affect (insert-if-absent set) ──
117    //
118    // The match section already contains the `not { $x has <dest> $d; }` guard,
119    // so this count == inserted.
120    let guarded_count_query = format!("{match_section}\nreduce $c = count;");
121
122    // ── Unguarded count: all candidate rows (matched) ───────────────────────
123    //
124    // Strip the `not { ... }` line from the match section.  The guard line is
125    // always of the form `  not { ... };` in the planner output.
126    let unguarded_match = strip_not_guard(match_section);
127    let total_count_query = format!("{unguarded_match}\nreduce $c = count;");
128
129    // ── Run guarded count (Read tx) ─────────────────────────────────────────
130    let inserted: u64 = {
131        let ctx = db.transaction_context(TxType::Read).await.map_err(|e| {
132            MigrationError::BackfillQuery {
133                message: format!(
134                    "backfill step {step_index}: failed to open read tx for guarded count: {e}"
135                ),
136            }
137        })?;
138        let result =
139            ctx.query(&guarded_count_query)
140                .await
141                .map_err(|e| MigrationError::BackfillQuery {
142                    message: format!("backfill step {step_index}: guarded count query failed: {e}"),
143                })?;
144        // Rollback / close (read txs don't need commit; best-effort close).
145        let _ = ctx.rollback().await;
146        extract_count(result, step_index, "guarded")?
147    };
148
149    // ── Run total count (Read tx) ───────────────────────────────────────────
150    let matched: u64 = {
151        let ctx = db.transaction_context(TxType::Read).await.map_err(|e| {
152            MigrationError::BackfillQuery {
153                message: format!(
154                    "backfill step {step_index}: failed to open read tx for total count: {e}"
155                ),
156            }
157        })?;
158        let result =
159            ctx.query(&total_count_query)
160                .await
161                .map_err(|e| MigrationError::BackfillQuery {
162                    message: format!("backfill step {step_index}: total count query failed: {e}"),
163                })?;
164        let _ = ctx.rollback().await;
165        extract_count(result, step_index, "total")?
166    };
167
168    // ── Prepare the backfill write (Write tx, deliberately uncommitted) ─────
169    let transaction =
170        db.transaction_context(TxType::Write)
171            .await
172            .map_err(|e| MigrationError::BackfillQuery {
173                message: format!("backfill step {step_index}: failed to open write tx: {e}"),
174            })?;
175    if let Err(error) = require_legacy_writer_open_in_transaction(&transaction).await {
176        let _ = transaction.rollback().await;
177        return Err(MigrationError::BackfillQuery {
178            message: format!(
179                "backfill step {step_index}: legacy cutover guard rejected the write: {error}"
180            ),
181        });
182    }
183    if let Err(error) = transaction.query(&step.forward).await {
184        let _ = transaction.rollback().await;
185        return Err(MigrationError::BackfillQuery {
186            message: format!("backfill step {step_index}: backfill write query failed: {error}"),
187        });
188    }
189
190    let skipped = matched.saturating_sub(inserted);
191
192    Ok(PreparedBackfill {
193        transaction,
194        result: BackfillResult {
195            step_index,
196            matched,
197            inserted,
198            skipped,
199            conflicts: 0,
200        },
201    })
202}
203
204/// Extract the count value from a `reduce $c = count;` answer.
205///
206/// TypeDB returns the answer keyed by the reduce variable (`c`) carrying a value
207/// envelope; the in-memory mock may use any key and a bare number. Both shapes
208/// are reduced via [`scalar_to_u64`]; an empty result set falls back to `0`.
209fn extract_count(
210    result: QueryResult,
211    step_index: usize,
212    label: &str,
213) -> Result<u64, MigrationError> {
214    // Both Rows and Documents carry the reduce answer as a single object keyed
215    // by the reduce variable (`c`); the mock may use any key, so we fall back to
216    // the first value. The answer itself is a TypeDB value envelope that
217    // `scalar_to_u64` unwraps.
218    let answer = match result {
219        QueryResult::Rows(items) | QueryResult::Documents(items) => match items.first() {
220            Some(item) => item.clone(),
221            None => return Ok(0),
222        },
223        // QueryResult::Ok from a read-reduce query is unexpected; treat as 0.
224        QueryResult::Ok => return Ok(0),
225    };
226
227    let value = answer
228        .get("c")
229        .or_else(|| answer.as_object().and_then(|m| m.values().next()))
230        .ok_or_else(|| MigrationError::BackfillQuery {
231            message: format!(
232                "backfill step {step_index}: {label} count answer has no recognizable key: {answer}"
233            ),
234        })?;
235
236    scalar_to_u64(value).ok_or_else(|| MigrationError::BackfillQuery {
237        message: format!(
238            "backfill step {step_index}: {label} count value is not a number: {value}"
239        ),
240    })
241}
242
243/// Reduce a TypeDB `reduce $c = count;` answer to a `u64`.
244///
245/// TypeDB returns the count as a value document
246/// `{"category":"Value","label":"integer","value":N,"value_type":"integer"}`;
247/// the in-memory mock returns a bare number. Descend into a `"value"` field when
248/// the answer is an object so both shapes parse identically (the same envelope
249/// the state backend unwraps in `state::typedb::extract_scalar`).
250fn scalar_to_u64(value: &serde_json::Value) -> Option<u64> {
251    match value {
252        serde_json::Value::Number(_) => value.as_u64().or_else(|| value.as_f64().map(|f| f as u64)),
253        serde_json::Value::Object(map) => map.get("value").and_then(scalar_to_u64),
254        _ => None,
255    }
256}
257
258/// Remove the `not { ... };` guard line from a match section.
259///
260/// The planner always writes the guard on a single line starting with
261/// `  not { ` (two-space indent).  We strip any line whose trimmed form starts
262/// with `not {`.
263fn strip_not_guard(match_section: &str) -> String {
264    match_section
265        .lines()
266        .filter(|line| !line.trim_start().starts_with("not {"))
267        .collect::<Vec<_>>()
268        .join("\n")
269}
270
271// ── Tests ─────────────────────────────────────────────────────────────────────
272
273#[cfg(test)]
274mod tests {
275    use super::*;
276    use crate::plan::StepKind;
277    use crate::testing::{MockEvent, MockMigrationBackend};
278    use type_bridge_orm::{Database, TxType};
279
280    // Build a backfill ExecutionStep carrying the TypeQL the planner passes
281    // through from `CopyAttribute.to_typeql()` (value-copy form `has <dest> == $v`).
282    fn backfill_step() -> ExecutionStep {
283        ExecutionStep {
284            tx_type: TxType::Write,
285            kind: StepKind::Backfill,
286            operation_kind: crate::plan::OperationKind::CopyAttribute,
287            forward: "match\n  $x isa person, has old-name $v;\n  not { $x has new-name $d; };\ninsert\n  $x has new-name == $v;".to_string(),
288            reverse: Some("match $x isa person, has new-name $v;\ndelete $v of $x;".to_string()),
289        }
290    }
291
292    // ── test: correct counts returned from scripted responses ─────────────────
293
294    #[tokio::test]
295    async fn backfill_derives_counts_from_scripted_mock_responses() {
296        use serde_json::json;
297        use type_bridge_orm::session::backend::QueryResult;
298
299        // Script: guarded count = 7, total count = 10 (3 already have dest → skipped=3).
300        // Transaction order: Read (guarded), Read (total), Write.
301        let scripted = vec![
302            // Tx 1: Read — guarded count query → 7 rows to insert
303            QueryResult::Rows(vec![json!({"c": 7})]),
304            // Tx 2: Read — total count query → 10 total candidates
305            QueryResult::Rows(vec![json!({"c": 10})]),
306            // Tx 3: Write — the backfill insert → Ok
307            QueryResult::Ok,
308        ];
309
310        let (backend, log) = MockMigrationBackend::with_responses(scripted);
311        let db = Database::with_backend(Box::new(backend), "test");
312        let step = backfill_step();
313
314        let result = execute_backfill(&db, &step, 0)
315            .await
316            .expect("execute_backfill should succeed");
317
318        assert_eq!(result.step_index, 0);
319        assert_eq!(result.inserted, 7, "inserted = guarded count");
320        assert_eq!(result.matched, 10, "matched = total count");
321        assert_eq!(result.skipped, 3, "skipped = matched - inserted");
322        assert_eq!(result.conflicts, 0, "conflicts always 0 in v1");
323
324        // Verify the write ran under TxType::Write and counts under TxType::Read.
325        let events = log.lock().unwrap();
326        // Expect: OpenTx(Read), Query(Read, ...), Rollback,
327        //         OpenTx(Read), Query(Read, ...), Rollback,
328        //         OpenTx(Write), Query(Write, ...), Commit
329        assert!(
330            matches!(events[0], MockEvent::OpenTx(TxType::Read)),
331            "first tx must be Read (guarded count)"
332        );
333        assert!(
334            matches!(events[3], MockEvent::OpenTx(TxType::Read)),
335            "second tx must be Read (total count)"
336        );
337        assert!(
338            matches!(events[6], MockEvent::OpenTx(TxType::Write)),
339            "third tx must be Write (backfill insert)"
340        );
341        // Commit closes the write tx.
342        assert!(matches!(events[8], MockEvent::Commit));
343    }
344
345    // ── test: count queries are composed from the step's carried match (invariant 2) ──
346
347    #[tokio::test]
348    async fn count_queries_are_built_from_carried_match_not_re_derived() {
349        use serde_json::json;
350        use type_bridge_orm::session::backend::QueryResult;
351
352        let scripted = vec![
353            QueryResult::Rows(vec![json!({"c": 0})]),
354            QueryResult::Rows(vec![json!({"c": 0})]),
355            QueryResult::Ok,
356        ];
357
358        let (backend, log) = MockMigrationBackend::with_responses(scripted);
359        let db = Database::with_backend(Box::new(backend), "test");
360        let step = backfill_step();
361
362        execute_backfill(&db, &step, 1)
363            .await
364            .expect("execute_backfill should succeed");
365
366        let events = log.lock().unwrap();
367
368        // The first count query (guarded) must contain the match body text from the step.
369        let guarded_query = events.iter().find_map(|e| {
370            if let MockEvent::Query(TxType::Read, q) = e {
371                Some(q.as_str())
372            } else {
373                None
374            }
375        });
376        assert!(
377            guarded_query.is_some(),
378            "expected at least one Read query in the event log"
379        );
380        let q = guarded_query.unwrap();
381        // The guarded count query must contain the core match pattern from the
382        // step's forward text — proving no semantic re-derivation happened.
383        assert!(
384            q.contains("$x isa person, has old-name $v"),
385            "guarded count query must contain the step's match body; got: {q}"
386        );
387        assert!(
388            q.contains("reduce $c = count"),
389            "guarded count query must end with reduce count; got: {q}"
390        );
391    }
392
393    // ── test: write runs under TxType::Write ──────────────────────────────────
394
395    #[tokio::test]
396    async fn backfill_write_runs_under_write_tx() {
397        use serde_json::json;
398        use type_bridge_orm::session::backend::QueryResult;
399
400        let scripted = vec![
401            QueryResult::Rows(vec![json!({"c": 5})]),
402            QueryResult::Rows(vec![json!({"c": 5})]),
403            QueryResult::Ok,
404        ];
405
406        let (backend, log) = MockMigrationBackend::with_responses(scripted);
407        let db = Database::with_backend(Box::new(backend), "test");
408        let step = backfill_step();
409
410        execute_backfill(&db, &step, 2)
411            .await
412            .expect("execute_backfill should succeed");
413
414        let events = log.lock().unwrap();
415
416        // Find the Write-typed query event.
417        let write_query = events.iter().find_map(|e| {
418            if let MockEvent::Query(TxType::Write, q) = e {
419                Some(q.as_str())
420            } else {
421                None
422            }
423        });
424        assert!(write_query.is_some(), "expected a Write-typed query event");
425        // The write query must be the step's full forward text (not a count query).
426        let wq = write_query.unwrap();
427        assert!(
428            wq.contains("insert"),
429            "write query must contain 'insert'; got: {wq}"
430        );
431        assert!(
432            !wq.contains("reduce"),
433            "write query must not contain 'reduce' (it is the insert, not a count query); got: {wq}"
434        );
435    }
436
437    #[tokio::test]
438    async fn cutover_rejects_backfill_before_the_write_query() {
439        let (backend, log) = MockMigrationBackend::with_legacy_cutover();
440        let db = Database::with_backend(Box::new(backend), "test");
441
442        let error = execute_backfill(&db, &backfill_step(), 0)
443            .await
444            .expect_err("cutover must reject the backfill write");
445        assert!(
446            error
447                .to_string()
448                .contains(crate::LEGACY_WRITER_CUTOVER_MESSAGE)
449        );
450        let events = log.lock().unwrap();
451        assert!(
452            !events
453                .iter()
454                .any(|event| matches!(event, MockEvent::Query(TxType::Write, _))),
455            "no user write query may run after cutover: {events:?}"
456        );
457        assert!(
458            events
459                .iter()
460                .any(|event| matches!(event, MockEvent::Rollback)),
461            "the rejected write transaction must be rolled back: {events:?}"
462        );
463    }
464
465    // ── test: strip_not_guard removes the guard line ──────────────────────────
466
467    #[test]
468    fn strip_not_guard_removes_not_line() {
469        let match_section =
470            "match\n  $x isa person, has old-name $v;\n  not { $x has new-name $d; };";
471        let stripped = strip_not_guard(match_section);
472        assert!(
473            !stripped.contains("not {"),
474            "stripped result must not contain 'not {{': {stripped}"
475        );
476        assert!(
477            stripped.contains("$x isa person"),
478            "stripped result must preserve the main match line: {stripped}"
479        );
480    }
481
482    // ── test: zero counts (empty database) ───────────────────────────────────
483
484    #[tokio::test]
485    async fn backfill_with_zero_counts_returns_zero_result() {
486        use serde_json::json;
487        use type_bridge_orm::session::backend::QueryResult;
488
489        let scripted = vec![
490            QueryResult::Rows(vec![json!({"c": 0})]),
491            QueryResult::Rows(vec![json!({"c": 0})]),
492            QueryResult::Ok,
493        ];
494
495        let (backend, _log) = MockMigrationBackend::with_responses(scripted);
496        let db = Database::with_backend(Box::new(backend), "test");
497        let step = backfill_step();
498
499        let result = execute_backfill(&db, &step, 0)
500            .await
501            .expect("execute_backfill should succeed with zero counts");
502
503        assert_eq!(result.matched, 0);
504        assert_eq!(result.inserted, 0);
505        assert_eq!(result.skipped, 0);
506        assert_eq!(result.conflicts, 0);
507    }
508}