1use super::Agent;
11use crate::channel::Channel;
12use zeph_llm::provider::{LlmProvider, Message, MessageMetadata, Role};
13
14impl<C: Channel> Agent<C> {
15 async fn call_llm_for_session_summary(
24 &self,
25 chat_messages: &[Message],
26 ) -> Option<zeph_memory::StructuredSummary> {
27 let provider = self.resolve_background_provider(
28 &self.services.memory.compaction.shutdown_summary_provider,
29 );
30 let timeout_dur = std::time::Duration::from_secs(
31 self.services
32 .memory
33 .compaction
34 .shutdown_summary_timeout_secs,
35 );
36 match tokio::time::timeout(
37 timeout_dur,
38 provider.chat_typed_erased::<zeph_memory::StructuredSummary>(chat_messages),
39 )
40 .await
41 {
42 Ok(Ok(s)) => Some(s),
43 Ok(Err(e)) => {
44 tracing::warn!(
45 "shutdown summary: structured LLM call failed, falling back to plain: {e:#}"
46 );
47 self.plain_text_summary_fallback(&provider, chat_messages, timeout_dur)
48 .await
49 }
50 Err(_) => {
51 tracing::warn!(
52 "shutdown summary: structured LLM call timed out after {}s, falling back to plain",
53 self.services
54 .memory
55 .compaction
56 .shutdown_summary_timeout_secs
57 );
58 self.plain_text_summary_fallback(&provider, chat_messages, timeout_dur)
59 .await
60 }
61 }
62 }
63 async fn plain_text_summary_fallback(
64 &self,
65 provider: &zeph_llm::any::AnyProvider,
66 chat_messages: &[Message],
67 timeout_dur: std::time::Duration,
68 ) -> Option<zeph_memory::StructuredSummary> {
69 match tokio::time::timeout(timeout_dur, provider.chat(chat_messages)).await {
70 Ok(Ok(plain)) => Some(zeph_memory::StructuredSummary {
71 summary: plain,
72 key_facts: vec![],
73 entities: vec![],
74 }),
75 Ok(Err(e)) => {
76 tracing::warn!("shutdown summary: plain LLM fallback failed: {e:#}");
77 None
78 }
79 Err(_) => {
80 tracing::warn!("shutdown summary: plain LLM fallback timed out");
81 None
82 }
83 }
84 }
85 pub(super) async fn flush_orphaned_tool_use_on_shutdown(&mut self) {
90 use zeph_llm::provider::{MessagePart, Role};
91
92 let msgs = &self.msg.messages;
96 let Some(asst_idx) = msgs.iter().rposition(|m| m.role == Role::Assistant) else {
98 return;
99 };
100 let asst_msg = &msgs[asst_idx];
101 let tool_use_ids: Vec<(&str, &str, &serde_json::Value)> = asst_msg
102 .parts
103 .iter()
104 .filter_map(|p| {
105 if let MessagePart::ToolUse { id, name, input } = p {
106 Some((id.as_str(), name.as_str(), input))
107 } else {
108 None
109 }
110 })
111 .collect();
112 if tool_use_ids.is_empty() {
113 return;
114 }
115
116 let paired_ids: std::collections::HashSet<&str> = msgs
118 .get(asst_idx + 1..)
119 .into_iter()
120 .flatten()
121 .filter(|m| m.role == Role::User)
122 .flat_map(|m| m.parts.iter())
123 .filter_map(|p| {
124 if let MessagePart::ToolResult { tool_use_id, .. } = p {
125 Some(tool_use_id.as_str())
126 } else {
127 None
128 }
129 })
130 .collect();
131
132 let unpaired: Vec<zeph_llm::provider::ToolUseRequest> = tool_use_ids
133 .iter()
134 .filter(|(id, _, _)| !paired_ids.contains(*id))
135 .map(|(id, name, input)| zeph_llm::provider::ToolUseRequest {
136 id: (*id).to_owned(),
137 name: (*name).to_owned().into(),
138 input: (*input).clone(),
139 })
140 .collect();
141
142 if unpaired.is_empty() {
143 return;
144 }
145
146 tracing::info!(
147 count = unpaired.len(),
148 "shutdown: persisting tombstone ToolResults for unpaired in-flight tool calls"
149 );
150 self.persist_cancelled_tool_results(&unpaired, Some(asst_idx + 1))
155 .await;
156 }
157 pub(super) async fn maybe_store_shutdown_summary(&mut self) {
168 if self.runtime.config.bare {
169 return;
170 }
171 if !self.services.memory.compaction.shutdown_summary {
172 return;
173 }
174 let Some(memory) = self.services.memory.persistence.memory.clone() else {
175 return;
176 };
177 let Some(conversation_id) = self.services.memory.persistence.conversation_id else {
178 return;
179 };
180
181 match memory.has_session_summary(conversation_id).await {
183 Ok(true) => {
184 tracing::debug!("shutdown summary: session already has a summary, skipping");
185 return;
186 }
187 Ok(false) => {}
188 Err(e) => {
189 tracing::warn!("shutdown summary: failed to check existing summary: {e:#}");
190 return;
191 }
192 }
193
194 let user_count = self
196 .msg
197 .messages
198 .iter()
199 .skip(1)
200 .filter(|m| m.role == Role::User)
201 .count();
202 let min_messages = self
203 .services
204 .memory
205 .compaction
206 .shutdown_summary_min_messages;
207 if user_count < min_messages {
208 tracing::debug!(
209 user_count,
210 min = min_messages,
211 "shutdown summary: too few user messages, skipping"
212 );
213 return;
214 }
215
216 self.channel
217 .send_status_best_effort("Saving session summary...")
218 .await;
219
220 let max = self
222 .services
223 .memory
224 .compaction
225 .shutdown_summary_max_messages;
226 if max == 0 {
227 tracing::debug!("shutdown summary: max_messages=0, skipping");
228 return;
229 }
230 let non_system: Vec<_> = self.msg.messages.iter().skip(1).collect();
231 let slice = if non_system.len() > max {
232 &non_system[non_system.len() - max..]
233 } else {
234 &non_system[..]
235 };
236
237 let msgs_for_prompt: Vec<(zeph_memory::MessageId, String, String)> = slice
238 .iter()
239 .map(|m| {
240 let role = match m.role {
241 Role::Assistant => "assistant".to_owned(),
242 Role::System => "system".to_owned(),
243 Role::User | _ => "user".to_owned(),
244 };
245 (zeph_memory::MessageId(0), role, m.content.clone())
246 })
247 .collect();
248
249 let prompt = zeph_memory::build_summarization_prompt(&msgs_for_prompt);
250 let chat_messages = vec![Message {
251 role: Role::User,
252 content: prompt,
253 parts: vec![],
254 metadata: MessageMetadata::default(),
255 }];
256
257 let Some(structured) = self.call_llm_for_session_summary(&chat_messages).await else {
258 self.channel.send_status_best_effort("").await;
259 return;
260 };
261
262 if let Err(e) = memory
263 .store_shutdown_summary(conversation_id, &structured.summary, &structured.key_facts)
264 .await
265 {
266 tracing::warn!("shutdown summary: storage failed: {e:#}");
267 } else {
268 tracing::info!(
269 conversation_id = conversation_id.0,
270 "shutdown summary stored"
271 );
272 }
273
274 self.channel.send_status_best_effort("").await;
275 }
276 #[tracing::instrument(name = "core.agent.shutdown", skip_all, level = "debug")]
292 #[allow(clippy::too_many_lines)]
293 pub async fn shutdown(&mut self) {
294 self.channel
295 .send_status_best_effort("Shutting down...")
296 .await;
297
298 self.provider.save_router_state().await;
300
301 if let Some(ref advisor) = self.services.orchestration.topology_advisor
303 && let Err(e) = advisor.save().await
304 {
305 tracing::warn!(error = %e, "adaptorch: failed to persist state");
306 }
307
308 if let Some(ref mut mgr) = self.services.orchestration.subagent_manager {
309 mgr.shutdown_all();
310 }
311
312 if let Some(ref manager) = self.services.mcp.manager {
313 manager.shutdown_all_shared().await;
314 }
315
316 if let Some(ref sink) = self.services.session.session_sink
321 && let Err(e) = sink.finalize().await
322 {
323 tracing::warn!(error = %e, "session anchor finalize failed");
324 }
325
326 if let Some(turns) = self.context_manager.turns_since_last_hard_compaction() {
330 self.update_metrics(|m| {
331 m.compaction_turns_after_hard.push(turns);
332 });
333 self.context_manager
334 .set_turns_since_last_hard_compaction(None);
335 }
336
337 if let Some(ref tx) = self.runtime.metrics.metrics_tx {
338 let m = tx.borrow();
339 if m.filter_applications > 0 {
340 #[allow(clippy::cast_precision_loss)]
341 let pct = if m.filter_raw_tokens > 0 {
342 m.filter_saved_tokens as f64 / m.filter_raw_tokens as f64 * 100.0
343 } else {
344 0.0
345 };
346 tracing::info!(
347 raw_tokens = m.filter_raw_tokens,
348 saved_tokens = m.filter_saved_tokens,
349 applications = m.filter_applications,
350 "tool output filtering saved ~{} tokens ({pct:.0}%)",
351 m.filter_saved_tokens,
352 );
353 }
354 if m.compaction_hard_count > 0 {
355 tracing::info!(
356 hard_compactions = m.compaction_hard_count,
357 turns_after_hard = ?m.compaction_turns_after_hard,
358 "hard compaction trajectory"
359 );
360 }
361 }
362
363 self.flush_orphaned_tool_use_on_shutdown().await;
367
368 if let Some(ref token) = self.services.experiments.cancel {
371 token.cancel();
372 }
373 if let Some(h) = self.services.experiments.handle.take() {
374 h.abort();
375 }
376
377 if let Some(memory) = self.services.memory.persistence.memory.as_ref() {
381 memory.cancel_graph_extraction();
382 }
383
384 self.runtime.lifecycle.supervisor.abort_all();
386
387 if let Some(h) = self.services.compression.pending_task_goal.take() {
390 h.abort();
391 }
392 if let Some(h) = self.services.compression.pending_sidequest_result.take() {
393 h.abort();
394 }
395 if let Some(h) = self.services.compression.pending_subgoal.take() {
396 h.abort();
397 }
398 self.flush_durable_writer().await;
399
400 self.services.learning_engine.learning_tasks.abort_all();
402
403 if let Some(h) = self.services.learning_engine.trace_extraction_handle.take() {
406 let deadline = std::time::Duration::from_mins(2);
407 match tokio::time::timeout(deadline, h.join()).await {
408 Ok(Ok(())) => {}
409 Ok(Err(e)) => tracing::warn!("trace_extraction: task error at shutdown: {e}"),
410 Err(_) => tracing::warn!(
411 "trace_extraction: timed out at shutdown ({}s), aborting",
412 deadline.as_secs()
413 ),
414 }
415 }
416
417 if let Some(h) = self
420 .services
421 .learning_engine
422 .heuristic_promotion_handle
423 .take()
424 {
425 h.abort();
426 }
427
428 if let Some(ref sentinel) = self.services.security.shadow_sentinel {
430 sentinel.drain_pending().await;
431 }
432
433 for _ in 0..4 {
438 tokio::task::yield_now().await;
439 }
440
441 self.maybe_store_shutdown_summary().await;
442 self.maybe_store_session_digest().await;
443
444 tracing::info!("agent shutdown complete");
445 }
446
447 async fn flush_durable_writer(&mut self) {
454 let flush_deadline = std::time::Duration::from_secs(2);
455 if let Some(ref writer) = self.services.orchestration.durable_writer {
456 match tokio::time::timeout(flush_deadline, writer.flush()).await {
457 Ok(Ok(())) => {}
458 Ok(Err(e)) => {
459 tracing::warn!(error = %e, "durable writer: flush on shutdown failed");
460 }
461 Err(_) => tracing::warn!("durable writer: flush timed out on shutdown"),
462 }
463 }
464 if let Some(h) = self.services.orchestration.durable_writer_task.take() {
465 h.abort();
466 }
467 if let Some(ref writer) = self.services.session.durable_writer {
468 match tokio::time::timeout(flush_deadline, writer.flush()).await {
469 Ok(Ok(())) => {}
470 Ok(Err(e)) => {
471 tracing::warn!(error = %e, "durable agent_turns writer: flush on shutdown failed");
472 }
473 Err(_) => tracing::warn!("durable agent_turns writer: flush timed out on shutdown"),
474 }
475 }
476 if let Some(ref ctx) = self.services.session.durable_ctx {
485 match tokio::time::timeout(
486 flush_deadline,
487 ctx.finalize(zeph_durable::ExecutionStatus::Completed),
488 )
489 .await
490 {
491 Ok(Ok(())) => {}
492 Ok(Err(e)) => {
493 tracing::warn!(
494 error = %e,
495 "durable agent_turns: failed to finalize execution on shutdown"
496 );
497 }
498 Err(_) => tracing::warn!("durable agent_turns: finalize timed out on shutdown"),
499 }
500 }
501 if let Some(h) = self.services.session.durable_writer_task.take() {
502 h.abort();
503 }
504 }
505}
506
507#[cfg(test)]
508mod tests {
509 use crate::agent::agent_tests::*;
510
511 fn agent_with_conversation() -> crate::agent::Agent<MockChannel> {
512 let provider = mock_provider(vec!["ok".into()]);
513 let channel = MockChannel::new(vec![]);
514 let registry = create_test_registry();
515 let executor = MockToolExecutor::no_tools();
516 let mut agent = crate::agent::Agent::new(provider, channel, registry, None, 5, executor);
517 agent.services.memory.persistence.conversation_id = Some(zeph_memory::ConversationId(1));
518 agent
519 }
520
521 #[tokio::test]
522 async fn flush_durable_writer_finalizes_the_p1_execution_as_completed() {
523 let dir = tempfile::tempdir().unwrap();
529 let db_url = dir.path().join("durable.db").to_string_lossy().into_owned();
530
531 let mut agent = agent_with_conversation();
532 agent.services.session.durable_agent_turns_config = Some(zeph_config::DurableConfig {
533 enabled: true,
534 agent_turns: true,
535 ..zeph_config::DurableConfig::default()
536 });
537 agent.services.session.durable_agent_turns_db_url = Some(db_url.clone());
538
539 agent.ensure_session_durable_ctx().await;
540 let exec_id = agent
541 .services
542 .session
543 .durable_ctx
544 .as_ref()
545 .expect("durable_ctx should be populated")
546 .execution_id();
547
548 agent.flush_durable_writer().await;
549
550 let backend = zeph_durable::LocalBackend::open(&db_url, 1_048_576)
551 .await
552 .unwrap();
553 let summaries = backend.list_executions(None, None, 10).await.unwrap();
554 let row = summaries
555 .iter()
556 .find(|s| s.execution_id == exec_id)
557 .expect("the execution's row must still exist");
558 assert_eq!(
559 row.status,
560 zeph_durable::ExecutionStatus::Completed,
561 "the P1 execution must finalize as Completed on graceful shutdown"
562 );
563 }
564
565 #[tokio::test]
566 async fn flush_durable_writer_finalize_is_bounded_by_the_2s_timeout() {
567 let dir = tempfile::tempdir().unwrap();
575 let db_url = dir.path().join("durable.db").to_string_lossy().into_owned();
576
577 let mut agent = agent_with_conversation();
578 agent.services.session.durable_agent_turns_config = Some(zeph_config::DurableConfig {
579 enabled: true,
580 agent_turns: true,
581 ..zeph_config::DurableConfig::default()
582 });
583 agent.services.session.durable_agent_turns_db_url = Some(db_url.clone());
584
585 agent.ensure_session_durable_ctx().await;
586 let exec_id = agent
587 .services
588 .session
589 .durable_ctx
590 .as_ref()
591 .expect("durable_ctx should be populated")
592 .execution_id();
593
594 let lock_holder = zeph_durable::LocalBackend::open(&db_url, 1_048_576)
598 .await
599 .unwrap();
600 let blocking_tx = zeph_db::begin_write(lock_holder.pool()).await.unwrap();
601
602 let start = std::time::Instant::now();
603 agent.flush_durable_writer().await;
604 let elapsed = start.elapsed();
605
606 drop(blocking_tx); assert!(
609 elapsed < std::time::Duration::from_secs(4),
610 "flush_durable_writer must return well within its 2s finalize timeout \
611 (plus the writer.flush() call's own bound), not the 5s sqlite busy_timeout; took \
612 {elapsed:?}"
613 );
614
615 let backend = zeph_durable::LocalBackend::open(&db_url, 1_048_576)
616 .await
617 .unwrap();
618 let summaries = backend.list_executions(None, None, 10).await.unwrap();
619 let row = summaries
620 .iter()
621 .find(|s| s.execution_id == exec_id)
622 .expect("the execution's row must still exist");
623 assert_eq!(
624 row.status,
625 zeph_durable::ExecutionStatus::Running,
626 "finalize must not have committed while the write lock was held elsewhere"
627 );
628
629 agent
634 .services
635 .session
636 .durable_ctx
637 .as_ref()
638 .unwrap()
639 .finalize(zeph_durable::ExecutionStatus::Completed)
640 .await
641 .unwrap();
642 let summaries = backend.list_executions(None, None, 10).await.unwrap();
643 let row = summaries
644 .iter()
645 .find(|s| s.execution_id == exec_id)
646 .expect("the execution's row must still exist");
647 assert_eq!(
648 row.status,
649 zeph_durable::ExecutionStatus::Completed,
650 "finalize succeeds once the write lock is free"
651 );
652 }
653}