Skip to content
Merged
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
20 changes: 20 additions & 0 deletions RELEASE_NOTES_2026.09.2.md
Original file line number Diff line number Diff line change
Expand Up @@ -344,6 +344,26 @@ fails the build rather than leaving the note quietly wrong.

## Bug fixes

### No primary writer was ever elected, so retention and continuous queries silently never ran ([#850](https://github.com/Basekick-Labs/arc/issues/850))

**Affects clusters with `cluster.failover_enabled=true` and `cluster.shared_storage_mode=false`, which is the Enterprise Helm chart's default for local-storage deployments.**

Singleton work — the retention and continuous-query schedulers, and the non-dry-run retention, CQ and delete endpoints — is gated on `IsPrimaryWriter()`. With writer failover enabled, that gate was false on every node in the cluster, forever, so retention never deleted anything, continuous queries never ran, and those endpoints answered 503 `is not primary writer`. Nothing logged an error, because each node simply believed it was not the writer. Turning failover off avoided it, since the gate then falls back to a plain role check.

Two defects combined. The writer failover manager could only fail over *from* an existing primary: its health check triggered a promotion only when it had already recorded one, and nothing else ever issued a promotion, so the first one could never happen. That is fixed by electing an initial primary when a cluster has none, mirroring what the compactor manager already did for its own lease. Second, the promotion updated the node registry, which hands out copies, while the gate reads the coordinator's own node object, so even a promotion that did happen never reached it. The promotion and the snapshot-restore path now both update that object.

Shared-storage multi-writer clusters (`cluster.shared_storage_mode=true`) were never affected: writer failover is suppressed there by design and the gate keys off Raft leadership instead.

A writer that simply restarted also used to lose the designation, because a re-join replaces the node's cluster record and the join payload carries no writer state. The cluster then had a primary it could not name, and in a single-writer deployment it never recovered one, because the only candidate was the node the failover logic had just excluded. The node table now keeps the designation across a re-join, and a healthy writer that is the only candidate can be re-elected rather than skipped.

One behaviour changes as a result. A write that arrives at a reader is proxied to a writer, and with no primary ever designated those proxied writes were spread across all healthy writers. They now go to the elected primary, which is what Pattern 1 intends, since only that node should be ingesting. Writers still serve their own traffic directly, and shared-storage clusters are unaffected.

This changes how two replicated commands are applied, so upgrade a cluster fully rather than leaving it mixed-version for long: an old binary elected as Raft leader will not elect a primary, and a new one will once it takes over.

`GET /api/v1/cluster/nodes` now reports each node's `writer_state`, which it did not before. That absence is a large part of why this went unnoticed: there was no way to ask the cluster which node held the primary role. `GET /api/v1/cluster/local` reports the same field for the node serving the request, and that is the one the scheduler gate actually reads.

