From 097ccfe552b59506223782e12064f03095ffbc84 Mon Sep 17 00:00:00 2001 From: bradjin8 Date: Fri, 31 Jul 2026 13:51:08 -0400 Subject: [PATCH 1/2] fix: queue follow-ups from #100 review (#103) Emit JOB_CANCELLED in cancel_all, dispatch listeners outside the lock, and treat FAILED/CANCELLED dependencies as unmet in all_dependencies_met. --- CHANGELOG.md | 3 + cli/localci/core/queue.py | 114 +++++++++++++++++++++++++------------- cli/tests/test_queue.py | 63 +++++++++++++++++++-- 3 files changed, 138 insertions(+), 42 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1668306..1c51b48 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,6 +28,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **`derive_image_tag()` prefix** honours `project.native_image_prefix` instead of hardcoding `capy-`, completing the generic-profile default behaviour documented for Week 30. - `PriorityJobQueue`: jobs in `WAITING_DEPS` return to `QUEUED` when their dependencies finish. Same-priority `needs` chains were leaving dependents stuck. - `PriorityJobQueue.cancel_all` drops cancelled keys from `_running_keys`. A `READY` job from `next_ready()` no longer counts as running after cancel. +- `PriorityJobQueue.cancel_all` emits `JOB_CANCELLED` for each cancelled job, matching `cancel()`. +- `DependencyResolver.all_dependencies_met` treats `FAILED` and `CANCELLED` dependencies as unmet; dependents no longer start after a failed or cancelled `needs` job. +- `PriorityJobQueue` event listeners run after `self._lock` is released, so callbacks can safely call back into the queue. ## [0.1.0] - TBD diff --git a/cli/localci/core/queue.py b/cli/localci/core/queue.py index 9a10db7..4ed8bee 100644 --- a/cli/localci/core/queue.py +++ b/cli/localci/core/queue.py @@ -27,6 +27,11 @@ logger = logging.getLogger(__name__) +Listener = Callable[[JobEvent], None] +PendingEmit = tuple[JobEvent, list[Listener]] + +_SATISFIED_DEP_STATUSES = frozenset({QueuedJobStatus.PASSED, QueuedJobStatus.SKIPPED}) + # --------------------------------------------------------------------------- # Priority assignment @@ -145,8 +150,12 @@ def get_dependencies(self, job_id: str) -> list[str]: def get_dependents(self, job_id: str) -> list[str]: return self._reverse.get(job_id, []) - def all_dependencies_met(self, job_id: str, completed: set[str]) -> bool: - return all(dep in completed for dep in self.get_dependencies(job_id)) + def all_dependencies_met(self, job_id: str, jobs: dict[str, QueuedJob]) -> bool: + return all( + (dep_job := jobs.get(dep_id)) is not None + and dep_job.status in _SATISFIED_DEP_STATUSES + for dep_id in self.get_dependencies(job_id) + ) # --------------------------------------------------------------------------- @@ -172,20 +181,27 @@ def __init__(self) -> None: self._failed_keys: set[str] = set() self._running_keys: set[str] = set() self._dep_resolver = DependencyResolver() - self._listeners: list[Callable[[JobEvent], None]] = [] + self._listeners: list[Listener] = [] - def add_listener(self, callback: Callable[[JobEvent], None]) -> None: + def add_listener(self, callback: Listener) -> None: self._listeners.append(callback) - def _emit(self, event_type: JobEventType, job: QueuedJob, **data: object) -> None: + def _prepare_emit( + self, event_type: JobEventType, job: QueuedJob, **data: object + ) -> PendingEmit: event = JobEvent(event_type=event_type, job=job, data=dict(data)) - for listener in self._listeners: - try: - listener(event) - except Exception as e: - logger.warning("Event listener error: %s", e) + return event, list(self._listeners) + + def _dispatch_emits(self, pending: list[PendingEmit]) -> None: + for event, listeners in pending: + for listener in listeners: + try: + listener(event) + except Exception as e: + logger.warning("Event listener error: %s", e) def enqueue(self, job: QueuedJob) -> None: + pending: list[PendingEmit] = [] with self._lock: key = job.queue_key self._jobs[key] = job @@ -201,13 +217,14 @@ def enqueue(self, job: QueuedJob) -> None: job.status = QueuedJobStatus.QUEUED else: job.status = QueuedJobStatus.WAITING_PRIORITY - self._emit(JobEventType.JOB_QUEUED, job) + pending.append(self._prepare_emit(JobEventType.JOB_QUEUED, job)) logger.debug( "Enqueued: %s (priority=%s, deps=%s)", job.matrix_entry.name, job.priority, job.dependencies, ) + self._dispatch_emits(pending) def enqueue_batch(self, jobs: list[QueuedJob]) -> None: for job in jobs: @@ -219,6 +236,8 @@ def enqueue_batch(self, jobs: list[QueuedJob]) -> None: ) def next_ready(self) -> QueuedJob | None: + pending: list[PendingEmit] = [] + result: QueuedJob | None = None with self._lock: if self._current_priority is None: return None @@ -227,18 +246,19 @@ def next_ready(self) -> QueuedJob | None: job = self._jobs[key] if job.status != QueuedJobStatus.QUEUED: continue - if not self._dep_resolver.all_dependencies_met( - key, self._completed_keys - ): + if not self._dep_resolver.all_dependencies_met(key, self._jobs): job.status = QueuedJobStatus.WAITING_DEPS continue job.status = QueuedJobStatus.READY self._running_keys.add(key) - self._emit(JobEventType.JOB_READY, job) - return job - return None + pending.append(self._prepare_emit(JobEventType.JOB_READY, job)) + result = job + break + self._dispatch_emits(pending) + return result def mark_completed(self, job: QueuedJob, success: bool = True) -> None: + pending: list[PendingEmit] = [] with self._lock: key = job.queue_key if success: @@ -249,27 +269,36 @@ def mark_completed(self, job: QueuedJob, success: bool = True) -> None: self._failed_keys.add(key) self._completed_keys.add(key) self._running_keys.discard(key) - self._after_job_terminal_state_change() + pending.extend(self._after_job_terminal_state_change()) + self._dispatch_emits(pending) def mark_skipped(self, job: QueuedJob) -> None: + pending: list[PendingEmit] = [] with self._lock: key = job.queue_key job.status = QueuedJobStatus.SKIPPED self._completed_keys.add(key) self._running_keys.discard(key) - self._after_job_terminal_state_change() + pending.extend(self._after_job_terminal_state_change()) + self._dispatch_emits(pending) def mark_running(self, job: QueuedJob) -> None: + pending: list[PendingEmit] = [] with self._lock: job.status = QueuedJobStatus.RUNNING - self._emit(JobEventType.JOB_STARTED, job) + pending.append(self._prepare_emit(JobEventType.JOB_STARTED, job)) + self._dispatch_emits(pending) def mark_preparing(self, job: QueuedJob) -> None: + pending: list[PendingEmit] = [] with self._lock: job.status = QueuedJobStatus.PREPARING - self._emit(JobEventType.JOB_PREPARING, job) + pending.append(self._prepare_emit(JobEventType.JOB_PREPARING, job)) + self._dispatch_emits(pending) def cancel(self, key: str) -> bool: + pending: list[PendingEmit] = [] + cancelled = False with self._lock: job = self._jobs.get(key) if job is None: @@ -282,11 +311,14 @@ def cancel(self, key: str) -> bool: return False job.status = QueuedJobStatus.CANCELLED self._completed_keys.add(key) - self._emit(JobEventType.JOB_CANCELLED, job) - self._after_job_terminal_state_change() - return True + pending.append(self._prepare_emit(JobEventType.JOB_CANCELLED, job)) + pending.extend(self._after_job_terminal_state_change()) + cancelled = True + self._dispatch_emits(pending) + return cancelled def cancel_all(self) -> int: + pending: list[PendingEmit] = [] count = 0 with self._lock: for key, job in list(self._jobs.items()): @@ -299,14 +331,16 @@ def cancel_all(self) -> int: job.status = QueuedJobStatus.CANCELLED self._completed_keys.add(key) self._running_keys.discard(key) + pending.append(self._prepare_emit(JobEventType.JOB_CANCELLED, job)) count += 1 if count > 0: - self._after_job_terminal_state_change() + pending.extend(self._after_job_terminal_state_change()) + self._dispatch_emits(pending) return count - def _after_job_terminal_state_change(self) -> None: + def _after_job_terminal_state_change(self) -> list[PendingEmit]: self._promote_waiting_deps() - self._check_priority_advance() + return self._check_priority_advance() def _promote_waiting_deps(self) -> None: if self._current_priority is None: @@ -315,12 +349,13 @@ def _promote_waiting_deps(self) -> None: job = self._jobs[key] if job.status != QueuedJobStatus.WAITING_DEPS: continue - if self._dep_resolver.all_dependencies_met(key, self._completed_keys): + if self._dep_resolver.all_dependencies_met(key, self._jobs): job.status = QueuedJobStatus.QUEUED - def _check_priority_advance(self) -> None: + def _check_priority_advance(self) -> list[PendingEmit]: + pending: list[PendingEmit] = [] if self._current_priority is None: - return + return pending current_keys = self._by_priority.get(self._current_priority, []) terminal = ( QueuedJobStatus.PASSED, @@ -332,18 +367,20 @@ def _check_priority_advance(self) -> None: ) all_done = all(self._jobs[k].status in terminal for k in current_keys) if not all_done: - return + return pending logger.info( "Priority level %s complete (%d jobs)", self._current_priority, len(current_keys), ) if current_keys: - self._emit( - JobEventType.PRIORITY_LEVEL_COMPLETE, - self._jobs[current_keys[0]], - priority=self._current_priority, - count=len(current_keys), + pending.append( + self._prepare_emit( + JobEventType.PRIORITY_LEVEL_COMPLETE, + self._jobs[current_keys[0]], + priority=self._current_priority, + count=len(current_keys), + ) ) idx = self._priority_levels.index(self._current_priority) if idx + 1 < len(self._priority_levels): @@ -357,7 +394,10 @@ def _check_priority_advance(self) -> None: self._current_priority = None if self._jobs: first_key = next(iter(self._jobs)) - self._emit(JobEventType.ALL_COMPLETE, self._jobs[first_key]) + pending.append( + self._prepare_emit(JobEventType.ALL_COMPLETE, self._jobs[first_key]) + ) + return pending @property def is_empty(self) -> bool: diff --git a/cli/tests/test_queue.py b/cli/tests/test_queue.py index 8038f16..0c817ad 100644 --- a/cli/tests/test_queue.py +++ b/cli/tests/test_queue.py @@ -8,6 +8,7 @@ from localci.core.config import LocalCIConfig from localci.core.models import ( + JobEvent, JobEventType, QueuedJob, QueuedJobStatus, @@ -142,11 +143,21 @@ def test_cyclic_dependency(self): def test_all_dependencies_met(self): resolver = DependencyResolver() - resolver.add_job("a", []) - resolver.add_job("b", ["a"]) - assert resolver.all_dependencies_met("a", set()) is True - assert resolver.all_dependencies_met("b", set()) is False - assert resolver.all_dependencies_met("b", {"a"}) is True + upstream = make_job("Upstream", priority=1, index=0) + downstream = make_job( + "Downstream", priority=1, index=1, deps=[upstream.queue_key] + ) + resolver.add_job(upstream.queue_key, []) + resolver.add_job(downstream.queue_key, [upstream.queue_key]) + jobs = {upstream.queue_key: upstream, downstream.queue_key: downstream} + assert resolver.all_dependencies_met(upstream.queue_key, jobs) is True + assert resolver.all_dependencies_met(downstream.queue_key, jobs) is False + upstream.status = QueuedJobStatus.PASSED + assert resolver.all_dependencies_met(downstream.queue_key, jobs) is True + upstream.status = QueuedJobStatus.FAILED + assert resolver.all_dependencies_met(downstream.queue_key, jobs) is False + upstream.status = QueuedJobStatus.CANCELLED + assert resolver.all_dependencies_met(downstream.queue_key, jobs) is False # --------------------------------------------------------------------------- @@ -227,6 +238,37 @@ def test_waiting_deps_promoted_when_dependency_completes(self): queue.mark_completed(second, success=True) assert queue.is_done is True + def test_dependent_stays_waiting_when_dependency_failed(self): + queue = PriorityJobQueue() + upstream = make_job("Upstream", priority=1, index=0) + downstream = make_job( + "Downstream", priority=1, index=1, deps=[upstream.queue_key] + ) + queue.enqueue(upstream) + queue.enqueue(downstream) + + first = queue.next_ready() + assert first is not None + queue.mark_running(first) + queue.mark_completed(first, success=False) + + assert queue.next_ready() is None + jobs = {j.queue_key: j for j in queue.get_all_jobs()} + assert jobs[downstream.queue_key].status == QueuedJobStatus.WAITING_DEPS + + def test_listener_can_acquire_lock_during_callback(self) -> None: + queue = PriorityJobQueue() + queue.enqueue(make_job("Job 1", priority=1, index=0)) + acquired = threading.Event() + + def listener(_event: JobEvent) -> None: + with queue._lock: + acquired.set() + + queue.add_listener(listener) + queue.enqueue(make_job("Job 2", priority=1, index=1)) + assert acquired.wait(timeout=5) + def test_is_done_acquires_lock(self): queue = PriorityJobQueue() queue.enqueue(make_job("Job", priority=1)) @@ -316,6 +358,17 @@ def test_cancel_all_after_next_ready_clears_running_keys(self): assert queue.pending_count >= 0 assert queue.is_done is True + def test_cancel_all_emits_cancelled_for_each_job(self): + events: list[JobEvent] = [] + queue = PriorityJobQueue() + queue.add_listener(lambda e: events.append(e)) + for i in range(3): + queue.enqueue(make_job(f"Job {i}", priority=1, index=i)) + count = queue.cancel_all() + assert count == 3 + cancelled = [e for e in events if e.event_type == JobEventType.JOB_CANCELLED] + assert len(cancelled) == 3 + def test_event_emission(self): events = [] queue = PriorityJobQueue() From f93c3803a948907d8313c39f76bd2cbc62c91d11 Mon Sep 17 00:00:00 2001 From: bradjin8 Date: Fri, 31 Jul 2026 14:14:44 -0400 Subject: [PATCH 2/2] fix: addressed ai review findings --- cli/localci/core/queue.py | 83 ++++++++++++++++++++------------------- cli/tests/test_queue.py | 15 ++++++- 2 files changed, 57 insertions(+), 41 deletions(-) diff --git a/cli/localci/core/queue.py b/cli/localci/core/queue.py index 4ed8bee..31af7a9 100644 --- a/cli/localci/core/queue.py +++ b/cli/localci/core/queue.py @@ -10,8 +10,9 @@ import fnmatch import logging import threading +from collections import deque from collections.abc import Callable -from dataclasses import dataclass, field +from dataclasses import dataclass, field, replace from typing import TYPE_CHECKING from localci.core.models import ( @@ -182,18 +183,28 @@ def __init__(self) -> None: self._running_keys: set[str] = set() self._dep_resolver = DependencyResolver() self._listeners: list[Listener] = [] + self._emit_fifo: deque[PendingEmit] = deque() def add_listener(self, callback: Listener) -> None: - self._listeners.append(callback) + with self._lock: + self._listeners.append(callback) def _prepare_emit( self, event_type: JobEventType, job: QueuedJob, **data: object ) -> PendingEmit: - event = JobEvent(event_type=event_type, job=job, data=dict(data)) + snapshot = replace(job) + event = JobEvent(event_type=event_type, job=snapshot, data=dict(data)) return event, list(self._listeners) - def _dispatch_emits(self, pending: list[PendingEmit]) -> None: - for event, listeners in pending: + def _enqueue_emit(self, pending: PendingEmit) -> None: + self._emit_fifo.append(pending) + + def _drain_emits(self) -> None: + while True: + with self._lock: + if not self._emit_fifo: + return + event, listeners = self._emit_fifo.popleft() for listener in listeners: try: listener(event) @@ -201,7 +212,6 @@ def _dispatch_emits(self, pending: list[PendingEmit]) -> None: logger.warning("Event listener error: %s", e) def enqueue(self, job: QueuedJob) -> None: - pending: list[PendingEmit] = [] with self._lock: key = job.queue_key self._jobs[key] = job @@ -217,14 +227,14 @@ def enqueue(self, job: QueuedJob) -> None: job.status = QueuedJobStatus.QUEUED else: job.status = QueuedJobStatus.WAITING_PRIORITY - pending.append(self._prepare_emit(JobEventType.JOB_QUEUED, job)) + self._enqueue_emit(self._prepare_emit(JobEventType.JOB_QUEUED, job)) logger.debug( "Enqueued: %s (priority=%s, deps=%s)", job.matrix_entry.name, job.priority, job.dependencies, ) - self._dispatch_emits(pending) + self._drain_emits() def enqueue_batch(self, jobs: list[QueuedJob]) -> None: for job in jobs: @@ -236,7 +246,6 @@ def enqueue_batch(self, jobs: list[QueuedJob]) -> None: ) def next_ready(self) -> QueuedJob | None: - pending: list[PendingEmit] = [] result: QueuedJob | None = None with self._lock: if self._current_priority is None: @@ -251,14 +260,13 @@ def next_ready(self) -> QueuedJob | None: continue job.status = QueuedJobStatus.READY self._running_keys.add(key) - pending.append(self._prepare_emit(JobEventType.JOB_READY, job)) + self._enqueue_emit(self._prepare_emit(JobEventType.JOB_READY, job)) result = job break - self._dispatch_emits(pending) + self._drain_emits() return result def mark_completed(self, job: QueuedJob, success: bool = True) -> None: - pending: list[PendingEmit] = [] with self._lock: key = job.queue_key if success: @@ -269,35 +277,31 @@ def mark_completed(self, job: QueuedJob, success: bool = True) -> None: self._failed_keys.add(key) self._completed_keys.add(key) self._running_keys.discard(key) - pending.extend(self._after_job_terminal_state_change()) - self._dispatch_emits(pending) + self._after_job_terminal_state_change() + self._drain_emits() def mark_skipped(self, job: QueuedJob) -> None: - pending: list[PendingEmit] = [] with self._lock: key = job.queue_key job.status = QueuedJobStatus.SKIPPED self._completed_keys.add(key) self._running_keys.discard(key) - pending.extend(self._after_job_terminal_state_change()) - self._dispatch_emits(pending) + self._after_job_terminal_state_change() + self._drain_emits() def mark_running(self, job: QueuedJob) -> None: - pending: list[PendingEmit] = [] with self._lock: job.status = QueuedJobStatus.RUNNING - pending.append(self._prepare_emit(JobEventType.JOB_STARTED, job)) - self._dispatch_emits(pending) + self._enqueue_emit(self._prepare_emit(JobEventType.JOB_STARTED, job)) + self._drain_emits() def mark_preparing(self, job: QueuedJob) -> None: - pending: list[PendingEmit] = [] with self._lock: job.status = QueuedJobStatus.PREPARING - pending.append(self._prepare_emit(JobEventType.JOB_PREPARING, job)) - self._dispatch_emits(pending) + self._enqueue_emit(self._prepare_emit(JobEventType.JOB_PREPARING, job)) + self._drain_emits() def cancel(self, key: str) -> bool: - pending: list[PendingEmit] = [] cancelled = False with self._lock: job = self._jobs.get(key) @@ -311,14 +315,13 @@ def cancel(self, key: str) -> bool: return False job.status = QueuedJobStatus.CANCELLED self._completed_keys.add(key) - pending.append(self._prepare_emit(JobEventType.JOB_CANCELLED, job)) - pending.extend(self._after_job_terminal_state_change()) + self._enqueue_emit(self._prepare_emit(JobEventType.JOB_CANCELLED, job)) + self._after_job_terminal_state_change() cancelled = True - self._dispatch_emits(pending) + self._drain_emits() return cancelled def cancel_all(self) -> int: - pending: list[PendingEmit] = [] count = 0 with self._lock: for key, job in list(self._jobs.items()): @@ -331,16 +334,18 @@ def cancel_all(self) -> int: job.status = QueuedJobStatus.CANCELLED self._completed_keys.add(key) self._running_keys.discard(key) - pending.append(self._prepare_emit(JobEventType.JOB_CANCELLED, job)) + self._enqueue_emit( + self._prepare_emit(JobEventType.JOB_CANCELLED, job) + ) count += 1 if count > 0: - pending.extend(self._after_job_terminal_state_change()) - self._dispatch_emits(pending) + self._after_job_terminal_state_change() + self._drain_emits() return count - def _after_job_terminal_state_change(self) -> list[PendingEmit]: + def _after_job_terminal_state_change(self) -> None: self._promote_waiting_deps() - return self._check_priority_advance() + self._check_priority_advance() def _promote_waiting_deps(self) -> None: if self._current_priority is None: @@ -352,10 +357,9 @@ def _promote_waiting_deps(self) -> None: if self._dep_resolver.all_dependencies_met(key, self._jobs): job.status = QueuedJobStatus.QUEUED - def _check_priority_advance(self) -> list[PendingEmit]: - pending: list[PendingEmit] = [] + def _check_priority_advance(self) -> None: if self._current_priority is None: - return pending + return current_keys = self._by_priority.get(self._current_priority, []) terminal = ( QueuedJobStatus.PASSED, @@ -367,14 +371,14 @@ def _check_priority_advance(self) -> list[PendingEmit]: ) all_done = all(self._jobs[k].status in terminal for k in current_keys) if not all_done: - return pending + return logger.info( "Priority level %s complete (%d jobs)", self._current_priority, len(current_keys), ) if current_keys: - pending.append( + self._enqueue_emit( self._prepare_emit( JobEventType.PRIORITY_LEVEL_COMPLETE, self._jobs[current_keys[0]], @@ -394,10 +398,9 @@ def _check_priority_advance(self) -> list[PendingEmit]: self._current_priority = None if self._jobs: first_key = next(iter(self._jobs)) - pending.append( + self._enqueue_emit( self._prepare_emit(JobEventType.ALL_COMPLETE, self._jobs[first_key]) ) - return pending @property def is_empty(self) -> bool: diff --git a/cli/tests/test_queue.py b/cli/tests/test_queue.py index 0c817ad..9d8ca6e 100644 --- a/cli/tests/test_queue.py +++ b/cli/tests/test_queue.py @@ -158,6 +158,13 @@ def test_all_dependencies_met(self): assert resolver.all_dependencies_met(downstream.queue_key, jobs) is False upstream.status = QueuedJobStatus.CANCELLED assert resolver.all_dependencies_met(downstream.queue_key, jobs) is False + upstream.status = QueuedJobStatus.SKIPPED + assert resolver.all_dependencies_met(downstream.queue_key, jobs) is True + jobs_missing_upstream = {downstream.queue_key: downstream} + assert ( + resolver.all_dependencies_met(downstream.queue_key, jobs_missing_upstream) + is False + ) # --------------------------------------------------------------------------- @@ -266,7 +273,13 @@ def listener(_event: JobEvent) -> None: acquired.set() queue.add_listener(listener) - queue.enqueue(make_job("Job 2", priority=1, index=1)) + worker = threading.Thread( + target=lambda: queue.enqueue(make_job("Job 2", priority=1, index=1)), + daemon=True, + ) + worker.start() + worker.join(timeout=_THREAD_JOIN_TIMEOUT) + assert not worker.is_alive(), "enqueue worker thread hung" assert acquired.wait(timeout=5) def test_is_done_acquires_lock(self):