Skip to content

orchestratord: use leader election and run multiple operator replicas - #37806

Merged
alex-hunt-materialize merged 2 commits into
MaterializeInc:mainfrom
alex-hunt-materialize:orchestratord-leader-election
Aug 3, 2026
Merged

alex-hunt-materialize merged 2 commits into
MaterializeInc:mainfrom
alex-hunt-materialize:orchestratord-leader-election

Conversation

@alex-hunt-materialize

@alex-hunt-materialize alex-hunt-materialize commented Jul 22, 2026 •

Copy link
Copy Markdown
Contributor

Motivation

Rolling out orchestratord updates causes short downtime of the CRD conversion webhook when running a single replica. Running multiple replicas fixes that, but requires that only one replica reconciles at a time.

Changes

  • Run orchestratord's controllers (materialize, balancer, console) under k8s-controller's new lease-based leader election (Add lease-based leader election k8s-controller#51, released in 0.12.0). A single coordination.k8s.io/v1 Lease in the operator's namespace guards all of the controllers, so that they can't end up scattered across replicas. Standby replicas keep serving the conversion webhook while waiting on the lease. The election identity must be unique per replica, otherwise replicas mistake each other's lease renewals for their own and all act as leader at once. It comes from the new --leader-election-identity, which the chart sets from the pod name via the downward API (as $ORCHESTRATORD_LEADER_ELECTION_IDENTITY), falling back to $HOSTNAME and then to a random identity so that a replica is still uniquely identified without the chart's help.
  • Losing the lease cancels in-flight reconciliation and the replica rejoins the election in process rather than restarting, so it keeps serving the conversion webhook throughout. Each controller, and the node upgrade watcher, runs as its own task, and what the lease guards is their abort-on-drop handles, so that losing it aborts all of them. Running them as tasks rather than joining their futures keeps them independently scheduled, so that none of them can hold up the poll of the others, or of the lease renewal that with_lease polls alongside them. A stalled renewal is the worse case, since it can miss the renew deadline and cost this replica its leadership. An aborted task runs until its next await point, so a reconciliation can still finish a request it had already issued. That is within the margin the lease timings leave, since the renew deadline is shorter than the lease duration.
  • The GCP node upgrade watcher moves under the lease too. It triggers rollouts by writing the materialize.cloud/force-rollout annotation, so running it on every replica would let one node pool upgrade trigger several rollouts of the same instance, each replica writing its own annotation value. Its in-memory dedup map doesn't coordinate across replicas, and its "a rollout is already in progress" check is a read followed by a write, so that doesn't serialize them either. Its Pub/Sub subscriber and its scan loop are spawned as tasks whose handles abort on drop, so that neither outlives the lease. A subscriber that did would keep pulling and acking notifications out from under the next holder. Its configuration is still validated at startup, so a misconfiguration fails immediately rather than only once the replica wins the election.
  • On SIGTERM the leader stops reconciling, releases the lease so that a standby takes leadership over immediately instead of waiting for the lease to expire, drains the conversion webhook server, and exits. The webhook keeps accepting requests for a few seconds before draining, since Kubernetes removes the pod from the webhook service's endpoints asynchronously.
  • Every step of that shutdown sequence is individually bounded, the lease release included. The kube client's default read timeout is far longer than the termination grace period, so an unreachable API server could otherwise consume the whole grace period in release() and leave no time to drain. The three steps add up to 15 seconds, and the chart sets terminationGracePeriodSeconds: 30 explicitly so that the budget the binary assumes cannot drift away from what Kubernetes grants it.
  • The SIGTERM handler is installed before the rest of startup, rather than alongside the controllers that act on it. The webhook server starts early and the readiness probe only checks that server, so the pod can be in the webhook service's endpoints while the rest of startup is still running, and with no handler installed the default action for SIGTERM would kill the process there and drop the conversion requests in flight. Tokio records a signal received before the first recv(), so one arriving during startup is still handled once the shutdown path is reached. This does not abort startup, so a SIGTERM during CRD registration is only acted on once startup finishes.
  • Helm chart: default to 2 operator replicas, spread them across nodes with a soft pod anti-affinity, create a PodDisruptionBudget tolerating one unavailable replica when running more than one (operator.podDisruptionBudget values), and grant the operator get/create/update on leases in coordination.k8s.io.
  • New orchestratord_is_leader gauge, set when a replica takes the lease and cleared when it stops holding it. Summed across the replicas it should be 1, so a sustained 0 means no replica can take the lease and the operator is reconciling nothing. That was previously visible only as a log line, which matters because the readiness probe checks the conversion webhook, so a replica that cannot reconcile at all still looks healthy. The likeliest cause is a service account without permission on leases, which is what an installation managing its own RBAC (rbac.create: false) ends up with after upgrading to this chart.
  • environmentd_needs_update is reset when leadership is lost. The gauge is derived from what reconciliation observed and lives as long as the process, which now outlives its own leadership, since a replica keeps serving the conversion webhook after losing the lease. Left set, a former leader would publish its last observation forever, and summing the gauge across the operator's replicas would count the same organizations once per past leader.

Testing

  • New leader-failover workflow in test/orchestratord/mzcompose.py, enabled in the Nightly pipeline. It verifies that the controller lease is held by an operator pod, that leadership fails over when the leader pod is deleted, after a rolling restart of the operator deployment, and when another candidate takes the lease over (also asserting that no operator pod restarted, i.e. that the replica rejoined the election in process), and that reconciliation completes after each failover. The candidate case additionally requires the lease's transition count to grow, because the incumbent leader is still running and eligible there, so it would otherwise satisfy "an operator pod holds the lease" without ever having given leadership up. It also checks that exactly the lease holder reports itself leader through orchestratord_is_leader (scraped through the API server's pod proxy, since the operator image has nothing to exec into), both when leadership is first established and after the candidate case, which is the one where the replica that gave leadership up is still running and so has to stop reporting itself leader rather than being replaced by a fresh pod.
  • get_orchestratord_data now returns only the pods of the operator deployment's current revision, found by matching the deployment's revision annotation against the ReplicaSet that carries it. A rollout leaves pods of the previous generation running until their replacements are up, so the existing callers that inspect a single pod's spec could otherwise read the outgoing pod. It rejects an empty result, since that filtering can transiently leave nothing behind (a new ReplicaSet carries the deployment's revision before its pods exist) and its callers index into the pods. The orchestratord upgrade test also no longer pins an exact operator pod count, since the chart versions it upgrades from install a single replica.
  • New helm-unittest coverage for the PodDisruptionBudget template, the leader election identity environment variable, the default pod anti-affinity and its interaction with a user-configured affinity, and the termination grace period.

