1#![allow(deprecated)]
44
45use std::collections::{HashMap, HashSet};
46
47use crate::{
48 ActivityDeltaEvent, ActivityMessage, ActivitySnapshotEvent, AssistantMessage, DeveloperMessage,
49 Event, InputContent, Interrupt, JsonObject, Message, MessageId, PatchOperation,
50 ReasoningEncryptedValueEvent, ReasoningEncryptedValueSubtype, ReasoningMessage,
51 ReasoningMessageChunkEvent, RunId, RunOutcome, SubagentErrorEvent, SubagentFinishedEvent,
52 SubagentOutcome, SubagentRunId, SubagentStartedEvent, SystemMessage, TextInputContent,
53 TextMessageChunkEvent, TextMessageRole, ThreadId, ToolCall, ToolCallChunkEvent, ToolCallId,
54 ToolMessage, UserContent, UserMessage,
55};
56use serde::Deserialize;
57use serde_json::Value;
58
59use crate::client::error::{Error, Result};
60use crate::metadata::merge_metadata_into;
61
62#[derive(Clone, Debug, PartialEq)]
68#[non_exhaustive]
69pub enum Changed {
70 Nothing,
73 Message(MessageChange),
75 MessagesReplaced,
78 State,
80 Reasoning(ReasoningChange),
83 Subagent(SubagentChange),
87 RunStarted {
89 thread_id: ThreadId,
91 run_id: RunId,
93 },
94 RunFinished {
97 outcome: RunOutcome,
99 result: Option<Value>,
101 },
102 RunError {
104 message: String,
106 code: Option<String>,
108 },
109}
110
111#[derive(Clone, Debug, PartialEq)]
113pub struct MessageChange {
114 pub index: usize,
116 pub id: MessageId,
118 pub kind: MessageChangeKind,
120}
121
122#[derive(Clone, Debug, PartialEq)]
124#[non_exhaustive]
125pub enum MessageChangeKind {
126 Started,
128 Content {
130 delta: String,
132 },
133 Ended,
135 ToolCallStarted {
138 tool_call_id: ToolCallId,
140 name: String,
142 },
143 ToolCallArgs {
150 tool_call_id: ToolCallId,
152 delta: String,
154 },
155 ToolCallEnded {
157 tool_call_id: ToolCallId,
159 },
160 ToolResult {
162 tool_call_id: ToolCallId,
164 },
165 Activity,
167 EncryptedValue,
170}
171
172#[derive(Clone, Debug, PartialEq)]
174pub struct ReasoningChange {
175 pub id: MessageId,
177 pub kind: ReasoningChangeKind,
179}
180
181#[derive(Clone, Debug, PartialEq)]
190#[non_exhaustive]
191pub enum ReasoningChangeKind {
192 Started,
194 Content {
196 delta: String,
198 },
199 Ended,
201 EncryptedValue,
203}
204
205#[derive(Clone, Debug, PartialEq)]
212pub struct Subagent {
213 pub run_id: SubagentRunId,
215 pub name: String,
217 pub description: Option<String>,
219 pub parent_subagent_run_id: Option<SubagentRunId>,
221 pub parent_tool_call_id: Option<ToolCallId>,
223 pub parent_message_id: Option<MessageId>,
225 pub status: SubagentStatus,
227}
228
229#[derive(Clone, Debug, PartialEq)]
231#[non_exhaustive]
232pub enum SubagentStatus {
233 Running,
235 Finished {
237 result: Option<Value>,
239 },
240 Suspended {
245 result: Option<Value>,
247 interrupt_ids: Vec<String>,
251 },
252 Failed {
254 message: String,
256 code: Option<String>,
258 },
259}
260
261#[derive(Clone, Debug, PartialEq)]
263pub struct SubagentChange {
264 pub index: usize,
266 pub run_id: SubagentRunId,
268 pub kind: SubagentChangeKind,
270}
271
272#[derive(Clone, Debug, PartialEq)]
274#[non_exhaustive]
275pub enum SubagentChangeKind {
276 Started,
278 Resumed,
281 Finished,
283 Suspended,
285 Failed,
287}
288
289#[derive(Clone, Debug)]
293pub struct Applier {
294 messages: Vec<Message>,
295 by_id: HashMap<MessageId, usize>,
296 tool_calls: HashMap<ToolCallId, usize>,
298 open_text: Vec<(Option<SubagentRunId>, MessageId)>,
302 state: Value,
303 reasoning: Vec<ReasoningMessage>,
304 reasoning_by_id: HashMap<MessageId, usize>,
305 open_reasoning: HashSet<MessageId>,
309 current_reasoning: Option<MessageId>,
312 thinking_counter: u64,
314 reasoning_streams: Vec<(Option<SubagentRunId>, MessageId)>,
317 thread_id: Option<ThreadId>,
318 run_id: Option<RunId>,
319 interrupts: Vec<Interrupt>,
320 subagents: Vec<Subagent>,
322 subagent_by_id: HashMap<SubagentRunId, usize>,
323}
324
325impl Default for Applier {
326 fn default() -> Self {
327 Self::new()
328 }
329}
330
331impl Applier {
332 pub fn new() -> Self {
337 Self {
338 messages: Vec::new(),
339 by_id: HashMap::new(),
340 tool_calls: HashMap::new(),
341 open_text: Vec::new(),
342 state: Value::Object(JsonObject::new()),
343 reasoning: Vec::new(),
344 reasoning_by_id: HashMap::new(),
345 open_reasoning: HashSet::new(),
346 current_reasoning: None,
347 thinking_counter: 0,
348 reasoning_streams: Vec::new(),
349 thread_id: None,
350 run_id: None,
351 interrupts: Vec::new(),
352 subagents: Vec::new(),
353 subagent_by_id: HashMap::new(),
354 }
355 }
356
357 #[must_use]
359 pub fn with_messages(mut self, messages: impl Into<Vec<Message>>) -> Self {
360 self.replace_messages(messages.into());
361 self
362 }
363
364 #[must_use]
366 pub fn with_state(mut self, state: impl Into<Value>) -> Self {
367 self.state = state.into();
368 self
369 }
370
371 pub fn messages(&self) -> &[Message] {
373 &self.messages
374 }
375
376 pub fn message(&self, id: &MessageId) -> Option<&Message> {
378 self.by_id
379 .get(id)
380 .and_then(|index| self.messages.get(*index))
381 }
382
383 pub fn text_of(&self, id: impl Into<MessageId>) -> Option<&str> {
387 match self.message(&id.into())? {
388 Message::Assistant(m) => m.content.as_deref(),
389 Message::System(m) => Some(&m.content),
390 Message::Developer(m) => Some(&m.content),
391 Message::Tool(m) => Some(&m.content),
392 Message::Reasoning(m) => Some(&m.content),
393 Message::User(m) => match &m.content {
394 UserContent::Text(text) => Some(text),
395 UserContent::Parts(parts) => parts.iter().find_map(|part| match part {
396 InputContent::Text(text) => Some(text.text.as_str()),
397 _ => None,
398 }),
399 },
400 Message::Activity(_) => None,
401 }
402 }
403
404 pub fn state(&self) -> &Value {
409 &self.state
410 }
411
412 pub fn state_as<T: for<'de> Deserialize<'de>>(&self) -> Result<T> {
429 T::deserialize(&self.state).map_err(Error::State)
430 }
431
432 pub fn set_state(&mut self, state: impl Into<Value>) {
434 self.state = state.into();
435 }
436
437 pub fn reasoning(&self) -> &[ReasoningMessage] {
443 &self.reasoning
444 }
445
446 pub fn reasoning_text(&self, id: &MessageId) -> Option<&str> {
448 self.reasoning_by_id
449 .get(id)
450 .and_then(|index| self.reasoning.get(*index))
451 .map(|message| message.content.as_str())
452 }
453
454 pub fn thread_id(&self) -> Option<&ThreadId> {
456 self.thread_id.as_ref()
457 }
458
459 pub fn run_id(&self) -> Option<&RunId> {
461 self.run_id.as_ref()
462 }
463
464 pub fn interrupts(&self) -> &[Interrupt] {
467 &self.interrupts
468 }
469
470 pub fn subagents(&self) -> &[Subagent] {
476 &self.subagents
477 }
478
479 pub fn subagent(&self, run_id: &SubagentRunId) -> Option<&Subagent> {
481 self.subagent_by_id
482 .get(run_id)
483 .and_then(|index| self.subagents.get(*index))
484 }
485
486 pub fn push_message(&mut self, message: Message) -> usize {
491 let index = self.messages.len();
492 self.by_id.insert(message.id().clone(), index);
493 if let Message::Assistant(assistant) = &message {
494 for call in assistant.tool_calls.iter().flatten() {
495 self.tool_calls.insert(call.id.clone(), index);
496 }
497 }
498 self.messages.push(message);
499 index
500 }
501
502 pub fn apply(&mut self, event: &Event) -> Result<Changed> {
509 let changed = self.apply_inner(event)?;
510 if let Some(metadata) = event.metadata() {
513 self.merge_event_metadata(event, &changed, metadata);
514 }
515 Ok(changed)
516 }
517
518 fn apply_inner(&mut self, event: &Event) -> Result<Changed> {
519 match event {
520 Event::TextMessageStart(e) => Ok(self.text_start(
521 e.message_id.clone(),
522 e.role,
523 e.name.clone(),
524 e.subagent_run_id.clone(),
525 )),
526 Event::TextMessageContent(e) => {
527 self.text_content(&e.message_id, &e.delta, &e.subagent_run_id)
528 }
529 Event::TextMessageEnd(e) => Ok(self.text_end(&e.message_id)),
530 Event::TextMessageChunk(e) => self.text_chunk(e),
531
532 Event::ToolCallStart(e) => Ok(self.tool_call_start(
533 e.tool_call_id.clone(),
534 e.tool_call_name.clone(),
535 e.parent_message_id.clone(),
536 e.subagent_run_id.clone(),
537 )),
538 Event::ToolCallArgs(e) => self.tool_call_args(&e.tool_call_id, &e.delta),
539 Event::ToolCallEnd(e) => Ok(self.tool_call_end(&e.tool_call_id)),
540 Event::ToolCallChunk(e) => self.tool_call_chunk(e),
541 Event::ToolCallResult(e) => Ok(self.tool_call_result(
542 e.message_id.clone(),
543 e.tool_call_id.clone(),
544 e.content.clone(),
545 e.subagent_run_id.clone(),
546 )),
547
548 Event::StateSnapshot(e) => {
549 self.state = e.snapshot.clone();
550 Ok(Changed::State)
551 }
552 Event::StateDelta(e) => {
553 apply_patch(&mut self.state, &e.delta, "state")?;
554 Ok(Changed::State)
555 }
556 Event::MessagesSnapshot(e) => {
557 self.merge_snapshot(e.messages.clone());
558 Ok(Changed::MessagesReplaced)
559 }
560
561 Event::ActivitySnapshot(e) => Ok(self.activity_snapshot(e)),
562 Event::ActivityDelta(e) => self.activity_delta(e),
563
564 Event::ReasoningStart(e) => {
565 Ok(self.reasoning_start(e.message_id.clone(), e.subagent_run_id.clone()))
566 }
567 Event::ReasoningMessageStart(e) => {
568 Ok(self.reasoning_start(e.message_id.clone(), e.subagent_run_id.clone()))
569 }
570 Event::ReasoningMessageContent(e) => {
571 Ok(self.reasoning_content(&e.message_id, &e.delta, &e.subagent_run_id))
572 }
573 Event::ReasoningMessageEnd(e) => Ok(self.reasoning_end(&e.message_id)),
574 Event::ReasoningEnd(e) => Ok(self.reasoning_end(&e.message_id)),
575 Event::ReasoningMessageChunk(e) => self.reasoning_chunk(e),
576 Event::ReasoningEncryptedValue(e) => Ok(self.encrypted_value(e)),
577
578 Event::ThinkingStart(_) => {
581 let id = self.mint_thinking_id();
582 Ok(self.reasoning_start(id, None))
583 }
584 Event::ThinkingTextMessageStart(_) => {
585 let id = self.thinking_id();
586 Ok(self.reasoning_start(id, None))
587 }
588 Event::ThinkingTextMessageContent(e) => {
589 let id = self.thinking_id();
590 Ok(self.reasoning_content(&id, &e.delta, &None))
591 }
592 Event::ThinkingTextMessageEnd(_) | Event::ThinkingEnd(_) => {
593 Ok(match self.current_reasoning.clone() {
594 Some(id) => self.reasoning_end(&id),
595 None => Changed::Nothing,
596 })
597 }
598
599 Event::RunStarted(e) => {
600 self.thread_id = Some(e.thread_id.clone());
601 self.run_id = Some(e.run_id.clone());
602 self.interrupts.clear();
603 Ok(Changed::RunStarted {
604 thread_id: e.thread_id.clone(),
605 run_id: e.run_id.clone(),
606 })
607 }
608 Event::RunFinished(e) => {
609 let outcome = e.outcome.clone().unwrap_or(RunOutcome::Success);
610 outcome.validate()?;
611 self.interrupts = outcome.interrupts().to_vec();
612 Ok(Changed::RunFinished {
613 outcome,
614 result: e.result.clone(),
615 })
616 }
617 Event::RunError(e) => Ok(Changed::RunError {
618 message: e.message.clone(),
619 code: e.code.clone(),
620 }),
621
622 Event::StepStarted(_) | Event::StepFinished(_) | Event::Raw(_) | Event::Custom(_) => {
623 Ok(Changed::Nothing)
624 }
625
626 Event::SubagentStarted(e) => Ok(self.subagent_started(e)),
627 Event::SubagentFinished(e) => Ok(self.subagent_finished(e)),
628 Event::SubagentError(e) => Ok(self.subagent_error(e)),
629 }
630 }
631
632 fn replace_messages(&mut self, messages: Vec<Message>) {
635 self.by_id.clear();
636 self.tool_calls.clear();
637 for (index, message) in messages.iter().enumerate() {
638 self.by_id.insert(message.id().clone(), index);
639 if let Message::Assistant(assistant) = message {
640 for call in assistant.tool_calls.iter().flatten() {
641 self.tool_calls.insert(call.id.clone(), index);
642 }
643 }
644 }
645 self.messages = messages;
646 let by_id = &self.by_id;
648 self.open_text.retain(|(_, open)| by_id.contains_key(open));
649 }
650
651 fn merge_snapshot(&mut self, snapshot: Vec<Message>) {
675 let snapshot_owns_activity = snapshot.iter().any(|m| matches!(m, Message::Activity(_)));
676
677 let mut incoming: Vec<Option<Message>> = snapshot.into_iter().map(Some).collect();
680 let mut position: HashMap<MessageId, usize> = HashMap::with_capacity(incoming.len());
681 for (index, message) in incoming.iter().enumerate() {
682 let id = message.as_ref().expect("every slot starts full").id();
683 position.insert(id.clone(), index);
684 }
685
686 let previous = std::mem::take(&mut self.messages);
687 let mut merged = Vec::with_capacity(previous.len().max(incoming.len()));
688 for message in previous {
689 let keep_local = !snapshot_owns_activity && matches!(message, Message::Activity(_));
690 if let Some(index) = position.get(message.id()).copied() {
691 if keep_local {
692 incoming[index] = None;
695 } else if let Some(replacement) = incoming[index].take() {
696 merged.push(replacement);
697 continue;
698 }
699 }
700 if keep_local {
701 merged.push(message);
702 }
703 }
704 merged.extend(incoming.into_iter().flatten());
705
706 self.replace_messages(merged);
707 }
708
709 fn message_change(&self, id: &MessageId, kind: MessageChangeKind) -> Changed {
710 match self.by_id.get(id) {
711 Some(index) => Changed::Message(MessageChange {
712 index: *index,
713 id: id.clone(),
714 kind,
715 }),
716 None => Changed::Nothing,
717 }
718 }
719
720 fn text_start(
721 &mut self,
722 id: MessageId,
723 role: TextMessageRole,
724 name: Option<String>,
725 owner: Option<SubagentRunId>,
726 ) -> Changed {
727 self.open_text.retain(|(_, open)| open != &id);
728 self.open_text.push((owner.clone(), id.clone()));
729 if let Some(index) = self.by_id.get(&id) {
730 return Changed::Message(MessageChange {
734 index: *index,
735 id,
736 kind: MessageChangeKind::Started,
737 });
738 }
739 let mut message = empty_message(id.clone(), role, name);
743 message.set_subagent_run_id(owner);
744 let index = self.push_message(message);
745 Changed::Message(MessageChange {
746 index,
747 id,
748 kind: MessageChangeKind::Started,
749 })
750 }
751
752 fn text_content(
753 &mut self,
754 id: &MessageId,
755 delta: &str,
756 owner: &Option<SubagentRunId>,
757 ) -> Result<Changed> {
758 if !self.by_id.contains_key(id) {
759 self.text_start(id.clone(), TextMessageRole::Assistant, None, owner.clone());
761 }
762 let Some(index) = self.by_id.get(id).copied() else {
763 return Ok(Changed::Nothing);
764 };
765 let Some(message) = self.messages.get_mut(index) else {
766 return Ok(Changed::Nothing);
767 };
768 append_text(message, delta)?;
769 Ok(Changed::Message(MessageChange {
770 index,
771 id: id.clone(),
772 kind: MessageChangeKind::Content {
773 delta: delta.to_owned(),
774 },
775 }))
776 }
777
778 fn text_end(&mut self, id: &MessageId) -> Changed {
779 self.open_text.retain(|(_, open)| open != id);
780 self.message_change(id, MessageChangeKind::Ended)
781 }
782
783 fn text_chunk(&mut self, event: &TextMessageChunkEvent) -> Result<Changed> {
792 let id = match &event.message_id {
793 Some(id) => id.clone(),
794 None => resolve_open(
795 &self.open_text,
796 &event.subagent_run_id,
797 "TEXT_MESSAGE_CHUNK",
798 "message",
799 )?,
800 };
801 if !self.open_text.iter().any(|(_, open)| open == &id) {
802 self.text_start(
803 id.clone(),
804 event.role.unwrap_or_default(),
805 event.name.clone(),
806 event.subagent_run_id.clone(),
807 );
808 }
809 match &event.delta {
810 Some(delta) => self.text_content(&id, delta, &event.subagent_run_id),
811 None => Ok(self.message_change(&id, MessageChangeKind::Started)),
812 }
813 }
814
815 fn tool_call_start(
818 &mut self,
819 tool_call_id: ToolCallId,
820 name: String,
821 parent_message_id: Option<MessageId>,
822 owner: Option<SubagentRunId>,
823 ) -> Changed {
824 let parent =
828 parent_message_id.unwrap_or_else(|| MessageId::new(format!("{tool_call_id}-message")));
829 let index = match self.by_id.get(&parent) {
830 Some(index) => *index,
831 None => self.push_message(Message::Assistant(AssistantMessage {
832 id: parent.clone(),
833 subagent_run_id: owner,
834 ..Default::default()
835 })),
836 };
837 let Some(Message::Assistant(assistant)) = self.messages.get_mut(index) else {
840 return Changed::Nothing;
841 };
842 let calls = assistant.tool_calls.get_or_insert_with(Vec::new);
843 if !calls.iter().any(|call| call.id == tool_call_id) {
844 calls.push(ToolCall::new(tool_call_id.clone(), name.clone(), ""));
845 }
846 self.tool_calls.insert(tool_call_id.clone(), index);
847 Changed::Message(MessageChange {
848 index,
849 id: parent,
850 kind: MessageChangeKind::ToolCallStarted { tool_call_id, name },
851 })
852 }
853
854 fn tool_call_args(&mut self, tool_call_id: &ToolCallId, delta: &str) -> Result<Changed> {
855 let Some(index) = self.tool_calls.get(tool_call_id).copied() else {
856 return Err(Error::protocol(format!(
857 "TOOL_CALL_ARGS for unknown tool call {tool_call_id:?}"
858 )));
859 };
860 let Some(Message::Assistant(assistant)) = self.messages.get_mut(index) else {
861 return Ok(Changed::Nothing);
862 };
863 let Some(call) = assistant
864 .tool_calls
865 .as_mut()
866 .and_then(|calls| calls.iter_mut().find(|call| &call.id == tool_call_id))
867 else {
868 return Ok(Changed::Nothing);
869 };
870 call.function.arguments.push_str(delta);
871 Ok(Changed::Message(MessageChange {
872 index,
873 id: assistant.id.clone(),
874 kind: MessageChangeKind::ToolCallArgs {
875 tool_call_id: tool_call_id.clone(),
876 delta: delta.to_owned(),
877 },
878 }))
879 }
880
881 fn tool_call_end(&mut self, tool_call_id: &ToolCallId) -> Changed {
882 match self.tool_calls.get(tool_call_id).copied() {
883 Some(index) => match self.messages.get(index) {
884 Some(message) => Changed::Message(MessageChange {
885 index,
886 id: message.id().clone(),
887 kind: MessageChangeKind::ToolCallEnded {
888 tool_call_id: tool_call_id.clone(),
889 },
890 }),
891 None => Changed::Nothing,
892 },
893 None => Changed::Nothing,
894 }
895 }
896
897 fn tool_call_chunk(&mut self, event: &ToolCallChunkEvent) -> Result<Changed> {
898 let known = event
899 .tool_call_id
900 .as_ref()
901 .is_some_and(|id| self.tool_calls.contains_key(id));
902 match (&event.tool_call_id, known) {
903 (Some(id), false) => {
904 let Some(name) = event.tool_call_name.clone() else {
905 return Err(Error::protocol(format!(
906 "TOOL_CALL_CHUNK opens tool call {id:?} without a toolCallName"
907 )));
908 };
909 let started = self.tool_call_start(
910 id.clone(),
911 name,
912 event.parent_message_id.clone(),
913 event.subagent_run_id.clone(),
914 );
915 match &event.delta {
916 Some(delta) => self.tool_call_args(id, delta),
917 None => Ok(started),
918 }
919 }
920 (Some(id), true) => match &event.delta {
921 Some(delta) => self.tool_call_args(id, delta),
922 None => Ok(self.tool_call_end(id)),
923 },
924 (None, _) => Err(Error::protocol(
925 "TOOL_CALL_CHUNK carries no toolCallId and no call is open",
926 )),
927 }
928 }
929
930 fn tool_call_result(
931 &mut self,
932 message_id: MessageId,
933 tool_call_id: ToolCallId,
934 content: String,
935 owner: Option<SubagentRunId>,
936 ) -> Changed {
937 let index = match self.by_id.get(&message_id).copied() {
938 Some(index) => {
939 if let Some(Message::Tool(tool)) = self.messages.get_mut(index) {
940 tool.content = content;
941 tool.subagent_run_id = owner;
943 }
944 index
945 }
946 None => self.push_message(Message::Tool(ToolMessage {
949 id: message_id.clone(),
950 content,
951 tool_call_id: tool_call_id.clone(),
952 subagent_run_id: owner,
953 ..Default::default()
954 })),
955 };
956 Changed::Message(MessageChange {
957 index,
958 id: message_id,
959 kind: MessageChangeKind::ToolResult { tool_call_id },
960 })
961 }
962
963 fn activity_snapshot(&mut self, event: &ActivitySnapshotEvent) -> Changed {
966 let is_new = !self.by_id.contains_key(&event.message_id);
967 let index = self.activity_index(&event.message_id, &event.activity_type);
968 if let Some(Message::Activity(activity)) = self.messages.get_mut(index) {
969 activity.activity_type = event.activity_type.clone();
970 if is_new || event.replace {
975 activity.subagent_run_id = event.subagent_run_id.clone();
976 }
977 if event.replace {
978 activity.content = event.content.clone();
979 } else {
980 let mut merged = Value::Object(std::mem::take(&mut activity.content));
983 json_patch::merge(&mut merged, &Value::Object(event.content.clone()));
984 if let Value::Object(object) = merged {
985 activity.content = object;
986 }
987 }
988 }
989 Changed::Message(MessageChange {
990 index,
991 id: event.message_id.clone(),
992 kind: MessageChangeKind::Activity,
993 })
994 }
995
996 fn activity_delta(&mut self, event: &ActivityDeltaEvent) -> Result<Changed> {
997 let index = self.activity_index(&event.message_id, &event.activity_type);
998 if let Some(Message::Activity(activity)) = self.messages.get_mut(index) {
999 let what = format!("activity {}", event.message_id);
1000 let mut content = Value::Object(activity.content.clone());
1006 apply_patch(&mut content, &event.patch, &what)?;
1007 let Value::Object(object) = content else {
1008 return Err(Error::Patch {
1009 target: what,
1010 message: format!(
1011 "patch replaced the whole activity with {}, which is not an object",
1012 kind_of(&content)
1013 ),
1014 });
1015 };
1016 activity.content = object;
1017 }
1018 Ok(Changed::Message(MessageChange {
1019 index,
1020 id: event.message_id.clone(),
1021 kind: MessageChangeKind::Activity,
1022 }))
1023 }
1024
1025 fn activity_index(&mut self, id: &MessageId, activity_type: &str) -> usize {
1028 match self.by_id.get(id).copied() {
1029 Some(index) => index,
1030 None => self.push_message(Message::Activity(ActivityMessage {
1031 id: id.clone(),
1032 activity_type: activity_type.to_owned(),
1033 content: JsonObject::new(),
1034 ..Default::default()
1035 })),
1036 }
1037 }
1038
1039 fn reasoning_start(&mut self, id: MessageId, owner: Option<SubagentRunId>) -> Changed {
1047 self.current_reasoning = Some(id.clone());
1048 self.reasoning_streams.retain(|(_, open)| open != &id);
1049 self.reasoning_streams.push((owner.clone(), id.clone()));
1050 if !self.reasoning_by_id.contains_key(&id) {
1051 self.reasoning_by_id
1052 .insert(id.clone(), self.reasoning.len());
1053 self.reasoning.push(ReasoningMessage {
1054 id: id.clone(),
1055 subagent_run_id: owner,
1056 ..Default::default()
1057 });
1058 }
1059 if !self.open_reasoning.insert(id.clone()) {
1060 return Changed::Nothing;
1061 }
1062 Changed::Reasoning(ReasoningChange {
1063 id,
1064 kind: ReasoningChangeKind::Started,
1065 })
1066 }
1067
1068 fn mint_thinking_id(&mut self) -> MessageId {
1074 self.thinking_counter += 1;
1075 MessageId::new(format!("thinking-{}", self.thinking_counter))
1076 }
1077
1078 fn thinking_id(&mut self) -> MessageId {
1084 match self.current_reasoning.clone() {
1085 Some(id) => id,
1086 None => self.mint_thinking_id(),
1087 }
1088 }
1089
1090 fn reasoning_content(
1091 &mut self,
1092 id: &MessageId,
1093 delta: &str,
1094 owner: &Option<SubagentRunId>,
1095 ) -> Changed {
1096 if !self.reasoning_by_id.contains_key(id) {
1097 self.reasoning_start(id.clone(), owner.clone());
1098 }
1099 if let Some(message) = self
1100 .reasoning_by_id
1101 .get(id)
1102 .and_then(|index| self.reasoning.get_mut(*index))
1103 {
1104 message.content.push_str(delta);
1105 }
1106 Changed::Reasoning(ReasoningChange {
1107 id: id.clone(),
1108 kind: ReasoningChangeKind::Content {
1109 delta: delta.to_owned(),
1110 },
1111 })
1112 }
1113
1114 fn reasoning_chunk(&mut self, event: &ReasoningMessageChunkEvent) -> Result<Changed> {
1117 let id = match &event.message_id {
1118 Some(id) => id.clone(),
1119 None => resolve_open(
1120 &self.reasoning_streams,
1121 &event.subagent_run_id,
1122 "REASONING_MESSAGE_CHUNK",
1123 "reasoning message",
1124 )?,
1125 };
1126 let started = if self.open_reasoning.contains(&id) {
1127 Changed::Nothing
1128 } else {
1129 self.reasoning_start(id.clone(), event.subagent_run_id.clone())
1130 };
1131 Ok(match &event.delta {
1132 Some(delta) => self.reasoning_content(&id, delta, &event.subagent_run_id),
1133 None => started,
1136 })
1137 }
1138
1139 fn reasoning_end(&mut self, id: &MessageId) -> Changed {
1145 if self.current_reasoning.as_ref() == Some(id) {
1146 self.current_reasoning = None;
1147 }
1148 self.reasoning_streams.retain(|(_, open)| open != id);
1149 if !self.open_reasoning.remove(id) {
1150 return Changed::Nothing;
1151 }
1152 Changed::Reasoning(ReasoningChange {
1153 id: id.clone(),
1154 kind: ReasoningChangeKind::Ended,
1155 })
1156 }
1157
1158 fn encrypted_value(&mut self, event: &ReasoningEncryptedValueEvent) -> Changed {
1159 let blob = event.encrypted_value.clone();
1160 match event.subtype {
1161 ReasoningEncryptedValueSubtype::ToolCall => {
1162 let tool_call_id = ToolCallId::new(event.entity_id.clone());
1163 let Some(index) = self.tool_calls.get(&tool_call_id).copied() else {
1164 return Changed::Nothing;
1165 };
1166 let Some(Message::Assistant(assistant)) = self.messages.get_mut(index) else {
1167 return Changed::Nothing;
1168 };
1169 if let Some(call) = assistant
1170 .tool_calls
1171 .as_mut()
1172 .and_then(|calls| calls.iter_mut().find(|call| call.id == tool_call_id))
1173 {
1174 call.encrypted_value = Some(blob);
1175 }
1176 Changed::Message(MessageChange {
1177 index,
1178 id: assistant.id.clone(),
1179 kind: MessageChangeKind::EncryptedValue,
1180 })
1181 }
1182 ReasoningEncryptedValueSubtype::Message => {
1183 let id = MessageId::new(event.entity_id.clone());
1184 if let Some(message) = self
1185 .reasoning_by_id
1186 .get(&id)
1187 .and_then(|index| self.reasoning.get_mut(*index))
1188 {
1189 message.encrypted_value = Some(blob);
1190 return Changed::Reasoning(ReasoningChange {
1191 id,
1192 kind: ReasoningChangeKind::EncryptedValue,
1193 });
1194 }
1195 let Some(index) = self.by_id.get(&id).copied() else {
1196 return Changed::Nothing;
1197 };
1198 if let Some(message) = self.messages.get_mut(index) {
1199 set_encrypted_value(message, blob);
1200 }
1201 Changed::Message(MessageChange {
1202 index,
1203 id,
1204 kind: MessageChangeKind::EncryptedValue,
1205 })
1206 }
1207 }
1208 }
1209
1210 fn subagent_started(&mut self, event: &SubagentStartedEvent) -> Changed {
1213 let run_id = event.subagent_run_id.clone();
1214 if let Some(index) = self.subagent_by_id.get(&run_id).copied() {
1215 let subagent = &mut self.subagents[index];
1220 let kind = if matches!(subagent.status, SubagentStatus::Suspended { .. }) {
1221 SubagentChangeKind::Resumed
1222 } else {
1223 SubagentChangeKind::Started
1224 };
1225 subagent.name.clone_from(&event.name);
1226 if event.description.is_some() {
1227 subagent.description.clone_from(&event.description);
1228 }
1229 if event.parent_subagent_run_id.is_some() {
1230 subagent
1231 .parent_subagent_run_id
1232 .clone_from(&event.parent_subagent_run_id);
1233 }
1234 if event.parent_tool_call_id.is_some() {
1235 subagent
1236 .parent_tool_call_id
1237 .clone_from(&event.parent_tool_call_id);
1238 subagent
1239 .parent_message_id
1240 .clone_from(&event.parent_message_id);
1241 }
1242 subagent.status = SubagentStatus::Running;
1243 return Changed::Subagent(SubagentChange {
1244 index,
1245 run_id,
1246 kind,
1247 });
1248 }
1249 let index = self.push_subagent(Subagent {
1250 run_id: run_id.clone(),
1251 name: event.name.clone(),
1252 description: event.description.clone(),
1253 parent_subagent_run_id: event.parent_subagent_run_id.clone(),
1254 parent_tool_call_id: event.parent_tool_call_id.clone(),
1255 parent_message_id: event.parent_message_id.clone(),
1256 status: SubagentStatus::Running,
1257 });
1258 Changed::Subagent(SubagentChange {
1259 index,
1260 run_id,
1261 kind: SubagentChangeKind::Started,
1262 })
1263 }
1264
1265 fn subagent_finished(&mut self, event: &SubagentFinishedEvent) -> Changed {
1266 let run_id = event.subagent_run_id.clone();
1267 let index = self.subagent_index(&run_id);
1268 let (status, kind) = match &event.outcome {
1269 Some(SubagentOutcome::Suspended { interrupt_ids }) => (
1270 SubagentStatus::Suspended {
1271 result: event.result.clone(),
1272 interrupt_ids: interrupt_ids.clone().unwrap_or_default(),
1273 },
1274 SubagentChangeKind::Suspended,
1275 ),
1276 Some(SubagentOutcome::Success) | None => (
1278 SubagentStatus::Finished {
1279 result: event.result.clone(),
1280 },
1281 SubagentChangeKind::Finished,
1282 ),
1283 };
1284 self.subagents[index].status = status;
1285 Changed::Subagent(SubagentChange {
1286 index,
1287 run_id,
1288 kind,
1289 })
1290 }
1291
1292 fn subagent_error(&mut self, event: &SubagentErrorEvent) -> Changed {
1293 let run_id = event.subagent_run_id.clone();
1294 let index = self.subagent_index(&run_id);
1295 self.subagents[index].status = SubagentStatus::Failed {
1296 message: event.message.clone(),
1297 code: event.code.clone(),
1298 };
1299 Changed::Subagent(SubagentChange {
1300 index,
1301 run_id,
1302 kind: SubagentChangeKind::Failed,
1303 })
1304 }
1305
1306 fn subagent_index(&mut self, run_id: &SubagentRunId) -> usize {
1311 match self.subagent_by_id.get(run_id).copied() {
1312 Some(index) => index,
1313 None => self.push_subagent(Subagent {
1314 run_id: run_id.clone(),
1315 name: run_id.to_string(),
1316 description: None,
1317 parent_subagent_run_id: None,
1318 parent_tool_call_id: None,
1319 parent_message_id: None,
1320 status: SubagentStatus::Running,
1321 }),
1322 }
1323 }
1324
1325 fn push_subagent(&mut self, subagent: Subagent) -> usize {
1326 let index = self.subagents.len();
1327 self.subagent_by_id.insert(subagent.run_id.clone(), index);
1328 self.subagents.push(subagent);
1329 index
1330 }
1331
1332 fn merge_event_metadata(&mut self, event: &Event, changed: &Changed, metadata: &JsonObject) {
1339 match event {
1340 Event::TextMessageStart(_)
1341 | Event::TextMessageContent(_)
1342 | Event::TextMessageEnd(_)
1343 | Event::TextMessageChunk(_)
1344 | Event::ToolCallResult(_)
1345 | Event::ActivitySnapshot(_)
1346 | Event::ActivityDelta(_) => {
1347 let Changed::Message(change) = changed else {
1348 return;
1349 };
1350 if let Some(message) = self.messages.get_mut(change.index) {
1351 merge_metadata_into(message.metadata_mut(), Some(metadata));
1352 }
1353 }
1354 Event::ToolCallStart(_)
1355 | Event::ToolCallArgs(_)
1356 | Event::ToolCallEnd(_)
1357 | Event::ToolCallChunk(_) => {
1358 let Changed::Message(change) = changed else {
1359 return;
1360 };
1361 let tool_call_id = match &change.kind {
1362 MessageChangeKind::ToolCallStarted { tool_call_id, .. }
1363 | MessageChangeKind::ToolCallArgs { tool_call_id, .. }
1364 | MessageChangeKind::ToolCallEnded { tool_call_id } => tool_call_id.clone(),
1365 _ => return,
1366 };
1367 let Some(Message::Assistant(assistant)) = self.messages.get_mut(change.index)
1368 else {
1369 return;
1370 };
1371 if let Some(call) = assistant
1372 .tool_calls
1373 .as_mut()
1374 .and_then(|calls| calls.iter_mut().find(|call| call.id == tool_call_id))
1375 {
1376 merge_metadata_into(&mut call.metadata, Some(metadata));
1377 }
1378 }
1379 Event::ReasoningMessageStart(_)
1380 | Event::ReasoningMessageContent(_)
1381 | Event::ReasoningMessageEnd(_)
1382 | Event::ReasoningMessageChunk(_) => {
1383 let id = match (event, changed) {
1384 (Event::ReasoningMessageStart(e), _) => e.message_id.clone(),
1385 (Event::ReasoningMessageContent(e), _) => e.message_id.clone(),
1386 (Event::ReasoningMessageEnd(e), _) => e.message_id.clone(),
1387 (_, Changed::Reasoning(change)) => change.id.clone(),
1388 _ => return,
1389 };
1390 if let Some(message) = self
1391 .reasoning_by_id
1392 .get(&id)
1393 .and_then(|index| self.reasoning.get_mut(*index))
1394 {
1395 merge_metadata_into(&mut message.metadata, Some(metadata));
1396 }
1397 }
1398 _ => {}
1399 }
1400 }
1401}
1402
1403fn resolve_open(
1407 streams: &[(Option<SubagentRunId>, MessageId)],
1408 tag: &Option<SubagentRunId>,
1409 kind: &str,
1410 what: &str,
1411) -> Result<MessageId> {
1412 if tag.is_some() {
1413 return streams
1414 .iter()
1415 .rev()
1416 .find(|(owner, _)| owner == tag)
1417 .map(|(_, id)| id.clone())
1418 .ok_or_else(|| {
1419 Error::protocol(format!(
1420 "{kind} carries no messageId and subagent {:?} has no {what} open",
1421 tag.as_deref().unwrap_or_default()
1422 ))
1423 });
1424 }
1425 if let Some((_, id)) = streams.iter().rev().find(|(owner, _)| owner.is_none()) {
1426 return Ok(id.clone());
1427 }
1428 match streams {
1429 [] => Err(Error::protocol(format!(
1430 "{kind} carries no messageId and no {what} is open"
1431 ))),
1432 [(_, id)] => Ok(id.clone()),
1433 _ => Err(Error::protocol(format!(
1434 "{kind} carries no messageId and several subagents have a {what} open; \
1435 attribute the chunk"
1436 ))),
1437 }
1438}
1439
1440fn apply_patch(target: &mut Value, operations: &[PatchOperation], what: &str) -> Result<()> {
1442 let document = serde_json::to_value(operations)?;
1446 let patch: json_patch::Patch =
1447 serde_json::from_value(document).map_err(|error| Error::Patch {
1448 target: what.to_owned(),
1449 message: format!("invalid patch document: {error}"),
1450 })?;
1451 json_patch::patch(target, &patch).map_err(|error| Error::Patch {
1452 target: what.to_owned(),
1453 message: error.to_string(),
1454 })
1455}
1456
1457fn kind_of(value: &Value) -> &'static str {
1459 match value {
1460 Value::Null => "null",
1461 Value::Bool(_) => "a boolean",
1462 Value::Number(_) => "a number",
1463 Value::String(_) => "a string",
1464 Value::Array(_) => "an array",
1465 Value::Object(_) => "an object",
1466 }
1467}
1468
1469fn empty_message(id: MessageId, role: TextMessageRole, name: Option<String>) -> Message {
1471 match role {
1472 TextMessageRole::Assistant => Message::Assistant(AssistantMessage {
1473 id,
1474 content: Some(String::new()),
1475 name,
1476 ..Default::default()
1477 }),
1478 TextMessageRole::User => Message::User(UserMessage {
1479 id,
1480 content: UserContent::Text(String::new()),
1481 name,
1482 ..Default::default()
1483 }),
1484 TextMessageRole::System => Message::System(SystemMessage {
1485 id,
1486 content: String::new(),
1487 name,
1488 ..Default::default()
1489 }),
1490 TextMessageRole::Developer => Message::Developer(DeveloperMessage {
1491 id,
1492 content: String::new(),
1493 name,
1494 ..Default::default()
1495 }),
1496 }
1497}
1498
1499fn append_text(message: &mut Message, delta: &str) -> Result<()> {
1501 match message {
1502 Message::Assistant(m) => m.content.get_or_insert_with(String::new).push_str(delta),
1503 Message::System(m) => m.content.push_str(delta),
1504 Message::Developer(m) => m.content.push_str(delta),
1505 Message::Reasoning(m) => m.content.push_str(delta),
1506 Message::Tool(m) => m.content.push_str(delta),
1507 Message::User(m) => match &mut m.content {
1508 UserContent::Text(text) => text.push_str(delta),
1509 UserContent::Parts(parts) => match parts.last_mut() {
1510 Some(InputContent::Text(text)) => text.text.push_str(delta),
1511 _ => parts.push(InputContent::Text(TextInputContent {
1512 text: delta.to_owned(),
1513 })),
1514 },
1515 },
1516 Message::Activity(m) => {
1517 return Err(Error::protocol(format!(
1518 "text streamed into activity message {:?}, which has no text",
1519 m.id
1520 )));
1521 }
1522 }
1523 Ok(())
1524}
1525
1526fn set_encrypted_value(message: &mut Message, blob: String) {
1528 match message {
1529 Message::Assistant(m) => m.encrypted_value = Some(blob),
1530 Message::System(m) => m.encrypted_value = Some(blob),
1531 Message::Developer(m) => m.encrypted_value = Some(blob),
1532 Message::User(m) => m.encrypted_value = Some(blob),
1533 Message::Tool(m) => m.encrypted_value = Some(blob),
1534 Message::Reasoning(m) => m.encrypted_value = Some(blob),
1535 Message::Activity(_) => {}
1537 }
1538}