@@ -34,6 +34,10 @@ pub struct Hub {
3434 sessions : Mutex < HashMap < String , Hosted > > ,
3535 /// Sessions whose engine is running a turn or has queued messages.
3636 busy : Mutex < HashSet < String > > ,
37+ /// Delivered messages a session thread has not processed yet. They hold
38+ /// a running slot: the engine still reports "ready" until it has read
39+ /// the message, and that must not free the slot early.
40+ claims : Mutex < HashMap < String , usize > > ,
3741 /// When each session last received a message (ms), for routing.
3842 last_active : Mutex < HashMap < String , u64 > > ,
3943 /// Serializes message delivery (scheduler ticks, sends, tool calls).
@@ -70,6 +74,7 @@ impl Hub {
7074 version,
7175 sessions : Mutex :: new ( HashMap :: new ( ) ) ,
7276 busy : Mutex :: new ( HashSet :: new ( ) ) ,
77+ claims : Mutex :: new ( HashMap :: new ( ) ) ,
7378 last_active : Mutex :: new ( HashMap :: new ( ) ) ,
7479 delivering : Mutex :: new ( ( ) ) ,
7580 next_id : AtomicU64 :: new ( 0 ) ,
@@ -277,6 +282,7 @@ impl Hub {
277282 let _ = self . store . record_session_closed ( id) ;
278283 }
279284 lock ( & self . busy ) . remove ( id) ;
285+ lock ( & self . claims ) . remove ( id) ;
280286 self . broadcast ( & json ! ( { "type" : "session_closed" , "session" : id } ) ) ;
281287 }
282288
@@ -315,8 +321,13 @@ impl Hub {
315321 }
316322 }
317323
318- /// Called by a session thread when its engine starts or finishes work.
324+ /// Called by a session thread with its engine's state every tick. A
325+ /// session with an unprocessed delivery stays busy.
319326 pub fn set_busy ( & self , session : & str , busy : bool ) {
327+ let busy = busy
328+ || lock ( & self . claims )
329+ . get ( session)
330+ . is_some_and ( |count| * count > 0 ) ;
320331 let changed = if busy {
321332 lock ( & self . busy ) . insert ( session. to_string ( ) )
322333 } else {
@@ -327,6 +338,17 @@ impl Hub {
327338 }
328339 }
329340
341+ /// The session thread has handed a delivered message to its engine.
342+ pub fn release_claim ( & self , session : & str ) {
343+ let mut claims = lock ( & self . claims ) ;
344+ if let Some ( count) = claims. get_mut ( session) {
345+ * count = count. saturating_sub ( 1 ) ;
346+ if * count == 0 {
347+ claims. remove ( session) ;
348+ }
349+ }
350+ }
351+
330352 pub fn agents_json ( & self ) -> Value {
331353 let records = self . store . sessions ( ) ;
332354 let busy = lock ( & self . busy ) . clone ( ) ;
@@ -610,13 +632,17 @@ impl Hub {
610632 }
611633 None => self . create_session ( None , Some ( & message. to ) ) ?,
612634 } ;
613- self . forward (
614- & session,
615- json ! ( { "op" : "user_message" , "content" : delivery_text( message) } ) ,
616- ) ?;
617- // A delivered message starts (or queues) a run; count it before the
618- // engine reports its status so the next message sees the slot taken.
635+ // A delivered message starts (or queues) a run: claim the slot before
636+ // forwarding, so the next message sees it taken.
637+ * lock ( & self . claims ) . entry ( session. clone ( ) ) . or_default ( ) += 1 ;
619638 lock ( & self . busy ) . insert ( session. clone ( ) ) ;
639+ if let Err ( error) = self . forward (
640+ & session,
641+ json ! ( { "op" : "user_message" , "content" : delivery_text( message) , "claimed" : true } ) ,
642+ ) {
643+ self . release_claim ( & session) ;
644+ return Err ( error) ;
645+ }
620646 lock ( & self . last_active ) . insert ( session. clone ( ) , now ( ) ) ;
621647 self . store
622648 . record_delivered ( & message. id , & session)
0 commit comments