Skip to content

Add lease-based leader election - #51

Merged
alex-hunt-materialize merged 3 commits into
mainfrom
leader-election
Jul 28, 2026
Merged

alex-hunt-materialize merged 3 commits into
mainfrom
leader-election

Conversation

@alex-hunt-materialize

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

Copy link
Copy Markdown
Contributor

Motivation

Rolling out controller updates causes short webhook downtime when running a single replica. Running multiple replicas behind a PodDisruptionBudget fixes that, but requires that only one replica reconciles at a time.

Changes

Adds LeaderElection, following the client-go leaderelection protocol on a coordination.k8s.io/v1 Lease:

  • The entry point is LeaderElection::with_lease(fut), which waits to acquire the lease, runs fut while renewing the lease in the background, and returns None when leadership is lost (dropping fut, which cancels a controller's in-flight reconciliations). A single lease can guard several controllers by passing a future that runs all of them, e.g. with_lease(future::join(a.run(), b.run())), so related controllers can't end up scattered across replicas. On loss, callers can either exit promptly and let Kubernetes restart the process, or rejoin the election in a loop (exiting is preferable if reconcilers spawn tasks or do blocking work, which dropping fut can't cancel).
  • Expiry is judged by observing the lease go unchanged for lease_duration on the candidate's own clock, so it is robust to clock skew between candidates. Renewal deadlines are measured from just before each renew request is sent (not from its response), so a slow response can't extend the leader's self-judged deadline past another candidate's takeover time. K8S leases are still best-effort, so controllers must still be robust to simultaneous reconciliations.
  • Writes use replace with resourceVersion, so takeover races lose with a conflict.
  • The leader renews on a tokio::time::Interval, so a slow attempt doesn't eat into the renew deadline budget between attempts.
  • LeaderElection::release hands leadership over immediately during graceful shutdown, and with_lease releases automatically if fut completes on its own.
  • Timings default to the client-go defaults (15s lease duration, 10s renew deadline, 2s retry period) and are configurable via builder methods, which panic on inconsistent values.

The controller's service account needs get/create/update on leases in coordination.k8s.io.

Also bumps the K8S version to 1.34, which is the oldest still supported version.

Testing

Exercised end-to-end by the orchestratord failover test in MaterializeInc/materialize#37806, which deploys 2 operator replicas with all three controllers guarded by a single lease, kills the leader pod, does a rolling restart of the operator deployment, and verifies lease takeover and continued reconciliation.

The test was run locally, pointing at this unpublished k8s-controller lib. The tests in MaterializeInc/materialize#37806 will fail until this crate is published and bumped in that PR.

🤖 Generated with Claude Code

Rolling out controller updates causes short webhook downtime with a
single replica. Running multiple replicas behind a PDB fixes that, but
requires that only one replica reconciles at a time.

LeaderElection follows the client-go leaderelection protocol on a
coordination.k8s.io/v1 Lease: expiry is judged by observing the lease
go unchanged for lease_duration on the candidate's own clock (clock-skew
safe), and writes use replace with resourceVersion so takeover races
lose with a conflict. Controller::run_with_leader_election returns when
leadership is lost; callers should exit and let Kubernetes restart the
process.

Requires get/create/update on leases in coordination.k8s.io.
@alex-hunt-materialize
alex-hunt-materialize marked this pull request as ready for review July 23, 2026 13:52
Comment thread src/controller.rs Outdated
/// leadership is lost (because the lease could not be renewed in time,
/// or was taken over by another candidate), the controller is stopped
/// and this method returns. Unlike [`run`](Self::run), this method
/// *does* return, and when it does, the process should be exited

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

is stopping the process actually necessary here? wouldn't dropping the self.run() future be sufficient to stop reconciliation? i'm not really sure i understand why something like loop { Controller::cluster(client.clone(), MyContext::new(), Default::default()).run_with_leader_election().await } would be unsafe.

if this is because this is what the go library recommends, i think we should try to do better here - idiomatic go involves a lot of spawning untrackable and uncancelable background tasks all over the place in a way that makes it difficult to be sure that everything is cleaned up properly without just restarting the process, but rust is a lot more explicit about when a future or task is running and encourages keeping track of that, and we shouldn't need to just throw our hands up here.

Comment thread src/controller.rs Outdated
leader_election.acquire().await;
let run = self.run();
let lost = leader_election.hold();
futures::pin_mut!(run, lost);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

nit: let run = pin!(self.run()); etc using std::pin::pin is generally preferable here (futures::pin_mut! is ~deprecated)

Comment thread src/controller.rs Outdated
// `run` never completes, so this returns only when leadership is
// lost; dropping `run` stops the controller
futures::future::select(run, lost).await;
warn!("lost leadership lease; the controller has been stopped");

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

we should indicate which controller this is referring to here, since a process may have multiple controllers (via Ctx::FINALIZER_NAME and Ctx::Resource::kind(Default::default()))

Comment thread src/leader_election.rs

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

this is reasonable for now, but it would really make more sense long term to upstream this to kube-runtime

Comment thread src/controller.rs Outdated
/// controller, then call [`LeaderElection::release`] on a clone of
/// `leader_election` to hand leadership over immediately rather than
/// making the other replicas wait for the lease to expire.
pub async fn run_with_leader_election(self, leader_election: LeaderElection) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

i wonder if we actually want to provide a more generic with_lease function on LeaderElection itself which implements the acquire/hold logic here for a generic future? then this function can just be implemented as leader_election.with_lease(self.run()).await, but it'd be easier to also do something like

leader_election.with_lease(
    [
        mz_ore::task::spawn(|| "foo", controller_foo.run()).abort_on_drop(),
        mz_ore::task::spawn(|| "bar", controller_bar.run()).abort_on_drop(),
    ].collect::<FuturesUnordered<_>>()
)

or whatever to allow a single lease to cover multiple related controllers (which i think is probably a pattern that will be more common in a lot of cases - allowing for the possibility of a process that runs multiple controllers where some are handled by one pod and others are handled by the other pod feels pretty hard to reason about)

Comment thread src/leader_election.rs Outdated
///
/// The `identity` must be non-empty and unique to each running instance of
/// the controller; the pod name (available in the `HOSTNAME` environment
/// variable, or via the downward API) is a good choice.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

the pod name is only a good choice if the process only runs a single controller, but i'm not sure that that will be the common case - maybe worth at least mentioning that caveat

Comment thread src/leader_election.rs Outdated
}

