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