Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
93 changes: 85 additions & 8 deletions coordinator/capture_write_lease.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,19 @@ import (
)

const (
witnessNonceSize = 16
witnessChallengeTimeout = time.Second
witnessNonceSize = 16
witnessChallengeTimeout = time.Second
nodeResourceUsageStaleThreshold = 5 * time.Second
)

type captureLeaseNodeState struct {
nodeEpoch uint64
lastRequestSeq uint64
nodeEpoch uint64
lastRequestSeq uint64
resourceUsageProtocolVersion uint32
eventStoreWriteBytes uint64
eventStoreWriteBytesPerSecond uint64
resourceUsageRateAvailable bool
resourceUsageUpdated time.Time
}

type pendingWitnessChallenge struct {
Expand All @@ -54,6 +60,7 @@ type captureWriteLeaseController struct {
nonce func([]byte) (int, error)

nodes map[node.ID]*captureLeaseNodeState
activeNodes map[node.ID]struct{}
p2pCapableNodes map[node.ID]struct{}
p2pLeaseEnabled bool
pendingWitness *pendingWitnessChallenge
Expand All @@ -67,6 +74,7 @@ func newCaptureWriteLeaseController(version int64, selfNodeID node.ID) *captureW
now: time.Now,
nonce: rand.Read,
nodes: make(map[node.ID]*captureLeaseNodeState),
activeNodes: make(map[node.ID]struct{}),
p2pCapableNodes: make(map[node.ID]struct{}),
}
}
Expand All @@ -82,6 +90,11 @@ func (c *captureWriteLeaseController) observeNodeCapability(id node.ID, version
// updateClusterMode enables P2P only when every active capture has reported
// support for the current protocol. A missing capability is treated as legacy.
func (c *captureWriteLeaseController) updateClusterMode(activeNodes []node.ID) {
c.activeNodes = make(map[node.ID]struct{}, len(activeNodes))
for _, id := range activeNodes {
c.activeNodes[id] = struct{}{}
}

p2pEnabled := len(activeNodes) > 0
for _, id := range activeNodes {
if _, ok := c.p2pCapableNodes[id]; !ok {
Expand Down Expand Up @@ -120,6 +133,29 @@ func (c *captureWriteLeaseController) handleHeartbeat(
return nil
}
state.lastRequestSeq = heartbeat.GetWriteLeaseRequestSeq()
state.resourceUsageProtocolVersion = heartbeat.GetNodeResourceUsageProtocolVersion()
usage := heartbeat.GetNodeResourceUsage()
if state.resourceUsageProtocolVersion == heartbeatpb.CurrentNodeResourceUsageProtocolVersion && usage != nil {
now := c.now()
writeBytes := usage.GetEventStoreWriteBytes()
if !state.resourceUsageUpdated.IsZero() && now.After(state.resourceUsageUpdated) {
if writeBytes >= state.eventStoreWriteBytes {
state.eventStoreWriteBytesPerSecond = uint64(
float64(writeBytes-state.eventStoreWriteBytes) /
now.Sub(state.resourceUsageUpdated).Seconds())
state.resourceUsageRateAvailable = true
} else {
state.resourceUsageRateAvailable = false
}
}
if state.resourceUsageUpdated.IsZero() || now.After(state.resourceUsageUpdated) {
state.eventStoreWriteBytes = writeBytes
state.resourceUsageUpdated = now
}
} else {
state.resourceUsageUpdated = time.Time{}
state.resourceUsageRateAvailable = false
}

messages := c.handleWitnessAck(from, heartbeat)
if from != c.selfNodeID {
Expand Down Expand Up @@ -248,20 +284,61 @@ func (c *captureWriteLeaseController) newGrant(
if c.p2pLeaseEnabled {
leaseDurationMs = uint64(writelease.P2PLeaseDuration.Milliseconds())
}
nodeResourceUsages, nodeResourceUsageStatus := c.nodeResourceUsageSnapshot()
return messaging.NewSingleTargetMessage(
target,
messaging.MaintainerManagerTopic,
&heartbeatpb.NodeHeartbeatResponse{
CoordinatorVersion: c.coordinatorVersion,
TargetNodeEpoch: targetNodeEpoch,
RequestSeq: requestSeq,
LeaseDurationMs: leaseDurationMs,
CoordinatorVersion: c.coordinatorVersion,
TargetNodeEpoch: targetNodeEpoch,
RequestSeq: requestSeq,
LeaseDurationMs: leaseDurationMs,
NodeResourceUsages: nodeResourceUsages,
NodeResourceUsageStatus: nodeResourceUsageStatus,
},
)
}

func (c *captureWriteLeaseController) nodeResourceUsageSnapshot() (
[]*heartbeatpb.NodeResourceUsage,
heartbeatpb.NodeResourceUsageStatus,
) {
if len(c.activeNodes) == 0 {
return nil, heartbeatpb.NodeResourceUsageStatus_UNSUPPORTED
}
for nodeID := range c.activeNodes {
state := c.nodes[nodeID]
if state == nil ||
state.resourceUsageProtocolVersion != heartbeatpb.CurrentNodeResourceUsageProtocolVersion {
return nil, heartbeatpb.NodeResourceUsageStatus_UNSUPPORTED
}
}

now := c.now()
nodeIDs := make([]node.ID, 0, len(c.activeNodes))
for nodeID := range c.activeNodes {
state := c.nodes[nodeID]
if !state.resourceUsageRateAvailable || state.resourceUsageUpdated.IsZero() ||
now.Sub(state.resourceUsageUpdated) > nodeResourceUsageStaleThreshold {
return nil, heartbeatpb.NodeResourceUsageStatus_INCOMPLETE
}
nodeIDs = append(nodeIDs, nodeID)
}
slices.Sort(nodeIDs)

result := make([]*heartbeatpb.NodeResourceUsage, 0, len(nodeIDs))
for _, nodeID := range nodeIDs {
result = append(result, &heartbeatpb.NodeResourceUsage{
NodeId: nodeID.String(),
EventStoreWriteBytesPerSecond: c.nodes[nodeID].eventStoreWriteBytesPerSecond,
})
}
return result, heartbeatpb.NodeResourceUsageStatus_AVAILABLE
}

func (c *captureWriteLeaseController) removeNode(id node.ID) {
delete(c.nodes, id)
delete(c.activeNodes, id)
delete(c.p2pCapableNodes, id)
if c.pendingWitness != nil &&
(c.pendingWitness.witnessNodeID == id || id == c.selfNodeID) {
Expand Down
97 changes: 93 additions & 4 deletions coordinator/capture_write_lease_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,94 @@ func TestCaptureWriteLeaseGrantsRemoteNode(t *testing.T) {
require.Equal(t, uint64(12), requireWriteLeaseResponse(t, messages[0]).TargetNodeEpoch)
}

func TestCaptureWriteLeaseSharesFreshNodeResourceUsage(t *testing.T) {
now := time.Unix(100, 0)
controller := newCaptureWriteLeaseController(10, node.ID("coordinator"))
controller.now = func() time.Time { return now }
enableP2PForNodes(controller, node.ID("capture-1"), node.ID("capture-2"))

capture2Heartbeat := newWriteLeaseHeartbeat(21, 1)
capture2Heartbeat.NodeResourceUsage = &heartbeatpb.NodeResourceUsage{
EventStoreWriteBytes: 1000,
}
controller.handleHeartbeat(node.ID("capture-2"), capture2Heartbeat, nil)

capture1Heartbeat := newWriteLeaseHeartbeat(11, 1)
capture1Heartbeat.NodeResourceUsage = &heartbeatpb.NodeResourceUsage{
EventStoreWriteBytes: 100,
}
messages := controller.handleHeartbeat(node.ID("capture-1"), capture1Heartbeat, nil)
require.Len(t, messages, 1)
require.Equal(t,
heartbeatpb.NodeResourceUsageStatus_INCOMPLETE,
requireWriteLeaseResponse(t, messages[0]).NodeResourceUsageStatus)

// Each node's rate is calculated from its own heartbeat interval.
now = now.Add(time.Second)
capture2Heartbeat.WriteLeaseRequestSeq = 2
capture2Heartbeat.NodeResourceUsage.EventStoreWriteBytes = 3000
controller.handleHeartbeat(node.ID("capture-2"), capture2Heartbeat, nil)
capture1Heartbeat.WriteLeaseRequestSeq = 2
capture1Heartbeat.NodeResourceUsage.EventStoreWriteBytes = 200
messages = controller.handleHeartbeat(node.ID("capture-1"), capture1Heartbeat, nil)
require.Len(t, messages, 1)
require.Equal(t,
heartbeatpb.NodeResourceUsageStatus_AVAILABLE,
requireWriteLeaseResponse(t, messages[0]).NodeResourceUsageStatus)
require.Equal(t, []*heartbeatpb.NodeResourceUsage{
{NodeId: "capture-1", EventStoreWriteBytesPerSecond: 100},
{NodeId: "capture-2", EventStoreWriteBytesPerSecond: 2000},
}, requireWriteLeaseResponse(t, messages[0]).NodeResourceUsages)

// Only capture-1 reports again. A response reuses capture-2's last valid
// rate instead of interpreting the unchanged snapshot as zero traffic.
now = now.Add(time.Second)
capture1Heartbeat.WriteLeaseRequestSeq = 3
capture1Heartbeat.NodeResourceUsage.EventStoreWriteBytes = 300
messages = controller.handleHeartbeat(node.ID("capture-1"), capture1Heartbeat, nil)
require.Equal(t, []*heartbeatpb.NodeResourceUsage{
{NodeId: "capture-1", EventStoreWriteBytesPerSecond: 100},
{NodeId: "capture-2", EventStoreWriteBytesPerSecond: 2000},
}, requireWriteLeaseResponse(t, messages[0]).NodeResourceUsages)

// A fresh report must not keep another node's stale sample in the cluster
// snapshot. All nodes support reporting, so this is an interruption rather
// than a rolling-upgrade fallback.
now = now.Add(nodeResourceUsageStaleThreshold + time.Nanosecond)
capture1Heartbeat.WriteLeaseRequestSeq = 4
capture1Heartbeat.NodeResourceUsage.EventStoreWriteBytes = 400
messages = controller.handleHeartbeat(node.ID("capture-1"), capture1Heartbeat, nil)
require.Len(t, messages, 1)
response := requireWriteLeaseResponse(t, messages[0])
require.Equal(t,
heartbeatpb.NodeResourceUsageStatus_INCOMPLETE,
response.NodeResourceUsageStatus)
require.Empty(t, response.NodeResourceUsages)

// Missing usage means the sender no longer supports or cannot provide the
// counter, so its previous value is removed immediately.
capture1Heartbeat.WriteLeaseRequestSeq = 5
capture1Heartbeat.NodeResourceUsage = nil
messages = controller.handleHeartbeat(node.ID("capture-1"), capture1Heartbeat, nil)
require.Len(t, messages, 1)
response = requireWriteLeaseResponse(t, messages[0])
require.Equal(t,
heartbeatpb.NodeResourceUsageStatus_INCOMPLETE,
response.NodeResourceUsageStatus)
require.Empty(t, response.NodeResourceUsages)

// A node that does not declare the resource protocol is a rolling-upgrade
// compatibility case, distinct from interrupted telemetry.
capture1Heartbeat.WriteLeaseRequestSeq = 6
capture1Heartbeat.NodeResourceUsageProtocolVersion = heartbeatpb.LegacyNodeResourceUsageProtocolVersion
messages = controller.handleHeartbeat(node.ID("capture-1"), capture1Heartbeat, nil)
require.Len(t, messages, 1)
response = requireWriteLeaseResponse(t, messages[0])
require.Equal(t,
heartbeatpb.NodeResourceUsageStatus_UNSUPPORTED,
response.NodeResourceUsageStatus)
}

func TestCaptureWriteLeaseRequiresRemoteWitnessForCoordinatorNode(t *testing.T) {
now := time.Unix(100, 0)
controller := newCaptureWriteLeaseController(10, node.ID("coordinator"))
Expand Down Expand Up @@ -246,10 +334,11 @@ func TestCaptureWriteLeaseRejectsInvalidHeartbeatAndLateWitness(t *testing.T) {

func newWriteLeaseHeartbeat(nodeEpoch, requestSeq uint64) *heartbeatpb.NodeHeartbeat {
return &heartbeatpb.NodeHeartbeat{
Liveness: heartbeatpb.NodeLiveness_ALIVE,
NodeEpoch: nodeEpoch,
WriteLeaseRequestSeq: requestSeq,
WriteLeaseProtocolVersion: heartbeatpb.CurrentWriteLeaseProtocolVersion,
Liveness: heartbeatpb.NodeLiveness_ALIVE,
NodeEpoch: nodeEpoch,
WriteLeaseRequestSeq: requestSeq,
WriteLeaseProtocolVersion: heartbeatpb.CurrentWriteLeaseProtocolVersion,
NodeResourceUsageProtocolVersion: heartbeatpb.CurrentNodeResourceUsageProtocolVersion,
}
}

Expand Down
1 change: 0 additions & 1 deletion downstreamadapter/dispatchermanager/dispatcher_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -924,7 +924,6 @@ func (e *DispatcherManager) aggregateDispatcherHeartbeats(needCompleteStatus boo
Watermark: heartbeatpb.NewMaxWatermark(),
RedoWatermark: heartbeatpb.NewMaxWatermark(),
}

toCleanMap := make([]*cleanMap, 0)
dispatcherCount := 0

Expand Down
Loading
Loading