Skip to main content

mj_controller/database/
committed.rs

1//! The writer publishes immutable records before acknowledging an operation.
2//!
3//! Connection-local triggers collect affected keys, including cascading deletes
4//! and writes made through another connection on the writer thread. Keys are
5//! hints, never commit receipts: after the operation returns we read committed
6//! rows and publish only actual differences. A rolled-back transaction therefore
7//! cannot publish a change. Nothing is added to the durable schema.
8
9use 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
25/// Every writable connection gets the same observer. The collector is active
26/// only inside an accepted writer job and only for that job's database path.
27pub(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        // Observe tables that exist; opening a connection must not depend on
65        // an unrelated optional table. Its own read/write still reports damage.
66        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    /// Each session's latest startup group, keyed by session. A session with
130    /// no group has no entry. API status readers ask this record instead of
131    /// the store, and a change to it publishes a revision.
132    pub startup_groups: SnapshotMap<String, Vec<StartupDelivery>>,
133    /// Each sub-agent child's recorded report, keyed by child session. A
134    /// child with nothing recorded has no entry.
135    pub subagent_reports: SnapshotMap<String, mj_core::subagent::SubagentReport>,
136    pub turns: SnapshotMap<String, CommittedTurn>,
137    /// Only this session's committed wait inputs advance this token.
138    pub wait_revisions: SnapshotMap<String, u64>,
139}
140
141#[derive(Clone, Debug, PartialEq, Eq)]
142pub struct CommittedTurn {
143    pub state: materialized::MaterializedTurnState,
144    /// Read once when a failed turn changes, never on each observation.
145    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
249/// Reads all changed records from one committed WAL snapshot. The caller owns
250/// the sole write lane, so no later job can overtake this publication.
251pub(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                // A formerly unsupported parent/child may now be readable.
354                // Re-evaluate only its indexed relationships, not all history.
355                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    /// A wait reads a session's startup status and a child's report from the
677    /// published records, so every write to them must publish, keyed to the
678    /// session it changed and equal to what the store holds.
679    #[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        // A writer that starts over the same store publishes the same records.
745        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}