Skip to main content

type_bridge_migration/
backfill.rs

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