1use 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
34pub struct BackfillResult {
35 pub step_index: usize,
37 pub matched: u64,
40 pub inserted: u64,
45 pub skipped: u64,
48 pub conflicts: u64,
51}
52
53pub(crate) struct PreparedBackfill {
58 pub(crate) transaction: TransactionContext,
59 pub(crate) result: BackfillResult,
60}
61
62pub 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
95pub(crate) async fn prepare_backfill(
97 db: &Database,
98 step: &ExecutionStep,
99 step_index: usize,
100) -> Result<PreparedBackfill, MigrationError> {
101 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 let guarded_count_query = format!("{match_section}\nreduce $c = count;");
122
123 let unguarded_match = strip_not_guard(match_section);
128 let total_count_query = format!("{unguarded_match}\nreduce $c = count;");
129
130 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 let _ = ctx.rollback().await;
147 extract_count(result, step_index, "guarded")?
148 };
149
150 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 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
205fn extract_count(
211 result: QueryResult,
212 step_index: usize,
213 label: &str,
214) -> Result<u64, MigrationError> {
215 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 => 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
244fn 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
259fn 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#[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 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 #[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 let scripted = vec![
303 QueryResult::Rows(vec![json!({"c": 7})]),
305 QueryResult::Rows(vec![json!({"c": 10})]),
307 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 let events = log.lock().unwrap();
327 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 assert!(matches!(events[8], MockEvent::Commit));
344 }
345
346 #[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 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 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 #[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 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 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]
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 #[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}