Skip to content
Merged
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
107 changes: 75 additions & 32 deletions cli/localci/core/queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand All @@ -27,6 +28,11 @@

logger = logging.getLogger(__name__)

Listener = Callable[[JobEvent], None]
PendingEmit = tuple[JobEvent, list[Listener]]

_SATISFIED_DEP_STATUSES = frozenset({QueuedJobStatus.PASSED, QueuedJobStatus.SKIPPED})


# ---------------------------------------------------------------------------
# Priority assignment
Expand Down Expand Up @@ -145,8 +151,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)
)


# ---------------------------------------------------------------------------
Expand All @@ -172,18 +182,34 @@ 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]] = []

def add_listener(self, callback: Callable[[JobEvent], None]) -> None:
self._listeners.append(callback)
self._listeners: list[Listener] = []
self._emit_fifo: deque[PendingEmit] = deque()

def _emit(self, event_type: JobEventType, job: QueuedJob, **data: object) -> None:
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)
def add_listener(self, callback: Listener) -> None:
with self._lock:
self._listeners.append(callback)

def _prepare_emit(
self, event_type: JobEventType, job: QueuedJob, **data: object
) -> PendingEmit:
snapshot = replace(job)
event = JobEvent(event_type=event_type, job=snapshot, data=dict(data))
return event, list(self._listeners)
Comment thread
coderabbitai[bot] marked this conversation as resolved.

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)
except Exception as e:
logger.warning("Event listener error: %s", e)

def enqueue(self, job: QueuedJob) -> None:
with self._lock:
Expand All @@ -201,13 +227,14 @@ def enqueue(self, job: QueuedJob) -> None:
job.status = QueuedJobStatus.QUEUED
else:
job.status = QueuedJobStatus.WAITING_PRIORITY
self._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._drain_emits()

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

def next_ready(self) -> QueuedJob | None:
result: QueuedJob | None = None
with self._lock:
if self._current_priority is None:
return None
Expand All @@ -227,16 +255,16 @@ 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
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_READY, job))
result = job
break
self._drain_emits()
return result

def mark_completed(self, job: QueuedJob, success: bool = True) -> None:
with self._lock:
Expand All @@ -250,6 +278,7 @@ def mark_completed(self, job: QueuedJob, success: bool = True) -> None:
self._completed_keys.add(key)
self._running_keys.discard(key)
self._after_job_terminal_state_change()
self._drain_emits()

def mark_skipped(self, job: QueuedJob) -> None:
with self._lock:
Expand All @@ -258,18 +287,22 @@ def mark_skipped(self, job: QueuedJob) -> None:
self._completed_keys.add(key)
self._running_keys.discard(key)
self._after_job_terminal_state_change()
self._drain_emits()

def mark_running(self, job: QueuedJob) -> None:
with self._lock:
job.status = QueuedJobStatus.RUNNING
self._emit(JobEventType.JOB_STARTED, job)
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_STARTED, job))
self._drain_emits()

def mark_preparing(self, job: QueuedJob) -> None:
with self._lock:
job.status = QueuedJobStatus.PREPARING
self._emit(JobEventType.JOB_PREPARING, job)
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_PREPARING, job))
self._drain_emits()

def cancel(self, key: str) -> bool:
cancelled = False
with self._lock:
job = self._jobs.get(key)
if job is None:
Expand All @@ -282,9 +315,11 @@ def cancel(self, key: str) -> bool:
return False
job.status = QueuedJobStatus.CANCELLED
self._completed_keys.add(key)
self._emit(JobEventType.JOB_CANCELLED, job)
self._enqueue_emit(self._prepare_emit(JobEventType.JOB_CANCELLED, job))
self._after_job_terminal_state_change()
return True
cancelled = True
self._drain_emits()
return cancelled

def cancel_all(self) -> int:
count = 0
Expand All @@ -299,9 +334,13 @@ def cancel_all(self) -> int:
job.status = QueuedJobStatus.CANCELLED
self._completed_keys.add(key)
self._running_keys.discard(key)
self._enqueue_emit(
self._prepare_emit(JobEventType.JOB_CANCELLED, job)
)
count += 1
if count > 0:
self._after_job_terminal_state_change()
self._drain_emits()
return count

def _after_job_terminal_state_change(self) -> None:
Expand All @@ -315,7 +354,7 @@ 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:
Expand All @@ -339,11 +378,13 @@ def _check_priority_advance(self) -> None:
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),
self._enqueue_emit(
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):
Expand All @@ -357,7 +398,9 @@ 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])
self._enqueue_emit(
self._prepare_emit(JobEventType.ALL_COMPLETE, self._jobs[first_key])
)

@property
def is_empty(self) -> bool:
Expand Down
76 changes: 71 additions & 5 deletions cli/tests/test_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@

from localci.core.config import LocalCIConfig
from localci.core.models import (
JobEvent,
JobEventType,
QueuedJob,
QueuedJobStatus,
Expand Down Expand Up @@ -142,11 +143,28 @@ 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
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
)


# ---------------------------------------------------------------------------
Expand Down Expand Up @@ -227,6 +245,43 @@ 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)
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):
queue = PriorityJobQueue()
queue.enqueue(make_job("Job", priority=1))
Expand Down Expand Up @@ -316,6 +371,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()
Expand Down
Loading