1use 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
33pub struct BackfillResult {
34 pub step_index: usize,
36 pub matched: u64,
39 pub inserted: u64,
44 pub skipped: u64,
47 pub conflicts: u64,
50}
51
52pub(crate) struct PreparedBackfill {
57 pub(crate) transaction: TransactionContext,
58 pub(crate) result: BackfillResult,
59}
60
61pub 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
94pub(crate) async fn prepare_backfill(
96 db: &Database,
97 step: &ExecutionStep,
98 step_index: usize,
99) -> Result<PreparedBackfill, MigrationError> {
100 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 let guarded_count_query = format!("{match_section}\nreduce $c = count;");
121
122 let unguarded_match = strip_not_guard(match_section);
127 let total_count_query = format!("{unguarded_match}\nreduce $c = count;");
128
129 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 let _ = ctx.rollback().await;
146 extract_count(result, step_index, "guarded")?
147 };
148
149 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 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
204fn extract_count(
210 result: QueryResult,
211 step_index: usize,
212 label: &str,
213) -> Result<u64, MigrationError> {
214 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 => 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
243fn 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
258fn 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#[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 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 #[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 let scripted = vec![
302 QueryResult::Rows(vec![json!({"c": 7})]),
304 QueryResult::Rows(vec![json!({"c": 10})]),
306 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 let events = log.lock().unwrap();
326 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 assert!(matches!(events[8], MockEvent::Commit));
343 }
344
345 #[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 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 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 #[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 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 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]
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 #[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}