use super::*;
use crate::row_locks::{PendingRowChange, RowLockKey};
const TABLE: [u8; 16] = [3; 16];
const OTHER: [u8; 16] = [7; 16];
#[derive(Default)]
struct Fixture {
frames: Vec<Vec<TransactionRowChange>>,
fail: bool,
}
impl RewriteRowChanges for Fixture {
fn storage_generation(&self, table: &str) -> Result<[u8; 16], SQLError> {
match table {
"target" => Ok(TABLE),
_ => Err(SQLError::UnknownTable(table.into())),
}
}
fn visit_changes(
&self,
visit: &mut dyn FnMut(&[TransactionRowChange]) -> Result<(), SQLError>,
) -> Result<(), SQLError> {
for frame in &self.frames {
visit(frame)?;
}
if self.fail {
return Err(SQLError::Internal("fixture read failed".into()));
}
Ok(())
}
}
fn change(generation: [u8; 16], id: DocId, kind: PendingRowChangeKind) -> TransactionRowChange {
TransactionRowChange {
pending: PendingRowChange {
key: RowLockKey {
table: 1,
doc_id: id,
},
kind,
},
source_generation: generation,
successor_generation: None,
query_origin: None,
}
}
fn rewritten(from: ([u8; 16], DocId), to: Option<([u8; 16], DocId)>) -> TransactionRowChange {
let mut change = change(
from.0,
from.1,
PendingRowChangeKind::Rewrite(RowLockKey {
table: 2,
doc_id: to.map_or(99, |(_, id)| id),
}),
);
change.successor_generation = to.map(|(generation, _)| generation);
change
}
fn entries(rows: &ChangedRows) -> Vec<(DocId, bool)> {
let mut result = Vec::new();
rows.visit(&mut |id, present| {
result.push((id, present));
Ok(())
})
.unwrap();
result
}
#[test]
fn surviving_changes_after_the_marker_keep_the_latest_presence() {
let memory = MemoryBudget::new(1 << 20);
let cancellation = CancellationToken::new();
let mut owner = Fixture {
frames: vec![
vec![change(TABLE, 2, PendingRowChangeKind::Insert)],
vec![change(TABLE, 8, PendingRowChangeKind::Delete)],
],
..Fixture::default()
};
let start = marker(&owner).unwrap();
assert_eq!(start, 2);
owner.frames[1].extend([
change(TABLE, 2, PendingRowChangeKind::Update),
change(TABLE, 2, PendingRowChangeKind::Update),
change(TABLE, 2, PendingRowChangeKind::Delete),
change(TABLE, 2, PendingRowChangeKind::Insert),
change(TABLE, 4, PendingRowChangeKind::Update),
change(TABLE, 4, PendingRowChangeKind::Delete),
change(OTHER, 2, PendingRowChangeKind::Delete),
change(OTHER, 9, PendingRowChangeKind::Insert),
]);
let rows = capture_since(
&owner,
start,
owner.storage_generation("target").unwrap(),
&memory,
&cancellation,
)
.unwrap();
assert!(!rows.is_empty());
assert_eq!(entries(&rows), [(2, true), (4, false)]);
assert_eq!(rows.get(2).unwrap(), Some(true));
assert_eq!(rows.get(4).unwrap(), Some(false));
assert_eq!(rows.get(8).unwrap(), None);
assert_eq!(rows.get(9).unwrap(), None);
drop(rows);
assert_eq!(memory.used(), 0);
}
#[test]
fn rewrites_follow_both_generations_and_keep_full_width_numeric_order() {
let owner = Fixture {
frames: vec![vec![
change(TABLE, 0, PendingRowChangeKind::Insert),
rewritten((TABLE, DocId::MAX), Some((OTHER, 255))),
rewritten((OTHER, 0), Some((TABLE, 256))),
rewritten((TABLE, 1), Some((TABLE, 2))),
rewritten((TABLE, 9), Some((TABLE, 9))),
rewritten((TABLE, 333), None),
]],
..Fixture::default()
};
let memory = MemoryBudget::new(1 << 20);
let cancellation = CancellationToken::new();
let rows = capture_since(&owner, 0, TABLE, &memory, &cancellation).unwrap();
assert_eq!(
entries(&rows),
[
(0, true),
(1, false),
(2, true),
(9, true),
(256, true),
(333, false),
(DocId::MAX, false),
]
);
let other = capture_since(&owner, 0, OTHER, &memory, &cancellation).unwrap();
assert_eq!(entries(&other), [(0, false), (255, true)]);
assert_eq!(other.get(DocId::MAX).unwrap(), None);
let deleted = capture_since(&owner, 5, TABLE, &memory, &cancellation).unwrap();
assert!(!deleted.is_empty());
assert_eq!(entries(&deleted), [(333, false)]);
}
#[test]
fn frame_commit_and_rollback_preserve_the_marker_and_retained_views() {
let memory = MemoryBudget::new(1 << 20);
let cancellation = CancellationToken::new();
let mut owner = Fixture {
frames: vec![
vec![change(TABLE, 1, PendingRowChangeKind::Insert)],
vec![change(TABLE, 2, PendingRowChangeKind::Insert)],
],
..Fixture::default()
};
let start = marker(&owner).unwrap();
owner.frames[1].push(change(TABLE, 3, PendingRowChangeKind::Update));
owner.frames.push(vec![
change(TABLE, 4, PendingRowChangeKind::Insert),
change(TABLE, 3, PendingRowChangeKind::Delete),
]);
let before_rollback = capture_since(&owner, start, TABLE, &memory, &cancellation).unwrap();
owner.frames.pop();
let after_rollback = capture_since(&owner, start, TABLE, &memory, &cancellation).unwrap();
assert_eq!(entries(&before_rollback), [(3, false), (4, true)]);
assert_eq!(entries(&after_rollback), [(3, true)]);
let committed = owner.frames.pop().unwrap();
owner.frames[0].extend(committed);
let after_commit = capture_since(&owner, start, TABLE, &memory, &cancellation).unwrap();
assert_eq!(entries(&after_commit), [(3, true)]);
owner.frames[0].truncate(start);
let empty = capture_since(&owner, start, TABLE, &memory, &cancellation).unwrap();
assert!(empty.is_empty());
owner.frames[0].push(change(TABLE, 5, PendingRowChangeKind::Insert));
let next = capture_since(&owner, start, TABLE, &memory, &cancellation).unwrap();
assert_eq!(entries(&next), [(5, true)]);
assert_eq!(entries(&before_rollback), [(3, false), (4, true)]);
owner.frames[0].clear();
assert!(matches!(
capture_since(&owner, start, TABLE, &memory, &cancellation),
Err(SQLError::Internal(_))
));
}
#[test]
fn an_unchanged_generation_needs_no_record_allocation() {
let owner = Fixture {
frames: vec![vec![change(OTHER, 1, PendingRowChangeKind::Insert)]],
..Fixture::default()
};
let memory = MemoryBudget::new(0);
let rows = capture_since(&owner, 0, TABLE, &memory, &CancellationToken::new()).unwrap();
assert!(rows.is_empty());
assert_eq!(rows.get(1).unwrap(), None);
assert!(entries(&rows).is_empty());
assert_eq!(memory.used(), 0);
}
#[test]
fn spilled_changes_share_the_original_allowance_and_release_replay_resources() {
verify_spilled_changes(64 << 10, 0, 8_192);
verify_spilled_changes(256 << 10, 128 << 10, 16_384);
}
fn verify_spilled_changes(limit: usize, caller_bytes: usize, count: u64) {
let memory = MemoryBudget::new(limit);
let held = memory.reserve(caller_bytes).unwrap();
assert!(usize::try_from(count).unwrap() * 9 > memory.limit() - held.bytes());
let owner = Fixture {
frames: vec![(0..count)
.rev()
.map(|id| {
change(
TABLE,
id,
if id.is_multiple_of(3) {
PendingRowChangeKind::Delete
} else {
PendingRowChangeKind::Update
},
)
})
.collect()],
..Fixture::default()
};
let rows = capture_since(&owner, 0, TABLE, &memory, &CancellationToken::new()).unwrap();
let retained = memory.used();
let mut next = 0;
rows.visit(&mut |id, present| {
assert_eq!(id, next);
assert_eq!(present, !id.is_multiple_of(3));
next += 1;
Ok(())
})
.unwrap();
assert_eq!(next, count);
assert_eq!(memory.used(), retained);
assert_eq!(
rows.get(count - 1).unwrap(),
Some(!(count - 1).is_multiple_of(3))
);
assert_eq!(
rows.get(count - 2).unwrap(),
Some(!(count - 2).is_multiple_of(3))
);
assert_eq!(rows.get(count).unwrap(), None);
let mut visited = 0;
let error = rows
.visit(&mut |_, _| {
visited += 1;
Err(SQLError::Routine {
sqlstate: "23514".into(),
message: "callback rejected the row".into(),
})
})
.unwrap_err();
assert_eq!(error.sqlstate(), Some("23514"));
assert_eq!(visited, 1);
assert_eq!(memory.used(), retained);
assert!(memory.peak() <= memory.limit());
drop(rows);
assert_eq!(memory.used(), held.bytes());
}
#[test]
fn source_errors_and_cancellation_release_partially_retained_changes() {
let memory = MemoryBudget::new(1 << 20);
let cancellation = CancellationToken::new();
let mut owner = Fixture {
frames: vec![vec![change(TABLE, 1, PendingRowChangeKind::Insert)]],
fail: true,
};
assert!(marker(&owner).is_err());
let error = capture_since(&owner, 0, TABLE, &memory, &cancellation)
.err()
.expect("source error");
assert!(matches!(error, SQLError::Internal(message) if message == "fixture read failed"));
assert_eq!(memory.used(), 0);
owner.fail = false;
let rows = capture_since(&owner, 0, TABLE, &memory, &cancellation).unwrap();
cancellation.cancel();
assert!(matches!(rows.get(1), Err(SQLError::Cancelled(_))));
assert!(matches!(
rows.visit(&mut |_, _| Ok(())),
Err(SQLError::Cancelled(_))
));
drop(rows);
assert!(matches!(
capture_since(&owner, 0, TABLE, &memory, &cancellation),
Err(SQLError::Cancelled(_))
));
assert_eq!(memory.used(), 0);
}