🤖 Generated with Claude Code

@alex-hunt-materialize
alex-hunt-materialize force-pushed the orchestratord-leader-election branch 4 times, most recently from 1a9b8df to 365f756 Compare July 28, 2026 08:59
@alex-hunt-materialize
alex-hunt-materialize marked this pull request as ready for review July 28, 2026 11:10
@alex-hunt-materialize
alex-hunt-materialize requested a review from a team as a code owner July 28, 2026 11:10
@alex-hunt-materialize
alex-hunt-materialize marked this pull request as draft July 28, 2026 11:21
@alex-hunt-materialize
alex-hunt-materialize force-pushed the orchestratord-leader-election branch 10 times, most recently from 59f4af5 to c209da4 Compare July 29, 2026 11:56
Run orchestratord's controllers (materialize, balancer, console) under
k8s-controller's lease-based leader election, with a single lease
guarding all of them, so that multiple replicas of orchestratord can
run while only one replica reconciles at a time. The other replicas
keep serving the conversion webhook, which avoids webhook downtime
during operator rollouts and node drains.

The leader election identity must be unique per replica, otherwise
replicas mistake each other's lease renewals for their own and all act
as leader at once. It comes from --leader-election-identity, which the
chart sets from the pod name via the downward API, and falls back to
$HOSTNAME so that a replica is still uniquely identified without the
chart's help.

Losing the lease drops the controllers, cancelling their in-flight
reconciliations, and the replica then rejoins the election in process
rather than restarting. Reconciliation is polled inline rather than
spawned, so none of it outlives the lease, and staying up keeps the
replica serving the conversion webhook.

The GCP node upgrade watcher moves under the lease as well. It triggers
rollouts by writing the force-rollout annotation, so running it on every
replica would let a single node pool upgrade trigger several rollouts of
the same instance, each replica writing its own annotation value. Its
in-memory dedup map does not coordinate across replicas, and the
"rollout already in progress" check it makes is a read followed by a
write, so it does not serialize them either. Its Pub/Sub subscriber is
polled inline rather than spawned for the same reason the controllers
are: a subscriber that outlived the lease would keep pulling and acking
notifications out from under the next holder. Its configuration is still
validated at startup, so a misconfiguration fails immediately rather
than only once the replica wins the election.

On SIGTERM the leader stops reconciling, releases the lease so that a
standby takes over immediately instead of waiting for the lease to
expire, drains the conversion webhook server, and exits. The webhook
keeps accepting requests for a few seconds before draining, since
Kubernetes removes the pod from the webhook service's endpoints
asynchronously. Every step of that sequence is bounded, the lease
release included, because its API call could otherwise consume the
whole termination grace period and leave no time to drain. The chart
sets terminationGracePeriodSeconds explicitly so that the budget the
binary assumes cannot drift away from what Kubernetes grants it.

The SIGTERM handler is installed before the rest of startup, rather than
alongside the controllers that act on it. The webhook server starts
early and the readiness probe only checks that server, so the pod can be
in the webhook service's endpoints while the rest of startup is still
running. With no handler installed yet, the default action for SIGTERM
would kill the process there and drop the conversion requests in flight.
Tokio records a signal received before the first recv(), so one arriving
during startup is still handled once the shutdown path is reached. Note
that this does not abort startup, so a SIGTERM during CRD registration
is only acted on once startup finishes.

