From bba33aeac7b4bd8e66146b1960cfa258ccaf74d4 Mon Sep 17 00:00:00 2001 From: Elias Hernandis Date: Sun, 6 Sep 2026 23:00:49 +0200 Subject: [PATCH 1/2] Add operational lifecycle signals --- CHANGELOG.md | 5 + README.md | 18 ++- docs/alternatives.rst | 6 +- docs/configuration.rst | 21 +++ steady_queue/models/queue.py | 15 +- steady_queue/processes/runnable.py | 27 +++- steady_queue/processes/supervisor.py | 64 +++++++-- steady_queue/signals.py | 37 +++++ tests/test_observability_signals.py | 202 +++++++++++++++++++++++++++ 9 files changed, 367 insertions(+), 28 deletions(-) create mode 100644 steady_queue/signals.py create mode 100644 tests/test_observability_signals.py diff --git a/CHANGELOG.md b/CHANGELOG.md index fb605de..1653423 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,11 @@ ## Unreleased +**Added:** + +- Django signals for process start, stop, and restart events, and queue pause + and resume actions (#23). + ## v0.2.1 - 2026-09-06 **Added:** diff --git a/README.md b/README.md index e2c345c..c8e8c73 100644 --- a/README.md +++ b/README.md @@ -506,8 +506,17 @@ instance that is handling the task, as well as the `task_result` (a `django.tasks.TaskResult` instance) with information on how the task was called and its status. -Unlike Solid Queue, steady queue doesn't yet emit signals related to the -lifecycle of its processes. +Steady Queue also emits operational signals from ``steady_queue.signals``: + +- `process_started` and `process_stopped` when a supervisor or child process + enters or leaves its run loop. +- `process_restarted` when the supervisor replaces a terminated child. +- `queue_paused` and `queue_resumed` after queue control actions. + +Process signals include process identity and metadata. Queue signals include a +`changed` flag so receivers can distinguish state transitions from idempotent +actions. See the [configuration documentation](https://steady-queue.readthedocs.io/en/latest/configuration.html#signals) +for the complete payloads. ## Logging @@ -795,8 +804,9 @@ there are a few differences which we outline below. will ever be supported, but we've kept the database column for compatibility. - Steady Queue worker processes do not set the process name (or procline) because doing so requires introducing an external dependency. -- Steady Queue does not expose rich instrumentation like Solid Queue does due to - the lack of a framework-native equivalent to `ActiveSupport::Notifications`. +- Steady Queue exposes Django signals for task and operational lifecycle events, + but does not yet provide the full timed instrumentation event set emitted by + Solid Queue through `ActiveSupport::Notifications`. - **Priority ordering:** Steady Queue follows Django's convention where larger numbers indicate higher priority (e.g., a task with priority 10 runs before priority 0), whereas Solid Queue uses the inverse (smaller numbers = higher diff --git a/docs/alternatives.rst b/docs/alternatives.rst index 713a120..9f9b375 100644 --- a/docs/alternatives.rst +++ b/docs/alternatives.rst @@ -83,6 +83,6 @@ differences in the external interface are: - **Recurring tasks.** Solid Queue supports command-based recurring tasks (arbitrary shell commands on a schedule). Steady Queue only supports recurring Python task functions. -- **Instrumentation.** Solid Queue emits rich ``ActiveSupport::Notifications`` - events. Steady Queue uses standard Python logging and the ``django.tasks`` - signals instead. +- **Instrumentation.** Solid Queue emits a broader set of timed + ``ActiveSupport::Notifications`` events. Steady Queue uses standard Python + logging plus Django signals for task, process, and queue lifecycle events. diff --git a/docs/configuration.rst b/docs/configuration.rst index ce33aeb..8f982fc 100644 --- a/docs/configuration.rst +++ b/docs/configuration.rst @@ -222,6 +222,27 @@ Steady Queue emits the standard `django.tasks signals All signals include ``sender`` (the ``SteadyQueueBackend`` instance) and ``task_result`` (a ``django.tasks.TaskResult``). +Steady Queue also emits operational lifecycle signals from +``steady_queue.signals``: + +- ``process_started`` — after a supervisor, worker, dispatcher, or scheduler + has booted and registered. +- ``process_stopped`` — after a process run loop exits. The ``error`` argument + is ``None`` for a normal exit and contains the exception otherwise. +- ``process_restarted`` — after a supervisor replaces a terminated child. +- ``queue_paused`` and ``queue_resumed`` — after a queue control action. The + ``changed`` argument distinguishes a real state transition from an + idempotent repeat. + +Process signals use the public ``steady_queue.signals.ProcessLifecycle`` type +as their sender and include ``process_kind``, ``process_name``, ``pid``, +``hostname``, and ``metadata``. ``process_restarted`` additionally includes +``exitcode``, ``replacement_pid``, and ``supervisor_pid``. Queue signals use +the public ``steady_queue.signals.QueueLifecycle`` type as their sender and +include ``queue_name`` and ``changed``. Signal payloads deliberately contain +no Steady Queue runtime or model objects. These synchronous Django signals can +feed application logging, metrics or tracing integrations. + Logging ------- diff --git a/steady_queue/models/queue.py b/steady_queue/models/queue.py index d4eb70a..f893497 100644 --- a/steady_queue/models/queue.py +++ b/steady_queue/models/queue.py @@ -2,6 +2,7 @@ from steady_queue.models.pause import Pause from steady_queue.models.ready_execution import ReadyExecution +from steady_queue.signals import QueueLifecycle, queue_paused, queue_resumed class QueueQuerySet(models.QuerySet): @@ -48,10 +49,20 @@ def is_running(self) -> bool: return not self.is_paused def pause(self) -> None: - Pause.objects.get_or_create(queue_name=self.queue_name) + _, changed = Pause.objects.get_or_create(queue_name=self.queue_name) + queue_paused.send( + sender=QueueLifecycle, + queue_name=self.queue_name, + changed=changed, + ) def resume(self) -> None: - Pause.objects.filter(queue_name=self.queue_name).delete() + deleted, _ = Pause.objects.filter(queue_name=self.queue_name).delete() + queue_resumed.send( + sender=QueueLifecycle, + queue_name=self.queue_name, + changed=deleted > 0, + ) def __str__(self) -> str: return self.queue_name diff --git a/steady_queue/processes/runnable.py b/steady_queue/processes/runnable.py index d2b6705..424965c 100644 --- a/steady_queue/processes/runnable.py +++ b/steady_queue/processes/runnable.py @@ -1,6 +1,12 @@ import logging from steady_queue.processes.supervised import Supervised +from steady_queue.signals import ( + _process_signal_context, + _send_process_signal, + process_started, + process_stopped, +) logger = logging.getLogger("steady_queue") @@ -10,11 +16,22 @@ class Runnable(Supervised): def start(self): self.boot() - - if self.is_running_async: - raise NotImplementedError - else: - self.run() + signal_context = _process_signal_context(self) + _send_process_signal(process_started, self, context=signal_context) + + error = None + try: + if self.is_running_async: + raise NotImplementedError + else: + self.run() + except BaseException as exception: + error = exception + raise + finally: + _send_process_signal( + process_stopped, self, context=signal_context, error=error + ) def stop(self): super().stop() diff --git a/steady_queue/processes/supervisor.py b/steady_queue/processes/supervisor.py index d423a77..13871ce 100644 --- a/steady_queue/processes/supervisor.py +++ b/steady_queue/processes/supervisor.py @@ -15,6 +15,14 @@ from steady_queue.processes.registrable import Registrable from steady_queue.processes.signals import Signals from steady_queue.processes.timer import wait_until +from steady_queue.signals import ( + ProcessLifecycle, + _process_signal_context, + _send_process_signal, + process_restarted, + process_started, + process_stopped, +) logger = logging.getLogger("steady_queue") @@ -40,19 +48,35 @@ def __init__(self, configuration: Configuration): super().__init__() def start(self) -> None: - logger.info("starting supervisor with PID %(pid)d", {"pid": self.pid}) + originating_pid = self.pid + logger.info("starting supervisor with PID %(pid)d", {"pid": originating_pid}) + started = False + error = None + signal_context = None try: - self.boot() - # Fork only after resetting DB state (connections + psycopg pools). - self.reset_database_connections() - self.start_processes() - self.launch_maintenance_task() - except SystemExit: - logger.info("supervisor interrupted during boot, shutting down") - self.restore_default_signal_handlers() - self.shutdown() - return - self.supervise() + try: + self.boot() + started = True + signal_context = _process_signal_context(self) + _send_process_signal(process_started, self, context=signal_context) + # Fork only after resetting DB state (connections + psycopg pools). + self.reset_database_connections() + self.start_processes() + self.launch_maintenance_task() + except SystemExit: + logger.info("supervisor interrupted during boot, shutting down") + self.restore_default_signal_handlers() + self.shutdown() + return + self.supervise() + except BaseException as exception: + error = exception + raise + finally: + if started and self.pid == originating_pid: + _send_process_signal( + process_stopped, self, context=signal_context, error=error + ) def boot(self) -> None: super().boot() @@ -80,7 +104,7 @@ def supervise(self): logger.debug("supervisor finally block") self.shutdown() - def start_process(self, process: Configuration.Process) -> None: + def start_process(self, process: Configuration.Process) -> int: logger.info("starting process %(process)s", {"process": process}) instance = process.instantiate() instance.supervisor = self.process @@ -98,6 +122,7 @@ def start_process(self, process: Configuration.Process) -> None: self.reset_database_connections() self.configured_processes[pid] = process self.forks[pid] = instance + return pid def set_procline(self) -> None: pass @@ -165,7 +190,18 @@ def replace_fork(self, pid: int, exitcode: int) -> None: logger.info("replacing fork %s due to exit code %s", pid, exitcode) if terminated_fork := self.forks.pop(pid, None): self.handle_claimed_jobs_by(terminated_fork, exitcode) - self.start_process(self.configured_processes.pop(pid)) + replacement_pid = self.start_process(self.configured_processes.pop(pid)) + process_restarted.send( + sender=ProcessLifecycle, + process_kind=terminated_fork.kind, + process_name=terminated_fork.name, + pid=pid, + hostname=terminated_fork.hostname, + metadata=terminated_fork.metadata, + exitcode=exitcode, + replacement_pid=replacement_pid, + supervisor_pid=self.pid, + ) def handle_claimed_jobs_by(self, terminated_fork: Base, exitcode: int) -> None: if not self.process: diff --git a/steady_queue/signals.py b/steady_queue/signals.py new file mode 100644 index 0000000..c410adc --- /dev/null +++ b/steady_queue/signals.py @@ -0,0 +1,37 @@ +from django.dispatch import Signal + + +class ProcessLifecycle: + """Public sender for process lifecycle signals.""" + + +class QueueLifecycle: + """Public sender for queue lifecycle signals.""" + + +# Process lifecycle signals carry stable data suitable for logs and metrics. +process_started = Signal() +process_stopped = Signal() +process_restarted = Signal() + +# Queue control signals. Receivers get ``queue_name`` and ``changed``. +queue_paused = Signal() +queue_resumed = Signal() + + +def _process_signal_context(process) -> dict: + return { + "process_kind": process.kind, + "process_name": process.name, + "pid": process.pid, + "hostname": process.hostname, + "metadata": process.metadata, + } + + +def _send_process_signal(signal: Signal, process, *, context=None, **kwargs) -> None: + signal.send( + sender=ProcessLifecycle, + **(context or _process_signal_context(process)), + **kwargs, + ) diff --git a/tests/test_observability_signals.py b/tests/test_observability_signals.py new file mode 100644 index 0000000..c523ce3 --- /dev/null +++ b/tests/test_observability_signals.py @@ -0,0 +1,202 @@ +import os +from unittest.mock import MagicMock, patch + +from django.test import SimpleTestCase, TestCase + +from steady_queue.configuration import Configuration +from steady_queue.models import Job, Queue +from steady_queue.processes.base import Base +from steady_queue.processes.runnable import Runnable +from steady_queue.processes.supervisor import Supervisor +from steady_queue.signals import ( + ProcessLifecycle, + QueueLifecycle, + process_restarted, + process_started, + process_stopped, + queue_paused, + queue_resumed, +) +from tests.dummy.tasks import dummy_task + + +class FakeRunnable(Runnable, Base): + mode = "inline" + + def run(self): + pass + + +class ProcessLifecycleSignalTests(SimpleTestCase): + def test_runnable_emits_started_and_stopped_with_operational_context(self): + process = FakeRunnable() + started = MagicMock() + stopped = MagicMock() + + process_started.connect(started) + process_stopped.connect(stopped) + self.addCleanup(process_started.disconnect, started) + self.addCleanup(process_stopped.disconnect, stopped) + + process.start() + + started.assert_called_once_with( + signal=process_started, + sender=ProcessLifecycle, + process_kind="fakerunnable", + process_name=process.name, + pid=process.pid, + hostname=process.hostname, + metadata={}, + ) + stopped.assert_called_once_with( + signal=process_stopped, + sender=ProcessLifecycle, + process_kind="fakerunnable", + process_name=process.name, + pid=process.pid, + hostname=process.hostname, + metadata={}, + error=None, + ) + + def test_stopped_signal_includes_run_error(self): + process = FakeRunnable() + error = RuntimeError("broken") + stopped = MagicMock() + process.run = MagicMock(side_effect=error) + + process_stopped.connect(stopped) + self.addCleanup(process_stopped.disconnect, stopped) + + with self.assertRaises(RuntimeError): + process.start() + + self.assertIs(stopped.call_args.kwargs["error"], error) + + def test_forked_child_does_not_emit_supervisor_stopped(self): + supervisor = Supervisor(Configuration(Configuration.Options())) + originating_pid = os.getpid() + read_fd, write_fd = os.pipe() + child_pids = [] + + def record_stop(**kwargs): + message = f"{os.getpid()}:{kwargs['pid']}\n".encode() + os.write(write_fd, message) + + def fork_child(): + pid = os.fork() + if pid == 0: + os.close(read_fd) + raise SystemExit + child_pids.append(pid) + + process_stopped.connect(record_stop) + try: + with ( + patch.object(supervisor, "boot"), + patch.object(supervisor, "reset_database_connections"), + patch.object(supervisor, "start_processes", side_effect=fork_child), + patch.object(supervisor, "launch_maintenance_task"), + patch.object(supervisor, "restore_default_signal_handlers"), + patch.object(supervisor, "shutdown"), + patch.object(supervisor, "supervise"), + ): + supervisor.start() + + if os.getpid() != originating_pid: + os.close(write_fd) + os._exit(0) + + os.waitpid(child_pids[0], 0) + os.close(write_fd) + write_fd = None + records = os.read(read_fd, 4096).decode().splitlines() + finally: + process_stopped.disconnect(record_stop) + os.close(read_fd) + if write_fd is not None: + os.close(write_fd) + + self.assertEqual(records, [f"{originating_pid}:{originating_pid}"]) + + def test_supervisor_emits_restart_after_replacing_child(self): + supervisor = Supervisor(Configuration(Configuration.Options())) + terminated = FakeRunnable() + configured = MagicMock() + supervisor.forks[123] = terminated + supervisor.configured_processes[123] = configured + restarted = MagicMock() + + process_restarted.connect(restarted) + self.addCleanup(process_restarted.disconnect, restarted) + + with ( + patch.object(supervisor, "handle_claimed_jobs_by"), + patch.object(supervisor, "start_process", return_value=456), + ): + supervisor.replace_fork(123, 9) + + restarted.assert_called_once_with( + signal=process_restarted, + sender=ProcessLifecycle, + process_kind="fakerunnable", + process_name=terminated.name, + pid=123, + hostname=terminated.hostname, + metadata={}, + exitcode=9, + replacement_pid=456, + supervisor_pid=supervisor.pid, + ) + + +class QueueLifecycleSignalTests(TestCase): + databases = {"default", "queue"} + + def setUp(self): + Job.objects.enqueue(dummy_task, [], {}) + self.queue = Queue.objects.get(queue_name="default") + + def test_pause_and_resume_report_real_and_idempotent_actions(self): + paused = MagicMock() + resumed = MagicMock() + queue_paused.connect(paused) + queue_resumed.connect(resumed) + self.addCleanup(queue_paused.disconnect, paused) + self.addCleanup(queue_resumed.disconnect, resumed) + + self.queue.pause() + self.queue.pause() + self.queue.resume() + self.queue.resume() + + self.assertEqual( + [call.kwargs["changed"] for call in paused.call_args_list], [True, False] + ) + self.assertEqual( + [call.kwargs["changed"] for call in resumed.call_args_list], [True, False] + ) + self.assertTrue( + all( + call.kwargs["queue_name"] == "default" for call in paused.call_args_list + ) + ) + self.assertTrue( + all("queue" not in call.kwargs for call in paused.call_args_list) + ) + self.assertTrue( + all("queue" not in call.kwargs for call in resumed.call_args_list) + ) + self.assertTrue( + all( + call.kwargs["sender"] is QueueLifecycle + for call in paused.call_args_list + ) + ) + self.assertTrue( + all( + call.kwargs["sender"] is QueueLifecycle + for call in resumed.call_args_list + ) + ) From 00bf757c9a4df3f87835413dfbbd672d7ab1b38c Mon Sep 17 00:00:00 2001 From: Elias Hernandis Date: Sun, 6 Sep 2026 23:08:04 +0200 Subject: [PATCH 2/2] Decode child wait status in restart signals --- steady_queue/processes/supervisor.py | 9 +++++---- tests/test_observability_signals.py | 7 ++++--- 2 files changed, 9 insertions(+), 7 deletions(-) diff --git a/steady_queue/processes/supervisor.py b/steady_queue/processes/supervisor.py index 13871ce..9289c42 100644 --- a/steady_queue/processes/supervisor.py +++ b/steady_queue/processes/supervisor.py @@ -158,14 +158,14 @@ def quit_forks(self) -> None: def reap_and_replace_terminated_forks(self) -> None: while True: try: - pid, exitcode = os.waitpid(-1, os.WNOHANG) + pid, wait_status = os.waitpid(-1, os.WNOHANG) except ChildProcessError: break else: if not pid: break - self.replace_fork(pid, exitcode) + self.replace_fork(pid, wait_status) def reap_terminated_forks(self) -> None: while True: @@ -186,10 +186,11 @@ def reap_terminated_forks(self) -> None: self.configured_processes.pop(pid, None) - def replace_fork(self, pid: int, exitcode: int) -> None: + def replace_fork(self, pid: int, wait_status: int) -> None: + exitcode = os.waitstatus_to_exitcode(wait_status) logger.info("replacing fork %s due to exit code %s", pid, exitcode) if terminated_fork := self.forks.pop(pid, None): - self.handle_claimed_jobs_by(terminated_fork, exitcode) + self.handle_claimed_jobs_by(terminated_fork, wait_status) replacement_pid = self.start_process(self.configured_processes.pop(pid)) process_restarted.send( sender=ProcessLifecycle, diff --git a/tests/test_observability_signals.py b/tests/test_observability_signals.py index c523ce3..c7d8188 100644 --- a/tests/test_observability_signals.py +++ b/tests/test_observability_signals.py @@ -132,11 +132,12 @@ def test_supervisor_emits_restart_after_replacing_child(self): self.addCleanup(process_restarted.disconnect, restarted) with ( - patch.object(supervisor, "handle_claimed_jobs_by"), + patch.object(supervisor, "handle_claimed_jobs_by") as recover_jobs, patch.object(supervisor, "start_process", return_value=456), ): - supervisor.replace_fork(123, 9) + supervisor.replace_fork(123, 7 << 8) + recover_jobs.assert_called_once_with(terminated, 7 << 8) restarted.assert_called_once_with( signal=process_restarted, sender=ProcessLifecycle, @@ -145,7 +146,7 @@ def test_supervisor_emits_restart_after_replacing_child(self): pid=123, hostname=terminated.hostname, metadata={}, - exitcode=9, + exitcode=7, replacement_pid=456, supervisor_pid=supervisor.pid, )