1use super::*;
10use mj_core::native_agent::{NativeAgentSummary, NativeAgentView};
11use mj_core::snapshot_map::SnapshotMap;
12use rusqlite::functions::FunctionFlags;
13use std::cell::RefCell;
14
15#[derive(Default)]
16struct PendingChanges {
17 path: PathBuf,
18 keys: BTreeSet<(String, String)>,
19}
20
21thread_local! {
22 static PENDING: RefCell<Option<PendingChanges>> = const { RefCell::new(None) };
23}
24
25pub(super) fn observe_connection(connection: &Connection, path: &Path) -> Result<()> {
28 let path = path.to_owned();
29 connection.create_scalar_function(
30 "mj_changed_record",
31 2,
32 FunctionFlags::SQLITE_UTF8,
33 move |arguments| {
34 let kind: String = arguments.get(0)?;
35 let key: String = arguments.get(1)?;
36 PENDING.with(|pending| {
37 if let Some(pending) = pending.borrow_mut().as_mut()
38 && pending.path == path
39 {
40 pending.keys.insert((kind, key));
41 }
42 });
43 Ok(0)
44 },
45 )?;
46 for (table, kind, key) in [
47 ("sessions", "session", "session_id"),
48 ("session_contexts", "session", "session_id"),
49 ("session_targets", "session", "session_id"),
50 ("session_mounts", "session", "session_id"),
51 ("session_mount_access", "session", "session_id"),
52 ("session_checkpoints", "session", "session_id"),
53 ("subagent_sessions", "relation", "child_session_id"),
54 ("subagent_preference", "preference", "singleton"),
55 ("mount_history", "mount_history", "host"),
56 ("project_locations", "mount_history", "host"),
57 ("host_container_sizes", "container_size", "host"),
58 ("session_moves", "move", "session_id"),
59 ("native_agents", "native_agent", "owner"),
60 ("startup_steps", "startup", "session_id"),
61 ("subagent_handbacks", "report", "child_session_id"),
62 ("materialized_sessions", "turn", "session_id"),
63 ] {
64 let exists: bool = connection.query_row(
67 "SELECT EXISTS(SELECT 1 FROM main.sqlite_schema WHERE type='table' AND name=?1)",
68 [table],
69 |row| row.get(0),
70 )?;
71 if !exists {
72 continue;
73 }
74 for (event, references) in [
75 ("INSERT", &["NEW"][..]),
76 ("DELETE", &["OLD"][..]),
77 ("UPDATE", &["OLD", "NEW"][..]),
78 ] {
79 let calls = references
80 .iter()
81 .map(|reference| {
82 let key = if table == "project_locations" {
83 format!("'project:' || CAST({reference}.host AS TEXT)")
84 } else if kind == "native_agent" {
85 format!("json_array({reference}.owner, {reference}.child)")
86 } else {
87 format!("CAST({reference}.{key} AS TEXT)")
88 };
89 format!("SELECT mj_changed_record('{kind}', {key});")
90 })
91 .collect::<String>();
92 let condition = if kind == "turn" && event == "UPDATE" {
93 " WHEN OLD.session_id IS NOT NEW.session_id
94 OR OLD.execution_state IS NOT NEW.execution_state
95 OR OLD.running_started_at_ms IS NOT NEW.running_started_at_ms
96 OR OLD.active_turn_json IS NOT NEW.active_turn_json
97 OR OLD.last_turn_outcome_json IS NOT NEW.last_turn_outcome_json"
98 } else {
99 ""
100 };
101 connection.execute_batch(&format!(
102 "CREATE TEMP TRIGGER mj_observe_{table}_{event} AFTER {event} ON main.{table}{condition}
103 BEGIN {calls} END;"
104 ))?;
105 }
106 }
107 Ok(())
108}
109
110pub(super) fn begin_operation(path: &Path) {
111 PENDING.with(|pending| {
112 assert!(
113 pending.borrow().is_none(),
114 "nested database writer operation"
115 );
116 *pending.borrow_mut() = Some(PendingChanges {
117 path: path.to_owned(),
118 keys: BTreeSet::new(),
119 });
120 });
121}
122
123#[derive(Clone)]
124pub struct CommittedState {
125 pub sequence: u64,
126 pub state: State,
127 pub moves: SnapshotMap<String, mj_core::state::MoveOperation>,
128 pub native_agents: SnapshotMap<String, SnapshotMap<String, NativeAgentSummary>>,
129 pub startup_groups: SnapshotMap<String, Vec<StartupDelivery>>,
133 pub subagent_reports: SnapshotMap<String, mj_core::subagent::SubagentReport>,
136 pub turns: SnapshotMap<String, CommittedTurn>,
137 pub wait_revisions: SnapshotMap<String, u64>,
139}
140
141#[derive(Clone, Debug, PartialEq, Eq)]
142pub struct CommittedTurn {
143 pub state: materialized::MaterializedTurnState,
144 pub failed_message: Option<String>,
146}
147
148impl CommittedTurn {
149 fn read(connection: &Connection, id: &str, previous: Option<&Self>) -> Result<Option<Self>> {
150 let Some(state) = materialized::read_materialized_turn_state(connection, id)? else {
151 return Ok(None);
152 };
153 let failed_message = match state.2.as_ref() {
154 Some(turn) if mj_core::subagent::failed_turn(turn, None).is_some() => {
155 if let Some(previous) = previous.filter(|previous| previous.state.2 == state.2) {
156 previous.failed_message.clone()
157 } else if let Some(start) = turn.turn_start_position {
158 materialized::last_materialized_agent_message_within(
159 connection,
160 id,
161 start,
162 turn.completed_ordinal,
163 )?
164 } else {
165 None
166 }
167 }
168 _ => None,
169 };
170 Ok(Some(Self {
171 state,
172 failed_message,
173 }))
174 }
175}
176
177impl CommittedState {
178 pub(super) fn bootstrap(connection: &mut Connection) -> Result<Self> {
179 let transaction =
180 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Deferred)?;
181 let state = state_io::load_state_with(&transaction)?;
182 let ids = transaction
183 .prepare("SELECT session_id FROM materialized_sessions")?
184 .query_map([], |row| row.get::<_, String>(0))?
185 .collect::<rusqlite::Result<Vec<_>>>()?;
186 let mut turns = SnapshotMap::new();
187 for id in ids {
188 if let Some(turn) = CommittedTurn::read(&transaction, &id, None)? {
189 turns.insert_shared(id, turn);
190 }
191 }
192 let mut startup_groups = SnapshotMap::new();
193 let session_ids = transaction
194 .prepare("SELECT DISTINCT session_id FROM startup_steps WHERE group_id IS NOT NULL")?
195 .query_map([], |row| row.get::<_, String>(0))?
196 .collect::<rusqlite::Result<Vec<_>>>()?;
197 for session_id in session_ids {
198 let group = startup::load_latest_startup_group_with(&transaction, &session_id)?;
199 if !group.is_empty() {
200 startup_groups.insert(session_id, group);
201 }
202 }
203 let mut subagent_reports = SnapshotMap::new();
204 let children = transaction
205 .prepare("SELECT child_session_id FROM subagent_handbacks")?
206 .query_map([], |row| row.get::<_, String>(0))?
207 .collect::<rusqlite::Result<Vec<_>>>()?;
208 for child in children {
209 if let Some(report) = sessions::load_subagent_report_with(&transaction, &child)? {
210 subagent_reports.insert(child, report);
211 }
212 }
213 let moves = session_move::load_move_operations_with(&transaction)?
214 .into_iter()
215 .map(|operation| (operation.selection.session_id.clone(), operation))
216 .collect();
217 let mut native_agents =
218 SnapshotMap::<String, SnapshotMap<String, NativeAgentSummary>>::new();
219 let mut statement =
220 transaction.prepare("SELECT owner, child, body FROM native_agents WHERE staging=0")?;
221 let rows = statement.query_map([], |row| {
222 Ok((
223 row.get::<_, String>(0)?,
224 row.get::<_, String>(1)?,
225 row.get::<_, String>(2)?,
226 ))
227 })?;
228 for row in rows {
229 let (owner, child, body) = row?;
230 let view: NativeAgentView = serde_json::from_str(&body)?;
231 native_agents
232 .entry(owner)
233 .or_insert_with(SnapshotMap::new)
234 .insert(child, NativeAgentSummary::of(&view));
235 }
236 Ok(Self {
237 sequence: 0,
238 state,
239 moves,
240 native_agents,
241 startup_groups,
242 subagent_reports,
243 turns,
244 wait_revisions: SnapshotMap::new(),
245 })
246 }
247}
248
249pub(super) fn finish_operation(
252 connection: &mut Connection,
253 previous: &CommittedState,
254) -> Result<Option<CommittedState>> {
255 let changes = PENDING
256 .with(|pending| pending.borrow_mut().take())
257 .context("database writer operation has no change collector")?;
258 ensure!(
259 connection.is_autocommit(),
260 "writer job left a transaction open"
261 );
262 if changes.keys.is_empty() {
263 return Ok(None);
264 }
265 let transaction =
266 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Deferred)?;
267 let mut state = previous.state.clone();
268 let mut moves = previous.moves.clone();
269 let mut native_agents = previous.native_agents.clone();
270 let mut startup_groups = previous.startup_groups.clone();
271 let mut subagent_reports = previous.subagent_reports.clone();
272 let mut turns = previous.turns.clone();
273 let mut changed_history = State::default();
274 let mut changed = false;
275 let mut relations = BTreeSet::new();
276 for (kind, key) in &changes.keys {
277 match kind.as_str() {
278 "turn" => {
279 let turn = CommittedTurn::read(&transaction, key, turns.get(key))?;
280 if turns.get(key) != turn.as_ref() {
281 match turn {
282 Some(turn) => {
283 turns.insert_shared(key.clone(), turn);
284 }
285 None => {
286 turns.remove_shared(key);
287 }
288 }
289 changed = true;
290 }
291 }
292 "move" => {
293 let operation = session_move::load_move_operation_with(&transaction, key)?;
294 if moves.get(key) != operation.as_ref() {
295 match operation {
296 Some(operation) => {
297 moves.insert(key.clone(), operation);
298 }
299 None => {
300 moves.remove(key);
301 }
302 }
303 changed = true;
304 }
305 }
306 "native_agent" => {
307 let (owner, child): (String, String) = serde_json::from_str(key)?;
308 let body: Option<String> = transaction
309 .query_row(
310 "SELECT body FROM native_agents WHERE owner=?1 AND child=?2 AND staging=0",
311 params![owner, child],
312 |row| row.get(0),
313 )
314 .optional()?;
315 let summary = body
316 .map(|body| {
317 serde_json::from_str::<NativeAgentView>(&body)
318 .map(|view| NativeAgentSummary::of(&view))
319 })
320 .transpose()?;
321 let old = native_agents
322 .get(&owner)
323 .and_then(|children| children.get(&child));
324 if old != summary.as_ref() {
325 if let Some(summary) = summary {
326 native_agents
327 .entry(owner)
328 .or_insert_with(SnapshotMap::new)
329 .insert(child, summary);
330 } else if let Some(children) = native_agents.get_mut(&owner) {
331 children.remove(&child);
332 if children.is_empty() {
333 native_agents.remove(&owner);
334 }
335 }
336 changed = true;
337 }
338 }
339 "session" => {
340 let record = state_io::load_session_with(&transaction, key)?;
341 let membership_changed = state.sessions.contains_key(key) != record.is_some();
342 if state.sessions.get(key) != record.as_ref() {
343 match record {
344 Some(record) => {
345 state.sessions.insert(key.clone(), record);
346 }
347 None => {
348 state.sessions.remove(key);
349 }
350 }
351 changed = true;
352 }
353 if membership_changed {
356 let mut statement = transaction.prepare(
357 "SELECT child_session_id FROM subagent_sessions
358 WHERE parent_session_id=?1 OR child_session_id=?1",
359 )?;
360 relations.extend(
361 statement
362 .query_map([key], |row| row.get::<_, String>(0))?
363 .collect::<rusqlite::Result<Vec<_>>>()?,
364 );
365 }
366 }
367 "relation" => {
368 relations.insert(key.clone());
369 let mut statement = transaction.prepare(
370 "SELECT child_session_id FROM subagent_sessions WHERE parent_session_id=?1",
371 )?;
372 relations.extend(
373 statement
374 .query_map([key], |row| row.get::<_, String>(0))?
375 .collect::<rusqlite::Result<Vec<_>>>()?,
376 );
377 }
378 "preference" => {
379 let json: Option<String> = transaction
380 .query_row(
381 "SELECT policy FROM subagent_preference WHERE singleton=1",
382 [],
383 |row| row.get(0),
384 )
385 .optional()?;
386 let policy = json
387 .map(|json| serde_json::from_str(&json))
388 .transpose()?
389 .unwrap_or_default();
390 if state.last_subagent_policy != policy {
391 state.last_subagent_policy = policy;
392 changed = true;
393 }
394 }
395 "mount_history" => {
396 let paths = state_io::read_mount_history(&transaction)?
397 .remove(key)
398 .unwrap_or_default();
399 let paths = (!paths.is_empty()).then_some(paths);
400 if state.mount_history.get(key) != paths.as_ref() {
401 match paths {
402 Some(paths) => {
403 changed_history
404 .mount_history
405 .insert(key.clone(), paths.clone());
406 state.mount_history.insert(key.clone(), paths);
407 }
408 None => {
409 state.mount_history.remove(key);
410 }
411 }
412 changed = true;
413 }
414 }
415 "container_size" => {
416 let size = transaction
417 .query_row(
418 "SELECT cpus, memory_bytes FROM host_container_sizes WHERE host=?1",
419 [key],
420 |row| {
421 Ok(HostContainerSize {
422 cpus: row.get::<_, i64>(0)? as u64,
423 memory_bytes: row.get::<_, i64>(1)? as u64,
424 })
425 },
426 )
427 .optional()?;
428 if state.container_sizes.get(key) != size.as_ref() {
429 match size {
430 Some(size) => {
431 changed_history.container_sizes.insert(key.clone(), size);
432 state.container_sizes.insert(key.clone(), size);
433 }
434 None => {
435 state.container_sizes.remove(key);
436 }
437 }
438 changed = true;
439 }
440 }
441 "startup" => {
442 let group = startup::load_latest_startup_group_with(&transaction, key)?;
443 let group = (!group.is_empty()).then_some(group);
444 if startup_groups.get(key) != group.as_ref() {
445 match group {
446 Some(group) => {
447 startup_groups.insert(key.clone(), group);
448 }
449 None => {
450 startup_groups.remove(key);
451 }
452 }
453 changed = true;
454 }
455 }
456 "report" => {
457 let report = sessions::load_subagent_report_with(&transaction, key)?;
458 if subagent_reports.get(key) != report.as_ref() {
459 match report {
460 Some(report) => {
461 subagent_reports.insert(key.clone(), report);
462 }
463 None => {
464 subagent_reports.remove(key);
465 }
466 }
467 changed = true;
468 }
469 }
470 _ => bail!("unknown committed record kind {kind}"),
471 }
472 }
473 for key in &relations {
474 let json: Option<String> = transaction
475 .query_row(
476 "SELECT record_json FROM subagent_sessions WHERE child_session_id=?1",
477 [key],
478 |row| row.get(0),
479 )
480 .optional()?;
481 let relation: Option<SubagentRecord> =
482 json.map(|json| serde_json::from_str(&json)).transpose()?;
483 let relation = relation.filter(|relation| {
484 state.sessions.contains_key(key)
485 && state.sessions.contains_key(&relation.parent_session_id)
486 });
487 if state.subagents.get(key) != relation.as_ref() {
488 match relation {
489 Some(relation) => {
490 state.subagents.insert(key.clone(), relation);
491 }
492 None => {
493 state.subagents.remove(key);
494 }
495 }
496 changed = true;
497 }
498 }
499 for key in &relations {
500 state.validate_subagent(key)?;
501 }
502 changed_history.validate()?;
503 let mut wait_revisions = previous.wait_revisions.clone();
504 let affected = previous
505 .state
506 .sessions
507 .changes(&state.sessions)
508 .map(|(id, _)| id)
509 .chain(
510 previous
511 .state
512 .subagents
513 .changes(&state.subagents)
514 .map(|(id, _)| id),
515 )
516 .chain(
517 previous
518 .startup_groups
519 .changes(&startup_groups)
520 .map(|(id, _)| id),
521 )
522 .chain(
523 previous
524 .subagent_reports
525 .changes(&subagent_reports)
526 .map(|(id, _)| id),
527 )
528 .chain(previous.turns.changes(&turns).map(|(id, _)| id));
529 for id in affected {
530 wait_revisions.insert_shared(id.clone(), previous.sequence + 1);
531 }
532 transaction.commit()?;
533 Ok(changed.then(|| CommittedState {
534 sequence: previous.sequence + 1,
535 state,
536 moves,
537 native_agents,
538 startup_groups,
539 subagent_reports,
540 turns,
541 wait_revisions,
542 }))
543}
544
545#[cfg(test)]
546mod tests {
547 use super::*;
548
549 #[test]
550 fn compact_turn_publications_are_committed_and_session_specific() {
551 let directory = tempfile::tempdir().unwrap();
552 let path = directory.path().join("controller.sqlite");
553 for id in ["first", "second"] {
554 save_session_to(&path, &super::super::tests::session(id, "project")).unwrap();
555 }
556 let owner = start_database_writer_at(&path, false).unwrap();
557 let before = owner.writer.committed_state().unwrap();
558 let second = before.turns.get_shared("second").unwrap();
559 owner.writer.execute("transcript frontier only", |connection| {
560 connection.execute("UPDATE materialized_sessions SET applied_event_ordinal=99 WHERE session_id='first'", [])?;
561 Ok(())
562 }).unwrap();
563 assert!(
564 Arc::ptr_eq(&before, &owner.writer.committed_state().unwrap()),
565 "transcript updates do not change compact wait facts"
566 );
567 let rollback: Result<()> = owner.writer.execute("rollback turn", |connection| {
568 let transaction = connection.transaction()?;
569 transaction.execute("UPDATE materialized_sessions SET execution_state='running',running_started_at_ms=7 WHERE session_id='first'", [])?;
570 bail!("rolled back");
571 });
572 assert!(rollback.is_err());
573 assert!(Arc::ptr_eq(
574 &before,
575 &owner.writer.committed_state().unwrap()
576 ));
577 owner.writer.execute("start first turn", |connection| {
578 connection.execute("UPDATE materialized_sessions SET execution_state='running',running_started_at_ms=7 WHERE session_id='first'", [])?;
579 Ok(())
580 }).unwrap();
581 let running = owner.writer.committed_state().unwrap();
582 assert_ne!(
583 running.wait_revisions["first"],
584 before.wait_revisions.get("first").copied().unwrap_or(0)
585 );
586 assert_eq!(
587 running.wait_revisions.get("second"),
588 before.wait_revisions.get("second")
589 );
590 assert!(Arc::ptr_eq(
591 &second,
592 &running.turns.get_shared("second").unwrap()
593 ));
594 owner.shutdown().unwrap();
595 let owner = start_database_writer_at(&path, false).unwrap();
596 assert_eq!(owner.writer.committed_state().unwrap().turns, running.turns);
597 owner
598 .writer
599 .execute("delete projection", |connection| {
600 connection.execute(
601 "DELETE FROM materialized_sessions WHERE session_id='first'",
602 [],
603 )?;
604 Ok(())
605 })
606 .unwrap();
607 assert!(
608 !owner
609 .writer
610 .committed_state()
611 .unwrap()
612 .turns
613 .contains_key("first")
614 );
615 assert!(running.turns.contains_key("first"));
616 }
617
618 #[test]
619 fn publication_failure_stops_mutations_without_replaying_the_committed_write() {
620 let directory = tempfile::tempdir().unwrap();
621 let path = directory.path().join("controller.sqlite");
622 save_session_to(&path, &super::super::tests::session("selected", "project")).unwrap();
623 let owner = start_database_writer_at(&path, false).unwrap();
624 let error = owner
625 .writer
626 .execute("invalid committed record", |connection| {
627 connection.execute("UPDATE sessions SET resource_allocation='[]'", [])?;
628 Ok(())
629 })
630 .unwrap_err();
631 assert!(error.to_string().contains("do not replay"));
632 assert!(owner.writer.committed_state().is_err());
633 assert!(
634 owner
635 .writer
636 .execute("must not execute", |_| -> Result<()> {
637 panic!("a failed publication must close mutation service");
638 })
639 .is_err()
640 );
641 let connection = open_reader(&path).unwrap();
642 let stored: String = connection
643 .query_row(
644 "SELECT resource_allocation FROM sessions WHERE session_id='selected'",
645 [],
646 |row| row.get(0),
647 )
648 .unwrap();
649 assert_eq!(
650 stored, "[]",
651 "publication failure cannot undo or replay a commit"
652 );
653 assert!(owner.shutdown().is_err());
654 }
655
656 #[test]
657 fn secondary_connections_publish_committed_records_before_the_write_reply() {
658 let directory = tempfile::tempdir().unwrap();
659 let path = directory.path().join("controller.sqlite");
660 let owner = start_database_writer_at(&path, false).unwrap();
661 let before = owner.writer.committed_state().unwrap();
662 let record = super::super::tests::session("created", "project");
663 let saved = record.clone();
664 owner
665 .writer
666 .execute("create on secondary connection", move |_| {
667 save_session_to(&path, &saved)
668 })
669 .unwrap();
670 let after = owner.writer.committed_state().unwrap();
671 assert!(before.state.sessions.is_empty());
672 assert_eq!(after.state.sessions["created"], record);
673 assert_eq!(after.sequence, before.sequence + 1);
674 }
675
676 #[test]
680 fn startup_groups_and_subagent_reports_are_published_per_session() {
681 let directory = tempfile::tempdir().unwrap();
682 let path = directory.path().join("controller.sqlite");
683 let owner = start_database_writer_at(&path, false).unwrap();
684 let before = owner.writer.committed_state().unwrap();
685 assert!(before.startup_groups.is_empty() && before.subagent_reports.is_empty());
686
687 owner
688 .writer
689 .execute("queue startup", |connection| {
690 connection.execute(
691 "INSERT INTO startup_steps(session_id,group_id,command_id,step_json,phase)
692 VALUES ('first','group-1','first:prompt','{}','pending'),
693 ('second','group-2','second:prompt','{}','pending')",
694 [],
695 )?;
696 Ok(())
697 })
698 .unwrap();
699 let queued = owner.writer.committed_state().unwrap();
700 assert_eq!(queued.sequence, before.sequence + 1);
701 let reader = open_reader(&path).unwrap();
702 for session in ["first", "second"] {
703 assert_eq!(
704 queued.startup_groups[session],
705 startup::load_latest_startup_group_with(&reader, session).unwrap()
706 );
707 }
708
709 let handback = mj_core::subagent::SubagentHandback {
710 command_id: "task".into(),
711 message: "the report".into(),
712 recorded_at_ms: 7,
713 };
714 let recorded = handback.clone();
715 let report_path = path.clone();
716 owner
717 .writer
718 .execute("finish one session", move |connection| {
719 connection.execute(
720 "UPDATE startup_steps SET phase='failed',error='refused'
721 WHERE session_id='first'",
722 [],
723 )?;
724 sessions::record_subagent_handback_to(&report_path, "first", &recorded)?;
725 Ok(())
726 })
727 .unwrap();
728 let after = owner.writer.committed_state().unwrap();
729 assert_eq!(after.startup_groups["first"][0].phase, "failed");
730 assert_eq!(
731 after.startup_groups["first"][0].error.as_deref(),
732 Some("refused")
733 );
734 assert_eq!(
735 after.subagent_reports["first"].handback.as_ref(),
736 Some(&handback)
737 );
738 assert_eq!(
739 after.startup_groups["second"], queued.startup_groups["second"],
740 "the other session's group is untouched"
741 );
742 assert!(!after.subagent_reports.contains_key("second"));
743
744 owner.shutdown().unwrap();
746 let owner = start_database_writer_at(&path, false).unwrap();
747 let bootstrapped = owner.writer.committed_state().unwrap();
748 assert_eq!(bootstrapped.startup_groups, after.startup_groups);
749 assert_eq!(bootstrapped.subagent_reports, after.subagent_reports);
750
751 owner
752 .writer
753 .execute("drop startup", |connection| {
754 connection.execute("DELETE FROM startup_steps WHERE session_id='first'", [])?;
755 connection.execute("DELETE FROM subagent_handbacks", [])?;
756 Ok(())
757 })
758 .unwrap();
759 let cleared = owner.writer.committed_state().unwrap();
760 assert!(!cleared.startup_groups.contains_key("first"));
761 assert!(cleared.startup_groups.contains_key("second"));
762 assert!(cleared.subagent_reports.is_empty());
763 }
764
765 #[test]
766 fn rollback_and_no_op_updates_do_not_publish_changes() {
767 let directory = tempfile::tempdir().unwrap();
768 let path = directory.path().join("controller.sqlite");
769 save_session_to(&path, &super::super::tests::session("selected", "project")).unwrap();
770 let owner = start_database_writer_at(&path, false).unwrap();
771 let before = owner.writer.committed_state().unwrap();
772 let result: Result<()> = owner.writer.execute("rollback", |connection| {
773 let transaction = connection.transaction()?;
774 transaction.execute("UPDATE sessions SET title='rolled back'", [])?;
775 bail!("operation failed before commit");
776 });
777 assert!(result.is_err());
778 owner
779 .writer
780 .execute("no-op update", |connection| {
781 connection.execute("UPDATE sessions SET title=title", [])?;
782 Ok(())
783 })
784 .unwrap();
785 let after = owner.writer.committed_state().unwrap();
786 assert_eq!(after.sequence, before.sequence);
787 assert!(Arc::ptr_eq(&before, &after));
788 assert_eq!(after.state, before.state);
789 }
790
791 #[test]
792 fn a_committed_write_is_published_even_when_later_work_in_the_operation_fails() {
793 let directory = tempfile::tempdir().unwrap();
794 let path = directory.path().join("controller.sqlite");
795 save_session_to(&path, &super::super::tests::session("selected", "project")).unwrap();
796 let owner = start_database_writer_at(&path, false).unwrap();
797 let result: Result<()> = owner.writer.execute("failure after commit", |connection| {
798 connection.execute("UPDATE sessions SET title='committed'", [])?;
799 bail!("later work failed");
800 });
801 assert!(
802 result
803 .unwrap_err()
804 .to_string()
805 .contains("later work failed")
806 );
807 assert_eq!(
808 owner.writer.committed_state().unwrap().state.sessions["selected"].title,
809 "committed"
810 );
811 }
812
813 #[test]
814 fn deleting_a_session_publishes_its_absence_and_keeps_a_held_snapshot() {
815 let directory = tempfile::tempdir().unwrap();
816 let path = directory.path().join("controller.sqlite");
817 save_session_to(&path, &super::super::tests::session("selected", "project")).unwrap();
818 let owner = start_database_writer_at(&path, false).unwrap();
819 let before = owner.writer.committed_state().unwrap();
820 owner
821 .writer
822 .execute("delete with cascading related rows", |connection| {
823 connection.execute("DELETE FROM sessions WHERE session_id='selected'", [])?;
824 Ok(())
825 })
826 .unwrap();
827 assert!(
828 owner
829 .writer
830 .committed_state()
831 .unwrap()
832 .state
833 .sessions
834 .is_empty()
835 );
836 assert!(before.state.sessions.contains_key("selected"));
837 }
838}