Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 24 additions & 7 deletions src/store/apply.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,12 @@ impl MemoryDb {
records_applied: 0,
records_skipped: 0,
};
let mut check_seen = transaction.prepare_cached(
"SELECT EXISTS(SELECT 1 FROM memory_applied_records WHERE tenant_id = ?1 AND person_id = ?2 AND record_kind = ?3 AND record_id = ?4 AND payload_hash = ?5)",
)?;
let mut insert_seen = transaction.prepare_cached(
"INSERT INTO memory_applied_records(tenant_id, person_id, record_kind, record_id, payload_hash, applied_at) VALUES(?1, ?2, ?3, ?4, ?5, ?6)",
)?;
for commit in &input.commits {
let mut accepted = vec![false; commit.records.len()];
for pass in PASSES {
Expand All @@ -99,20 +105,29 @@ impl MemoryDb {
}
let (record_kind, record_id) = record_identity(record);
let payload_hash = record_hash(record)?;
let seen: bool = transaction.query_row(
"SELECT EXISTS(SELECT 1 FROM memory_applied_records WHERE tenant_id = ?1 AND person_id = ?2 AND record_kind = ?3 AND record_id = ?4 AND payload_hash = ?5)",
params![input.tenant_id.0, input.person_id.0, record_kind, record_id, payload_hash],
let seen: bool = check_seen.query_row(
params![
input.tenant_id.0,
input.person_id.0,
record_kind,
record_id,
payload_hash
],
|row| row.get(0),
)?;
if seen {
applied.records_skipped += 1;
continue;
}
apply_record(&transaction, record, applied_at)?;
transaction.execute(
"INSERT INTO memory_applied_records(tenant_id, person_id, record_kind, record_id, payload_hash, applied_at) VALUES(?1, ?2, ?3, ?4, ?5, ?6)",
params![input.tenant_id.0, input.person_id.0, record_kind, record_id, payload_hash, applied_at],
)?;
insert_seen.execute(params![
input.tenant_id.0,
input.person_id.0,
record_kind,
record_id,
payload_hash,
applied_at
])?;
accepted[index] = true;
applied.records_applied += 1;
}
Expand All @@ -139,6 +154,8 @@ impl MemoryDb {
)?;
applied.commits_applied += 1;
}
drop(check_seen);
drop(insert_seen);
record_operation(
&transaction,
&input.tenant_id,
Expand Down
Loading