Skip to content

Commit 71e6137

Browse files
committed
Fix four concurrency races exposed by race-clean integration tests
- daemon.WebhookClient: Close() no longer closes the event channel, which could race with a concurrent Emit() ending in send-on-closed-channel panic; run() now drains and exits via the closed signal instead - daemon.ManagedEngine.persist / PolicyRunner.persist: snapshot peers via a new clonePeersLocked helper so json.Marshal cannot observe concurrent map writes after the RLock is released - registry.Server.flushSave: deep-copy mutable node and network slices (Networks, Tags, LANAddrs, Members, MemberRoles, MemberTags, AllowedPorts) under RLock; replace rawNetCopy pointer-to-NetworkInfo with a value type so Phase 2 encoding no longer races with handleRegister/handleDeregister
1 parent 2cfb652 commit 71e6137

4 files changed

Lines changed: 113 additions & 36 deletions

File tree

pkg/daemon/managed.go

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,27 @@ type managedSnapshot struct {
5050
CycleNum int `json:"cycle_num"`
5151
}
5252

53+
// clonePeersLocked returns a deep copy of a peers map. Caller must hold the
54+
// owning mutex. The copy is safe to hand to json.Marshal after the lock is
55+
// released, because it shares no mutable state with the original.
56+
func clonePeersLocked(src map[uint32]*managedPeer) map[uint32]*managedPeer {
57+
dst := make(map[uint32]*managedPeer, len(src))
58+
for k, p := range src {
59+
pc := *p
60+
if p.Topics != nil {
61+
pc.Topics = make(map[string]int, len(p.Topics))
62+
for tk, tv := range p.Topics {
63+
pc.Topics[tk] = tv
64+
}
65+
}
66+
if p.Tags != nil {
67+
pc.Tags = append([]string(nil), p.Tags...)
68+
}
69+
dst[k] = &pc
70+
}
71+
return dst
72+
}
73+
5374
// NewManagedEngine creates a managed engine for a network.
5475
// It loads persisted state if available, or bootstraps from the member list.
5576
func NewManagedEngine(netID uint16, rules *registry.NetworkRules, d *Daemon) *ManagedEngine {
@@ -416,7 +437,7 @@ func (me *ManagedEngine) persist() {
416437
me.mu.RLock()
417438
snap := managedSnapshot{
418439
NetworkID: me.netID,
419-
Peers: me.peers,
440+
Peers: clonePeersLocked(me.peers),
420441
JoinedAt: me.joinedAt.Format(time.RFC3339),
421442
}
422443
me.mu.RUnlock()

pkg/daemon/policy_runner.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -780,7 +780,7 @@ func (pr *PolicyRunner) persist() {
780780
pr.mu.RLock()
781781
snap := policySnapshot{
782782
NetworkID: pr.netID,
783-
Peers: pr.peers,
783+
Peers: clonePeersLocked(pr.peers),
784784
JoinedAt: pr.joinedAt.Format(time.RFC3339),
785785
CycleNum: pr.cycleNum,
786786
}

pkg/daemon/webhook.go

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -138,7 +138,6 @@ func (wc *WebhookClient) Close() {
138138
}
139139
wc.closeOnce.Do(func() {
140140
close(wc.closed)
141-
close(wc.ch)
142141
})
143142
select {
144143
case <-wc.done:
@@ -149,8 +148,20 @@ func (wc *WebhookClient) Close() {
149148

150149
func (wc *WebhookClient) run() {
151150
defer close(wc.done)
152-
for ev := range wc.ch {
153-
wc.post(ev)
151+
for {
152+
select {
153+
case ev := <-wc.ch:
154+
wc.post(ev)
155+
case <-wc.closed:
156+
for {
157+
select {
158+
case ev := <-wc.ch:
159+
wc.post(ev)
160+
default:
161+
return
162+
}
163+
}
164+
}
154165
}
155166
}
156167

pkg/registry/server.go

Lines changed: 76 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -5755,35 +5755,81 @@ func (s *Server) flushSave() error {
57555755
nextNode := s.nextNode
57565756
nextNet := s.nextNet
57575757

5758-
// Copy node raw values (no encoding under lock)
5758+
// Copy node raw values (no encoding under lock). Slice fields that can be
5759+
// mutated in place elsewhere (Networks/Tags/LANAddrs are append-grown) are
5760+
// deep-copied so Phase 2 sees a stable snapshot after the lock is released.
57595761
rawNodes := make([]rawNodeCopy, 0, len(s.nodes))
57605762
for _, n := range s.nodes {
57615763
rawNodes = append(rawNodes, rawNodeCopy{
57625764
id: n.ID,
57635765
owner: n.Owner,
57645766
publicKey: n.PublicKey,
57655767
realAddr: n.RealAddr,
5766-
networks: n.Networks,
5768+
networks: append([]uint16(nil), n.Networks...),
57675769
lastSeen: n.getLastSeen(),
57685770
public: n.Public,
57695771
hostname: n.Hostname,
5770-
tags: n.Tags,
5772+
tags: append([]string(nil), n.Tags...),
57715773
poloScore: n.PoloScore,
57725774
taskExec: n.TaskExec,
5773-
lanAddrs: n.LANAddrs,
5775+
lanAddrs: append([]string(nil), n.LANAddrs...),
57745776
keyMeta: n.KeyMeta,
57755777
externalID: n.ExternalID,
57765778
version: n.Version,
57775779
})
57785780
}
57795781

5780-
// Copy network data (Created.Format is the only costly op — defer it)
5782+
// Copy network data. Members/MemberRoles/MemberTags mutate in place
5783+
// (handleRegister/handleDeregister append-grow Members, map writes add/remove
5784+
// roles and tags), so Phase 2 must see deep copies rather than live pointers.
57815785
type rawNetCopy struct {
5782-
info *NetworkInfo
5786+
id uint16
5787+
name string
5788+
joinRule string
5789+
token string
5790+
members []uint32
5791+
memberRoles map[uint32]Role
5792+
memberTags map[uint32][]string
5793+
adminToken string
5794+
policy NetworkPolicy
5795+
rules *NetworkRules
5796+
exprPolicy json.RawMessage
5797+
enterprise bool
5798+
created time.Time
5799+
requestCount int64
57835800
}
57845801
rawNets := make([]rawNetCopy, 0, len(s.networks))
57855802
for _, n := range s.networks {
5786-
rawNets = append(rawNets, rawNetCopy{info: n})
5803+
rc := rawNetCopy{
5804+
id: n.ID,
5805+
name: n.Name,
5806+
joinRule: n.JoinRule,
5807+
token: n.Token,
5808+
members: append([]uint32(nil), n.Members...),
5809+
adminToken: n.AdminToken,
5810+
policy: n.Policy,
5811+
rules: n.Rules,
5812+
exprPolicy: n.ExprPolicy,
5813+
enterprise: n.Enterprise,
5814+
created: n.Created,
5815+
requestCount: n.requestCount.Load(),
5816+
}
5817+
if len(n.Policy.AllowedPorts) > 0 {
5818+
rc.policy.AllowedPorts = append([]uint16(nil), n.Policy.AllowedPorts...)
5819+
}
5820+
if len(n.MemberRoles) > 0 {
5821+
rc.memberRoles = make(map[uint32]Role, len(n.MemberRoles))
5822+
for k, v := range n.MemberRoles {
5823+
rc.memberRoles[k] = v
5824+
}
5825+
}
5826+
if len(n.MemberTags) > 0 {
5827+
rc.memberTags = make(map[uint32][]string, len(n.MemberTags))
5828+
for k, v := range n.MemberTags {
5829+
rc.memberTags[k] = append([]string(nil), v...)
5830+
}
5831+
}
5832+
rawNets = append(rawNets, rc)
57875833
}
57885834

57895835
// Copy index maps
@@ -5918,39 +5964,38 @@ func (s *Server) flushSave() error {
59185964
}
59195965
}
59205966

5921-
for _, rn := range rawNets {
5922-
n := rn.info
5967+
for i := range rawNets {
5968+
rn := &rawNets[i]
59235969
sn := &snapshotNet{
5924-
ID: n.ID,
5925-
Name: n.Name,
5926-
JoinRule: n.JoinRule,
5927-
Token: n.Token,
5928-
Members: n.Members,
5929-
AdminToken: n.AdminToken,
5930-
Enterprise: n.Enterprise,
5931-
RequestCount: n.requestCount.Load(),
5932-
Created: n.Created.Format(time.RFC3339),
5933-
}
5934-
if len(n.MemberRoles) > 0 {
5935-
sn.MemberRoles = make(map[string]string, len(n.MemberRoles))
5936-
for nodeID, role := range n.MemberRoles {
5970+
ID: rn.id,
5971+
Name: rn.name,
5972+
JoinRule: rn.joinRule,
5973+
Token: rn.token,
5974+
Members: rn.members,
5975+
AdminToken: rn.adminToken,
5976+
Enterprise: rn.enterprise,
5977+
RequestCount: rn.requestCount,
5978+
Created: rn.created.Format(time.RFC3339),
5979+
}
5980+
if len(rn.memberRoles) > 0 {
5981+
sn.MemberRoles = make(map[string]string, len(rn.memberRoles))
5982+
for nodeID, role := range rn.memberRoles {
59375983
sn.MemberRoles[fmt.Sprintf("%d", nodeID)] = string(role)
59385984
}
59395985
}
5940-
if len(n.MemberTags) > 0 {
5941-
sn.MemberTags = make(map[string][]string, len(n.MemberTags))
5942-
for nodeID, tags := range n.MemberTags {
5986+
if len(rn.memberTags) > 0 {
5987+
sn.MemberTags = make(map[string][]string, len(rn.memberTags))
5988+
for nodeID, tags := range rn.memberTags {
59435989
sn.MemberTags[fmt.Sprintf("%d", nodeID)] = tags
59445990
}
59455991
}
5946-
// Persist policy if any field is set
5947-
if n.Policy.MaxMembers != 0 || len(n.Policy.AllowedPorts) > 0 || n.Policy.Description != "" {
5948-
pol := n.Policy // copy
5992+
if rn.policy.MaxMembers != 0 || len(rn.policy.AllowedPorts) > 0 || rn.policy.Description != "" {
5993+
pol := rn.policy
59495994
sn.Policy = &pol
59505995
}
5951-
sn.Rules = n.Rules
5952-
sn.ExprPolicy = n.ExprPolicy
5953-
snap.Networks[fmt.Sprintf("%d", n.ID)] = sn
5996+
sn.Rules = rn.rules
5997+
sn.ExprPolicy = rn.exprPolicy
5998+
snap.Networks[fmt.Sprintf("%d", rn.id)] = sn
59545999
}
59556000

59566001
snap.PubKeyIdx = pubKeyIdx

0 commit comments

Comments
 (0)