Two known gaps remain and are tracked separately. Readers are still not promotion candidates ([#856](https://github.com/Basekick-Labs/arc/issues/856)), so a cluster needs more than one writer-role node to survive losing one; the Enterprise Helm chart's Pattern 1 example is corrected accordingly. And a load balancer still has no way to target the elected primary, because `/ready` does not distinguish roles ([#857](https://github.com/Basekick-Labs/arc/issues/857)).

### Cluster shutdown no longer holds the coordinator and Raft locks while it waits for its subsystems ([#813](https://github.com/Basekick-Labs/arc/issues/813))

Stopping a clustered node joined every subsystem — the file puller, the writer and compactor failover managers, the Raft node, the delete worker — while holding the coordinator's lock, and the Raft node's own `Stop` held its lock across the Raft shutdown, which waits for the goroutine that applies entries to the FSM. Anything on one of those joined goroutines that took either lock deadlocked the shutdown, and a deadlocked shutdown hook hangs the whole process until the supervisor kills it. #797 fixed one such caller, the FSM delete callback, and left the structure in place with a contract comment and a test that guards the FSM callbacks only; a puller gate check had already been working around it with a try-lock.
Expand Down
17 changes: 9 additions & 8 deletions helm/arc-enterprise/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -198,18 +198,19 @@ automatically by the chart (see `templates/_helpers.tpl`). This:

#### Local storage (Pattern 1) — `storage.mode: local`

Each writer pod owns its own PersistentVolume. In current Arc Enterprise
v26.06.x this pattern still runs in **single-writer + multi-reader** mode:
one writer takes ingest, readers replicate the WAL and can be promoted on
writer failure via Arc's writer-failover controller. (Pattern 1 multi-writer
— per-shard writer ownership — is a future initiative; see the multi-writer
plan doc.)
Each writer pod owns its own PersistentVolume. This pattern runs in
**single-writer + multi-reader** mode: one writer takes ingest and the readers
replicate its WAL. Arc's writer-failover controller elects the primary writer
among the nodes whose role is `writer`; readers are **not** promotion
candidates today, so give the cluster more than one writer replica if you want
it to survive losing one. (Pattern 1 multi-writer — per-shard writer ownership
— is a future initiative; see the multi-writer plan doc.)

```yaml
writer:
replicas: 1 # one active writer
replicas: 2 # one is elected primary, the other can take over
reader:
replicas: 3 # readers double as failover pool
replicas: 3 # readers serve queries and replicate the WAL
cluster:
failover:
enabled: true # automatic writer + compactor failover
Expand Down
14 changes: 10 additions & 4 deletions internal/api/cluster.go
Original file line number Diff line number Diff line change
Expand Up @@ -247,10 +247,16 @@ func (h *ClusterHandler) respondNotEnabled(c *fiber.Ctx) error {
// nodeToMap converts a Node to a map for JSON serialization.
func (h *ClusterHandler) nodeToMap(node *cluster.Node) map[string]interface{} {
return map[string]interface{}{
"id": node.ID,
"name": node.Name,
"role": node.Role,
"state": node.GetState(),
"id": node.ID,
"name": node.Name,
"role": node.Role,
"state": node.GetState(),
// Which writer is the primary is what gates every singleton task
// (retention, CQ, delete). It was absent here, so an operator had no
// way to see that no node held it — which is how #850 stayed hidden.
// Empty for readers, compactors, and for writers in shared-storage
// mode, where there is no primary/standby distinction.
"writer_state": node.GetWriterState(),
"address": node.Address,
"api_address": node.APIAddress,
"cluster_name": node.ClusterName,
Expand Down
29 changes: 28 additions & 1 deletion internal/cluster/coordinator.go
Original file line number Diff line number Diff line change
Expand Up @@ -2070,7 +2070,7 @@ func (c *Coordinator) GetRole() NodeRole {
// wins the election would run retention/CQ/delete.
//
// Leader-change semantics: each scheduler (retention, CQ, delete,
// reconciliation) checks IsPrimaryWriter() ONCE at the start of each
// delete endpoints) checks IsPrimaryWriter() ONCE at the start of each
// tick and runs all work for that tick if true. A leader change
// mid-tick will let the demoted node complete its current tick's
// work; the new leader's next tick picks up from there. This is
Expand Down Expand Up @@ -2378,6 +2378,30 @@ func (c *Coordinator) onRaftNodeAdded(n *raft.NodeInfo) {
if err := c.registry.Register(node); err != nil {
c.logger.Error().Err(err).Str("node_id", n.ID).Msg("Failed to register node from Raft")
}
// A snapshot restore replays the membership without firing a promotion
// callback, so the writer state the FSM carries has to reach the local
// node here or this node would forget it was the primary across a
// restart (#850). Mirrored unconditionally, including the empty value:
// the FSM record is authoritative, and applyAddNode now preserves a live
// designation across a re-join, so an empty value here means this node
// genuinely holds none and must not keep claiming one.
c.setLocalWriterState(n.ID, WriterState(n.WriterState))
}

// setLocalWriterState mirrors a writer-state change onto c.localNode when the
// id names this node.
//
// SECURITY of the invariant, not of access: Registry.Get hands out a CLONE, so
// onWriterPromoted below updates a copy and re-registers it, replacing the map
// entry. c.localNode is a different object, and it is the one
// Coordinator.IsPrimaryWriter consults through GetLocalNode(). Before #850 the
// promotion therefore never reached the gate: the registry said "primary" while
// the node's own scheduler gate said "not primary" and silently skipped every
// tick.
func (c *Coordinator) setLocalWriterState(nodeID string, state WriterState) {
if c.localNode != nil && nodeID == c.localNode.ID {
c.localNode.SetWriterState(state)
}
}

// onRaftNodeRemoved is called when a node is removed via Raft consensus.
Expand Down Expand Up @@ -2408,13 +2432,16 @@ func (c *Coordinator) onWriterPromoted(newPrimaryID, oldPrimaryID string) {
oldNode.SetWriterState(WriterStateStandby)
c.registry.Register(oldNode)
}
c.setLocalWriterState(oldPrimaryID, WriterStateStandby)
}

// Promote new primary in registry
if newNode, exists := c.registry.Get(newPrimaryID); exists {
newNode.SetWriterState(WriterStatePrimary)
c.registry.Register(newNode)
}
// The registry entry is a clone; the gate reads c.localNode (#850).
c.setLocalWriterState(newPrimaryID, WriterStatePrimary)

c.logger.Info().
Str("new_primary", newPrimaryID).
Expand Down
32 changes: 21 additions & 11 deletions internal/cluster/raft/fsm.go
Original file line number Diff line number Diff line change
Expand Up @@ -932,6 +932,15 @@ func (f *ClusterFSM) applyAddNode(payload []byte) interface{} {
}

f.mu.Lock()
// A re-join replaces the record wholesale, and the join payload carries no
// writer state, so a writer that merely restarted used to come back
// undesignated while primaryWriterID still named it. The failover manager
// then saw no primary, took the failover branch, and excluded the only
// writer from selection: the cluster never regained a primary and every
// singleton task stayed off (#850). Keep the table self-consistent.
if p.Node.WriterState == "" && f.primaryWriterID == p.Node.ID && p.Node.Role == "writer" {
p.Node.WriterState = "primary"
}
f.nodes[p.Node.ID] = &p.Node
callback := f.onNodeAdded
f.mu.Unlock()
Expand Down Expand Up @@ -1033,10 +1042,18 @@ func (f *ClusterFSM) applyPromoteWriter(payload []byte) interface{} {
}

f.mu.Lock()
// Validate the node exists and is a writer
if node, exists := f.nodes[p.NodeID]; exists && node.Role != "writer" {
// Validate the node exists and is a writer. Both checks happen before any
// mutation: demoting the old primary and then refusing the new one left
// the table with a standby old primary and a designation naming a node
// that is not there, which no callback ever announces.
newNode, exists := f.nodes[p.NodeID]
if !exists {
f.mu.Unlock()
return fmt.Errorf("promote writer: node %s has role %s, expected writer", p.NodeID, node.Role)
return fmt.Errorf("node %s not found", p.NodeID)
}
if newNode.Role != "writer" {
f.mu.Unlock()
return fmt.Errorf("promote writer: node %s has role %s, expected writer", p.NodeID, newNode.Role)
}

// Warn if OldPrimaryID doesn't match actual primary (informational only — FSM uses its own tracking)
Expand All @@ -1054,18 +1071,11 @@ func (f *ClusterFSM) applyPromoteWriter(payload []byte) interface{} {
}

// Promote new primary
newNode, exists := f.nodes[p.NodeID]
if exists {
newNode.WriterState = "primary"
}
newNode.WriterState = "primary"
f.primaryWriterID = p.NodeID
callback := f.onWriterPromoted
f.mu.Unlock()

if !exists {
return fmt.Errorf("node %s not found", p.NodeID)
}

f.logger.Info().
Str("new_primary", p.NodeID).
Str("old_primary", oldPrimaryID).
Expand Down
4 changes: 2 additions & 2 deletions internal/cluster/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -280,7 +280,7 @@ func (r *Registry) GetPrimaryWriter() *Node {
r.mu.RLock()
defer r.mu.RUnlock()
for _, node := range r.nodes {
if node.Role == RoleWriter && node.WriterSt == WriterStatePrimary && node.IsHealthy() {
if node.Role == RoleWriter && node.GetWriterState() == WriterStatePrimary && node.IsHealthy() {
return node.Clone()
}
}
Expand All @@ -290,7 +290,7 @@ func (r *Registry) GetPrimaryWriter() *Node {
// GetStandbyWriters returns healthy standby writer nodes.
func (r *Registry) GetStandbyWriters() []*Node {
return r.filterNodes(func(n *Node) bool {
return n.Role == RoleWriter && n.WriterSt == WriterStateStandby && n.IsHealthy()
return n.Role == RoleWriter && n.GetWriterState() == WriterStateStandby && n.IsHealthy()
})
}

Expand Down
120 changes: 109 additions & 11 deletions internal/cluster/writer_failover.go
Original file line number Diff line number Diff line change
Expand Up @@ -157,12 +157,24 @@ func (m *WriterFailoverManager) checkPrimaryHealth() {
if primary == nil {
// No primary — check if there are any writers at all
writers := m.cfg.Registry.GetWriters()
if len(writers) > 0 && m.primaryID != "" {
// We had a primary but it's gone — trigger failover
m.consecutiveFails++
if m.consecutiveFails >= m.cfg.UnhealthyThreshold {
m.triggerFailoverLocked()
}
if len(writers) == 0 {
return
}
if m.primaryID == "" {
// No primary has ever existed in this cluster: elect one.
// Without this the manager deadlocks at boot — the failover
// branch below needs a previous primary to fail over FROM, and
// nothing else ever issues CommandPromoteWriter, so every node
// stayed WriterState-less and IsPrimaryWriter() was false
// cluster-wide, silently disabling the retention and CQ
// schedulers and the delete/retention/CQ endpoints (#850).
m.tryInitialElectionLocked()
return
}
// We had a primary but it's gone — trigger failover
m.consecutiveFails++
if m.consecutiveFails >= m.cfg.UnhealthyThreshold {
m.triggerFailoverLocked()
}
return
}
Expand Down Expand Up @@ -210,6 +222,62 @@ func (m *WriterFailoverManager) HandleWriterUnhealthy(node *Node) {
}
}

// tryInitialElectionLocked elects the first primary writer of a cluster that
// has never had one (must hold lock). It is the writer-side counterpart of
// CompactorFailoverManager.tryInitialAssignment and mirrors it: the same
// in-progress guard, the same asynchronous Raft apply so the health-check loop
// is never blocked, and the same completion bookkeeping.
//
// Only the Raft leader reaches here (checkPrimaryHealth returns early
// otherwise), so exactly one node elects. An election is not a failover: there
// is no old primary to demote, no cooldown to respect on the way in, and none
// armed on the way out — see completeElection. A failing election therefore
// retries on the next tick rather than backing off, which is what the
// compactor's initial assignment does too; each attempt costs one Raft apply
// bounded by FailoverTimeout, so a quorum outage logs at most one error per
// timeout rather than per tick.
func (m *WriterFailoverManager) tryInitialElectionLocked() {
if m.failoverInProg {
return
}
newPrimaryID := m.selectNewPrimary("")
if newPrimaryID == "" {
return
}
m.failoverInProg = true
m.logger.Info().
Str("node_id", newPrimaryID).
Msg("Electing initial primary writer")
m.wg.Add(1)
go func() {
defer m.wg.Done()
err := m.cfg.RaftNode.PromoteWriter(newPrimaryID, "", m.cfg.FailoverTimeout)
if err != nil {
m.logger.Error().Err(err).
Str("node_id", newPrimaryID).
Msg("Failed to elect the initial primary writer")
}
m.completeElection(newPrimaryID, err == nil)
}()
}

// completeElection finishes an initial election. It is deliberately NOT
// completeFailover: that arms the failover cooldown unconditionally, and an
// election happening at boot would then swallow a genuine writer failure for
// the whole cooldown period. An election is not a failover and nothing was
// lost, so there is nothing to back off from.
// It deliberately does not invoke onFailoverComplete either: that callback
// announces a failover, and no failover happened.
func (m *WriterFailoverManager) completeElection(newPrimaryID string, success bool) {
m.mu.Lock()
m.failoverInProg = false
if success {
m.primaryID = newPrimaryID
m.consecutiveFails = 0
}
m.mu.Unlock()
}

// triggerFailoverLocked initiates failover (must hold lock).
func (m *WriterFailoverManager) triggerFailoverLocked() {
// Check cooldown
Expand All @@ -230,17 +298,21 @@ func (m *WriterFailoverManager) triggerFailoverLocked() {
m.wg.Add(1)
go func() {
defer m.wg.Done()
m.executeFailover(oldPrimary)
m.executeFailover(oldPrimary, true)
}()
}

// executeFailover performs the actual failover operation.
func (m *WriterFailoverManager) executeFailover(oldPrimaryID string) {
func (m *WriterFailoverManager) executeFailover(oldPrimaryID string, allowSelf bool) {
ctx, cancel := context.WithTimeout(m.ctx, m.cfg.FailoverTimeout)
defer cancel()

// Select best standby
newPrimaryID := m.selectNewPrimary(oldPrimaryID)
// allowSelf: the old primary is still a candidate when it is healthy and
// alone, which is how a writer that lost its designation to a restart gets
// it back instead of deadlocking. An operator-requested failover never
// takes that path — being asked to move off a node and then landing back
// on it is not a failover.
newPrimaryID := m.selectPrimary(oldPrimaryID, allowSelf)
if newPrimaryID == "" {
m.logger.Error().
Str("old_primary", oldPrimaryID).
Expand Down Expand Up @@ -291,6 +363,22 @@ func (m *WriterFailoverManager) executeFailover(oldPrimaryID string) {
// selectNewPrimary picks the best standby writer to promote.
// Prefers standby writers; falls back to any healthy writer excluding the failed primary.
func (m *WriterFailoverManager) selectNewPrimary(excludeNodeID string) string {
return m.selectPrimary(excludeNodeID, false)
}

// selectPrimary picks a writer to promote. It excludes excludeNodeID, the
// primary being failed away from, unless allowSelf is set and no other
// candidate exists.
//
// allowSelf exists because "no primary" does not always mean the primary
// died. A writer that merely restarted re-joins with a payload that carries
// no writer state, which clears its designation; the manager then looks for
// someone to fail over TO and, in the single-writer topology this feature
// documents, finds nobody, because the only candidate is the node it just
// excluded. The cluster would sit with no primary forever. GetWriters already
// filters to healthy nodes, so re-selecting the excluded node can never
// resurrect a dead one.
func (m *WriterFailoverManager) selectPrimary(excludeNodeID string, allowSelf bool) string {
// GetStandbyWriters and GetWriters already filter for healthy nodes
for _, node := range m.cfg.Registry.GetStandbyWriters() {
if node.ID != excludeNodeID {
Expand All @@ -302,6 +390,13 @@ func (m *WriterFailoverManager) selectNewPrimary(excludeNodeID string) string {
return node.ID
}
}
if allowSelf && excludeNodeID != "" {
for _, node := range m.cfg.Registry.GetWriters() {
if node.ID == excludeNodeID {
return node.ID
}
}
}
return ""
}

Expand Down Expand Up @@ -345,8 +440,11 @@ func (m *WriterFailoverManager) TriggerManualFailover() error {
Str("current_primary", primary.ID).
Msg("Manual writer failover initiated")

// Tracked on the WaitGroup like the other two, so Stop joins it.
m.wg.Add(1)
go func() {
m.executeFailover(primary.ID)
defer m.wg.Done()
m.executeFailover(primary.ID, false)
}()

return nil
Expand Down
Loading
Loading