Skip to content
Merged
Show file tree
Hide file tree
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
24 changes: 12 additions & 12 deletions aw-datastore/src/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,7 @@ pub enum Command {
/// filters inserts/heartbeats, so every write/delete of this key (and startup)
/// must reload the engine. RefreshPrivacyFilter exists for explicit reloads.
const PRIVACY_FILTERS_KEY: &str = "settings.privacy_filters";
const STOPWATCH_BUCKET_TYPE: &str = "general.stopwatch";

fn _unwrap_empty_response(response: Response) -> Result<(), DatastoreError> {
match response {
Expand Down Expand Up @@ -319,14 +320,8 @@ impl DatastoreWorker {

self.uncommitted_events = 0;
self.commit = false;
// ForceCommit and Close promise the caller that their data is
// committed, so their acks are held back until the transaction
// below has actually committed. Acking first (as before) let a
// caller reopen the database and read a pre-commit snapshot —
// harmless under the rollback journal's locking, but a real race
// in WAL mode where readers never block on the writer.
// All other commands are acked immediately: a watcher heartbeat
// must not wait up to 15 s for the batch commit.
// Commands that force a commit are acknowledged only after it
// succeeds. Other commands can return before the batch commits.
let mut deferred_ack = None;
loop {
let (request, response_sender) = match self.responder.poll() {
Expand All @@ -338,10 +333,8 @@ impl DatastoreWorker {
break;
}
};
let ack_after_commit = matches!(request, Command::ForceCommit() | Command::Close());
let response = self.handle_request(request, &mut ds, &tx);
if ack_after_commit {
// Both commands force a commit, so the loop ends here.
if self.commit || self.quit {
deferred_ack = Some((response_sender, response));
break;
}
Expand Down Expand Up @@ -441,6 +434,9 @@ impl DatastoreWorker {
Ok(events) => {
self.uncommitted_events += events.len();
self.last_heartbeat.insert(bucketname.to_string(), None); // invalidate last_heartbeat cache

// Manual timer changes must be durable before the UI confirms them.
self.commit |= ds.get_bucket(&bucketname)?._type == STOPWATCH_BUCKET_TYPE;
Ok(Response::EventList(events))
}
Err(e) => Err(e),
Expand Down Expand Up @@ -472,6 +468,7 @@ impl DatastoreWorker {
) {
Ok(e) => {
self.uncommitted_events += 1;
self.commit |= ds.get_bucket(&bucketname)?._type == STOPWATCH_BUCKET_TYPE;
Ok(Response::Event(e))
}
Err(e) => Err(e),
Expand Down Expand Up @@ -502,7 +499,10 @@ impl DatastoreWorker {
}
Command::DeleteEventsById(bucketname, event_ids) => {
match ds.delete_events_by_id(tx, &bucketname, event_ids) {
Ok(()) => Ok(Response::Empty()),
Ok(()) => {
self.commit |= ds.get_bucket(&bucketname)?._type == STOPWATCH_BUCKET_TYPE;
Ok(Response::Empty())
}
Err(e) => Err(e),
}
}
Expand Down
37 changes: 37 additions & 0 deletions aw-datastore/tests/datastore.rs
Original file line number Diff line number Diff line change
Expand Up @@ -832,6 +832,43 @@ mod datastore_tests {
}
}

#[test]
fn stopwatch_changes_are_committed_before_ack() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("stopwatch.db");
let ds = Datastore::new(path.to_str().unwrap().to_string(), false);
let mut bucket = test_bucket();
bucket._type = "general.stopwatch".to_string();
ds.create_bucket(&bucket).unwrap();
let conn = rusqlite::Connection::open(&path).unwrap();

let mut event = test_event(Utc::now(), Duration::zero());
event.data = json_map! {"running": json!(true)};
event = ds.heartbeat(&bucket.id, event, 1.0).unwrap();
let id = event.id.unwrap();
let read_data = || -> serde_json::Value {
let data: String = conn
.query_row("SELECT data FROM events WHERE id = ?1", [id], |row| {
row.get(0)
})
.unwrap();
serde_json::from_str(&data).unwrap()
};
assert_eq!(read_data()["running"], true);

event.data = json_map! {"running": json!(false)};
ds.insert_events(&bucket.id, &[event]).unwrap();
assert_eq!(read_data()["running"], false);

ds.delete_events_by_id(&bucket.id, vec![id]).unwrap();
let count: i64 = conn
.query_row("SELECT count(*) FROM events WHERE id = ?1", [id], |row| {
row.get(0)
})
.unwrap();
assert_eq!(count, 0);
}

#[test]
fn test_migration_v4_to_v5() {
let test_dir = tempfile::tempdir().unwrap();
Expand Down