From 178aac8e4245c6721ed3126d8ac4ab4f8c83ac85 Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 08:46:35 +0200 Subject: [PATCH 1/2] Updated documentation to include rules engine setup --- README.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/README.md b/README.md index 165921b..1d91036 100644 --- a/README.md +++ b/README.md @@ -153,12 +153,15 @@ poetry run manage.py createsuperuser sudo cp systemd/iot.api.service /etc/systemd/system/iot.api.service sudo cp systemd/iot.mqttsubscriber.service /etc/systemd/system/iot.mqttsubscriber.service +sudo cp systemd/iot.rulesengine.service /etc/systemd/system/iot.rulesengine.service sudo systemctl start iot.api.service sudo systemctl start iot.mqttsubscriber.service +sudo systemctl start iot.rulesengine.service sudo systemctl enable iot.api.service sudo systemctl enable iot.mqttsubscriber.service +sudo systemctl enable iot.rulesengine.service ``` ### Make a SD Card Backup From 2c26662b401bb2456643e64a11128255451aa461 Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 08:47:02 +0200 Subject: [PATCH 2/2] Refactored rules engine to handle new devices added to the databse in realtime --- .../apps/device/management/commands/rules.py | 77 ++++++++++------ .../tests/management/commands/test_rules.py | 90 +++++++++++++------ systemd/iot.api.service | 1 + systemd/iot.mqttsubscriber.service | 3 +- systemd/iot.rulesengine.service | 1 + 5 files changed, 118 insertions(+), 54 deletions(-) diff --git a/iotserver/apps/device/management/commands/rules.py b/iotserver/apps/device/management/commands/rules.py index 0c143e1..bbc44dc 100644 --- a/iotserver/apps/device/management/commands/rules.py +++ b/iotserver/apps/device/management/commands/rules.py @@ -6,46 +6,73 @@ from iotserver.apps.device.models import Device from iotserver.apps.device.utils import rules +# Seconds between checks for newly added/removed active non-managed-firmware devices. +DEVICE_POLL_INTERVAL = 30 + class Command(BaseCommand): help = 'Run the rules engine for non-managed-firmware devices.' + def add_arguments(self, parser): + parser.add_argument( + '--poll-interval', + type=int, + default=DEVICE_POLL_INTERVAL, + help='Seconds between checks for added/removed devices.', + ) + def handle(self, *args, **options): """ - Spawn one worker thread per active, non-managed-firmware device and run - until interrupted. + Continuously reconciles worker threads against active, non-managed- + firmware devices, so devices added/removed/(de)activated while running + are picked up without a restart, until interrupted. """ - stop_event = threading.Event() - - devices = Device.objects.filter(managed_firmware=False, active=True) - threads = [ - threading.Thread( - target=rules.run_device, args=(device, stop_event), daemon=True - ) - for device in devices - ] - - if not threads: - self.stdout.write( - self.style.WARNING('No active non-managed-firmware devices found.') - ) - return + self.shutdown_event = threading.Event() def handle_shutdown(signum, frame): self.stdout.write(self.style.WARNING('Shutting down rules engine...')) - stop_event.set() + self.shutdown_event.set() signal.signal(signal.SIGINT, handle_shutdown) signal.signal(signal.SIGTERM, handle_shutdown) - for thread in threads: - thread.start() - - self.stdout.write( - self.style.SUCCESS(f'Rules engine started for {len(threads)} device(s).') - ) + workers = {} + while not self.shutdown_event.is_set(): + self._reconcile_workers(workers) + self.shutdown_event.wait(options['poll_interval']) - for thread in threads: + for stop_event, _thread in workers.values(): + stop_event.set() + for _stop_event, thread in workers.values(): thread.join() self.stdout.write(self.style.SUCCESS('Rules engine stopped.')) + + def _reconcile_workers(self, workers): + """ + Starts workers for newly active non-managed-firmware devices and stops + workers for devices which no longer match, mutating `workers` in place + (keyed by device id, valued by (stop_event, thread)). + """ + current_device_ids = set() + + for device in Device.objects.filter(managed_firmware=False, active=True): + current_device_ids.add(device.id) + if device.id not in workers: + stop_event = threading.Event() + thread = threading.Thread( + target=rules.run_device, args=(device, stop_event), daemon=True + ) + thread.start() + workers[device.id] = (stop_event, thread) + self.stdout.write( + self.style.SUCCESS(f'Started worker for device {device.id}.') + ) + + for device_id in [key for key in workers if key not in current_device_ids]: + stop_event, thread = workers.pop(device_id) + stop_event.set() + thread.join() + self.stdout.write( + self.style.WARNING(f'Stopped worker for device {device_id}.') + ) diff --git a/iotserver/apps/device/tests/management/commands/test_rules.py b/iotserver/apps/device/tests/management/commands/test_rules.py index d6f5d9d..00ddfa6 100644 --- a/iotserver/apps/device/tests/management/commands/test_rules.py +++ b/iotserver/apps/device/tests/management/commands/test_rules.py @@ -1,47 +1,81 @@ import pytest -from django.core.management import call_command +from iotserver.apps.device.management.commands.rules import Command from iotserver.apps.device.tests import factories as device_factories @pytest.mark.django_db -class TestRulesCommand: - def test_spawns_one_thread_per_active_non_managed_device(self, mocker): - matching_device = device_factories.DeviceFactory( - managed_firmware=False, - active=True, - ip_address='192.168.0.1', - mac_address='0E:00:20:01:71:AE', - ) - device_factories.DeviceFactory( - managed_firmware=True, - active=True, - ip_address='192.168.0.2', - mac_address='0E:00:20:01:71:AF', - ) - device_factories.DeviceFactory( - managed_firmware=False, - active=False, - ip_address='192.168.0.3', - mac_address='0E:00:20:01:71:B0', - ) - +class TestReconcileWorkers: + def test_starts_worker_for_new_active_device(self, mocker): + device = device_factories.DeviceFactory(managed_firmware=False, active=True) mock_thread_cls = mocker.patch('threading.Thread') + command = Command() + workers = {} - call_command('rules') + command._reconcile_workers(workers) - assert mock_thread_cls.call_count == 1 + assert device.id in workers + mock_thread_cls.assert_called_once() _, kwargs = mock_thread_cls.call_args - assert kwargs['args'][0] == matching_device + assert kwargs['args'][0] == device assert kwargs['daemon'] is True mock_thread_cls.return_value.start.assert_called_once() - mock_thread_cls.return_value.join.assert_called_once() - def test_exits_without_spawning_when_no_devices_match(self, mocker): + def test_ignores_managed_or_inactive_devices(self, mocker): device_factories.DeviceFactory(managed_firmware=True, active=True) + device_factories.DeviceFactory( + managed_firmware=False, + active=False, + ip_address='192.168.0.2', + mac_address='0E:00:20:01:71:AF', + ) + mock_thread_cls = mocker.patch('threading.Thread') + command = Command() + + command._reconcile_workers({}) + + mock_thread_cls.assert_not_called() + def test_does_not_restart_an_already_running_device(self, mocker): + device_factories.DeviceFactory(managed_firmware=False, active=True) mock_thread_cls = mocker.patch('threading.Thread') + command = Command() + workers = {} + command._reconcile_workers(workers) + mock_thread_cls.reset_mock() - call_command('rules') + command._reconcile_workers(workers) mock_thread_cls.assert_not_called() + + def test_stops_worker_for_device_no_longer_matching(self, mocker): + device = device_factories.DeviceFactory(managed_firmware=False, active=True) + mocker.patch('threading.Thread') + command = Command() + workers = {} + command._reconcile_workers(workers) + stop_event, thread = workers[device.id] + + device.active = False + device.save() + command._reconcile_workers(workers) + + assert device.id not in workers + assert stop_event.is_set() + thread.join.assert_called_once() + + +class TestHandle: + def test_polls_until_shutdown_event_is_set(self, mocker): + command = Command() + + def stop_after_first_call(workers): + command.shutdown_event.set() + + mock_reconcile = mocker.patch.object( + command, '_reconcile_workers', side_effect=stop_after_first_call + ) + + command.handle(poll_interval=0) + + mock_reconcile.assert_called_once() diff --git a/systemd/iot.api.service b/systemd/iot.api.service index 6a5498f..e1b533a 100644 --- a/systemd/iot.api.service +++ b/systemd/iot.api.service @@ -6,6 +6,7 @@ After=network.target User=pi Group=pi WorkingDirectory=/home/pi/iotserver +Environment="PATH=/home/pi/.local/bin:$PATH" ExecStart=/home/pi/.local/bin/poetry run gunicorn iotserver.wsgi:application -w 4 -b :8000 [Install] diff --git a/systemd/iot.mqttsubscriber.service b/systemd/iot.mqttsubscriber.service index 5d2a03d..172ca60 100644 --- a/systemd/iot.mqttsubscriber.service +++ b/systemd/iot.mqttsubscriber.service @@ -6,7 +6,8 @@ After=network.target User=pi Group=pi WorkingDirectory=/home/pi/iotserver -ExecStart=/home/pi/.local/bin/poetry run /home/pi/iotserver/manage.py mqtt +Environment="PATH=/home/pi/.local/bin:$PATH" +ExecStart=/home/pi/.local/bin/poetry run manage.py mqtt [Install] WantedBy=multi-user.target diff --git a/systemd/iot.rulesengine.service b/systemd/iot.rulesengine.service index 91af082..a436837 100644 --- a/systemd/iot.rulesengine.service +++ b/systemd/iot.rulesengine.service @@ -6,6 +6,7 @@ After=network.target User=pi Group=pi WorkingDirectory=/home/pi/iotserver +Environment="PATH=/home/pi/.local/bin:$PATH" ExecStart=/home/pi/.local/bin/poetry run /home/pi/iotserver/manage.py rules [Install]