Skip to content

Commit 2c804a3

Browse files
authored
fix: PriorityJobQueue follow-ups from #100 review (#103) (#104)
* 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. * fix: addressed ai review findings
1 parent 09ebd73 commit 2c804a3

3 files changed

Lines changed: 149 additions & 37 deletions

File tree

CHANGELOG.md

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
2828
- **`derive_image_tag()` prefix** honours `project.native_image_prefix` instead of hardcoding `capy-`, completing the generic-profile default behaviour documented for Week 30.
2929
- `PriorityJobQueue`: jobs in `WAITING_DEPS` return to `QUEUED` when their dependencies finish. Same-priority `needs` chains were leaving dependents stuck.
3030
- `PriorityJobQueue.cancel_all` drops cancelled keys from `_running_keys`. A `READY` job from `next_ready()` no longer counts as running after cancel.
31+
- `PriorityJobQueue.cancel_all` emits `JOB_CANCELLED` for each cancelled job, matching `cancel()`.
32+
- `DependencyResolver.all_dependencies_met` treats `FAILED` and `CANCELLED` dependencies as unmet; dependents no longer start after a failed or cancelled `needs` job.
33+
- `PriorityJobQueue` event listeners run after `self._lock` is released, so callbacks can safely call back into the queue.
3134

3235
## [0.1.0] - TBD
3336

cli/localci/core/queue.py

Lines changed: 75 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,9 @@
1010
import fnmatch
1111
import logging
1212
import threading
13+
from collections import deque
1314
from collections.abc import Callable
14-
from dataclasses import dataclass, field
15+
from dataclasses import dataclass, field, replace
1516
from typing import TYPE_CHECKING
1617

1718
from localci.core.models import (
@@ -27,6 +28,11 @@
2728

2829
logger = logging.getLogger(__name__)
2930

31+
Listener = Callable[[JobEvent], None]
32+
PendingEmit = tuple[JobEvent, list[Listener]]
33+
34+
_SATISFIED_DEP_STATUSES = frozenset({QueuedJobStatus.PASSED, QueuedJobStatus.SKIPPED})
35+
3036

3137
# ---------------------------------------------------------------------------
3238
# Priority assignment
@@ -145,8 +151,12 @@ def get_dependencies(self, job_id: str) -> list[str]:
145151
def get_dependents(self, job_id: str) -> list[str]:
146152
return self._reverse.get(job_id, [])
147153

148-
def all_dependencies_met(self, job_id: str, completed: set[str]) -> bool:
149-
return all(dep in completed for dep in self.get_dependencies(job_id))
154+
def all_dependencies_met(self, job_id: str, jobs: dict[str, QueuedJob]) -> bool:
155+
return all(
156+
(dep_job := jobs.get(dep_id)) is not None
157+
and dep_job.status in _SATISFIED_DEP_STATUSES
158+
for dep_id in self.get_dependencies(job_id)
159+
)
150160

151161

152162
# ---------------------------------------------------------------------------
@@ -172,18 +182,34 @@ def __init__(self) -> None:
172182
self._failed_keys: set[str] = set()
173183
self._running_keys: set[str] = set()
174184
self._dep_resolver = DependencyResolver()
175-
self._listeners: list[Callable[[JobEvent], None]] = []
176-
177-
def add_listener(self, callback: Callable[[JobEvent], None]) -> None:
178-
self._listeners.append(callback)
185+
self._listeners: list[Listener] = []
186+
self._emit_fifo: deque[PendingEmit] = deque()
179187

180-
def _emit(self, event_type: JobEventType, job: QueuedJob, **data: object) -> None:
181-
event = JobEvent(event_type=event_type, job=job, data=dict(data))
182-
for listener in self._listeners:
183-
try:
184-
listener(event)
185-
except Exception as e:
186-
logger.warning("Event listener error: %s", e)
188+
def add_listener(self, callback: Listener) -> None:
189+
with self._lock:
190+
self._listeners.append(callback)
191+
192+
def _prepare_emit(
193+
self, event_type: JobEventType, job: QueuedJob, **data: object
194+
) -> PendingEmit:
195+
snapshot = replace(job)
196+
event = JobEvent(event_type=event_type, job=snapshot, data=dict(data))
197+
return event, list(self._listeners)
198+
199+
def _enqueue_emit(self, pending: PendingEmit) -> None:
200+
self._emit_fifo.append(pending)
201+
202+
def _drain_emits(self) -> None:
203+
while True:
204+
with self._lock:
205+
if not self._emit_fifo:
206+
return
207+
event, listeners = self._emit_fifo.popleft()
208+
for listener in listeners:
209+
try:
210+
listener(event)
211+
except Exception as e:
212+
logger.warning("Event listener error: %s", e)
187213

188214
def enqueue(self, job: QueuedJob) -> None:
189215
with self._lock:
@@ -201,13 +227,14 @@ def enqueue(self, job: QueuedJob) -> None:
201227
job.status = QueuedJobStatus.QUEUED
202228
else:
203229
job.status = QueuedJobStatus.WAITING_PRIORITY
204-
self._emit(JobEventType.JOB_QUEUED, job)
230+
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_QUEUED, job))
205231
logger.debug(
206232
"Enqueued: %s (priority=%s, deps=%s)",
207233
job.matrix_entry.name,
208234
job.priority,
209235
job.dependencies,
210236
)
237+
self._drain_emits()
211238

