Skip to content

Commit 3924f9e

Browse files
gaikwadabhishekalex-aizman
authored andcommitted
shared data-mover: demux/transport lock-order
* release rxmu (protecting receivers map and demux: xid => receiver) before calling Close() - no need to hold this particular mutex over unrelated function(s) - there's an r-lock on the receive (SDM.recv) side * when dropping Rx, differentiate SDM close vs xaction abort ---- * sperately, cmn.SparseWarn to show up in the source where it belongs Signed-off-by: Alex Aizman <alex.aizman@gmail.com>
1 parent 227ce54 commit 3924f9e

2 files changed

Lines changed: 17 additions & 9 deletions

File tree

cmn/sparse_warn.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,6 @@ import (
1515
// per-module verbose level is at or above 5.
1616
func SparseWarn(mod int, cnt int64, args ...any) {
1717
if Rom.V(5, mod) || cos.Sparse(cnt) {
18-
nlog.Warningln(args...)
18+
nlog.WarningDepth(1, args...)
1919
}
2020
}

transport/bundle/shared_dm.go

Lines changed: 16 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -213,11 +213,13 @@ func (sdm *sharedDM) Close(err ...error) error {
213213
nlog.Errorln(sdm.trname(), "closing despite", msg, "[", err[0], "]")
214214
}
215215

216-
sdm.dm.Close(nil)
217-
sdm.dm.UnregRecv()
216+
// release rxmu (intended only to protect sdm.receivers) before calling dm.Close()
218217
sdm.receivers = nil
219218
sdm.rxmu.Unlock()
220219

220+
sdm.dm.Close(nil)
221+
sdm.dm.UnregRecv()
222+
221223
sdm.ocmu.Unlock()
222224

223225
sdm.sbrs.mtx.Lock()
@@ -305,19 +307,25 @@ func (sdm *sharedDM) recv(hdr *transport.ObjHdr, r io.Reader, err error) error {
305307
return err
306308
}
307309

310+
var (
311+
en *rxent
312+
closed, ok bool
313+
)
308314
sdm.rxmu.RLock()
309-
en, ok := sdm.receivers[xid]
315+
closed = sdm.receivers == nil
316+
en, ok = sdm.receivers[xid]
317+
sdm.rxmu.RUnlock()
318+
310319
if !ok {
311-
sdm.rxmu.RUnlock()
312320
transport.DrainAndFreeReader(r)
313321
n := sdm.stats.drops.Inc()
314-
if n < 5 || n%100 == 0 || cmn.Rom.V(4, cos.ModTransport) {
315-
nlog.Warningf("%s: xid %s not found, dropping recv [oname: %s, drops: %d]",
316-
sdm.trname(), xid, hdr.ObjName, n)
322+
if closed {
323+
cmn.SparseWarn(cos.ModTransport, n, sdm.trname(), "is closed (dropping all Rx)")
324+
} else {
325+
cmn.SparseWarn(cos.ModTransport, n, sdm.trname(), "xid", xid, "not found - dropping", hdr.ObjName)
317326
}
318327
return nil
319328
}
320-
sdm.rxmu.RUnlock()
321329

322330
// (unlikely)
323331
if en.rx.ID() != xid {

0 commit comments

Comments
 (0)