fn lease_duration_seconds(&self) -> i32 {
i32::try_from(self.lease_duration.as_secs()).unwrap_or(i32::MAX)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

this feels more reasonable to just panic on if it fails to parse rather than use an unexpected default

Comment thread src/leader_election.rs Outdated
pub(crate) async fn hold(&self) {
let mut last_renew = Instant::now();
loop {
tokio::time::sleep(self.retry_period).await;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

this would probably be better expressed as a tokio::time::Interval which triggers at a consistent rate regardless of how long the rest of the loop takes (for example, if one of the api calls happens to take 7 seconds due to some network issue or whatever, the current implementation will result in a timeout unnecessarily)

Comment thread src/leader_election.rs Outdated
pub(crate) async fn hold(&self) {
let mut last_renew = Instant::now();
loop {
tokio::time::sleep(self.retry_period).await;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

also, i wonder if it'd make more sense to use a watcher on the Lease instead of polling? that feels like it might cut down on the amount of time it's possible to end up with two different controller instances running (like, if the renew deadline expires and another instance grabs the lease, but this one is stuck waiting for the sleep to finish before it will cancel its own reconciliations)

Comment thread src/lib.rs
//! If you run multiple replicas of your controller (for instance, to avoid
//! downtime of webhooks served by the same process during rollouts), you
//! can use [leader election](LeaderElection) to ensure that only one
//! replica reconciles at a time:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

this is useful, but it'd probably also be good to give an example of how to implement it if you want the LeaderElection to cover multiple different controllers

- replace Controller::run_with_leader_election with
  LeaderElection::with_lease, which runs any future while holding the
  lease so a single lease can guard several controllers
- measure renew deadlines from before each request is sent: another
  candidate starts its takeover clock when it observes the written
  lease, which can happen well before the response arrives, so
  measuring from the response could extend the leader's deadline past
  the takeover time
- renew on a tokio Interval so slow attempts don't eat into the renew
  deadline budget between attempts
- panic on lease_duration > i32::MAX seconds instead of silently
  clamping
- release the lease when the guarded future completes on its own
- document rejoining the election in a loop as an alternative to
  exiting the process, and how to guard multiple controllers

Refs #51

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@alex-hunt-materialize
alex-hunt-materialize marked this pull request as draft July 27, 2026 12:47
@alex-hunt-materialize

Copy link
Copy Markdown
Contributor Author

@doy-materialize I think I've addressed all your comments except for the watcher suggestion. That added several hundred lines of complexity, and didn't remove all cases where we needed to poll (ie: when lease expires without any events). If you feel really strongly, I can add it, but it just seemed to make the code harder to read.

@alex-hunt-materialize
alex-hunt-materialize marked this pull request as ready for review July 27, 2026 15:23
@alex-hunt-materialize
alex-hunt-materialize merged commit 8f2ba15 into main Jul 28, 2026
2 checks passed
@alex-hunt-materialize
alex-hunt-materialize deleted the leader-election branch July 28, 2026 08:49
alex-hunt-materialize added a commit to MaterializeInc/materialize that referenced this pull request Aug 3, 2026
…#37806)

### 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
(MaterializeInc/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](https://claude.com/claude-code)

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
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