212239
def enqueue_batch(self, jobs: list[QueuedJob]) -> None:
213240
for job in jobs:
@@ -219,6 +246,7 @@ def enqueue_batch(self, jobs: list[QueuedJob]) -> None:
219246
)
220247

221248
def next_ready(self) -> QueuedJob | None:
249+
result: QueuedJob | None = None
222250
with self._lock:
223251
if self._current_priority is None:
224252
return None
@@ -227,16 +255,16 @@ def next_ready(self) -> QueuedJob | None:
227255
job = self._jobs[key]
228256
if job.status != QueuedJobStatus.QUEUED:
229257
continue
230-
if not self._dep_resolver.all_dependencies_met(
231-
key, self._completed_keys
232-
):
258+
if not self._dep_resolver.all_dependencies_met(key, self._jobs):
233259
job.status = QueuedJobStatus.WAITING_DEPS
234260
continue
235261
job.status = QueuedJobStatus.READY
236262
self._running_keys.add(key)
237-
self._emit(JobEventType.JOB_READY, job)
238-
return job
239-
return None
263+
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_READY, job))
264+
result = job
265+
break
266+
self._drain_emits()
267+
return result
240268

241269
def mark_completed(self, job: QueuedJob, success: bool = True) -> None:
242270
with self._lock:
@@ -250,6 +278,7 @@ def mark_completed(self, job: QueuedJob, success: bool = True) -> None:
250278
self._completed_keys.add(key)
251279
self._running_keys.discard(key)
252280
self._after_job_terminal_state_change()
281+
self._drain_emits()
253282

254283
def mark_skipped(self, job: QueuedJob) -> None:
255284
with self._lock:
@@ -258,18 +287,22 @@ def mark_skipped(self, job: QueuedJob) -> None:
258287
self._completed_keys.add(key)
259288
self._running_keys.discard(key)
260289
self._after_job_terminal_state_change()
290+
self._drain_emits()
261291

262292
def mark_running(self, job: QueuedJob) -> None:
263293
with self._lock:
264294
job.status = QueuedJobStatus.RUNNING
265-
self._emit(JobEventType.JOB_STARTED, job)
295+
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_STARTED, job))
296+
self._drain_emits()
266297

267298
def mark_preparing(self, job: QueuedJob) -> None:
268299
with self._lock:
269300
job.status = QueuedJobStatus.PREPARING
270-
self._emit(JobEventType.JOB_PREPARING, job)
301+
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_PREPARING, job))
302+
self._drain_emits()
271303

