1use std::panic::{catch_unwind, AssertUnwindSafe};
13use std::sync::Arc;
14
15use futures::future::BoxFuture;
16use rpi_agent::events::{AgentEmitter, AgentEvent};
17use rpi_plugin_sdk::{EventTag, StablePluginEvent, StbString};
18
19use crate::host_free_string;
20use crate::registry::RegistrySnapshot;
21
22pub fn event_tag_for(event: &AgentEvent) -> Option<EventTag> {
28 match event {
29 AgentEvent::AgentStart => Some(EventTag::AgentStart),
30 AgentEvent::AgentEnd { .. } => Some(EventTag::AgentEnd),
31 AgentEvent::RetryScheduled { .. } => None,
32 AgentEvent::TurnStart => Some(EventTag::TurnStart),
33 AgentEvent::TurnEnd { .. } => Some(EventTag::TurnEnd),
34 AgentEvent::MessageStart { .. } => Some(EventTag::MessageStart),
35 AgentEvent::MessageUpdate { .. } => Some(EventTag::MessageUpdate),
36 AgentEvent::MessageEnd { .. } => Some(EventTag::MessageEnd),
37 AgentEvent::ToolExecutionStart { .. } => Some(EventTag::ToolExecutionStart),
38 AgentEvent::ToolExecutionUpdate { .. } => Some(EventTag::ToolExecutionUpdate),
39 AgentEvent::ToolExecutionEnd { .. } => Some(EventTag::ToolExecutionEnd),
40 }
41}
42
43fn message_to_stb(message: &rpi_agent::message::AgentMessage) -> StbString {
47 let text = serde_json::to_string(message).unwrap_or_else(|_| "null".to_string());
49 StbString::from_string(text)
50}
51
52pub fn translate(event: &AgentEvent) -> Option<StablePluginEvent> {
57 let tag = event_tag_for(event)?;
58 match event {
59 AgentEvent::MessageStart { message }
60 | AgentEvent::MessageUpdate { message, .. }
61 | AgentEvent::MessageEnd { message } => {
62 let stb = message_to_stb(message);
63 Some(StablePluginEvent::message(tag, stb))
64 }
65
66 AgentEvent::ToolExecutionStart {
67 tool_call_id,
68 tool_name,
69 args,
70 }
71 | AgentEvent::ToolExecutionUpdate {
72 tool_call_id,
73 tool_name,
74 args,
75 ..
76 } => Some(StablePluginEvent::tool_call(
77 tag,
78 StbString::from_string(tool_call_id.clone()),
79 StbString::from_string(tool_name.clone()),
80 StbString::from_string(serde_json::to_string(args).unwrap_or_else(|_| "null".into())),
81 )),
82
83 AgentEvent::ToolExecutionEnd {
84 tool_call_id,
85 tool_name,
86 result,
87 is_error,
88 } => {
89 let result_json = agent_tool_result_to_json(result);
91 Some(StablePluginEvent::tool_result(
92 tag,
93 StbString::from_string(tool_call_id.clone()),
94 StbString::from_string(tool_name.clone()),
95 StbString::from_string(result_json),
96 *is_error,
97 ))
98 }
99
100 AgentEvent::AgentStart
102 | AgentEvent::TurnStart
103 | AgentEvent::AgentEnd { .. }
104 | AgentEvent::TurnEnd { .. } => Some(StablePluginEvent::empty(tag)),
105
106 AgentEvent::RetryScheduled { .. } => None,
107 }
108}
109
110fn agent_tool_result_to_json(result: &rpi_agent::types::AgentToolResult) -> String {
115 let mut txt = String::new();
116 txt.push('{');
117 txt.push_str("\"content\":[");
118 for (i, c) in result.content.iter().enumerate() {
119 if i > 0 {
120 txt.push(',');
121 }
122 match c {
123 rpi_agent::types::TextContentOrImage::Text(t) => {
124 txt.push_str(
125 &serde_json::to_string(&serde_json::json!({ "type": "text", "text": t.text }))
126 .unwrap_or_else(|_| "\"\"".into()),
127 );
128 }
129 rpi_agent::types::TextContentOrImage::Image(img) => {
130 txt.push_str(
131 &serde_json::to_string(&serde_json::json!({
132 "type": "image",
133 "data": img.data,
134 "mimeType": img.mime_type,
135 }))
136 .unwrap_or_else(|_| "\"\"".into()),
137 );
138 }
139 }
140 }
141 txt.push(']');
142 txt.push_str(",\"details\":");
143 txt.push_str(&serde_json::to_string(&result.details).unwrap_or_else(|_| "null".into()));
144 txt.push_str(",\"terminate\":");
145 txt.push_str(if result.terminate { "true" } else { "false" });
146 txt.push_str(",\"addedToolNames\":");
147 txt.push_str(&serde_json::to_string(&result.added_tool_names).unwrap_or_else(|_| "[]".into()));
148 txt.push('}');
149 txt
150}
151
152#[allow(dead_code)]
169fn free_event_strings(_event: &StablePluginEvent) {
170 }
172
173pub struct ExtensionEmitter {
195 snapshot: Arc<RegistrySnapshot>,
196 #[allow(dead_code)]
201 keepalive: Arc<crate::PluginKeepalive>,
202}
203
204impl ExtensionEmitter {
205 pub fn new(snapshot: Arc<RegistrySnapshot>, keepalive: Arc<crate::PluginKeepalive>) -> Self {
208 Self {
209 snapshot,
210 keepalive,
211 }
212 }
213
214 fn dispatch(&self, event: &StablePluginEvent) {
219 dispatch_to_handlers(&self.snapshot, event);
220 }
221}
222
223pub fn dispatch_to_handlers(snapshot: &RegistrySnapshot, event: &StablePluginEvent) {
228 if !crate::registry::assert_active(snapshot.active_flag()) {
230 return;
231 }
232 let handlers = snapshot.handlers_for(event.tag);
233 if handlers.is_empty() {
234 free_dispatched_event(event);
237 return;
238 }
239 for h in handlers {
240 let outcome = catch_unwind(AssertUnwindSafe(|| (h.handler)(*event, h.user_data)));
243 match outcome {
244 Ok(rc) if rc != 0 => {
245 tracing::warn!(tag = ?event.tag, rc, "extension event handler returned nonzero");
246 }
247 Ok(_) => {}
248 Err(_) => {
249 tracing::error!(tag = ?event.tag, "extension event handler panicked — skipped");
250 }
251 }
252 }
253 free_dispatched_event(event);
255}
256
257pub fn dispatch_data_event(snapshot: &RegistrySnapshot, tag: EventTag, data: &str) -> bool {
262 let handlers = snapshot.handlers_for(tag);
263 if handlers.is_empty() {
264 return false;
265 }
266 let event = StablePluginEvent::data(tag, StbString::from_string(data.to_string()));
267 dispatch_to_handlers(snapshot, &event);
268 true
269}
270
271impl AgentEmitter for ExtensionEmitter {
272 fn emit(&self, event: AgentEvent) -> BoxFuture<'static, ()> {
273 if let Some(stable) = translate(&event) {
277 self.dispatch(&stable);
278 }
279 Box::pin(async {})
280 }
281
282 fn try_emit(&self, event: AgentEvent) {
283 if let Some(stable) = translate(&event) {
284 self.dispatch(&stable);
285 }
286 }
287}
288
289pub struct TeeEmitter {
308 emitters: Vec<Arc<dyn AgentEmitter>>,
309}
310
311impl TeeEmitter {
312 pub fn new(emitters: Vec<Arc<dyn AgentEmitter>>) -> Self {
317 Self { emitters }
318 }
319}
320
321impl AgentEmitter for TeeEmitter {
322 fn emit(&self, event: AgentEvent) -> BoxFuture<'static, ()> {
323 let emitters = self.emitters.clone();
329 Box::pin(async move {
330 for e in &emitters {
331 e.emit(event.clone()).await;
332 }
333 })
334 }
335
336 fn try_emit(&self, event: AgentEvent) {
337 for e in &self.emitters {
338 e.try_emit(event.clone());
339 }
340 }
341}
342
343fn free_dispatched_event(event: &StablePluginEvent) {
347 use rpi_plugin_sdk::EventTag as T;
348 match event.tag {
349 T::MessageStart | T::MessageUpdate | T::MessageEnd => {
350 unsafe { host_free_string(event.payload.message.message) };
352 }
353 T::ToolCall | T::ToolExecutionStart | T::ToolExecutionUpdate => {
354 unsafe {
356 let tc = &event.payload.tool_call;
357 host_free_string(tc.tool_call_id);
358 host_free_string(tc.tool_name);
359 host_free_string(tc.params);
360 }
361 }
362 T::ToolResult | T::ToolExecutionEnd => {
363 unsafe {
365 let tr = &event.payload.tool_result;
366 host_free_string(tr.tool_call_id);
367 host_free_string(tr.tool_name);
368 host_free_string(tr.result);
369 }
370 }
371 T::ProjectTrust
372 | T::ResourcesDiscover
373 | T::SessionStart
374 | T::SessionInfoChanged
375 | T::SessionBeforeSwitch
376 | T::SessionBeforeFork
377 | T::SessionBeforeCompact
378 | T::SessionCompact
379 | T::SessionShutdown
380 | T::SessionBeforeTree
381 | T::SessionTree
382 | T::Context
383 | T::BeforeAgentStart
384 | T::AgentStart
385 | T::AgentEnd
386 | T::AgentSettled
387 | T::TurnStart
388 | T::TurnEnd
389 | T::ModelSelect
390 | T::ThinkingLevelSelect
391 | T::UserBash
392 | T::Input => {
393 }
395 T::BeforeProviderRequest | T::BeforeProviderHeaders | T::AfterProviderResponse => {
398 unsafe { host_free_string(event.payload.data.data) };
400 }
401 }
402}
403
404#[cfg(test)]
412mod tests {
413 use super::*;
414 use rpi_agent::message::AgentMessage;
415 use rpi_agent::types::AgentToolResult;
416 use rpi_ai::types::{AssistantMessage, Usage};
417 use rpi_plugin_sdk::EventTag;
418 use std::sync::atomic::{AtomicUsize, Ordering};
419 use std::sync::Mutex;
420
421 static HANDLER_TEST_LOCK: Mutex<()> = Mutex::new(());
429
430 #[test]
431 fn all_ten_agent_events_map_to_a_tag() {
432 let am = AgentMessage::Assistant(Box::new(AssistantMessage {
433 role: rpi_ai::types::AssistantRole,
434 content: vec![rpi_ai::types::Content::text("hi")],
435 api: rpi_ai::types::Api::AnthropicMessages,
436 provider: "anthropic".to_string(),
437 model: "m".into(),
438 response_model: None,
439 response_id: None,
440 usage: Usage::zero(),
441 stop_reason: rpi_ai::types::StopReason::Stop,
442 deferred: None,
443 error_message: None,
444 raw_stop_reason: None,
445 end_turn: None,
446 timestamp: 0,
447 }));
448 let events = vec![
449 AgentEvent::AgentStart,
450 AgentEvent::AgentEnd { messages: vec![] },
451 AgentEvent::TurnStart,
452 AgentEvent::TurnEnd {
453 message: am.clone(),
454 tool_results: vec![],
455 },
456 AgentEvent::MessageStart {
457 message: am.clone(),
458 },
459 AgentEvent::MessageUpdate {
460 message: am.clone(),
461 assistant_message_event: rpi_ai::types::AssistantMessageEvent::Start {
462 partial: std::sync::Arc::new((*am.as_assistant().unwrap()).clone()),
463 },
464 },
465 AgentEvent::MessageEnd { message: am },
466 AgentEvent::ToolExecutionStart {
467 tool_call_id: "c1".into(),
468 tool_name: "echo".into(),
469 args: serde_json::json!({}),
470 },
471 AgentEvent::ToolExecutionUpdate {
472 tool_call_id: "c1".into(),
473 tool_name: "echo".into(),
474 args: serde_json::json!({}),
475 partial_result: std::sync::Arc::new(AgentToolResult::text("...")),
476 },
477 AgentEvent::ToolExecutionEnd {
478 tool_call_id: "c1".into(),
479 tool_name: "echo".into(),
480 result: AgentToolResult::text("done"),
481 is_error: false,
482 },
483 ];
484 for e in &events {
485 assert!(
486 event_tag_for(e).is_some(),
487 "event {:?} should map",
488 e.type_tag()
489 );
490 }
491 assert_eq!(event_tag_for(&events[0]), Some(EventTag::AgentStart));
493 assert_eq!(event_tag_for(&events[2]), Some(EventTag::TurnStart));
494 assert_eq!(
495 event_tag_for(&events[7]),
496 Some(EventTag::ToolExecutionStart)
497 );
498 assert!(event_tag_for(&AgentEvent::RetryScheduled {
499 attempt: 1,
500 max_retries: 10,
501 delay_ms: 2_000,
502 error: "503".into(),
503 })
504 .is_none());
505 }
506
507 #[test]
508 fn translate_message_end_produces_message_payload() {
509 let am = AgentMessage::Assistant(Box::new(AssistantMessage {
510 role: rpi_ai::types::AssistantRole,
511 content: vec![rpi_ai::types::Content::text("hi")],
512 api: rpi_ai::types::Api::AnthropicMessages,
513 provider: "anthropic".to_string(),
514 model: "m".into(),
515 response_model: None,
516 response_id: None,
517 usage: Usage::zero(),
518 stop_reason: rpi_ai::types::StopReason::Stop,
519 deferred: None,
520 error_message: None,
521 raw_stop_reason: None,
522 end_turn: None,
523 timestamp: 0,
524 }));
525 let ev = AgentEvent::MessageEnd { message: am };
526 let stable = translate(&ev).expect("maps");
527 assert_eq!(stable.tag, EventTag::MessageEnd);
528 unsafe { host_free_string(stable.payload.message.message) };
531 }
532
533 static HANDLER_HITS: AtomicUsize = AtomicUsize::new(0);
536
537 extern "C" fn counting_handler(_ev: StablePluginEvent, _ud: *mut std::ffi::c_void) -> i32 {
538 HANDLER_HITS.fetch_add(1, Ordering::SeqCst);
539 0
540 }
541
542 #[test]
543 fn emitter_dispatches_to_registered_handlers() {
544 let _guard = HANDLER_TEST_LOCK.lock().unwrap();
545 HANDLER_HITS.store(0, Ordering::SeqCst);
546 let mut reg = crate::registry::ExtensionRegistry::new();
547 reg.register_event_handler(EventTag::MessageEnd, counting_handler, std::ptr::null_mut());
548 let snap = Arc::new(reg.snapshot());
549 let emitter = ExtensionEmitter::new(snap, crate::loader::PluginKeepalive::empty());
550
551 let am = AgentMessage::Assistant(Box::new(AssistantMessage {
552 role: rpi_ai::types::AssistantRole,
553 content: vec![rpi_ai::types::Content::text("hi")],
554 api: rpi_ai::types::Api::AnthropicMessages,
555 provider: "anthropic".to_string(),
556 model: "m".into(),
557 response_model: None,
558 response_id: None,
559 usage: Usage::zero(),
560 stop_reason: rpi_ai::types::StopReason::Stop,
561 deferred: None,
562 error_message: None,
563 raw_stop_reason: None,
564 end_turn: None,
565 timestamp: 0,
566 }));
567 emitter.try_emit(AgentEvent::MessageEnd { message: am });
569 assert_eq!(HANDLER_HITS.load(Ordering::SeqCst), 1);
570
571 reg.invalidate();
573 let am2 = AgentMessage::Assistant(Box::new(AssistantMessage {
574 role: rpi_ai::types::AssistantRole,
575 content: vec![rpi_ai::types::Content::text("hi")],
576 api: rpi_ai::types::Api::AnthropicMessages,
577 provider: "anthropic".to_string(),
578 model: "m".into(),
579 response_model: None,
580 response_id: None,
581 usage: Usage::zero(),
582 stop_reason: rpi_ai::types::StopReason::Stop,
583 deferred: None,
584 error_message: None,
585 raw_stop_reason: None,
586 end_turn: None,
587 timestamp: 0,
588 }));
589 emitter.try_emit(AgentEvent::MessageEnd { message: am2 });
590 assert_eq!(
591 HANDLER_HITS.load(Ordering::SeqCst),
592 1,
593 "stale registry must not dispatch"
594 );
595 }
596
597 #[tokio::test]
600 async fn tee_emitter_fans_out_to_every_child() {
601 let _guard = HANDLER_TEST_LOCK.lock().unwrap();
602 use rpi_agent::events::{AgentEmitter, CollectorEmitter};
603
604 let (collector_a, events_a) = CollectorEmitter::new();
606 let (collector_b, events_b) = CollectorEmitter::new();
607 HANDLER_HITS.store(0, Ordering::SeqCst);
608 let mut reg = crate::registry::ExtensionRegistry::new();
609 reg.register_event_handler(EventTag::MessageEnd, counting_handler, std::ptr::null_mut());
610 let snap = Arc::new(reg.snapshot());
611 let ext = ExtensionEmitter::new(snap, crate::loader::PluginKeepalive::empty());
612
613 let tee = TeeEmitter::new(vec![
614 Arc::new(collector_a),
615 Arc::new(collector_b),
616 Arc::new(ext),
617 ]);
618
619 let am = AgentMessage::Assistant(Box::new(AssistantMessage {
620 role: rpi_ai::types::AssistantRole,
621 content: vec![rpi_ai::types::Content::text("hi")],
622 api: rpi_ai::types::Api::AnthropicMessages,
623 provider: "anthropic".to_string(),
624 model: "m".into(),
625 response_model: None,
626 response_id: None,
627 usage: Usage::zero(),
628 stop_reason: rpi_ai::types::StopReason::Stop,
629 deferred: None,
630 error_message: None,
631 raw_stop_reason: None,
632 end_turn: None,
633 timestamp: 0,
634 }));
635 tee.emit(AgentEvent::MessageEnd { message: am }).await;
638 assert_eq!(
639 events_a.lock().unwrap().len(),
640 1,
641 "collector A got the event"
642 );
643 assert_eq!(
644 events_b.lock().unwrap().len(),
645 1,
646 "collector B got the event"
647 );
648 assert_eq!(
649 HANDLER_HITS.load(Ordering::SeqCst),
650 1,
651 "plugin handler fired once"
652 );
653 }
654}