Add an orchestratord_is_leader gauge, set when this replica takes the
lease and cleared when it stops holding it. Summed across the replicas
it should be 1, so a sustained 0 means no replica can take the lease and
the operator is reconciling nothing. That was previously visible only as
a log line, which matters because the readiness probe checks the
conversion webhook, so a replica that cannot reconcile at all still
looks healthy. The most likely cause is a service account without
permission on leases, which is what an installation managing its own
RBAC ends up with after upgrading.

Reset environmentd_needs_update when leadership is lost. The gauge is
derived from what reconciliation observed and lives for the lifetime of
the process, which now outlives its own leadership, since a replica
keeps serving the conversion webhook after losing the lease. Left set,
a former leader would publish its last observation forever, so summing
the gauge across the replicas would count the same organizations once
per past leader.

Helm chart changes: default to 2 operator replicas, spread them across
nodes with a soft pod anti-affinity, create a PodDisruptionBudget
tolerating one unavailable replica when running more than one, and
grant the operator get/create/update on coordination.k8s.io leases.
The chart's unit tests cover the new budget, the anti-affinity and its
interaction with a user-configured affinity, the leader election
identity, and the termination grace period.

Add a leader-failover workflow to test/orchestratord that verifies the
lease is held by an operator pod, that leadership fails over when the
leader pod is deleted, after a rolling restart of the operator
deployment, and when another candidate takes the lease, and that
reconciliation completes after each. Enable it in the Nightly pipeline.
The candidate case requires the lease's transition count to grow,
because the incumbent is still running and eligible there and would
otherwise satisfy the check without ever having given leadership up.

It also checks that exactly the lease holder reports itself leader
through orchestratord_is_leader, scraped through the API server's pod
proxy since the operator image has nothing to exec into. Once when
leadership is first established, and again after the candidate case,
which is the one where the replica that gave leadership up is still
running and so has to stop reporting itself leader rather than being
replaced by a fresh pod.

The upgrade test asserted that exactly one operator pod exists, which
breaks now that the chart defaults to two replicas. The released chart
versions it upgrades from still install a single replica, so require
that every operator pod runs the expected image instead of pinning a
pod count.

get_orchestratord_data now returns only the pods of the operator
deployment's current revision, found by matching the deployment's
revision annotation against the ReplicaSet that carries it. A rollout
leaves pods of the previous generation running until their replacements
are up, so the callers that inspect a single pod's spec could otherwise
read the outgoing pod. It rejects an empty result, since that filtering
can transiently leave nothing behind and its callers index into the
pods.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
@alex-hunt-materialize
alex-hunt-materialize force-pushed the orchestratord-leader-election branch from c209da4 to d129fa7 Compare July 29, 2026 13:06
@alex-hunt-materialize
alex-hunt-materialize marked this pull request as ready for review July 29, 2026 14:38
@alex-hunt-materialize
alex-hunt-materialize requested a review from a team as a code owner July 29, 2026 14:38

@doy-materialize doy-materialize left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

mostly fine, just a few small questions

// this future stops it too. A subscriber that outlived it would keep
// pulling notifications, and acking them, out from under whichever
// replica holds the lease next.
future::join(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

i think this isn't the ideal pattern for these kind of long running loops - just using future::join on raw futures allows one future to starve out all of the others if it behaves badly. i think we should continue to spawn the components as tasks, but use future::join(mz_ore::task::spawn(...).abort_on_drop(), ...) instead, which keeps them as tasks which can be independently scheduled while still blocking the owning task until they exit.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done

let leader_election_identity = args
.leader_election_identity
.or_else(|| env::var("HOSTNAME").ok())
.unwrap_or_else(|| format!("orchestratord-{}", uuid::Uuid::new_v4()));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

if just defaulting to a random value works for this case, why do we go through the effort of passing the pod name/hostname/etc through? can we just always use Uuid::new_v4() here?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is nice to be able to identify which pod is the leader. I don't feel strongly about this, though, so if you think it isn't worth it we can drop all that.

// both them and the metric to the moment leadership is acquired.
let controllers = Box::pin(leader_election.with_lease(async {
metrics.leadership_acquired();
join4(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

similarly here, i think we should continue to use separate tasks

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done

The controllers and the node upgrade watcher are joined rather than
spawned, so that losing the leadership lease drops them. Joining raw
futures puts all of them, and the lease renewal that `with_lease` polls
alongside them, on a single task, where one of them blocking its poll
stalls the rest. A stalled renewal is the worse case, since it can miss
the renew deadline and cost this replica its leadership.

Spawn each of them as a task instead, and join the abort-on-drop
handles, which keeps them independently scheduled while still stopping
all of them when the lease is lost.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
@alex-hunt-materialize
alex-hunt-materialize merged commit 8dad2db into MaterializeInc:main Aug 3, 2026
160 checks passed
@alex-hunt-materialize
alex-hunt-materialize deleted the orchestratord-leader-election branch August 3, 2026 12:59
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants