Skip to main content

mj_controller/database/
native_agents.rs

1//! Native children are projections of their owner's relay, not worker sessions.
2use super::*;
3use mj_core::native_agent::{NativeAgent, NativeAgentEvent, NativeAgentState, NativeAgentView};
4use mj_core::relay::{RelayEvent, RelayObservation};
5
6/// Read identity and tail from one snapshot so replay cannot mix generations.
7pub fn load_native_agent_view(
8    owner: &str,
9    child: &str,
10    limit: usize,
11) -> Result<Option<NativeAgentView>> {
12    let mut connection = open_reader(&database_path())?;
13    let transaction = connection.transaction()?;
14    load_native_agent_view_from(&transaction, owner, child, limit)
15}
16
17fn load_native_agent_view_from(
18    connection: &Connection,
19    owner: &str,
20    child: &str,
21    limit: usize,
22) -> Result<Option<NativeAgentView>> {
23    let body: Option<String> = connection
24        .query_row(
25            "SELECT body FROM native_agents WHERE owner=?1 AND child=?2 AND staging=0",
26            params![owner, child],
27            |row| row.get(0),
28        )
29        .optional()?;
30    body.map(|body| {
31        let mut view: NativeAgentView = serde_json::from_str(&body)?;
32        view.projection.transcript = transcript(connection, owner, child, 0, limit)?;
33        Ok(view)
34    })
35    .transpose()
36}
37
38pub fn load_native_agents(owner: &str, limit: usize) -> Result<Vec<NativeAgentView>> {
39    let mut connection = open_reader(&database_path())?;
40    let transaction = connection.transaction()?;
41    load_native_agents_from(&transaction, owner, limit)
42}
43
44pub(super) fn load_native_agents_from(
45    connection: &Connection,
46    owner: &str,
47    limit: usize,
48) -> Result<Vec<NativeAgentView>> {
49    load_native_agents_at(connection, owner, limit, 0)
50}
51
52fn load_native_agents_at(
53    connection: &Connection,
54    owner: &str,
55    limit: usize,
56    staging: i64,
57) -> Result<Vec<NativeAgentView>> {
58    let mut statement = connection
59        .prepare("SELECT body FROM native_agents WHERE owner=?1 AND staging=?2 ORDER BY child")?;
60    let bodies = statement
61        .query_map(params![owner, staging], |row| row.get::<_, String>(0))?
62        .collect::<rusqlite::Result<Vec<_>>>()?;
63    bodies
64        .into_iter()
65        .map(|body| {
66            let mut view: NativeAgentView = serde_json::from_str(&body)?;
67            view.projection.transcript =
68                transcript(connection, owner, &view.agent.session_id, staging, limit)?;
69            Ok(view)
70        })
71        .collect()
72}
73
74fn transcript(
75    connection: &Connection,
76    owner: &str,
77    child: &str,
78    staging: i64,
79    limit: usize,
80) -> Result<Vec<Arc<TranscriptItem>>> {
81    let mut statement = connection.prepare("SELECT body FROM (SELECT position, stable_id, body FROM native_agent_transcript WHERE owner=?1 AND child=?2 AND staging=?3 ORDER BY position DESC, stable_id DESC LIMIT ?4) ORDER BY position, stable_id")?;
82    let bodies = statement
83        .query_map(params![owner, child, staging, limit as i64], |row| {
84            row.get::<_, String>(0)
85        })?
86        .collect::<rusqlite::Result<Vec<_>>>()?;
87    bodies
88        .into_iter()
89        .map(|body| Ok(Arc::new(serde_json::from_str(&body)?)))
90        .collect()
91}
92
93pub(super) fn apply_native_agent_event(
94    connection: &Connection,
95    owner: &str,
96    relay: &RelayEvent,
97) -> Result<()> {
98    let RelayObservation::NativeAgent { event } = &relay.observation else {
99        bail!("expected native agent observation")
100    };
101    match event {
102        NativeAgentEvent::Availability { reports, complete } => {
103            for mut view in load_native_agents_from(connection, owner, 0)? {
104                view.agent.apply_availability(reports, *complete);
105                save(connection, &view, 0)?;
106            }
107            return Ok(());
108        }
109        NativeAgentEvent::ReplayBegin => {
110            for mut view in load_native_agents_from(connection, owner, 0)? {
111                view.agent.invalidate_availability();
112                if view.agent.state == NativeAgentState::Running {
113                    view.agent.state = NativeAgentState::Disconnected;
114                }
115                save(connection, &view, 0)?;
116            }
117            connection.execute(
118                "DELETE FROM native_agents WHERE owner=?1 AND staging=1",
119                [owner],
120            )?;
121            connection.execute(
122                "INSERT OR IGNORE INTO native_agent_replay(owner) VALUES (?1)",
123                [owner],
124            )?;
125            return Ok(());
126        }
127        NativeAgentEvent::ReplayCommit => {
128            let staging: bool = connection.query_row(
129                "SELECT EXISTS(SELECT 1 FROM native_agent_replay WHERE owner=?1)",
130                [owner],
131                |row| row.get(0),
132            )?;
133            if staging {
134                for mut view in load_native_agents_at(connection, owner, 0, 1)? {
135                    view.agent.finish_replay();
136                    view.projection.execution = MaterializedExecutionState::Idle;
137                    finish_streaming(connection, owner, &view.agent.session_id, 1)?;
138                    save(connection, &view, 1)?;
139                }
140                connection.execute(
141                    "DELETE FROM native_agents WHERE owner=?1 AND staging=0 AND child IN (SELECT child FROM native_agents WHERE owner=?1 AND staging=1)",
142                    [owner],
143                )?;
144                connection.execute(
145                    "UPDATE native_agents SET staging=0 WHERE owner=?1 AND staging=1",
146                    [owner],
147                )?;
148                connection.execute("DELETE FROM native_agent_replay WHERE owner=?1", [owner])?;
149            }
150            return Ok(());
151        }
152        NativeAgentEvent::Disconnected => {
153            connection.execute(
154                "DELETE FROM native_agents WHERE owner=?1 AND staging=1",
155                [owner],
156            )?;
157            connection.execute("DELETE FROM native_agent_replay WHERE owner=?1", [owner])?;
158            for mut view in load_native_agents_from(connection, owner, 0)? {
159                view.agent.invalidate_availability();
160                if view.agent.state == NativeAgentState::Running {
161                    view.agent.state = NativeAgentState::Disconnected;
162                }
163                view.projection.execution = MaterializedExecutionState::Idle;
164                finish_streaming(connection, owner, &view.agent.session_id, 0)?;
165                view.projection.applied_event_ordinal = relay.ordinal;
166                save(connection, &view, 0)?;
167            }
168            return Ok(());
169        }
170        _ => {}
171    }
172    let staging: i64 = connection.query_row(
173        "SELECT EXISTS(SELECT 1 FROM native_agent_replay WHERE owner=?1)",
174        [owner],
175        |row| row.get(0),
176    )?;
177    let child = match event {
178        NativeAgentEvent::Spawned { session_id, .. }
179        | NativeAgentEvent::State { session_id, .. }
180        | NativeAgentEvent::Update { session_id, .. } => session_id,
181        _ => unreachable!(),
182    };
183    let stored: Option<String> = connection
184        .query_row(
185            "SELECT body FROM native_agents WHERE owner=?1 AND child=?2 AND staging=?3",
186            params![owner, child, staging],
187            |row| row.get(0),
188        )
189        .optional()?;
190    let mut view: NativeAgentView = match stored {
191        Some(body) => serde_json::from_str(&body)?,
192        None => {
193            let NativeAgentEvent::Spawned {
194                session_id,
195                parent_session_id,
196                name,
197                task,
198                capabilities,
199            } = event
200            else {
201                bail!("native child {child} update precedes spawn for owner {owner}");
202            };
203            let agent = NativeAgent {
204                availability: Default::default(),
205                availability_reason: None,
206                stable_id: None,
207                owner_session_id: owner.to_owned(),
208                session_id: session_id.clone(),
209                parent_session_id: parent_session_id.clone(),
210                name: name.clone(),
211                task: task.clone(),
212                capabilities: capabilities.clone(),
213                state: NativeAgentState::Running,
214            };
215            let mut projection = MaterializedSession::empty(agent.view_id());
216            projection.session_title = Some(name.clone());
217            projection.execution = MaterializedExecutionState::Running {
218                started_at_ms: relay.recorded_at_ms,
219            };
220            NativeAgentView {
221                generation_ordinal: relay.ordinal,
222                agent,
223                projection,
224            }
225        }
226    };
227    match event {
228        NativeAgentEvent::State { state, .. } => {
229            view.agent.state = *state;
230            if *state != NativeAgentState::Running {
231                finish_streaming(connection, owner, child, staging)?;
232            }
233            view.projection.execution = if *state == NativeAgentState::Running {
234                MaterializedExecutionState::Running {
235                    started_at_ms: relay.recorded_at_ms,
236                }
237            } else {
238                MaterializedExecutionState::Idle
239            };
240        }
241        NativeAgentEvent::Update { update, .. } => {
242            // Streaming text needs the tail; an old tool update additionally needs
243            // its named row, even when many newer messages have followed it.
244            view.projection.transcript = transcript(connection, owner, child, staging, 64)?;
245            let tool_id = match update.as_ref() {
246                agent_client_protocol::schema::v1::SessionUpdate::ToolCall(call) => {
247                    Some(call.tool_call_id.to_string())
248                }
249                agent_client_protocol::schema::v1::SessionUpdate::ToolCallUpdate(call) => {
250                    Some(call.tool_call_id.to_string())
251                }
252                _ => None,
253            };
254            if let Some(id) = tool_id {
255                let stable_id = format!("tool:{id}");
256                if !view
257                    .projection
258                    .transcript
259                    .iter()
260                    .any(|item| item.stable_id == stable_id)
261                {
262                    let body: Option<String> = connection.query_row("SELECT body FROM native_agent_transcript WHERE owner=?1 AND child=?2 AND staging=?3 AND stable_id=?4", params![owner, child, staging, stable_id], |row| row.get(0)).optional()?;
263                    if let Some(body) = body {
264                        view.projection
265                            .transcript
266                            .push(Arc::new(serde_json::from_str(&body)?));
267                    }
268                }
269            }
270            let mutation =
271                mj_transcript::projection::project_native_update(&view.projection, relay, update)?;
272            for change in mutation.transcript {
273                match change {
274                    TranscriptMutation::Upsert(item) => {
275                        connection.execute("INSERT INTO native_agent_transcript(owner,child,staging,stable_id,position,body) VALUES (?1,?2,?3,?4,?5,?6) ON CONFLICT(owner,child,staging,stable_id) DO UPDATE SET body=excluded.body", params![owner, child, staging, item.stable_id, item.position, serde_json::to_string(&item)?])?;
276                    }
277                    TranscriptMutation::Remove { stable_id } => {
278                        connection.execute("DELETE FROM native_agent_transcript WHERE owner=?1 AND child=?2 AND staging=?3 AND stable_id=?4", params![owner, child, staging, stable_id])?;
279                    }
280                }
281            }
282            view.projection.transcript.clear();
283        }
284        _ => {}
285    }
286    view.projection.applied_event_ordinal = relay.ordinal;
287    view.projection
288        .applied_event_digest
289        .clone_from(&relay.digest);
290    view.projection.last_activity_at_ms = Some(relay.recorded_at_ms);
291    save(connection, &view, staging)
292}
293
294fn finish_streaming(connection: &Connection, owner: &str, child: &str, staging: i64) -> Result<()> {
295    connection.execute("UPDATE native_agent_transcript SET body=json_set(body,'$.body.streaming',json('false')) WHERE owner=?1 AND child=?2 AND staging=?3 AND json_extract(body,'$.body.kind') IN ('agent','thought')", params![owner, child, staging])?;
296    Ok(())
297}
298
299fn save(connection: &Connection, view: &NativeAgentView, staging: i64) -> Result<()> {
300    connection.execute("INSERT INTO native_agents(owner,child,staging,body) VALUES (?1,?2,?3,?4) ON CONFLICT(owner,child,staging) DO UPDATE SET body=excluded.body", params![view.agent.owner_session_id, view.agent.session_id, staging, serde_json::to_string(view)?])?;
301    Ok(())
302}
303
304pub fn native_agent_history(
305    owner: &str,
306    child: &str,
307    before: Option<(u64, String)>,
308) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
309    let mut reader = open_reader(&database_path())?;
310    let connection = reader.transaction()?;
311    native_agent_history_from(&connection, owner, child, before)
312}
313
314fn native_agent_history_from(
315    connection: &Connection,
316    owner: &str,
317    child: &str,
318    before: Option<(u64, String)>,
319) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
320    let body: String = connection.query_row(
321        "SELECT body FROM native_agents WHERE owner=?1 AND child=?2 AND staging=0",
322        params![owner, child],
323        |row| row.get(0),
324    )?;
325    let view: NativeAgentView = serde_json::from_str(&body)?;
326    let (position, stable_id) = before.unwrap_or((i64::MAX as u64, String::new()));
327    let mut statement = connection.prepare("SELECT body FROM native_agent_transcript WHERE owner=?1 AND child=?2 AND staging=0 AND (position, stable_id) < (?3, ?4) ORDER BY position DESC, stable_id DESC LIMIT 201")?;
328    let bodies = statement
329        .query_map(params![owner, child, position, stable_id], |row| {
330            row.get::<_, String>(0)
331        })?
332        .collect::<rusqlite::Result<Vec<_>>>()?;
333    let has_more = bodies.len() > 200;
334    let mut items = bodies
335        .into_iter()
336        .take(200)
337        .map(|body| Ok(Arc::new(serde_json::from_str(&body)?)))
338        .collect::<Result<Vec<_>>>()?;
339    items.reverse();
340    Ok(mj_core::native_agent::NativeAgentHistoryPage {
341        generation_ordinal: view.generation_ordinal,
342        items,
343        has_more,
344    })
345}
346
347#[cfg(test)]
348mod tests {
349    use super::*;
350    use mj_core::native_agent::NativeAgentCapabilities;
351    use mj_core::relay::{RELAY_EVENT_FORMAT_V1, RELAY_EVENT_GENESIS_DIGEST, relay_event_digest};
352    use serde_json::json;
353
354    fn event(ordinal: u64, native: NativeAgentEvent) -> RelayEvent {
355        let mut event = RelayEvent {
356            format: RELAY_EVENT_FORMAT_V1,
357            ordinal,
358            previous_digest: RELAY_EVENT_GENESIS_DIGEST.into(),
359            digest: String::new(),
360            recorded_at_ms: ordinal as i64 * 100,
361            command_id: None,
362            observation: RelayObservation::NativeAgent { event: native },
363        };
364        event.digest = relay_event_digest(&event).unwrap();
365        event
366    }
367
368    fn spawn(child: &str, parent: Option<&str>) -> NativeAgentEvent {
369        NativeAgentEvent::Spawned {
370            session_id: child.into(),
371            parent_session_id: parent.map(str::to_owned),
372            name: child.into(),
373            task: "Inspect code".into(),
374            capabilities: NativeAgentCapabilities::default(),
375        }
376    }
377
378    fn text(child: &str, message: &str) -> NativeAgentEvent {
379        NativeAgentEvent::Update { session_id: child.into(), update: Box::new(serde_json::from_value(json!({"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":message}})).unwrap()) }
380    }
381
382    #[test]
383    fn committed_native_agents_keep_replay_private_and_publish_cascading_deletion() {
384        let directory = tempfile::tempdir().unwrap();
385        let path = directory.path().join("controller.sqlite");
386        save_session_to(&path, &super::super::tests::session("owner", "project")).unwrap();
387        let writer = start_database_writer_at(&path, false).unwrap();
388        writer
389            .writer
390            .execute("spawn native child", |connection| {
391                apply_native_agent_event(connection, "owner", &event(1, spawn("child", None)))
392            })
393            .unwrap();
394        let held = writer.writer.committed_state().unwrap();
395        assert!(held.native_agents["owner"].contains_key("child"));
396        writer.shutdown().unwrap();
397        let writer = start_database_writer_at(&path, false).unwrap();
398        assert_eq!(
399            writer.writer.committed_state().unwrap().native_agents,
400            held.native_agents
401        );
402        writer
403            .writer
404            .execute("begin native replay", |connection| {
405                apply_native_agent_event(
406                    connection,
407                    "owner",
408                    &event(2, NativeAgentEvent::ReplayBegin),
409                )
410            })
411            .unwrap();
412        let before = writer.writer.committed_state().unwrap();
413        assert_eq!(
414            before.native_agents["owner"]["child"].agent.state,
415            NativeAgentState::Disconnected
416        );
417        writer
418            .writer
419            .execute("stage native child", |connection| {
420                apply_native_agent_event(connection, "owner", &event(3, spawn("replacement", None)))
421            })
422            .unwrap();
423        let staged = writer.writer.committed_state().unwrap();
424        assert_eq!(staged.sequence, before.sequence);
425        assert_eq!(staged.native_agents, before.native_agents);
426        writer
427            .writer
428            .execute("publish native replay", |connection| {
429                apply_native_agent_event(
430                    connection,
431                    "owner",
432                    &event(4, NativeAgentEvent::ReplayCommit),
433                )
434            })
435            .unwrap();
436        let published = writer.writer.committed_state().unwrap();
437        assert!(published.native_agents["owner"].contains_key("child"));
438        assert!(published.native_agents["owner"].contains_key("replacement"));
439        writer
440            .writer
441            .execute("delete native owner", |connection| {
442                connection.execute("DELETE FROM sessions WHERE session_id='owner'", [])?;
443                Ok(())
444            })
445            .unwrap();
446        assert!(
447            writer
448                .writer
449                .committed_state()
450                .unwrap()
451                .native_agents
452                .is_empty()
453        );
454        assert!(held.native_agents["owner"].contains_key("child"));
455    }
456
457    #[test]
458    fn completed_children_remain_reusable_but_replay_does_not_prove_availability() {
459        use mj_core::native_agent::{NativeAgentAvailability, NativeAgentAvailabilityReport};
460        let directory = tempfile::tempdir().unwrap();
461        let path = directory.path().join("test.sqlite3");
462        save_session_to(&path, &super::super::tests::session("owner", "project")).unwrap();
463        let connection = open(&path).unwrap();
464        let apply = |ordinal, update| {
465            apply_native_agent_event(&connection, "owner", &event(ordinal, update)).unwrap()
466        };
467        apply(1, spawn("child", None));
468        apply(2, text("child", "retained context"));
469        apply(
470            3,
471            NativeAgentEvent::State {
472                session_id: "child".into(),
473                state: NativeAgentState::Completed,
474            },
475        );
476        apply(
477            4,
478            NativeAgentEvent::Availability {
479                reports: vec![NativeAgentAvailabilityReport {
480                    session_id: "child".into(),
481                    stable_id: Some("stable-child".into()),
482                    state: None,
483                    availability: NativeAgentAvailability::Available,
484                    reason: None,
485                }],
486                complete: true,
487            },
488        );
489        let reusable = load_native_agents_from(&connection, "owner", 200).unwrap();
490        assert_eq!(reusable[0].agent.state, NativeAgentState::Completed);
491        assert_eq!(
492            reusable[0].agent.availability,
493            NativeAgentAvailability::Available
494        );
495        apply(
496            5,
497            NativeAgentEvent::State {
498                session_id: "child".into(),
499                state: NativeAgentState::Running,
500            },
501        );
502        assert_eq!(
503            load_native_agents_from(&connection, "owner", 200).unwrap()[0]
504                .projection
505                .transcript,
506            reusable[0].projection.transcript
507        );
508        apply(6, NativeAgentEvent::ReplayBegin);
509        apply(7, NativeAgentEvent::ReplayCommit);
510        let retained = load_native_agents_from(&connection, "owner", 200).unwrap();
511        assert_eq!(
512            retained.len(),
513            1,
514            "missing replay must retain inspectable history"
515        );
516        assert_eq!(
517            retained[0].agent.availability,
518            NativeAgentAvailability::Unknown
519        );
520        assert_eq!(retained[0].agent.state, NativeAgentState::Disconnected);
521        apply(
522            8,
523            NativeAgentEvent::Availability {
524                reports: vec![],
525                complete: false,
526            },
527        );
528        assert_eq!(
529            load_native_agents_from(&connection, "owner", 200).unwrap()[0]
530                .agent
531                .availability,
532            NativeAgentAvailability::Unknown
533        );
534        apply(
535            9,
536            NativeAgentEvent::Availability {
537                reports: vec![],
538                complete: true,
539            },
540        );
541        let missing = load_native_agents_from(&connection, "owner", 200).unwrap();
542        assert_eq!(
543            missing[0].agent.availability,
544            NativeAgentAvailability::Unavailable
545        );
546        assert_eq!(missing[0].projection.transcript.len(), 1);
547    }
548
549    #[test]
550    fn native_replay_is_atomic_and_does_not_duplicate_child_history() {
551        let directory = tempfile::tempdir().unwrap();
552        let path = directory.path().join("test.sqlite3");
553        save_session_to(&path, &super::super::tests::session("owner", "project")).unwrap();
554        let connection = open(&path).unwrap();
555        let apply = |ordinal, update| {
556            apply_native_agent_event(&connection, "owner", &event(ordinal, update)).unwrap()
557        };
558        apply(1, spawn("child", None));
559        apply(2, text("child", &"x".repeat(80_000)));
560        let mut original = load_native_agents_from(&connection, "owner", 200).unwrap();
561        assert_eq!(original[0].projection.transcript.len(), 1);
562        apply(3, NativeAgentEvent::ReplayBegin);
563        original[0].agent.invalidate_availability();
564        original[0].agent.state = NativeAgentState::Disconnected;
565        apply(4, spawn("child", None));
566        apply(5, text("child", "replacement"));
567        assert_eq!(
568            load_native_agents_from(&connection, "owner", 200).unwrap(),
569            original
570        );
571        apply(6, NativeAgentEvent::ReplayCommit);
572        let replaced = load_native_agents_from(&connection, "owner", 200).unwrap();
573        assert_eq!(replaced[0].projection.transcript.len(), 1);
574        assert!(
575            serde_json::to_string(&replaced[0])
576                .unwrap()
577                .contains("replacement")
578        );
579        apply(7, NativeAgentEvent::ReplayBegin);
580        apply(8, spawn("child", None));
581        apply(9, text("child", "replacement"));
582        apply(10, NativeAgentEvent::ReplayCommit);
583        assert_eq!(
584            load_native_agents_from(&connection, "owner", 200).unwrap()[0]
585                .projection
586                .transcript
587                .len(),
588            1
589        );
590        apply(11, NativeAgentEvent::ReplayBegin);
591        apply(12, spawn("incomplete", None));
592        apply(13, NativeAgentEvent::Disconnected);
593        let recovered = load_native_agents_from(&connection, "owner", 200).unwrap();
594        assert_eq!(recovered.len(), 1);
595        assert_eq!(recovered[0].agent.session_id, "child");
596        assert_eq!(recovered[0].agent.state, NativeAgentState::Disconnected);
597        assert!(matches!(
598            recovered[0].projection.transcript[0].body,
599            TranscriptBody::Agent {
600                streaming: false,
601                ..
602            }
603        ));
604        assert!(
605            load_materialized_session_from(&path, "owner")
606                .unwrap()
607                .unwrap()
608                .transcript
609                .is_empty()
610        );
611    }
612
613    #[test]
614    fn native_children_isolate_tool_ids_and_preserve_nesting_and_terminal_states() {
615        let directory = tempfile::tempdir().unwrap();
616        let path = directory.path().join("test.sqlite3");
617        save_session_to(&path, &super::super::tests::session("owner", "project")).unwrap();
618        let connection = open(&path).unwrap();
619        for (ordinal, native) in [spawn("child", None), spawn("nested", Some("child"))]
620            .into_iter()
621            .enumerate()
622        {
623            apply_native_agent_event(&connection, "owner", &event(ordinal as u64 + 1, native))
624                .unwrap();
625        }
626        for (index, child) in ["child", "nested"].into_iter().enumerate() {
627            let update = serde_json::from_value(json!({"sessionUpdate":"tool_call","toolCallId":"same-id","title":child,"kind":"read","status":"in_progress"})).unwrap();
628            apply_native_agent_event(
629                &connection,
630                "owner",
631                &event(
632                    index as u64 + 3,
633                    NativeAgentEvent::Update {
634                        session_id: child.into(),
635                        update: Box::new(update),
636                    },
637                ),
638            )
639            .unwrap();
640        }
641        apply_native_agent_event(
642            &connection,
643            "owner",
644            &event(
645                5,
646                NativeAgentEvent::State {
647                    session_id: "nested".into(),
648                    state: NativeAgentState::Failed,
649                },
650            ),
651        )
652        .unwrap();
653        let views = load_native_agents_from(&connection, "owner", 200).unwrap();
654        assert_eq!(views.len(), 2);
655        assert_eq!(views[0].projection.transcript.len(), 1);
656        assert_eq!(views[1].projection.transcript.len(), 1);
657        assert_eq!(views[0].agent.state, NativeAgentState::Running);
658        assert_eq!(views[1].agent.state, NativeAgentState::Failed);
659        assert_eq!(views[1].agent.parent_view_id(), views[0].agent.view_id());
660    }
661    #[test]
662    fn native_history_pages_older_messages_without_duplicates() {
663        let directory = tempfile::tempdir().unwrap();
664        let path = directory.path().join("test.sqlite3");
665        save_session_to(&path, &super::super::tests::session("owner", "project")).unwrap();
666        let connection = open(&path).unwrap();
667        apply_native_agent_event(&connection, "owner", &event(1, spawn("child", None))).unwrap();
668        for ordinal in 2..=232 {
669            let update = serde_json::from_value(json!({"sessionUpdate":"user_message_chunk","content":{"type":"text","text":format!("message {ordinal}: {}", "x".repeat(1024))}})).unwrap();
670            apply_native_agent_event(
671                &connection,
672                "owner",
673                &event(
674                    ordinal,
675                    NativeAgentEvent::Update {
676                        session_id: "child".into(),
677                        update: Box::new(update),
678                    },
679                ),
680            )
681            .unwrap();
682        }
683        let page = native_agent_history_from(&connection, "owner", "child", None).unwrap();
684        assert_eq!(page.items.len(), 200);
685        assert!(serde_json::to_vec(&page).unwrap().len() > 64 * 1024);
686        assert!(page.has_more);
687        let first = &page.items[0];
688        let older = native_agent_history_from(
689            &connection,
690            "owner",
691            "child",
692            Some((first.position, first.stable_id.clone())),
693        )
694        .unwrap();
695        assert_eq!(older.items.len(), 31);
696        assert!(!older.has_more);
697        assert_eq!(older.generation_ordinal, page.generation_ordinal);
698        assert_eq!(older.items.last().unwrap().position + 1, first.position);
699        assert_eq!(older.items[0].position, 2);
700    }
701}