Skip to content

Commit 6073f7f

Browse files
chala2001rvlane
andcommitted
Add LeaseLock for leader election
Leader election could only use a ConfigMap. client-go moved to the coordination.k8s.io Lease resource and dropped ConfigMap support, so add a LeaseLock next to ConfigMapLock and use it in the example. Based on the implementation from #2314, with four corrections: use the package logger instead of calling logging.basicConfig at import, accept the times str(datetime) produces when the microseconds are zero, keep the created lease so a following update has a reference, and store the times as the real UTC instant rather than local wall clock labelled as UTC. Co-authored-by: Lane Richard <rick.lane@nokia.com>
1 parent d4b497b commit 6073f7f

3 files changed

Lines changed: 370 additions & 1 deletion

File tree

‎kubernetes/base/leaderelection/example.py‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
from kubernetes import client, config
1717
from kubernetes.leaderelection import leaderelection
1818
from kubernetes.leaderelection.resourcelock.configmaplock import ConfigMapLock
19+
from kubernetes.leaderelection.resourcelock.leaselock import LeaseLock
1920
from kubernetes.leaderelection import electionconfig
2021

2122

@@ -42,8 +43,13 @@ def example_func():
4243
# A user can choose not to provide any callbacks for what to do when a candidate fails to lead - onStoppedLeading()
4344
# In that case, a default callback function will be used
4445

