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
22 changes: 19 additions & 3 deletions kubernetes/base/leaderelection/leaderelection.py
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand All @@ -75,14 +86,19 @@ 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")

retry_period = self.election_config.retry_period
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

Expand Down
27 changes: 27 additions & 0 deletions kubernetes/base/leaderelection/leaderelection_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import unittest
import threading
import json
import sys
import time
import pytest
from unittest.mock import patch
Expand Down Expand Up @@ -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)
Expand Down
Loading