272304
def cancel(self, key: str) -> bool:
305+
cancelled = False
273306
with self._lock:
274307
job = self._jobs.get(key)
275308
if job is None:
@@ -282,9 +315,11 @@ def cancel(self, key: str) -> bool:
282315
return False
283316
job.status = QueuedJobStatus.CANCELLED
284317
self._completed_keys.add(key)
285-
self._emit(JobEventType.JOB_CANCELLED, job)
318+
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_CANCELLED, job))
286319
self._after_job_terminal_state_change()
287-
return True
320+
cancelled = True
321+
self._drain_emits()
322+
return cancelled
288323

289324
def cancel_all(self) -> int:
290325
count = 0
@@ -299,9 +334,13 @@ def cancel_all(self) -> int:
299334
job.status = QueuedJobStatus.CANCELLED
300335
self._completed_keys.add(key)
301336
self._running_keys.discard(key)
337+
self._enqueue_emit(
338+
self._prepare_emit(JobEventType.JOB_CANCELLED, job)
339+
)
302340
count += 1
303341
if count > 0:
304342
self._after_job_terminal_state_change()
343+
self._drain_emits()
305344
return count
306345

307346
def _after_job_terminal_state_change(self) -> None:
@@ -315,7 +354,7 @@ def _promote_waiting_deps(self) -> None:
315354
job = self._jobs[key]
316355
if job.status != QueuedJobStatus.WAITING_DEPS:
317356
continue
318-
if self._dep_resolver.all_dependencies_met(key, self._completed_keys):
357+
if self._dep_resolver.all_dependencies_met(key, self._jobs):
319358
job.status = QueuedJobStatus.QUEUED
320359

