@@ -16,7 +16,7 @@ use agentkit_acp::{
1616 } ,
1717} ;
1818use agentkit_core:: {
19- CancellationController , FinishReason , Item , ItemKind , Part , SessionId , ToolOutput ,
19+ CancellationController , FinishReason , Item , ItemKind , Part , SessionId , ToolOutput , Usage ,
2020} ;
2121use agentkit_loop:: {
2222 AgentEvent , LoopDriver , LoopError , LoopInterrupt , LoopObserver , LoopStep , ModelSession ,
@@ -257,6 +257,21 @@ struct ResponseReplacementObserver<S> {
257257 session_id : wire:: SessionId ,
258258}
259259
260+ fn usage_update ( usage : & Usage ) -> Option < wire:: UsageUpdate > {
261+ let tokens = usage. tokens . as_ref ( ) ?;
262+ let used = tokens. input_tokens . checked_add ( tokens. output_tokens ) ?;
263+ let size = [
264+ "context_window" ,
265+ "context_window_tokens" ,
266+ "model.context_window" ,
267+ "model.context_length" ,
268+ "openrouter.context_length" ,
269+ ]
270+ . iter ( )
271+ . find_map ( |key| usage. metadata . get ( * key) . and_then ( |value| value. as_u64 ( ) ) ) ?;
272+ Some ( wire:: UsageUpdate :: new ( used, size) )
273+ }
274+
260275impl < S > ResponseReplacementObserver < S > {
261276 fn new (
262277 inner : AcpIntegration ,
@@ -286,6 +301,19 @@ where
286301 S : AcpSessionUpdateSink + Clone ,
287302{
288303 fn handle_event ( & self , event : ObservedEvent ) {
304+ if let AgentEvent :: UsageUpdated ( usage) = & event. event {
305+ let Some ( update) = usage_update ( usage) else {
306+ return ;
307+ } ;
308+ let notification = wire:: UpdateSessionNotification :: new (
309+ self . session_id . clone ( ) ,
310+ wire:: SessionUpdate :: UsageUpdate ( update) ,
311+ ) ;
312+ if let Err ( error) = self . sink . update ( notification) {
313+ tracing:: debug!( %error, "failed to queue ACP v2 usage update" ) ;
314+ }
315+ return ;
316+ }
289317 if matches ! (
290318 & event. event,
291319 AgentEvent :: ContentDelta ( delta) if crate :: response_attempt:: is_marker( delta)
@@ -1762,6 +1790,42 @@ mod tests {
17621790 }
17631791 }
17641792
1793+ #[ test]
1794+ fn observer_reports_usage_with_a_known_context_window ( ) {
1795+ let recording = RecordingSink :: default ( ) ;
1796+ let sink = ResponseReplacementSink :: new ( recording. clone ( ) ) ;
1797+ let observer = ResponseReplacementObserver :: new (
1798+ AcpIntegration :: default ( ) ,
1799+ sink,
1800+ wire:: SessionId :: new ( "usage-session" ) ,
1801+ ) ;
1802+ let loop_session_id = SessionId :: new ( "usage-loop" ) ;
1803+ let emit = |usage| {
1804+ observer. handle_event ( ObservedEvent {
1805+ session_id : Arc :: new ( loop_session_id. clone ( ) ) ,
1806+ event : AgentEvent :: UsageUpdated ( usage) ,
1807+ } ) ;
1808+ } ;
1809+
1810+ emit ( agentkit_core:: Usage :: new ( agentkit_core:: TokenUsage :: new (
1811+ 10 , 2 ,
1812+ ) ) ) ;
1813+ emit (
1814+ agentkit_core:: Usage :: new ( agentkit_core:: TokenUsage :: new ( 50_000 , 3_000 ) ) . with_metadata (
1815+ MetadataMap :: from ( [ ( "context_window" . into ( ) , json ! ( 272_000 ) ) ] ) ,
1816+ ) ,
1817+ ) ;
1818+
1819+ let updates = recording. updates . lock ( ) . unwrap ( ) ;
1820+ assert_eq ! ( updates. len( ) , 1 ) ;
1821+ let wire:: SessionUpdate :: UsageUpdate ( usage) = & updates[ 0 ] . update else {
1822+ panic ! ( "expected usage update, got {:?}" , updates[ 0 ] . update) ;
1823+ } ;
1824+ assert_eq ! ( usage. used, 53_000 ) ;
1825+ assert_eq ! ( usage. size, 272_000 ) ;
1826+ assert ! ( usage. cost. is_none( ) ) ;
1827+ }
1828+
17651829 #[ test]
17661830 fn response_replacement_clears_and_remaps_message_ids_in_new_chunk_order ( ) {
17671831 let integration = AcpIntegration :: default ( ) ;
0 commit comments