use type_bridge_orm::Database;
use type_bridge_orm::session::TransactionContext;
use type_bridge_orm::session::backend::QueryResult;
use crate::error::MigrationError;
use crate::plan::ExecutionStep;
use serde::{Deserialize, Serialize};
use type_bridge_orm::TxType;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BackfillResult {
pub step_index: usize,
pub matched: u64,
pub inserted: u64,
pub skipped: u64,
pub conflicts: u64,
}
pub(crate) struct PreparedBackfill {
pub(crate) transaction: TransactionContext,
pub(crate) result: BackfillResult,
}
pub async fn execute_backfill(
db: &Database,
step: &ExecutionStep,
step_index: usize,
) -> Result<BackfillResult, MigrationError> {
let prepared = prepare_backfill(db, step, step_index).await?;
prepared
.transaction
.commit()
.await
.map_err(|e| MigrationError::BackfillQuery {
message: format!("backfill step {step_index}: backfill write commit failed: {e}"),
})?;
Ok(prepared.result)
}
pub(crate) async fn prepare_backfill(
db: &Database,
step: &ExecutionStep,
step_index: usize,
) -> Result<PreparedBackfill, MigrationError> {
let (match_section, _insert_section) = step
.forward
.split_once("\ninsert\n")
.ok_or_else(|| MigrationError::BackfillQuery {
message: format!(
"backfill step {step_index}: forward TypeQL does not contain '\\ninsert\\n' separator; \
cannot compose count queries without re-deriving semantics (invariant 2)"
),
})?;
let guarded_count_query = format!("{match_section}\nreduce $c = count;");
let unguarded_match = strip_not_guard(match_section);
let total_count_query = format!("{unguarded_match}\nreduce $c = count;");
let inserted: u64 = {
let ctx = db.transaction_context(TxType::Read).await.map_err(|e| {
MigrationError::BackfillQuery {
message: format!(
"backfill step {step_index}: failed to open read tx for guarded count: {e}"
),
}
})?;
let result =
ctx.query(&guarded_count_query)
.await
.map_err(|e| MigrationError::BackfillQuery {
message: format!("backfill step {step_index}: guarded count query failed: {e}"),
})?;
let _ = ctx.rollback().await;
extract_count(result, step_index, "guarded")?
};
let matched: u64 = {
let ctx = db.transaction_context(TxType::Read).await.map_err(|e| {
MigrationError::BackfillQuery {
message: format!(
"backfill step {step_index}: failed to open read tx for total count: {e}"
),
}
})?;
let result =
ctx.query(&total_count_query)
.await
.map_err(|e| MigrationError::BackfillQuery {
message: format!("backfill step {step_index}: total count query failed: {e}"),
})?;
let _ = ctx.rollback().await;
extract_count(result, step_index, "total")?
};
let transaction =
db.transaction_context(TxType::Write)
.await
.map_err(|e| MigrationError::BackfillQuery {
message: format!("backfill step {step_index}: failed to open write tx: {e}"),
})?;
if let Err(error) = transaction.query(&step.forward).await {
let _ = transaction.rollback().await;
return Err(MigrationError::BackfillQuery {
message: format!("backfill step {step_index}: backfill write query failed: {error}"),
});
}
let skipped = matched.saturating_sub(inserted);
Ok(PreparedBackfill {
transaction,
result: BackfillResult {
step_index,
matched,
inserted,
skipped,
conflicts: 0,
},
})
}
fn extract_count(
result: QueryResult,
step_index: usize,
label: &str,
) -> Result<u64, MigrationError> {
let answer = match result {
QueryResult::Rows(items) | QueryResult::Documents(items) => match items.first() {
Some(item) => item.clone(),
None => return Ok(0),
},
QueryResult::Ok => return Ok(0),
};
let value = answer
.get("c")
.or_else(|| answer.as_object().and_then(|m| m.values().next()))
.ok_or_else(|| MigrationError::BackfillQuery {
message: format!(
"backfill step {step_index}: {label} count answer has no recognizable key: {answer}"
),
})?;
scalar_to_u64(value).ok_or_else(|| MigrationError::BackfillQuery {
message: format!(
"backfill step {step_index}: {label} count value is not a number: {value}"
),
})
}
fn scalar_to_u64(value: &serde_json::Value) -> Option<u64> {
match value {
serde_json::Value::Number(_) => value.as_u64().or_else(|| value.as_f64().map(|f| f as u64)),
serde_json::Value::Object(map) => map.get("value").and_then(scalar_to_u64),
_ => None,
}
}
fn strip_not_guard(match_section: &str) -> String {
match_section
.lines()
.filter(|line| !line.trim_start().starts_with("not {"))
.collect::<Vec<_>>()
.join("\n")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::plan::StepKind;
use crate::testing::{MockEvent, MockMigrationBackend};
use type_bridge_orm::{Database, TxType};
fn backfill_step() -> ExecutionStep {
ExecutionStep {
tx_type: TxType::Write,
kind: StepKind::Backfill,
operation_kind: crate::plan::OperationKind::CopyAttribute,
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(),
reverse: Some("match $x isa person, has new-name $v;\ndelete $v of $x;".to_string()),
}
}
#[tokio::test]
async fn backfill_derives_counts_from_scripted_mock_responses() {
use serde_json::json;
use type_bridge_orm::session::backend::QueryResult;
let scripted = vec![
QueryResult::Rows(vec![json!({"c": 7})]),
QueryResult::Rows(vec![json!({"c": 10})]),
QueryResult::Ok,
];
let (backend, log) = MockMigrationBackend::with_responses(scripted);
let db = Database::with_backend(Box::new(backend), "test");
let step = backfill_step();
let result = execute_backfill(&db, &step, 0)
.await
.expect("execute_backfill should succeed");
assert_eq!(result.step_index, 0);
assert_eq!(result.inserted, 7, "inserted = guarded count");
assert_eq!(result.matched, 10, "matched = total count");
assert_eq!(result.skipped, 3, "skipped = matched - inserted");
assert_eq!(result.conflicts, 0, "conflicts always 0 in v1");
let events = log.lock().unwrap();
assert!(
matches!(events[0], MockEvent::OpenTx(TxType::Read)),
"first tx must be Read (guarded count)"
);
assert!(
matches!(events[3], MockEvent::OpenTx(TxType::Read)),
"second tx must be Read (total count)"
);
assert!(
matches!(events[6], MockEvent::OpenTx(TxType::Write)),
"third tx must be Write (backfill insert)"
);
assert!(matches!(events[8], MockEvent::Commit));
}
#[tokio::test]
async fn count_queries_are_built_from_carried_match_not_re_derived() {
use serde_json::json;
use type_bridge_orm::session::backend::QueryResult;
let scripted = vec![
QueryResult::Rows(vec![json!({"c": 0})]),
QueryResult::Rows(vec![json!({"c": 0})]),
QueryResult::Ok,
];
let (backend, log) = MockMigrationBackend::with_responses(scripted);
let db = Database::with_backend(Box::new(backend), "test");
let step = backfill_step();
execute_backfill(&db, &step, 1)
.await
.expect("execute_backfill should succeed");
let events = log.lock().unwrap();
let guarded_query = events.iter().find_map(|e| {
if let MockEvent::Query(TxType::Read, q) = e {
Some(q.as_str())
} else {
None
}
});
assert!(
guarded_query.is_some(),
"expected at least one Read query in the event log"
);
let q = guarded_query.unwrap();
assert!(
q.contains("$x isa person, has old-name $v"),
"guarded count query must contain the step's match body; got: {q}"
);
assert!(
q.contains("reduce $c = count"),
"guarded count query must end with reduce count; got: {q}"
);
}
#[tokio::test]
async fn backfill_write_runs_under_write_tx() {
use serde_json::json;
use type_bridge_orm::session::backend::QueryResult;
let scripted = vec![
QueryResult::Rows(vec![json!({"c": 5})]),
QueryResult::Rows(vec![json!({"c": 5})]),
QueryResult::Ok,
];
let (backend, log) = MockMigrationBackend::with_responses(scripted);
let db = Database::with_backend(Box::new(backend), "test");
let step = backfill_step();
execute_backfill(&db, &step, 2)
.await
.expect("execute_backfill should succeed");
let events = log.lock().unwrap();
let write_query = events.iter().find_map(|e| {
if let MockEvent::Query(TxType::Write, q) = e {
Some(q.as_str())
} else {
None
}
});
assert!(write_query.is_some(), "expected a Write-typed query event");
let wq = write_query.unwrap();
assert!(
wq.contains("insert"),
"write query must contain 'insert'; got: {wq}"
);
assert!(
!wq.contains("reduce"),
"write query must not contain 'reduce' (it is the insert, not a count query); got: {wq}"
);
}
#[test]
fn strip_not_guard_removes_not_line() {
let match_section =
"match\n $x isa person, has old-name $v;\n not { $x has new-name $d; };";
let stripped = strip_not_guard(match_section);
assert!(
!stripped.contains("not {"),
"stripped result must not contain 'not {{': {stripped}"
);
assert!(
stripped.contains("$x isa person"),
"stripped result must preserve the main match line: {stripped}"
);
}
#[tokio::test]
async fn backfill_with_zero_counts_returns_zero_result() {
use serde_json::json;
use type_bridge_orm::session::backend::QueryResult;
let scripted = vec![
QueryResult::Rows(vec![json!({"c": 0})]),
QueryResult::Rows(vec![json!({"c": 0})]),
QueryResult::Ok,
];
let (backend, _log) = MockMigrationBackend::with_responses(scripted);
let db = Database::with_backend(Box::new(backend), "test");
let step = backfill_step();
let result = execute_backfill(&db, &step, 0)
.await
.expect("execute_backfill should succeed with zero counts");
assert_eq!(result.matched, 0);
assert_eq!(result.inserted, 0);
assert_eq!(result.skipped, 0);
assert_eq!(result.conflicts, 0);
}
}