diff --git a/kubernetes/base/leaderelection/leaderelection.py b/kubernetes/base/leaderelection/leaderelection.py index ea20f0a570..fc72a1d95e 100644 --- a/kubernetes/base/leaderelection/leaderelection.py +++ b/kubernetes/base/leaderelection/leaderelection.py @@ -55,13 +55,24 @@ def run(self): logger.info("{} successfully acquired lease".format(self.election_config.lock.identity)) # Start leading and call OnStartedLeading() - threading.Thread(target=self.election_config.onstarted_leading, daemon=True).start() + leading_failed = threading.Event() + threading.Thread(target=self.run_onstarted_leading, args=(leading_failed,), daemon=True).start() - self.renew_loop() + self.renew_loop(leading_failed) # Failed to update lease, run OnStoppedLeading callback self.election_config.onstopped_leading() + def run_onstarted_leading(self, leading_failed): + # Run the callback in this thread, recording whether it raised. Without + # this the exception is swallowed by the worker thread and the lease keeps + # being renewed even though the work it protects is no longer running. + try: + self.election_config.onstarted_leading() + except Exception: + logger.exception("onstarted_leading raised an exception, stopping leading") + leading_failed.set() + def acquire(self): # Follower logger.info("{} is a follower".format(self.election_config.lock.identity)) @@ -75,7 +86,7 @@ def acquire(self): time.sleep(retry_period) - def renew_loop(self): + def renew_loop(self, leading_failed=None): # Leader logger.info("Leader has entered renew loop and will try to update lease continuously") @@ -83,6 +94,11 @@ def renew_loop(self): renew_deadline = self.election_config.renew_deadline * 1000 while True: + # onstarted_leading raised, so stop renewing a lease that no longer + # protects anything and let run() call onstopped_leading. + if leading_failed is not None and leading_failed.is_set(): + return + timeout = int(time.time() * 1000) + renew_deadline succeeded = False diff --git a/kubernetes/base/leaderelection/leaderelection_test.py b/kubernetes/base/leaderelection/leaderelection_test.py index 0cbcf00c88..ad9c7e7d1e 100644 --- a/kubernetes/base/leaderelection/leaderelection_test.py +++ b/kubernetes/base/leaderelection/leaderelection_test.py @@ -20,6 +20,7 @@ import unittest import threading import json +import sys import time import pytest from unittest.mock import patch @@ -222,6 +223,32 @@ def record_thread(*args, **kwargs): self.assertIn("daemon", captured) self.assertTrue(captured["daemon"]) + """Expected behavior: if onstarted_leading raises, the lease it protects is no + longer backed by any running work, so the candidate must stop renewing it and + run onstopped_leading. The lock below never refuses a renewal, so the renew + loop can only end because the callback failed.""" + def test_stops_leading_when_onstarted_leading_raises(self): + stopped = threading.Event() + + mock_lock = MockResourceLock("mock", "mock_namespace", "mock", thread_lock, + lambda: None, lambda: None, lambda: None, None) + mock_lock.renew_count_max = sys.maxsize + + def on_started_leading(): + raise RuntimeError("onstarted_leading failed") + + config = electionconfig.Config(lock=mock_lock, lease_duration=2, + renew_deadline=1.5, retry_period=1.1, + onstarted_leading=on_started_leading, + onstopped_leading=stopped.set) + + # Run in a daemon thread so a regression times out instead of hanging. + threading.Thread(target=leaderelection.LeaderElection(config).run, + daemon=True).start() + + self.assertTrue(stopped.wait(10), + "onstopped_leading was not called after onstarted_leading raised") + def assert_history(self, history, expected): self.assertIsNotNone(expected) self.assertIsNotNone(history)