321360
def _check_priority_advance(self) -> None:
@@ -339,11 +378,13 @@ def _check_priority_advance(self) -> None:
339378
len(current_keys),
340379
)
341380
if current_keys:
342-
self._emit(
343-
JobEventType.PRIORITY_LEVEL_COMPLETE,
344-
self._jobs[current_keys[0]],
345-
priority=self._current_priority,
346-
count=len(current_keys),
381+
self._enqueue_emit(
382+
self._prepare_emit(
383+
JobEventType.PRIORITY_LEVEL_COMPLETE,
384+
self._jobs[current_keys[0]],
385+
priority=self._current_priority,
386+
count=len(current_keys),
387+
)
347388
)
348389
idx = self._priority_levels.index(self._current_priority)
349390
if idx + 1 < len(self._priority_levels):
@@ -357,7 +398,9 @@ def _check_priority_advance(self) -> None:
357398
self._current_priority = None
358399
if self._jobs:
359400
first_key = next(iter(self._jobs))
360-
self._emit(JobEventType.ALL_COMPLETE, self._jobs[first_key])
401+
self._enqueue_emit(
402+
self._prepare_emit(JobEventType.ALL_COMPLETE, self._jobs[first_key])
403+
)
361404

362405
@property
363406
def is_empty(self) -> bool:

cli/tests/test_queue.py

Lines changed: 71 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88

99
from localci.core.config import LocalCIConfig
1010
from localci.core.models import (
11+
JobEvent,
1112
JobEventType,
1213
QueuedJob,
1314
QueuedJobStatus,
@@ -142,11 +143,28 @@ def test_cyclic_dependency(self):
142143

143144
def test_all_dependencies_met(self):
144145
resolver = DependencyResolver()
145-
resolver.add_job("a", [])
146-
resolver.add_job("b", ["a"])
147-
assert resolver.all_dependencies_met("a", set()) is True
148-
assert resolver.all_dependencies_met("b", set()) is False
149-
assert resolver.all_dependencies_met("b", {"a"}) is True
146+
upstream = make_job("Upstream", priority=1, index=0)
147+
downstream = make_job(
148+
"Downstream", priority=1, index=1, deps=[upstream.queue_key]
149+
)
150+
resolver.add_job(upstream.queue_key, [])
151+
resolver.add_job(downstream.queue_key, [upstream.queue_key])
152+
jobs = {upstream.queue_key: upstream, downstream.queue_key: downstream}
153+
assert resolver.all_dependencies_met(upstream.queue_key, jobs) is True
154+
assert resolver.all_dependencies_met(downstream.queue_key, jobs) is False
155+
upstream.status = QueuedJobStatus.PASSED
156+
assert resolver.all_dependencies_met(downstream.queue_key, jobs) is True
157+
upstream.status = QueuedJobStatus.FAILED
158+
assert resolver.all_dependencies_met(downstream.queue_key, jobs) is False
159+
upstream.status = QueuedJobStatus.CANCELLED
160+
assert resolver.all_dependencies_met(downstream.queue_key, jobs) is False
161+
upstream.status = QueuedJobStatus.SKIPPED
162+
assert resolver.all_dependencies_met(downstream.queue_key, jobs) is True
163+
jobs_missing_upstream = {downstream.queue_key: downstream}
164+
assert (
165+
resolver.all_dependencies_met(downstream.queue_key, jobs_missing_upstream)
166+
is False
167+
)
150168

151169

152170
# ---------------------------------------------------------------------------
@@ -227,6 +245,43 @@ def test_waiting_deps_promoted_when_dependency_completes(self):
227245
queue.mark_completed(second, success=True)
228246
assert queue.is_done is True
229247

248+
def test_dependent_stays_waiting_when_dependency_failed(self):
249+
queue = PriorityJobQueue()
250+
upstream = make_job("Upstream", priority=1, index=0)
251+
downstream = make_job(
252+
"Downstream", priority=1, index=1, deps=[upstream.queue_key]
253+
)
254+
queue.enqueue(upstream)
255+
queue.enqueue(downstream)
256+
257+
first = queue.next_ready()
258+
assert first is not None
259+
queue.mark_running(first)
260+
queue.mark_completed(first, success=False)
261+
262+
assert queue.next_ready() is None
263+
jobs = {j.queue_key: j for j in queue.get_all_jobs()}
264+
assert jobs[downstream.queue_key].status == QueuedJobStatus.WAITING_DEPS
265+
266+
def test_listener_can_acquire_lock_during_callback(self) -> None:
267+
queue = PriorityJobQueue()
268+
queue.enqueue(make_job("Job 1", priority=1, index=0))
269+
acquired = threading.Event()
270+
271+
def listener(_event: JobEvent) -> None:
272+
with queue._lock:
273+
acquired.set()
274+
275+
queue.add_listener(listener)
276+
worker = threading.Thread(
277+
target=lambda: queue.enqueue(make_job("Job 2", priority=1, index=1)),
278+
daemon=True,
279+
)
280+
worker.start()
281+
worker.join(timeout=_THREAD_JOIN_TIMEOUT)
282+
assert not worker.is_alive(), "enqueue worker thread hung"
283+
assert acquired.wait(timeout=5)
284+
230285
def test_is_done_acquires_lock(self):
231286
queue = PriorityJobQueue()
232287
queue.enqueue(make_job("Job", priority=1))
@@ -316,6 +371,17 @@ def test_cancel_all_after_next_ready_clears_running_keys(self):
316371
assert queue.pending_count >= 0
317372
assert queue.is_done is True
318373

374+
def test_cancel_all_emits_cancelled_for_each_job(self):
375+
events: list[JobEvent] = []
376+
queue = PriorityJobQueue()
377+
queue.add_listener(lambda e: events.append(e))
378+
for i in range(3):
379+
queue.enqueue(make_job(f"Job {i}", priority=1, index=i))
380+
count = queue.cancel_all()
381+
assert count == 3
382+
cancelled = [e for e in events if e.event_type == JobEventType.JOB_CANCELLED]
383+
assert len(cancelled) == 3
384+
319385
def test_event_emission(self):
320386
events = []
321387
queue = PriorityJobQueue()

0 commit comments

Comments
 (0)