46+
lock = LeaseLock(lock_name, lock_namespace, candidate_id)
47+
# Choose the lock to use. LeaseLock uses the coordination.k8s.io Lease
48+
# resource and is the recommended one; ConfigMapLock is kept for existing users.
49+
# lock = ConfigMapLock(lock_name, lock_namespace, candidate_id)
50+
4551
# Create config
46-
config = electionconfig.Config(ConfigMapLock(lock_name, lock_namespace, candidate_id), lease_duration=17,
52+
config = electionconfig.Config(lock, lease_duration=17,
4753
renew_deadline=15, retry_period=5, onstarted_leading=example_func,
4854
onstopped_leading=None)
4955

Lines changed: 155 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,155 @@
1+
# Copyright 2021 The Kubernetes Authors.
2+
#
3+
# Licensed under the Apache License, Version 2.0 (the "License");
4+
# you may not use this file except in compliance with the License.
5+
# You may obtain a copy of the License at
6+
#
7+
# http://www.apache.org/licenses/LICENSE-2.0
8+
#
9+
# Unless required by applicable law or agreed to in writing, software
10+
# distributed under the License is distributed on an "AS IS" BASIS,
11+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
# See the License for the specific language governing permissions and
13+
# limitations under the License.
14+
15+
from kubernetes.client.rest import ApiException
16+
from kubernetes import client
17+
from ..leaderelectionrecord import LeaderElectionRecord
18+
from datetime import datetime, timezone
19+
import logging
20+
logger = logging.getLogger("leaderelection")
21+
22+
# Formats produced by str(datetime). The microsecond component is omitted
23+
# when it is exactly zero, so both spellings have to be accepted.
24+
TIME_FORMATS = (
25+
"%Y-%m-%d %H:%M:%S.%f%z",
26+
"%Y-%m-%d %H:%M:%S.%f",
27+
"%Y-%m-%d %H:%M:%S%z",
28+
"%Y-%m-%d %H:%M:%S",
29+
)
30+
31+
32+
class LeaseLock:
33+
def __init__(self, name, namespace, identity):
34+
"""
35+
:param name: name of the lock
36+
:param namespace: namespace
37+
:param identity: A unique identifier that the candidate is using
38+
"""
39+
self.api_instance = client.CoordinationV1Api()
40+
self.name = name
41+
self.namespace = namespace
42+
self.identity = str(identity)
43+
self.lease_reference = None
44+
45+
# get returns the election record from a Lease spec
46+
def get(self, name, namespace):
47+
"""
48+
:param name: Name of the lease object information to get
49+
:param namespace: Namespace in which the lease object is to be searched
50+
:return: 'True, election record' if object found else 'False, exception response'
51+
"""
52+
try:
53+
lease = self.api_instance.read_namespaced_lease(name, namespace)
54+
except ApiException as e:
55+
return False, e
56+
57+
self.lease_reference = lease
58+
return True, self.get_lock_object(lease)
59+
60+
def create(self, name, namespace, election_record):
61+
"""
62+
:param name: Name of the lease object to be created
63+
:param namespace: Namespace in which the lease object is to be created
64+
:param election_record: The election record to store in the lease spec
65+
:return: 'True' if object is created else 'False' if failed
66+
"""
67+
body = client.V1Lease(metadata={"name": name},
68+
spec=self.get_lease_spec(election_record))
69+
70+
try:
71+
# Keep the created lease so that a following update has a
72+
# reference to work from without re-reading it.
73+
self.lease_reference = self.api_instance.create_namespaced_lease(
74+
namespace, body)
75+
return True
76+
except ApiException as e:
77+
logger.info("Failed to create lock as {}".format(e))
78+
return False
79+
80+
def update(self, name, namespace, updated_record):
81+
"""
82+
:param name: name of the lock to be updated
83+
:param namespace: namespace the lock is in
84+
:param updated_record: the updated election record
85+
:return: True if update is successful False if it fails
86+
"""
87+
if self.lease_reference is None:
88+
logger.info("Lease not initialized, call get or create first")
89+
return False
90+
91+
try:
92+
self.lease_reference.spec = self.get_lease_spec(
93+
updated_record, self.lease_reference.spec)
94+
self.api_instance.replace_namespaced_lease(
95+
name=name, namespace=namespace, body=self.lease_reference)
96+
return True
97+
except ApiException as e:
98+
logger.info("Failed to update lock as {}".format(e))
99+
return False
100+
101+
def get_lease_spec(self, leader_election_record, current_spec=None):
102+
"""Build the lease spec that holds the given election record."""
103+
spec = current_spec if current_spec else client.V1LeaseSpec()
104+
105+
spec.holder_identity = leader_election_record.holder_identity
106+
spec.lease_duration_seconds = int(leader_election_record.lease_duration)
107+
spec.acquire_time = self.time_str_to_iso(
108+
leader_election_record.acquire_time)
109+
spec.renew_time = self.time_str_to_iso(
110+
leader_election_record.renew_time)
111+
112+
return spec
113+
114+
def get_lock_object(self, lease):
115+
"""Build the election record held in the given lease spec."""
116+
leader_election_record = LeaderElectionRecord(None, None, None, None)
117+
118+
if not lease.spec:
119+
return leader_election_record
120+
121+
if lease.spec.holder_identity:
122+
leader_election_record.holder_identity = lease.spec.holder_identity
123+
if lease.spec.lease_duration_seconds:
124+
leader_election_record.lease_duration = str(
125+
lease.spec.lease_duration_seconds)
126+
if lease.spec.acquire_time:
127+
leader_election_record.acquire_time = self.time_from_utc(
128+
lease.spec.acquire_time)
129+
if lease.spec.renew_time:
130+
leader_election_record.renew_time = self.time_from_utc(
131+
lease.spec.renew_time)
132+
133+
return leader_election_record
134+
135+
def time_str_to_iso(self, str_time):
136+
"""Convert an election record time into the instant to store.
137+
138+
``leaderelection.py`` builds its times with
139+
``datetime.fromtimestamp()``, which is local and naive. The Lease is
140+
shared with other clients, so the value has to go on the wire as the
141+
real UTC instant rather than as local wall clock labelled as UTC.
142+
"""
143+
for fmt in TIME_FORMATS:
144+
try:
145+
parsed = datetime.strptime(str_time, fmt)
146+
except ValueError:
147+
continue
148+
if parsed.tzinfo is None:
149+
parsed = parsed.astimezone()
150+
return parsed.astimezone(timezone.utc)
151+
raise ValueError("Failed to parse time string: {}".format(str_time))
152+
153+
def time_from_utc(self, value):
154+
"""Inverse of :meth:`time_str_to_iso`, back to the record format."""
155+
return str(value.astimezone().replace(tzinfo=None))
Lines changed: 208 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,208 @@
1+
# Copyright 2021 The Kubernetes Authors.
2+
#
3+
# Licensed under the Apache License, Version 2.0 (the "License");
4+
# you may not use this file except in compliance with the License.
5+
# You may obtain a copy of the License at
6+
#
7+
# http://www.apache.org/licenses/LICENSE-2.0
8+
#
9+
# Unless required by applicable law or agreed to in writing, software
10+
# distributed under the License is distributed on an "AS IS" BASIS,
11+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
# See the License for the specific language governing permissions and
13+
# limitations under the License.
14+
15+
import datetime
16+
import os
17+
import time
18+
import unittest
19+
from unittest import mock
20+
21+
from kubernetes import client
22+
from kubernetes.client.rest import ApiException
23+
24+
from ..leaderelectionrecord import LeaderElectionRecord
25+
from .leaselock import LeaseLock
26+
27+
28+
def make_lock():
29+
with mock.patch.object(client, 'CoordinationV1Api'):
30+
lock = LeaseLock('lock', 'default', 'candidate')
31+
lock.api_instance = mock.MagicMock()
32+
return lock
33+
34+
35+
UTC = datetime.timezone.utc
36+
37+
38+
class LeaseLockTest(unittest.TestCase):
39+
40+
def test_create_writes_the_election_record_to_the_spec(self):
41+
lock = make_lock()
42+
record = LeaderElectionRecord('candidate', '17',
43+
'2026-08-29 01:02:03.456789',
44+
'2026-08-29 01:02:03.456789')
45+
46+
self.assertTrue(lock.create('lock', 'default', record))
47+
48+
namespace, body = lock.api_instance.create_namespaced_lease.call_args[0]
49+
self.assertEqual('default', namespace)
50+
self.assertEqual('lock', body.metadata.name)
51+
self.assertEqual('candidate', body.spec.holder_identity)
52+
self.assertEqual(17, body.spec.lease_duration_seconds)
53+
# stored as UTC on the wire, denoting the same instant as the
54+
# local wall clock the election record carries
55+
self.assertEqual(UTC, body.spec.acquire_time.tzinfo)
56+
self.assertEqual('2026-08-29 01:02:03.456789',
57+
str(body.spec.acquire_time.astimezone()
58+
.replace(tzinfo=None)))
59+
self.assertEqual(body.spec.acquire_time, body.spec.renew_time)
60+
61+
def test_create_returns_false_when_the_api_fails(self):
62+
lock = make_lock()
63+
lock.api_instance.create_namespaced_lease.side_effect = ApiException(
64+
status=409, reason='Conflict')
65+
record = LeaderElectionRecord('candidate', '17', '2026-08-29 01:02:03',
66+
'2026-08-29 01:02:03')
67+
68+
self.assertFalse(lock.create('lock', 'default', record))
69+
70+
def test_get_returns_the_exception_when_the_lease_is_missing(self):
71+
lock = make_lock()
72+
expected = ApiException(status=404, reason='Not Found')
73+
lock.api_instance.read_namespaced_lease.side_effect = expected
74+
75+
status, response = lock.get('lock', 'default')
76+
77+
self.assertFalse(status)
78+
self.assertIs(expected, response)
79+
80+
def test_get_reads_the_election_record_from_the_lease(self):
81+
lock = make_lock()
82+
acquired = datetime.datetime(2026, 8, 29, 1, 2, 3, 456789)
83+
lock.api_instance.read_namespaced_lease.return_value = client.V1Lease(
84+
metadata={'name': 'lock'},
85+
spec=client.V1LeaseSpec(holder_identity='candidate',
86+
lease_duration_seconds=17,
87+
acquire_time=acquired,
88+
renew_time=acquired))
89+
90+
status, record = lock.get('lock', 'default')
91+
92+
self.assertTrue(status)
93+
self.assertEqual('candidate', record.holder_identity)
94+
self.assertEqual('17', record.lease_duration)
95+
self.assertEqual('2026-08-29 01:02:03.456789', record.acquire_time)
96+
97+
def test_get_on_a_lease_without_a_spec_returns_an_empty_record(self):
98+
lock = make_lock()
99+
lock.api_instance.read_namespaced_lease.return_value = client.V1Lease(
100+
metadata={'name': 'lock'}, spec=None)
101+
102+
status, record = lock.get('lock', 'default')
103+
104+
self.assertTrue(status)
105+
self.assertIsNone(record.holder_identity)
106+
107+
def test_record_survives_a_write_and_read_unchanged(self):
108+
"""leaderelection.py compares the stored record with the observed one
109+
using __dict__, so a round trip has to come back identical."""
110+
lock = make_lock()
111+
now = datetime.datetime.fromtimestamp(1787000000.5)
112+
record = LeaderElectionRecord('candidate', str(17), str(now), str(now))
113+
114+
spec = lock.get_lease_spec(record)
115+
read_back = lock.get_lock_object(
116+
client.V1Lease(metadata={'name': 'lock'}, spec=spec))
117+
118+
self.assertEqual(record.__dict__, read_back.__dict__)
119+
120+
def test_record_without_microseconds_survives_a_write_and_read(self):
121+
"""str(datetime) drops the microseconds when they are exactly zero."""
122+
lock = make_lock()
123+
now = datetime.datetime(2026, 8, 29, 1, 2, 3)
124+
self.assertEqual('2026-08-29 01:02:03', str(now))
125+
record = LeaderElectionRecord('candidate', str(17), str(now), str(now))
126+
127+
spec = lock.get_lease_spec(record)
128+
read_back = lock.get_lock_object(
129+
client.V1Lease(metadata={'name': 'lock'}, spec=spec))
130+
131+
self.assertEqual(record.__dict__, read_back.__dict__)
132+
133+
def test_update_replaces_the_lease_it_read(self):
134+
lock = make_lock()
135+
acquired = datetime.datetime(2026, 8, 29, 1, 2, 3, 456789)
136+
lock.api_instance.read_namespaced_lease.return_value = client.V1Lease(
137+
metadata={'name': 'lock'},
138+
spec=client.V1LeaseSpec(holder_identity='other',
139+
lease_duration_seconds=17,
140+
acquire_time=acquired,
141+
renew_time=acquired))
142+
lock.get('lock', 'default')
143+
144+
record = LeaderElectionRecord('candidate', '17',
145+
'2026-08-29 01:02:03.456789',
146+
'2026-08-29 02:03:04.567890')
147+
self.assertTrue(lock.update('lock', 'default', record))
148+
149+
body = lock.api_instance.replace_namespaced_lease.call_args[1]['body']
150+
self.assertEqual('candidate', body.spec.holder_identity)
151+
self.assertEqual('2026-08-29 02:03:04.567890',
152+
str(body.spec.renew_time.astimezone()
153+
.replace(tzinfo=None)))
154+
155+
def test_update_returns_false_when_the_api_fails(self):
156+
lock = make_lock()
157+
lock.lease_reference = client.V1Lease(metadata={'name': 'lock'},
158+
spec=client.V1LeaseSpec())
159+
lock.api_instance.replace_namespaced_lease.side_effect = ApiException(
160+
status=409, reason='Conflict')
161+
record = LeaderElectionRecord('candidate', '17', '2026-08-29 01:02:03',
162+
'2026-08-29 01:02:03')
163+
164+
self.assertFalse(lock.update('lock', 'default', record))
165+
166+
@unittest.skipUnless(hasattr(time, 'tzset'), 'requires tzset')
167+
def test_times_are_written_as_the_real_utc_instant(self):
168+
"""The election record holds local wall clock. The Lease is shared
169+
with other clients, so it has to carry the real UTC instant."""
170+
lock = make_lock()
171+
previous = os.environ.get('TZ')
172+
os.environ['TZ'] = 'Asia/Kolkata' # UTC+05:30, no DST
173+
time.tzset()
174+
try:
175+
record = LeaderElectionRecord('candidate', '17',
176+
'2026-08-29 01:02:03.456789',
177+
'2026-08-29 01:02:03.456789')
178+
spec = lock.get_lease_spec(record)
179+
180+
self.assertEqual(
181+
datetime.datetime(2026, 8, 28, 19, 32, 3, 456789, tzinfo=UTC),
182+
spec.acquire_time)
183+
# and it still round trips back to the local wall clock
184+
read_back = lock.get_lock_object(
185+
client.V1Lease(metadata={'name': 'lock'}, spec=spec))
186+
self.assertEqual(record.__dict__, read_back.__dict__)
187+
finally:
188+
if previous is None:
189+
os.environ.pop('TZ', None)
190+
else:
191+
os.environ['TZ'] = previous
192+
time.tzset()
193+
194+
def test_update_before_get_or_create_does_not_raise(self):
195+
lock = make_lock()
196+
record = LeaderElectionRecord('candidate', '17', '2026-08-29 01:02:03',
197+
'2026-08-29 01:02:03')
198+
199+
self.assertFalse(lock.update('lock', 'default', record))
200+
201+
def test_an_unparsable_time_is_reported(self):
202+
lock = make_lock()
203+
with self.assertRaises(ValueError):
204+
lock.time_str_to_iso('not a time')
205+
206+
207+
if __name__ == '__main__':
208+
unittest.main()

0 commit comments

Comments
 (0)