Skip to content

Commit c99cdb4

Browse files
committed
fix: prevent duplicate concurrent pods on unresponsive activation restart
When the monitor detects an unresponsive activation and triggers a restart, _watch_pod_deletion falsely reports success if the watch stream times out without receiving a DELETED event. This allows _unresponsive_policy to restart the activation while the original pod is still running, producing ghost pods that compete for shared resources like Kafka consumer groups. Two fixes: 1. _watch_pod_deletion now tracks whether a DELETED event was actually received. If the watch times out without one, it retries instead of reporting success. After all retries, it raises ContainerCleanupError. 2. _unresponsive_policy now calls _cleanup explicitly via a new _handle_unresponsive_activation helper. If cleanup fails, the restart is skipped to avoid creating a duplicate pod. The activation is still marked FAILED regardless of cleanup outcome. 3. adds a list_then_watch flow for deleting a pod so that we don't loose track of the pods we are trying to delete Fixes: AAP-81201
1 parent 51086aa commit c99cdb4

4 files changed

Lines changed: 307 additions & 36 deletions

File tree

src/aap_eda/services/activation/activation_manager.py

Lines changed: 30 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -377,45 +377,52 @@ def _unresponsive_policy(self, check_type: str):
377377
"Unresponsive policy called for "
378378
f"activation id: {self.db_instance.id}",
379379
)
380-
container_logger = self.container_logger_class(self.latest_instance.id)
381-
if self.db_instance.restart_policy == RestartPolicy.NEVER:
382-
LOGGER.info(
383-
f"Activation id: {self.db_instance.id} "
384-
f"Restart policy is set to {self.db_instance.restart_policy}."
385-
"No restart policy is applied."
386-
)
387-
user_msg = (
388-
"Activation is unresponsive. "
389-
f"{check_type} check for ansible-rulebook timed out. "
390-
"Restart policy is not applicable."
380+
try:
381+
self._handle_unresponsive_activation(check_type=check_type)
382+
except engine_exceptions.ContainerCleanupError as e:
383+
LOGGER.error(
384+
f"Error occurred while handling unresponsive activation: {e}"
391385
)
392-
container_logger.write(user_msg, flush=True)
393-
self._fail_instance(user_msg)
394-
self.set_status(
395-
ActivationStatus.FAILED,
396-
user_msg,
386+
return
387+
388+
if not self.db_instance.restart_policy == RestartPolicy.NEVER:
389+
system_restart_activation(
390+
self.db_instance_type, self.db_instance.id, delay_seconds=1
397391
)
398392

399-
else:
400-
LOGGER.info(
393+
def _handle_unresponsive_activation(self, check_type: str):
394+
try:
395+
self._cleanup()
396+
except engine_exceptions.ContainerCleanupError:
397+
raise
398+
finally:
399+
container_logger = self.container_logger_class(
400+
self.latest_instance.id
401+
)
402+
log_message = (
401403
f"Activation id: {self.db_instance.id} "
402404
f"Restart policy is set to {self.db_instance.restart_policy}."
403-
"Restart policy is applied.",
404405
)
405406
user_msg = (
406407
"Activation is unresponsive. "
407408
f"{check_type} check for ansible-rulebook timed out. "
408-
"Activation is going to be restarted."
409409
)
410+
411+
if self.db_instance.restart_policy == RestartPolicy.NEVER:
412+
log_message += " Restart policy is not applicable."
413+
user_msg += " Restart policy is not applicable."
414+
else:
415+
log_message += " Restart policy is applied."
416+
user_msg += " Activation is going to be restarted."
417+
418+
LOGGER.info(log_message)
419+
410420
container_logger.write(user_msg, flush=True)
411421
self._fail_instance(msg=user_msg)
412422
self.set_status(
413423
ActivationStatus.FAILED,
414424
user_msg,
415425
)
416-
system_restart_activation(
417-
self.db_instance_type, self.db_instance.id, delay_seconds=1
418-
)
419426

420427
def _missing_container_policy(self):
421428
LOGGER.info(

src/aap_eda/services/activation/engine/kubernetes.py

Lines changed: 32 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@
5050

5151
K8S_API_RETRIES = 3
5252
K8S_API_RETRY_BACKOFF = 1.0
53-
K8S_API_TRANSIENT_STATUS_CODES = {401, 403, 500, 502, 503, 504}
53+
K8S_API_TRANSIENT_STATUS_CODES = {401, 403, 410, 500, 502, 503, 504}
5454

5555
INVALID_IMAGE_NAME = "InvalidImageName"
5656
IMAGE_PULL_BACK_OFF = "ImagePullBackOff"
@@ -566,9 +566,22 @@ def _delete_job(self, log_handler: LogHandler) -> None:
566566
activation_job_name = activation_job.items[0].metadata.name
567567
self._delete_job_resource(activation_job_name, log_handler)
568568

569+
def _get_resource_version(self) -> str:
570+
pod_list = self._call_k8s_api(
571+
self.client.core_api.list_namespaced_pod,
572+
namespace=self.namespace,
573+
label_selector=f"job-name={self.job_name}",
574+
error_cls=ContainerCleanupError,
575+
description=f"list namespaced jobs {self.job_name}",
576+
)
577+
resource_version = pod_list.metadata.resource_version
578+
return resource_version
579+
569580
def _delete_job_resource(
570581
self, job_name: str, log_handler: LogHandler
571582
) -> None:
583+
resource_version = self._get_resource_version()
584+
572585
result = self._call_k8s_api(
573586
self.client.batch_api.delete_namespaced_job,
574587
name=job_name,
@@ -581,32 +594,42 @@ def _delete_job_resource(
581594
if result.status == "Failure":
582595
raise ContainerCleanupError(f"{result}")
583596

584-
self._watch_pod_deletion(log_handler)
597+
self._watch_pod_deletion(
598+
log_handler, resource_version=resource_version
599+
)
585600

586-
def _watch_pod_deletion(self, log_handler: LogHandler) -> None:
601+
def _watch_pod_deletion(
602+
self, log_handler: LogHandler, resource_version: str
603+
) -> None:
587604
"""Watch for pod deletion with retry on transient errors."""
588605
desc = f"watch pod deletion {self.job_name}"
589606
last_exc = None
590607
for attempt in range(1, K8S_API_RETRIES + 1):
591608
watcher = watch.Watch()
609+
pod_deleted = False
610+
if attempt > 1:
611+
resource_version = self._get_resource_version()
592612
try:
593613
for event in watcher.stream(
594614
self.client.core_api.list_namespaced_pod,
595615
namespace=self.namespace,
596616
label_selector=f"job-name={self.job_name}",
597617
timeout_seconds=POD_DELETE_TIMEOUT,
618+
resource_version=resource_version,
598619
):
599620
if event["type"] == "DELETED":
600621
log_handler.write(
601622
f"Pod '{self.job_name}' is deleted.",
602623
flush=True,
603624
)
625+
pod_deleted = True
604626
break
605-
log_handler.write(
606-
f"Job {self.job_name} is cleaned up.",
607-
flush=True,
608-
)
609-
return
627+
if pod_deleted:
628+
log_handler.write(
629+
f"Job {self.job_name} is cleaned up.",
630+
flush=True,
631+
)
632+
return
610633
except ApiException as exc:
611634
last_exc = exc
612635
if exc.status == status.HTTP_404_NOT_FOUND:
@@ -660,7 +683,7 @@ def _process_pod_start_event(self, event) -> bool:
660683
)
661684
return False
662685

663-
def _wait_for_pod_to_start(self, _log_handler: LogHandler) -> None:
686+
def _wait_for_pod_to_start(self, log_handler: LogHandler) -> None:
664687
LOGGER.info("Waiting for pod to start")
665688
desc = f"watch pod start {self.job_name}"
666689
last_exc = None

tests/integration/services/activation/engine/test_kubernetes.py

Lines changed: 99 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -714,7 +714,11 @@ def test_delete_job(mock_watch, init_kubernetes_data, kubernetes_engine):
714714

715715
with mock.patch.object(engine.client, "batch_api") as batch_api_mock:
716716
batch_api_mock.list_namespaced_job.return_value.items = [job_mock]
717-
engine._delete_job(log_handler)
717+
mock_watch.return_value.stream.return_value = [{"type": "DELETED"}]
718+
with mock.patch.object(
719+
engine, "_get_resource_version", return_value="1"
720+
):
721+
engine._delete_job(log_handler)
718722

719723
batch_api_mock.delete_namespaced_job.assert_called_once()
720724

@@ -730,7 +734,10 @@ def test_delete_job(mock_watch, init_kubernetes_data, kubernetes_engine):
730734
with mock.patch.object(engine.client, "batch_api") as batch_api_mock:
731735
batch_api_mock.list_namespaced_job.return_value.items = [job_mock]
732736
mock_watch.return_value.stream.return_value = [event]
733-
engine._delete_job(log_handler)
737+
with mock.patch.object(
738+
engine, "_get_resource_version", return_value="1"
739+
):
740+
engine._delete_job(log_handler)
734741

735742
log_messages = [
736743
record.log for record in models.RulebookProcessLog.objects.all()
@@ -1153,8 +1160,11 @@ def test_watch_pod_deletion_retries_on_transient(
11531160

11541161
mock_watch.return_value.stream.side_effect = ApiException(status=503)
11551162

1156-
with pytest.raises(ContainerCleanupError, match="retries"):
1157-
engine._watch_pod_deletion(log_handler)
1163+
with mock.patch.object(
1164+
engine, "_get_resource_version", return_value="999"
1165+
):
1166+
with pytest.raises(ContainerCleanupError, match="retries"):
1167+
engine._watch_pod_deletion(log_handler, resource_version="1")
11581168

11591169
assert mock_watch.return_value.stream.call_count == K8S_API_RETRIES
11601170
assert mock_watch.return_value.stop.call_count == K8S_API_RETRIES
@@ -1204,6 +1214,91 @@ def test_wait_for_pod_to_start_raises_on_non_transient(
12041214
assert mock_watch.return_value.stream.call_count == 1
12051215

12061216

1217+
@mock.patch("aap_eda.services.activation.engine.kubernetes.watch.Watch")
1218+
@mock.patch("aap_eda.services.activation.engine.kubernetes.time.sleep")
1219+
@pytest.mark.django_db
1220+
def test_watch_pod_deletion_timeout_without_deleted_event(
1221+
sleep_mock,
1222+
mock_watch,
1223+
init_kubernetes_data,
1224+
kubernetes_engine,
1225+
):
1226+
"""_watch_pod_deletion raises ContainerCleanupError when watch times out
1227+
without receiving a DELETED event across all retries."""
1228+
engine = kubernetes_engine
1229+
engine.job_name = "test-job"
1230+
log_handler = DBLogger(init_kubernetes_data.activation_instance.id)
1231+
1232+
mock_watch.return_value.stream.return_value = iter(
1233+
[
1234+
{"type": "MODIFIED"},
1235+
]
1236+
)
1237+
1238+
with mock.patch.object(
1239+
engine, "_get_resource_version", return_value="999"
1240+
):
1241+
with pytest.raises(ContainerCleanupError, match="retries"):
1242+
engine._watch_pod_deletion(log_handler, resource_version="1")
1243+
1244+
assert mock_watch.return_value.stream.call_count == K8S_API_RETRIES
1245+
assert mock_watch.return_value.stop.call_count == K8S_API_RETRIES
1246+
1247+
1248+
@mock.patch("aap_eda.services.activation.engine.kubernetes.watch.Watch")
1249+
@pytest.mark.django_db
1250+
def test_watch_pod_deletion_empty_stream_retries(
1251+
mock_watch,
1252+
init_kubernetes_data,
1253+
kubernetes_engine,
1254+
):
1255+
"""_watch_pod_deletion retries when the watch stream ends without any
1256+
events (simulating a pure timeout with no events)."""
1257+
engine = kubernetes_engine
1258+
engine.job_name = "test-job"
1259+
log_handler = DBLogger(init_kubernetes_data.activation_instance.id)
1260+
1261+
mock_watch.return_value.stream.return_value = iter([])
1262+
1263+
with mock.patch.object(
1264+
engine, "_get_resource_version", return_value="999"
1265+
):
1266+
with pytest.raises(ContainerCleanupError, match="retries"):
1267+
engine._watch_pod_deletion(log_handler, resource_version="1")
1268+
1269+
assert mock_watch.return_value.stream.call_count == K8S_API_RETRIES
1270+
1271+
1272+
@mock.patch("aap_eda.services.activation.engine.kubernetes.watch.Watch")
1273+
@pytest.mark.django_db
1274+
def test_watch_pod_deletion_success_on_deleted_event(
1275+
mock_watch,
1276+
init_kubernetes_data,
1277+
kubernetes_engine,
1278+
):
1279+
"""_watch_pod_deletion returns successfully when a DELETED event
1280+
is received."""
1281+
engine = kubernetes_engine
1282+
engine.job_name = "test-job"
1283+
log_handler = DBLogger(init_kubernetes_data.activation_instance.id)
1284+
1285+
mock_watch.return_value.stream.return_value = iter(
1286+
[
1287+
{"type": "MODIFIED"},
1288+
{"type": "DELETED"},
1289+
]
1290+
)
1291+
1292+
engine._watch_pod_deletion(log_handler, resource_version="1")
1293+
1294+
assert mock_watch.return_value.stream.call_count == 1
1295+
log_messages = [
1296+
record.log for record in models.RulebookProcessLog.objects.all()
1297+
]
1298+
assert f"Pod '{engine.job_name}' is deleted." in log_messages
1299+
assert f"Job {engine.job_name} is cleaned up." in log_messages
1300+
1301+
12071302
@pytest.mark.django_db
12081303
def test_process_pod_start_event_pending(kubernetes_engine):
12091304
"""_process_pod_start_event returns False for Pending."""

0 commit comments

Comments
 (0)