1use super::*;
3use mj_core::native_agent::{NativeAgent, NativeAgentEvent, NativeAgentState, NativeAgentView};
4use mj_core::relay::{RelayEvent, RelayObservation};
5
6pub 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 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}