@@ -296,6 +296,17 @@ func (s *AgentService) GetAgentAuth(ctx context.Context, req *ConnectorAuthReque
296296 return & ConnectorAuthResponse {Key : agent .AgentKey , TenantId : agent .TenantID }, nil
297297}
298298
299+ // evictIfOwner deletes the AgentStreamMap entry for agentID only if it still
300+ // points to stream. Prevents a slow-exiting prior AgentStream goroutine from
301+ // clobbering the fresh entry a newly-reconnected agent installed.
302+ func (s * AgentService ) evictIfOwner (agentID uint , stream AgentService_AgentStreamServer ) {
303+ s .AgentStreamMutex .Lock ()
304+ if s .AgentStreamMap [agentID ] == stream {
305+ delete (s .AgentStreamMap , agentID )
306+ }
307+ s .AgentStreamMutex .Unlock ()
308+ }
309+
299310func (s * AgentService ) AgentStream (stream AgentService_AgentStreamServer ) error {
300311 id , _ , _ , err := utils .GetItemsFromContext (stream .Context ())
301312 if err != nil {
@@ -307,11 +318,12 @@ func (s *AgentService) AgentStream(stream AgentService_AgentStreamServer) error
307318 }
308319 idUint := uint (idInt )
309320
321+ // Replace any prior entry rather than rejecting the reconnect. A dead
322+ // prior stream's goroutine may still be looping on Recv (see
323+ // utils.WaitForReconnect) and would otherwise block the agent from
324+ // re-registering for minutes. evictIfOwner guards the map so the old
325+ // goroutine's eventual delete does not clobber the fresh entry.
310326 s .AgentStreamMutex .Lock ()
311- if _ , ok := s .AgentStreamMap [idUint ]; ok {
312- s .AgentStreamMutex .Unlock ()
313- return status .Error (codes .AlreadyExists , "stream already exists" )
314- }
315327 s .AgentStreamMap [idUint ] = stream
316328 s .AgentStreamMutex .Unlock ()
317329
@@ -324,18 +336,17 @@ func (s *AgentService) AgentStream(stream AgentService_AgentStreamServer) error
324336 if err == io .EOF {
325337 err = utils .WaitForReconnect (stream .Context (), stream )
326338 if err != nil {
327- s .AgentStreamMutex .Lock ()
328- delete (s .AgentStreamMap , idUint )
329- s .AgentStreamMutex .Unlock ()
330-
339+ catcher .Info ("AgentStream: WaitForReconnect failed, evicting stream" ,
340+ map [string ]any {"agent_id" : idUint , "err" : err .Error (), "process" : "agent-manager" })
341+ s .evictIfOwner (idUint , stream )
331342 return status .Error (codes .Internal , fmt .Sprintf ("failed to reconnect: %v" , err ))
332343 }
333344 continue
334345 }
335346 if err != nil {
336- s . AgentStreamMutex . Lock ()
337- delete ( s . AgentStreamMap , idUint )
338- s .AgentStreamMutex . Unlock ( )
347+ catcher . Info ( "AgentStream: Recv errored, evicting stream" ,
348+ map [ string ] any { "agent_id" : idUint , "err" : err . Error (), "process" : "agent-manager" } )
349+ s .evictIfOwner ( idUint , stream )
339350 return status .Error (codes .Internal , fmt .Sprintf ("failed to receive message: %v" , err ))
340351 }
341352
0 commit comments