@@ -224,6 +224,11 @@ func (transport *Transport) readLoopTCPConn(conn net.Conn, logger zerolog.Logger
224224 for {
225225 packet , err := transport .decoder .DecodeNext (conn )
226226 if err != nil {
227+ // net.ErrClosed means we closed the connection on purpose
228+ // (e.g. UDP keep-alive timeout), and a fault was already sent
229+ if errors .Is (err , net .ErrClosed ) {
230+ return
231+ }
227232 logger .Error ().Stack ().Err (err ).Msg ("decode" )
228233 transport .errChan <- err
229234 transport .SendFault ()
@@ -268,6 +273,10 @@ func (transport *Transport) SendMessage(message abstraction.TransportMessage) er
268273 return err
269274}
270275
276+ // faultWriteTimeout bounds each TCP write of the fault broadcast; boards on
277+ // the vehicle LAN ack in milliseconds, so exceeding this means a dead peer
278+ const faultWriteTimeout = time .Second
279+
271280// handlePacketEvent is used to send an order to one of the connected boards
272281func (transport * Transport ) handlePacketEvent (message PacketMessage ) error {
273282 eventLogger := transport .logger .With ().Str ("type" , fmt .Sprintf ("%T" , message .Packet )).Uint16 ("id" , uint16 (message .Id ())).Logger ()
@@ -287,17 +296,30 @@ func (transport *Transport) handlePacketEvent(message PacketMessage) error {
287296 defer transport .connectionsMx .RUnlock ()
288297 for target , conn := range transport .connections {
289298 targetName := string (target )
299+
300+ // Bound each write so a dead peer with a full send buffer cannot
301+ // hold the connections lock (and the caller) for minutes
302+ conn .SetWriteDeadline (time .Now ().Add (faultWriteTimeout ))
303+ var writeErr error
290304 totalWritten := 0
291305 for totalWritten < len (data ) {
292306 n , err := conn .Write (data [totalWritten :])
293307 eventLogger .Trace ().Str ("target" , targetName ).Int ("amount" , n ).Msg ("written chunk" )
294308 totalWritten += n
295309 if err != nil {
296- eventLogger .Error ().Str ("target" , targetName ).Stack ().Err (err ).Msg ("write" )
297- transport .errChan <- err
298- return err
310+ writeErr = err
311+ break
299312 }
300313 }
314+ conn .SetWriteDeadline (time.Time {})
315+
316+ // Keep broadcasting to the remaining boards even if one write
317+ // fails: the fault must reach every live board
318+ if writeErr != nil {
319+ eventLogger .Error ().Str ("target" , targetName ).Stack ().Err (writeErr ).Msg ("write" )
320+ transport .errChan <- writeErr
321+ continue
322+ }
301323 eventLogger .Info ().Str ("target" , targetName ).Msg ("sent" )
302324 }
303325 return nil
@@ -475,6 +497,22 @@ func (transport *Transport) TargetFromIp(ip string) (abstraction.TransportTarget
475497 return target , ok
476498}
477499
500+ // DisconnectTarget forcefully closes the TCP connection to target, if any.
501+ // The connection handler wakes up with reason, cleans up and notifies the
502+ // disconnection, and the client reconnection loop takes over.
503+ func (transport * Transport ) DisconnectTarget (target abstraction.TransportTarget , reason error ) bool {
504+ transport .connectionsMx .RLock ()
505+ conn , ok := transport .connections [target ]
506+ transport .connectionsMx .RUnlock ()
507+ if ! ok {
508+ return false
509+ }
510+
511+ transport .logger .Warn ().Str ("target" , string (target )).Err (reason ).Msg ("forcefully disconnecting target" )
512+ tcp .CloseWithError (conn , reason )
513+ return true
514+ }
515+
478516func (transport * Transport ) SendFault () {
479517 err := transport .SendMessage (NewPacketMessage (data .NewPacket (0 )))
480518 if err != nil {
0 commit comments