From 8d4f4253fd1e08485f66688728b3ea5397303149 Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Wed, 23 Sep 2026 19:28:45 +0200 Subject: [PATCH 01/12] Added field to device to indicate if the firmware is managed --- iotserver/apps/device/models.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/iotserver/apps/device/models.py b/iotserver/apps/device/models.py index 33badba..5506f3b 100644 --- a/iotserver/apps/device/models.py +++ b/iotserver/apps/device/models.py @@ -52,6 +52,7 @@ class Device(models.Model): updated_at = models.DateTimeField(auto_now=True) active = models.BooleanField(default=False) + managed_firmware = models.BooleanField(default=True) name = models.CharField(max_length=128) description = models.CharField(max_length=1024) @@ -205,7 +206,7 @@ def status(self): @receiver(pre_save, sender=Device) def handle_device_default_config(sender, instance, *args, **kwargs): """Get the config from the new device and update the config field.""" - if settings.AUTO_SYNC_DEVICE: + if settings.AUTO_SYNC_DEVICE and instance.managed_firmware and instance.active: if instance.config is None: temp_file_path = f'/tmp/config.{instance.id}.json' with open(temp_file_path, 'w') as input_file: @@ -232,7 +233,7 @@ def handle_device_default_config(sender, instance, *args, **kwargs): @receiver(post_save, sender=Device) def handle_device_config_update(sender, instance, *args, **kwargs): """Update config on the physical device via webrepl.""" - if settings.AUTO_SYNC_DEVICE: + if settings.AUTO_SYNC_DEVICE and instance.managed_firmware: temp_file_path = f'/tmp/config.{instance.id}.json' with open(temp_file_path, 'w') as input_file: input_file.write(json.dumps(instance.full_config, indent=4)) From 649d99505627e2b9ac8e03ff4668be29ea08628d Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Wed, 23 Sep 2026 19:30:06 +0200 Subject: [PATCH 02/12] Added device migration --- .../migrations/0012_device_managed_firmware.py | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) create mode 100644 iotserver/apps/device/migrations/0012_device_managed_firmware.py diff --git a/iotserver/apps/device/migrations/0012_device_managed_firmware.py b/iotserver/apps/device/migrations/0012_device_managed_firmware.py new file mode 100644 index 0000000..609ef74 --- /dev/null +++ b/iotserver/apps/device/migrations/0012_device_managed_firmware.py @@ -0,0 +1,17 @@ +# Generated by Django 5.2.17 on 2026-09-23 17:27 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + dependencies = [ + ('device', '0011_devicepintype_devicetype_identifier_devicepin_type'), + ] + + operations = [ + migrations.AddField( + model_name='device', + name='managed_firmware', + field=models.BooleanField(default=True), + ), + ] From c1abafcc55a900fd5c9f67cbd4be0e788511c71c Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Wed, 23 Sep 2026 19:46:00 +0200 Subject: [PATCH 03/12] Improved the json widget to mask known secrets when dispalyed in admin --- iotserver/apps/device/tests/test_widgets.py | 159 ++++++++++++++++++++ iotserver/apps/device/widgets.py | 99 +++++++++++- 2 files changed, 255 insertions(+), 3 deletions(-) create mode 100644 iotserver/apps/device/tests/test_widgets.py diff --git a/iotserver/apps/device/tests/test_widgets.py b/iotserver/apps/device/tests/test_widgets.py new file mode 100644 index 0000000..b5be02e --- /dev/null +++ b/iotserver/apps/device/tests/test_widgets.py @@ -0,0 +1,159 @@ +import json + +from iotserver.apps.device.widgets import MASK_VALUE, PrettyJSONWidget + + +class TestPrettyJSONWidget(object): + def setup_method(self, test_method): + self.widget = PrettyJSONWidget() + + # -- format_value ----------------------------------------------------- + + def test_format_value_none(self): + assert self.widget.format_value(None) is None + + def test_format_value_empty_string(self): + assert self.widget.format_value("") is None + + def test_format_value_json_string(self): + value = json.dumps({'b': 1, 'a': 2}) + + result = self.widget.format_value(value) + + assert result == json.dumps({'a': 2, 'b': 1}, indent=4, sort_keys=True) + + def test_format_value_dict(self): + result = self.widget.format_value({'b': 1, 'a': 2}) + + assert result == json.dumps({'a': 2, 'b': 1}, indent=4, sort_keys=True) + + def test_format_value_list(self): + result = self.widget.format_value([3, 1, 2]) + + assert result == json.dumps([3, 1, 2], indent=4, sort_keys=True) + + def test_format_value_invalid_json_string_returned_unchanged(self): + value = '{not valid json' + + assert self.widget.format_value(value) == value + + def test_format_value_non_json_scalar_returned_unchanged(self): + assert self.widget.format_value(123) == 123 + + def test_format_value_masks_top_level_secret_keys(self): + value = { + 'password': 'hunter2', + 'webrepl_password': 'secret', + 'auth_header': 'Bearer xyz', + } + + result = json.loads(self.widget.format_value(value)) + + assert result == { + 'password': MASK_VALUE, + 'webrepl_password': MASK_VALUE, + 'auth_header': MASK_VALUE, + } + + def test_format_value_masks_nested_secret_keys(self): + value = {'wifi': {'ssid': 'home', 'password': 'hunter2'}, 'name': 'device'} + + result = json.loads(self.widget.format_value(value)) + + assert result == { + 'wifi': {'ssid': 'home', 'password': MASK_VALUE}, + 'name': 'device', + } + + def test_format_value_masks_secret_keys_in_list_of_dicts(self): + value = { + 'endpoints': [{'auth_header': 'Bearer xyz'}, {'auth_header': 'Bearer abc'}] + } + + result = json.loads(self.widget.format_value(value)) + + assert result == { + 'endpoints': [{'auth_header': MASK_VALUE}, {'auth_header': MASK_VALUE}] + } + + def test_format_value_does_not_mask_empty_secret_value(self): + result = json.loads(self.widget.format_value({'password': ''})) + + assert result == {'password': ''} + + def test_format_value_is_case_insensitive_for_secret_keys(self): + result = json.loads(self.widget.format_value({'Password': 'hunter2'})) + + assert result == {'Password': MASK_VALUE} + + # -- render ------------------------------------------------------------- + + def test_render_includes_masked_visible_value(self): + html = self.widget.render('config', {'password': 'hunter2'}) + visible_textarea = html.split('__pretty_json_original')[0] + + assert MASK_VALUE in visible_textarea + assert 'hunter2' not in visible_textarea + + def test_render_includes_hidden_original_value(self): + html = self.widget.render('config', {'password': 'hunter2'}) + + assert 'name="config__pretty_json_original"' in html + assert 'hunter2' in html + + def test_render_hidden_original_empty_for_none_value(self): + html = self.widget.render('config', None) + + assert 'name="config__pretty_json_original" value=""' in html + + # -- value_from_datadict -------------------------------------------- + + def test_value_from_datadict_restores_untouched_masked_value(self): + original = {'ssid': 'home', 'password': 'hunter2'} + data = { + 'config': json.dumps({'ssid': 'home', 'password': MASK_VALUE}), + 'config__pretty_json_original': json.dumps(original), + } + + result = json.loads(self.widget.value_from_datadict(data, {}, 'config')) + + assert result == original + + def test_value_from_datadict_keeps_new_secret_value(self): + original = {'ssid': 'home', 'password': 'hunter2'} + data = { + 'config': json.dumps({'ssid': 'home', 'password': 'new-password'}), + 'config__pretty_json_original': json.dumps(original), + } + + result = json.loads(self.widget.value_from_datadict(data, {}, 'config')) + + assert result == {'ssid': 'home', 'password': 'new-password'} + + def test_value_from_datadict_restores_nested_untouched_masked_value(self): + original = {'wifi': {'ssid': 'home', 'password': 'hunter2'}} + data = { + 'config': json.dumps({'wifi': {'ssid': 'home', 'password': MASK_VALUE}}), + 'config__pretty_json_original': json.dumps(original), + } + + result = json.loads(self.widget.value_from_datadict(data, {}, 'config')) + + assert result == original + + def test_value_from_datadict_without_original_returns_submitted_unchanged(self): + data = {'config': json.dumps({'password': MASK_VALUE})} + + result = self.widget.value_from_datadict(data, {}, 'config') + + assert result == data['config'] + + def test_value_from_datadict_non_structured_submission_returned_unchanged(self): + data = { + 'config': 'not json', + 'config__pretty_json_original': json.dumps({'password': 'hunter2'}), + } + + result = self.widget.value_from_datadict(data, {}, 'config') + + assert result == 'not json' diff --git a/iotserver/apps/device/widgets.py b/iotserver/apps/device/widgets.py index 9b41204..4df4ee6 100644 --- a/iotserver/apps/device/widgets.py +++ b/iotserver/apps/device/widgets.py @@ -1,16 +1,109 @@ import json from django.forms import Textarea +from django.utils.html import format_html +from django.utils.safestring import mark_safe + +MASK_VALUE = '**********' +MASKED_KEYS = frozenset({'webrepl_password', 'password', 'auth_header'}) + + +def _is_masked_key(key): + return isinstance(key, str) and key.lower() in MASKED_KEYS + + +def _mask_secrets(value): + """Recursively replace values of masked keys with MASK_VALUE.""" + if isinstance(value, dict): + return { + key: MASK_VALUE + if _is_masked_key(key) and val not in (None, '') + else _mask_secrets(val) + for key, val in value.items() + } + if isinstance(value, list): + return [_mask_secrets(item) for item in value] + return value + + +def _restore_masked(submitted, original): + """Recursively replace untouched MASK_VALUE placeholders with their original values.""" + if isinstance(submitted, dict): + original_dict = original if isinstance(original, dict) else {} + restored = {} + for key, val in submitted.items(): + if _is_masked_key(key) and val == MASK_VALUE and key in original_dict: + restored[key] = original_dict[key] + else: + restored[key] = _restore_masked(val, original_dict.get(key)) + return restored + if isinstance(submitted, list): + original_list = original if isinstance(original, list) else [] + return [ + _restore_masked( + item, original_list[index] if index < len(original_list) else None + ) + for index, item in enumerate(submitted) + ] + return submitted class PrettyJSONWidget(Textarea): - """A widget that pretty-prints JSON data.""" + """A widget that pretty-prints JSON data and masks sensitive keys on display. + + Masked keys (see MASKED_KEYS) are never rendered in plain text. A hidden + companion field carries the real, unmasked value across the request so + that untouched masked fields are restored instead of being saved as the + literal mask placeholder. + """ + + ORIGINAL_FIELD_SUFFIX = '__pretty_json_original' + + def original_field_name(self, name): + return f'{name}{self.ORIGINAL_FIELD_SUFFIX}' def format_value(self, value): - if value == "" or value is None: + if value is None or value == '': return None + # value may already be a parsed dict/list (e.g. a field default) rather than a string. + if isinstance(value, (dict, list)): + return json.dumps(_mask_secrets(value), indent=4, sort_keys=True) try: parsed = json.loads(value) - return json.dumps(parsed, indent=4, sort_keys=True) except (TypeError, ValueError): return value + if isinstance(parsed, (dict, list)): + parsed = _mask_secrets(parsed) + return json.dumps(parsed, indent=4, sort_keys=True) + + def render(self, name, value, attrs=None, renderer=None): + rendered = super().render(name, value, attrs, renderer) + original = self._to_python(value) + original_json = '' if original is None else json.dumps(original) + hidden_input = format_html( + '', + self.original_field_name(name), + original_json, + ) + return mark_safe(f'{rendered}{hidden_input}') + + def value_from_datadict(self, data, files, name): + submitted = super().value_from_datadict(data, files, name) + original = self._to_python(data.get(self.original_field_name(name))) + submitted_parsed = self._to_python(submitted) + + if original is None or not isinstance(submitted_parsed, (dict, list)): + return submitted + + return json.dumps(_restore_masked(submitted_parsed, original)) + + @staticmethod + def _to_python(value): + if value is None or value == '': + return None + if isinstance(value, (dict, list)): + return value + try: + return json.loads(value) + except (TypeError, ValueError): + return None From fa2e794daf98a2895ab89c3532452f314045964f Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Wed, 23 Sep 2026 20:04:35 +0200 Subject: [PATCH 04/12] Added managed firmware flag to device admin --- iotserver/apps/device/admin.py | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/iotserver/apps/device/admin.py b/iotserver/apps/device/admin.py index ec6bc7a..4871a70 100644 --- a/iotserver/apps/device/admin.py +++ b/iotserver/apps/device/admin.py @@ -33,7 +33,19 @@ class DeviceTypeModelAdmin(admin.ModelAdmin): class DeviceModelAdmin(admin.ModelAdmin): actions = [toggle_devices_on, toggle_devices_off] fieldsets = ( - (None, {'fields': ('active', 'name', 'description', 'type', 'location')}), + ( + None, + { + 'fields': ( + 'active', + 'managed_firmware', + 'name', + 'description', + 'type', + 'location', + ) + }, + ), ( 'Advanced options', { @@ -53,12 +65,13 @@ class DeviceModelAdmin(admin.ModelAdmin): 'name', 'description', 'active', + 'managed_firmware', 'created_at', 'type', 'location', 'ip_address', ) - list_filter = ('active', 'type__name', 'location__name') + list_filter = ('active', 'managed_firmware', 'type__name', 'location__name') formfield_overrides = { JSONField: {'widget': widgets.PrettyJSONWidget(attrs={'rows': 20, 'cols': 120})} } From ec93c693b4fe8e68ddcad82ed3e1e6c049212054 Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Wed, 23 Sep 2026 20:16:40 +0200 Subject: [PATCH 05/12] Added tests for device model signals --- iotserver/apps/device/tests/test_models.py | 53 ++++++++++++++++++++++ 1 file changed, 53 insertions(+) diff --git a/iotserver/apps/device/tests/test_models.py b/iotserver/apps/device/tests/test_models.py index 5a38353..6dbafb4 100644 --- a/iotserver/apps/device/tests/test_models.py +++ b/iotserver/apps/device/tests/test_models.py @@ -1,5 +1,8 @@ +import json + import pytest +from iotserver.apps.device import models from iotserver.apps.device.tests import factories as device_factories @@ -94,6 +97,56 @@ def test_mqtt_toggle(self, mocker): self.device.mqtt_toggle('on') mock_mqtt_toggle.assert_called_once_with(self.device.id, '1') + def test_handle_device_default_config(self, settings, mocker): + settings.AUTO_SYNC_DEVICE = True + self.device.managed_firmware = True + self.device.active = True + self.device.config = None + + mock_web_socket = mocker.Mock() + mock_get_websocket = mocker.patch( + 'iotserver.apps.device.models.webrepl.get_websocket', + return_value=(mocker.Mock(), mock_web_socket), + ) + + def fake_get_file(web_socket, path, remote_path): + with open(path, 'w') as input_file: + input_file.write(json.dumps({'main': {}, 'pins': []})) + + mocker.patch( + 'iotserver.apps.device.models.webrepl.get_file', side_effect=fake_get_file + ) + + models.handle_device_default_config(sender=models.Device, instance=self.device) + + mock_get_websocket.assert_called_once_with( + self.device.ip_address, settings.WEBREPL_PORT, settings.WEBREPL_PASSWORD + ) + assert self.device.config == {'main': {'identifier': str(self.device.id)}} + + def test_handle_device_config_update(self, settings, mocker): + settings.AUTO_SYNC_DEVICE = True + self.device.managed_firmware = True + self.device.config = {'main': {}} + + mock_web_socket = mocker.Mock() + mock_get_websocket = mocker.patch( + 'iotserver.apps.device.models.webrepl.get_websocket', + return_value=(mocker.Mock(), mock_web_socket), + ) + mock_put_file = mocker.patch('iotserver.apps.device.models.webrepl.put_file') + + models.handle_device_config_update(sender=models.Device, instance=self.device) + + mock_get_websocket.assert_called_once_with( + self.device.ip_address, settings.WEBREPL_PORT, settings.WEBREPL_PASSWORD + ) + mock_put_file.assert_called_once_with( + mock_web_socket, + f'/tmp/config.{self.device.id}.json', + 'config/config.json', + ) + @pytest.mark.django_db class TestDevicePinModel(object): From 801e1c12f0538ea226ea20e0cc8aa6a72f6a8abc Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 07:50:14 +0200 Subject: [PATCH 06/12] Cleanup env example file --- .env.example | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/.env.example b/.env.example index 828c8ba..f23056e 100644 --- a/.env.example +++ b/.env.example @@ -2,9 +2,11 @@ IOTSERVER_SECRET_KEY='insecure-secret' IOTSERVER_CORS_ORIGIN_WHITELIST='http://localhost:3000' IOTSERVER_CORS_ALLOW_ALL_ORIGINS='0' + IOTSERVER_POSTGRES_USER='postgres' IOTSERVER_POSTGRES_PASSWORD='insecure-password' IOTSERVER_POSTGRES_DBNAME='iotserver' + IOTSERVER_SONOFF_AUTH_URL='https://eu-api.coolkit.cc:8080/api/user/login' IOTSERVER_SONOFF_EMAIL='test@example.com' IOTSERVER_SONOFF_PASSWORD='insecure-password' @@ -13,14 +15,18 @@ IOTSERVER_SONOFF_COUNTRY_CODE='+27' IOTSERVER_SONOFF_APP_ID='sonoff-app-id' IOTSERVER_SONOFF_APP_SECRET='sonoff-app-secret' IOTSERVER_SONOFF_DEVICE_URL='https://eu-api.coolkit.cc:8080/api/user/device/status' + IOTSERVER_SOLARMAN_BASE_URL='https://globalapi.solarmanpv.com' IOTSERVER_SOLARMAN_APP_ID='solarman-app-id' IOTSERVER_SOLARMAN_APP_SECRET='solarman-app-secret' IOTSERVER_SOLARMAN_EMAIL='solarman-email@example.com' IOTSERVER_SOLARMAN_PASSWORD='solarman-password' + IOTSERVER_OPENWEATHER_URL='https://api.openweathermap.org/data/3.0/onecall' IOTSERVER_OPENWEATHER_APIKEY='openweather-key' + IOTSERVER_GOOGLEMAPS_APIKEY='googlemaps-key' + IOTSERVER_WEBREPL_PORT='8266' IOTSERVER_WEBREPL_PASSWORD='insecure-password' IOTSERVER_AUTO_SYNC_DEVICE=1 @@ -30,3 +36,4 @@ DOCKER_POSTGIS_IMAGE='kartoza/postgis:15' # GDAL local path IOTSERVER_GDAL_LIBRARY_PATH='/path/to/your/libgdal.so' +IOTSERVER_GEOS_LIBRARY_PATH='/path/to/your/libgeos_c.so' From 48f0cb4acfacaa846ec3cdf52ec56c2aba00896c Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 07:50:35 +0200 Subject: [PATCH 07/12] Added initial rule utils --- iotserver/apps/device/utils/rules.py | 114 +++++++++++++++++++++++++++ 1 file changed, 114 insertions(+) create mode 100644 iotserver/apps/device/utils/rules.py diff --git a/iotserver/apps/device/utils/rules.py b/iotserver/apps/device/utils/rules.py new file mode 100644 index 0000000..e993488 --- /dev/null +++ b/iotserver/apps/device/utils/rules.py @@ -0,0 +1,114 @@ +import requests +from django.utils import timezone + +from iotserver.apps.device.integrations.sonoff import Sonoff + +CONDITION_OPERATORS = { + 'eq': lambda input, value: input == value, + 'gt': lambda input, value: input > value, + 'lt': lambda input, value: input < value, +} + + +def find_xpath_value(response, xpaths): + """ + Recursively finds a value in a nested dict/list structure based on a list of + xpaths. `xpaths` must be reversed before being passed in. + """ + xpath = xpaths.pop() + + try: + response = response[xpath] + except (KeyError, IndexError, TypeError): + return None + + if not xpaths: + return response + + return find_xpath_value(response, xpaths) + + +def evaluate_condition(input, operator, value): + """ + Evaluates a condition and returns a boolean. + """ + try: + return CONDITION_OPERATORS[operator](input, value) + except TypeError: + return False + + +def handle_conditions(rule_values, input_value): + """ + Returns a dict of `must`/`should` condition boolean lists to be evaluated. + """ + condition_values = {'must': [], 'should': []} + for condition_type, conditions in input_value['conditions'].items(): + if condition_type in condition_values: + for pin_identifier, condition in conditions.items(): + xpaths = pin_identifier.split('.') + xpaths.reverse() + condition_values[condition_type].append( + evaluate_condition( + find_xpath_value(rule_values, xpaths), **condition + ) + ) + + return condition_values + + +def value_to_bool(value): + """ + Converts a string, int, bool, or None to a boolean. + """ + if isinstance(value, str): + return value.strip().lower() in ['true', '1', 'yes', 'on'] + return bool(value) + + +def timer(**kwargs): + """ + Checks if the current local time is within a start/end time range, handling + ranges which span midnight (e.g. 18:00 - 05:00). + """ + now = timezone.localtime().strftime('%H%M') + start = kwargs.get('start_time').replace(':', '') + end = kwargs.get('end_time').replace(':', '') + + if start <= end: + return start <= now <= end + return now >= start or now <= end + + +def service(**kwargs): + """ + Calls a web service and returns its JSON response. + """ + url = kwargs.get('url') + auth_header = kwargs.get('auth_header') + headers = {'Authorization': f'Token {auth_header}'} if auth_header else None + + response = requests.get(url, headers=headers) + response.raise_for_status() + + return response.json() + + +def mqtt_toggle(**kwargs): + """ + Reads the latest retained MQTT message received for `topic` by the device + worker's MQTT client and returns it as a boolean. + """ + topic = kwargs.get('topic') + mqtt_values = kwargs.get('mqtt_values', {}) + return value_to_bool(mqtt_values.get(topic)) + + +def sonoff_toggle(**kwargs): + """ + Toggles a Sonoff cloud device on/off via the eWeLink/Coolkit API. + """ + on = kwargs.get('on') + device_id = kwargs.get('device_id') + state = Sonoff(device_id).toggle_device('on' if on else 'off') + return value_to_bool(state) From 2be42de016f55771ee7fb1ff4702f2d04137ecd0 Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 07:57:19 +0200 Subject: [PATCH 08/12] Implimented rule runner worker loop --- iotserver/apps/device/utils/rules.py | 161 +++++++++++++++++++++++++++ 1 file changed, 161 insertions(+) diff --git a/iotserver/apps/device/utils/rules.py b/iotserver/apps/device/utils/rules.py index e993488..72f4b6d 100644 --- a/iotserver/apps/device/utils/rules.py +++ b/iotserver/apps/device/utils/rules.py @@ -1,8 +1,20 @@ +import hashlib +import json +import logging +import threading + +import paho.mqtt.client as mqtt import requests from django.utils import timezone from iotserver.apps.device.integrations.sonoff import Sonoff +logger = logging.getLogger(__name__) + +# Ordering (not alphabetical) matches the IoTDevice firmware's LOG_LEVELS so the +# same `logging.level` config value has the same meaning on both sides. +LOG_LEVELS = ['info', 'debug', 'warning', 'error'] + CONDITION_OPERATORS = { 'eq': lambda input, value: input == value, 'gt': lambda input, value: input > value, @@ -112,3 +124,152 @@ def sonoff_toggle(**kwargs): device_id = kwargs.get('device_id') state = Sonoff(device_id).toggle_device('on' if on else 'off') return value_to_bool(state) + + +# Explicit whitelist of callable rule actions, rather than `getattr` on this +# module, so `rule['action']` (data, not code) can't invoke arbitrary attributes. +RULE_ACTIONS = { + 'timer': timer, + 'service': service, + 'mqtt_toggle': mqtt_toggle, + 'sonoff_toggle': sonoff_toggle, +} + + +def _build_mqtt_client(mqtt_config, on_connect, on_message): + """ + Builds, connects, and starts the network loop for a device's persistent MQTT + client. + """ + client = mqtt.Client( + mqtt.CallbackAPIVersion.VERSION2, client_id=mqtt_config['client_id'] + ) + client.on_connect = on_connect + client.on_message = on_message + + username = mqtt_config.get('username') + if username: + client.username_pw_set(username, mqtt_config.get('password')) + + if mqtt_config.get('ssl_enabled'): + client.tls_set() + + lastwill = mqtt_config.get('lastwill') + if lastwill: + client.will_set(topic=lastwill['topic'], payload=lastwill['message']) + + client.connect(host=mqtt_config['host'], port=mqtt_config.get('port', 1883)) + client.loop_start() + + return client + + +def _resolve_rule_params(rule, rule_values, mqtt_values): + """ + Resolves a rule's `input` into keyword arguments, evaluating any must/should + conditions against previously collected rule values. + """ + rule_params = {} + for key, value in rule['input'].items(): + if isinstance(value, dict) and 'conditions' in value: + condition_values = handle_conditions(rule_values, value) + rule_params[key] = any( + [all(condition_values['must']), any(condition_values['should'])] + ) + else: + rule_params[key] = value + + if rule['action'] == 'mqtt_toggle': + rule_params['mqtt_values'] = mqtt_values + + return rule_params + + +def _publish_status(client, device_id, rule_values, previous_status_hash): + """ + Publishes the current rule values as the device status if they've changed + since the last publish, returning the new dedup hash. + """ + payload = json.dumps(rule_values, sort_keys=True) + status_hash = hashlib.sha1(payload.encode()).digest() + + if status_hash != previous_status_hash: + client.publish(f'iot-devices/{device_id}/status', payload) + + return status_hash + + +def _publish_log(client, device_id, logging_config, level, message): + """ + Logs locally and, if the configured logging level permits it, publishes the + message to the device's log topic for ingestion by the `mqtt` command. + """ + getattr(logger, level, logger.info)(message) + + threshold = logging_config.get('level', 'warning') + if LOG_LEVELS.index(level) >= LOG_LEVELS.index(threshold): + client.publish(f'iot-devices/{device_id}/logs', message) + + +def run_device(device, stop_event: threading.Event) -> None: + """ + Runs the rules engine loop for a single non-managed-firmware device until + `stop_event` is set. Mirrors the IoTDevice firmware's main loop, but drives + the device over HTTP/MQTT (via its own config) instead of local GPIO pins. + """ + config = device.full_config + device_id = str(device.id) + + mqtt_values = {} + mqtt_topics = [ + pin['rule']['input']['topic'] + for pin in config['pins'] + if pin['rule']['action'] == 'mqtt_toggle' + ] + + def on_connect(client, userdata, connect_flags, reason_code, properties): + for topic in mqtt_topics: + client.subscribe(topic) + + def on_message(client, userdata, message): + mqtt_values[message.topic] = message.payload.decode('utf-8') + + client = _build_mqtt_client(config['mqtt'], on_connect, on_message) + + rule_values = {} + previous_status_hash = None + run_count = 0 + + try: + while not stop_event.is_set(): + try: + requests.get(config['health']['url'], timeout=10) + + for pin in config['pins']: + if run_count % pin.get('interval', 1): + continue + + rule = pin['rule'] + rule_params = _resolve_rule_params(rule, rule_values, mqtt_values) + rule_values[pin['identifier']] = RULE_ACTIONS[rule['action']]( + **rule_params + ) + + previous_status_hash = _publish_status( + client, device_id, rule_values, previous_status_hash + ) + except Exception: + logger.exception('Error running rules for device %s', device_id) + _publish_log( + client, + device_id, + config['logging'], + 'error', + f'Error running rules for device {device_id}', + ) + + run_count += 1 + stop_event.wait(config['main']['process_interval']) + finally: + client.loop_stop() + client.disconnect() From ebd29ac3482490582654deac3e32508dfdcbe95f Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 08:02:08 +0200 Subject: [PATCH 09/12] Added management command to run rule workers --- .../apps/device/management/commands/rules.py | 51 +++++++++++++++++++ 1 file changed, 51 insertions(+) create mode 100644 iotserver/apps/device/management/commands/rules.py diff --git a/iotserver/apps/device/management/commands/rules.py b/iotserver/apps/device/management/commands/rules.py new file mode 100644 index 0000000..0c143e1 --- /dev/null +++ b/iotserver/apps/device/management/commands/rules.py @@ -0,0 +1,51 @@ +import signal +import threading + +from django.core.management.base import BaseCommand + +from iotserver.apps.device.models import Device +from iotserver.apps.device.utils import rules + + +class Command(BaseCommand): + help = 'Run the rules engine for non-managed-firmware devices.' + + def handle(self, *args, **options): + """ + Spawn one worker thread per active, non-managed-firmware device and run + 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 + + def handle_shutdown(signum, frame): + self.stdout.write(self.style.WARNING('Shutting down rules engine...')) + stop_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).') + ) + + for thread in threads: + thread.join() + + self.stdout.write(self.style.SUCCESS('Rules engine stopped.')) From 0e005cabca69e83a2eee427424ff7c2899d78c4c Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 08:10:24 +0200 Subject: [PATCH 10/12] Added tests for mqtt and rules managemnt commands --- .../apps/device/tests/management/__init__.py | 0 .../tests/management/commands/__init__.py | 0 .../tests/management/commands/test_mqtt.py | 111 ++++++++++++++++++ .../tests/management/commands/test_rules.py | 47 ++++++++ 4 files changed, 158 insertions(+) create mode 100644 iotserver/apps/device/tests/management/__init__.py create mode 100644 iotserver/apps/device/tests/management/commands/__init__.py create mode 100644 iotserver/apps/device/tests/management/commands/test_mqtt.py create mode 100644 iotserver/apps/device/tests/management/commands/test_rules.py diff --git a/iotserver/apps/device/tests/management/__init__.py b/iotserver/apps/device/tests/management/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/iotserver/apps/device/tests/management/commands/__init__.py b/iotserver/apps/device/tests/management/commands/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/iotserver/apps/device/tests/management/commands/test_mqtt.py b/iotserver/apps/device/tests/management/commands/test_mqtt.py new file mode 100644 index 0000000..113d370 --- /dev/null +++ b/iotserver/apps/device/tests/management/commands/test_mqtt.py @@ -0,0 +1,111 @@ +import io +import json + +import pytest +from django.core.management.base import OutputWrapper + +from iotserver.apps.device.management.commands.mqtt import Command +from iotserver.apps.device.models import DeviceStatus +from iotserver.apps.device.tests import factories as device_factories + + +@pytest.fixture +def command(): + command = Command() + command.stdout = OutputWrapper(io.StringIO()) + return command + + +@pytest.mark.django_db +class TestHandleStatusQueue: + def test_creates_device_status_for_existing_device(self, command, mocker): + device = device_factories.DeviceFactory() + message = mocker.Mock( + topic=f'iot-devices/{device.id}/status', + payload=json.dumps({'sonoff-switch': True}).encode(), + ) + + command.handle_status_queue(message) + + status = DeviceStatus.objects.get(device=device) + assert status.status == {'sonoff-switch': True} + + def test_ignores_message_for_unknown_device(self, command, mocker): + message = mocker.Mock( + topic='iot-devices/00000000-0000-0000-0000-000000000000/status', + payload=b'{}', + ) + + command.handle_status_queue(message) + + assert DeviceStatus.objects.count() == 0 + assert 'does not exist' in command.stdout._out.getvalue() + + +class TestHandleLogQueue: + def test_writes_log_message_to_stdout(self, command, mocker): + message = mocker.Mock( + topic='iot-devices/device-1/logs', payload=b'Device disconnected' + ) + + command.handle_log_queue(message) + + output = command.stdout._out.getvalue() + assert 'device-1' in output + assert 'Device disconnected' in output + + +class TestMqttOnConnect: + def test_subscribes_to_all_device_topics(self, command, mocker): + client = mocker.Mock() + + command.mqtt_on_connect(client, None, None, None, None) + + client.subscribe.assert_called_once_with('iot-devices/#') + + +class TestMqttOnMessage: + def test_dispatches_status_messages(self, command, mocker): + mock_handle_status = mocker.patch.object(command, 'handle_status_queue') + message = mocker.Mock(topic='iot-devices/device-1/status') + + command.mqtt_on_message(None, None, message) + + mock_handle_status.assert_called_once_with(message) + + def test_dispatches_log_messages(self, command, mocker): + mock_handle_log = mocker.patch.object(command, 'handle_log_queue') + message = mocker.Mock(topic='iot-devices/device-1/logs') + + command.mqtt_on_message(None, None, message) + + mock_handle_log.assert_called_once_with(message) + + def test_ignores_unrelated_topics(self, command, mocker): + mock_handle_status = mocker.patch.object(command, 'handle_status_queue') + mock_handle_log = mocker.patch.object(command, 'handle_log_queue') + message = mocker.Mock(topic='iot-devices/device-1/toggle') + + command.mqtt_on_message(None, None, message) + + mock_handle_status.assert_not_called() + mock_handle_log.assert_not_called() + + +class TestHandle: + def test_connects_and_starts_loop_forever(self, command, mocker): + mock_client_cls = mocker.patch( + 'iotserver.apps.device.management.commands.mqtt.mqtt.Client' + ) + mocker.patch( + 'iotserver.apps.device.management.commands.mqtt.settings.MQTT', + {'host': 'broker.local', 'port': 1883}, + ) + + command.handle() + + client = mock_client_cls.return_value + assert client.on_connect == command.mqtt_on_connect + assert client.on_message == command.mqtt_on_message + client.connect.assert_called_once_with(host='broker.local', port=1883) + client.loop_forever.assert_called_once() diff --git a/iotserver/apps/device/tests/management/commands/test_rules.py b/iotserver/apps/device/tests/management/commands/test_rules.py new file mode 100644 index 0000000..d6f5d9d --- /dev/null +++ b/iotserver/apps/device/tests/management/commands/test_rules.py @@ -0,0 +1,47 @@ +import pytest +from django.core.management import call_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', + ) + + mock_thread_cls = mocker.patch('threading.Thread') + + call_command('rules') + + assert mock_thread_cls.call_count == 1 + _, kwargs = mock_thread_cls.call_args + assert kwargs['args'][0] == matching_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): + device_factories.DeviceFactory(managed_firmware=True, active=True) + + mock_thread_cls = mocker.patch('threading.Thread') + + call_command('rules') + + mock_thread_cls.assert_not_called() From e148a04e94814f43117eac7ed5086282d394169b Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 08:10:52 +0200 Subject: [PATCH 11/12] Added tests for rule utils --- .../apps/device/tests/utils/test_rules.py | 434 ++++++++++++++++++ 1 file changed, 434 insertions(+) create mode 100644 iotserver/apps/device/tests/utils/test_rules.py diff --git a/iotserver/apps/device/tests/utils/test_rules.py b/iotserver/apps/device/tests/utils/test_rules.py new file mode 100644 index 0000000..e535d06 --- /dev/null +++ b/iotserver/apps/device/tests/utils/test_rules.py @@ -0,0 +1,434 @@ +import hashlib +import json +import threading + +import pytest + +from iotserver.apps.device.utils import rules + + +class TestFindXpathValue: + def test_returns_top_level_value(self): + assert rules.find_xpath_value({'a': 1}, ['a']) == 1 + + def test_returns_nested_value(self): + # xpaths are passed in already reversed, e.g. 'a.b' -> ['b', 'a'] + assert rules.find_xpath_value({'a': {'b': 2}}, ['b', 'a']) == 2 + + def test_returns_none_for_missing_key(self): + assert rules.find_xpath_value({'a': 1}, ['missing']) is None + + def test_returns_none_when_intermediate_is_not_subscriptable(self): + assert rules.find_xpath_value({'a': 1}, ['b', 'a']) is None + + +class TestEvaluateCondition: + @pytest.mark.parametrize( + 'operator,input,value,expected', + [ + ('eq', True, True, True), + ('eq', True, False, False), + ('gt', 5, 3, True), + ('gt', 2, 3, False), + ('lt', 2, 3, True), + ], + ) + def test_evaluate(self, operator, input, value, expected): + assert rules.evaluate_condition(input, operator, value) is expected + + def test_returns_false_on_type_error(self): + assert rules.evaluate_condition(None, 'gt', 3) is False + + +def test_handle_conditions_returns_must_and_should_booleans(): + rule_values = { + 'timer': True, + 'weather-service-current': {'rain': False}, + 'mqtt-toggle': True, + } + input_value = { + 'conditions': { + 'must': { + 'timer': {'operator': 'eq', 'value': True}, + 'weather-service-current.rain': {'operator': 'eq', 'value': False}, + }, + 'should': { + 'mqtt-toggle': {'operator': 'eq', 'value': True}, + }, + } + } + + assert rules.handle_conditions(rule_values, input_value) == { + 'must': [True, True], + 'should': [True], + } + + +class TestValueToBool: + @pytest.mark.parametrize( + 'value,expected', + [ + ('true', True), + ('1', True), + ('yes', True), + ('on', True), + ('false', False), + ('0', False), + (1, True), + (0, False), + (None, False), + (True, True), + ], + ) + def test_value_to_bool(self, value, expected): + assert rules.value_to_bool(value) is expected + + +class TestTimer: + def test_within_range_same_day(self, mocker): + mock_localtime = mocker.patch( + 'iotserver.apps.device.utils.rules.timezone.localtime' + ) + mock_localtime.return_value.strftime.return_value = '1200' + + assert rules.timer(start_time='09:00', end_time='17:00') is True + + def test_outside_range_same_day(self, mocker): + mock_localtime = mocker.patch( + 'iotserver.apps.device.utils.rules.timezone.localtime' + ) + mock_localtime.return_value.strftime.return_value = '2000' + + assert rules.timer(start_time='09:00', end_time='17:00') is False + + def test_within_overnight_range(self, mocker): + mock_localtime = mocker.patch( + 'iotserver.apps.device.utils.rules.timezone.localtime' + ) + mock_localtime.return_value.strftime.return_value = '2300' + + assert rules.timer(start_time='18:00', end_time='05:00') is True + + def test_outside_overnight_range(self, mocker): + mock_localtime = mocker.patch( + 'iotserver.apps.device.utils.rules.timezone.localtime' + ) + mock_localtime.return_value.strftime.return_value = '1200' + + assert rules.timer(start_time='18:00', end_time='05:00') is False + + +class TestService: + def test_returns_json_with_auth_header(self, mocker): + mock_requests = mocker.patch('iotserver.apps.device.utils.rules.requests') + mock_requests.get.return_value.json.return_value = {'temperature': 20} + + result = rules.service(url='http://example.com', auth_header='abc123') + + mock_requests.get.assert_called_once_with( + 'http://example.com', headers={'Authorization': 'Token abc123'} + ) + mock_requests.get.return_value.raise_for_status.assert_called_once() + assert result == {'temperature': 20} + + def test_omits_headers_without_auth_header(self, mocker): + mock_requests = mocker.patch('iotserver.apps.device.utils.rules.requests') + + rules.service(url='http://example.com') + + mock_requests.get.assert_called_once_with('http://example.com', headers=None) + + +class TestMqttToggle: + def test_returns_true_for_retained_truthy_value(self): + assert rules.mqtt_toggle(topic='t', mqtt_values={'t': '1'}) is True + + def test_returns_false_for_missing_topic(self): + assert rules.mqtt_toggle(topic='t', mqtt_values={}) is False + + +class TestSonoffToggle: + def test_toggles_on(self, mocker): + mock_sonoff_cls = mocker.patch('iotserver.apps.device.utils.rules.Sonoff') + mock_sonoff_cls.return_value.toggle_device.return_value = 'on' + + result = rules.sonoff_toggle(on=True, device_id='abc123') + + mock_sonoff_cls.assert_called_once_with('abc123') + mock_sonoff_cls.return_value.toggle_device.assert_called_once_with('on') + assert result is True + + def test_toggles_off(self, mocker): + mock_sonoff_cls = mocker.patch('iotserver.apps.device.utils.rules.Sonoff') + mock_sonoff_cls.return_value.toggle_device.return_value = 'off' + + result = rules.sonoff_toggle(on=False, device_id='abc123') + + mock_sonoff_cls.return_value.toggle_device.assert_called_once_with('off') + assert result is False + + +def test_rule_actions_dispatch_table(): + assert rules.RULE_ACTIONS == { + 'timer': rules.timer, + 'service': rules.service, + 'mqtt_toggle': rules.mqtt_toggle, + 'sonoff_toggle': rules.sonoff_toggle, + } + + +class TestBuildMqttClient: + @pytest.fixture + def mock_client_cls(self, mocker): + return mocker.patch('iotserver.apps.device.utils.rules.mqtt.Client') + + def test_connects_and_starts_loop_with_default_port(self, mock_client_cls): + client = rules._build_mqtt_client( + {'client_id': 'device-1', 'host': 'broker.local'}, None, None + ) + + mock_client_cls.assert_called_once_with( + rules.mqtt.CallbackAPIVersion.VERSION2, client_id='device-1' + ) + client.connect.assert_called_once_with(host='broker.local', port=1883) + client.loop_start.assert_called_once() + + def test_uses_configured_port(self, mock_client_cls): + client = rules._build_mqtt_client( + {'client_id': 'd', 'host': 'h', 'port': 8883}, None, None + ) + + client.connect.assert_called_once_with(host='h', port=8883) + + def test_sets_credentials_when_configured(self, mock_client_cls): + client = rules._build_mqtt_client( + {'client_id': 'd', 'host': 'h', 'username': 'user', 'password': 'pass'}, + None, + None, + ) + + client.username_pw_set.assert_called_once_with('user', 'pass') + + def test_skips_credentials_when_not_configured(self, mock_client_cls): + client = rules._build_mqtt_client({'client_id': 'd', 'host': 'h'}, None, None) + + client.username_pw_set.assert_not_called() + + def test_enables_tls_when_configured(self, mock_client_cls): + client = rules._build_mqtt_client( + {'client_id': 'd', 'host': 'h', 'ssl_enabled': True}, None, None + ) + + client.tls_set.assert_called_once() + + def test_sets_lastwill_when_configured(self, mock_client_cls): + client = rules._build_mqtt_client( + { + 'client_id': 'd', + 'host': 'h', + 'lastwill': {'topic': 'iot-devices/d/logs', 'message': 'disconnected'}, + }, + None, + None, + ) + + client.will_set.assert_called_once_with( + topic='iot-devices/d/logs', payload='disconnected' + ) + + +class TestResolveRuleParams: + def test_passes_through_plain_values(self): + rule = { + 'action': 'timer', + 'input': {'start_time': '18:00', 'end_time': '05:00'}, + } + + assert rules._resolve_rule_params(rule, {}, {}) == { + 'start_time': '18:00', + 'end_time': '05:00', + } + + def test_evaluates_conditions(self): + rule = { + 'action': 'sonoff_toggle', + 'input': { + 'device_id': 'abc123', + 'on': { + 'conditions': { + 'must': {'timer': {'operator': 'eq', 'value': True}}, + 'should': {}, + } + }, + }, + } + + result = rules._resolve_rule_params(rule, {'timer': True}, {}) + + assert result == {'device_id': 'abc123', 'on': True} + + def test_injects_mqtt_values_for_mqtt_toggle_action(self): + rule = {'action': 'mqtt_toggle', 'input': {'topic': 'iot-devices/x/toggle'}} + mqtt_values = {'iot-devices/x/toggle': '1'} + + result = rules._resolve_rule_params(rule, {}, mqtt_values) + + assert result == {'topic': 'iot-devices/x/toggle', 'mqtt_values': mqtt_values} + + +class TestPublishStatus: + def test_publishes_when_status_has_changed(self, mocker): + client = mocker.Mock() + + status_hash = rules._publish_status(client, 'device-1', {'a': 1}, None) + + payload = json.dumps({'a': 1}, sort_keys=True) + client.publish.assert_called_once_with('iot-devices/device-1/status', payload) + assert status_hash == hashlib.sha1(payload.encode()).digest() + + def test_skips_publish_when_status_unchanged(self, mocker): + client = mocker.Mock() + payload = json.dumps({'a': 1}, sort_keys=True) + previous_hash = hashlib.sha1(payload.encode()).digest() + + result_hash = rules._publish_status(client, 'device-1', {'a': 1}, previous_hash) + + client.publish.assert_not_called() + assert result_hash == previous_hash + + +class TestPublishLog: + def test_publishes_when_at_or_above_threshold(self, mocker): + mock_logger = mocker.patch('iotserver.apps.device.utils.rules.logger') + client = mocker.Mock() + + rules._publish_log(client, 'device-1', {'level': 'warning'}, 'error', 'boom') + + mock_logger.error.assert_called_once_with('boom') + client.publish.assert_called_once_with('iot-devices/device-1/logs', 'boom') + + def test_suppresses_publish_below_threshold(self, mocker): + mock_logger = mocker.patch('iotserver.apps.device.utils.rules.logger') + client = mocker.Mock() + + rules._publish_log(client, 'device-1', {'level': 'warning'}, 'debug', 'noisy') + + mock_logger.debug.assert_called_once_with('noisy') + client.publish.assert_not_called() + + +class TestRunDevice: + def _device(self, mocker, **config_overrides): + config = { + 'main': {'process_interval': 0}, + 'health': {'url': 'http://example.com/health/'}, + 'mqtt': {'client_id': 'device-1', 'host': 'broker.local'}, + 'logging': {'level': 'warning'}, + 'pins': [], + } + config.update(config_overrides) + return mocker.Mock(id='device-1', full_config=config) + + def test_runs_health_check_dispatches_rule_and_publishes_status(self, mocker): + stop_event = threading.Event() + mock_client = mocker.Mock() + mocker.patch( + 'iotserver.apps.device.utils.rules._build_mqtt_client', + return_value=mock_client, + ) + mock_requests_get = mocker.patch( + 'iotserver.apps.device.utils.rules.requests.get' + ) + mock_action = mocker.Mock(return_value=True) + mocker.patch.dict( + 'iotserver.apps.device.utils.rules.RULE_ACTIONS', {'timer': mock_action} + ) + mock_publish_status = mocker.patch( + 'iotserver.apps.device.utils.rules._publish_status', + side_effect=lambda *a, **k: stop_event.set(), + ) + + device = self._device( + mocker, + pins=[ + { + 'identifier': 'day-timer', + 'interval': 1, + 'rule': { + 'action': 'timer', + 'input': {'start_time': '18:00', 'end_time': '05:00'}, + }, + } + ], + ) + + rules.run_device(device, stop_event) + + mock_requests_get.assert_called_once_with( + 'http://example.com/health/', timeout=10 + ) + mock_action.assert_called_once_with(start_time='18:00', end_time='05:00') + mock_publish_status.assert_called_once() + mock_client.loop_stop.assert_called_once() + mock_client.disconnect.assert_called_once() + + def test_catches_iteration_errors_and_keeps_looping(self, mocker): + stop_event = threading.Event() + mock_client = mocker.Mock() + mocker.patch( + 'iotserver.apps.device.utils.rules._build_mqtt_client', + return_value=mock_client, + ) + mocker.patch( + 'iotserver.apps.device.utils.rules.requests.get', + side_effect=[Exception('boom'), None], + ) + mock_publish_log = mocker.patch( + 'iotserver.apps.device.utils.rules._publish_log' + ) + mocker.patch( + 'iotserver.apps.device.utils.rules._publish_status', + side_effect=lambda *a, **k: stop_event.set(), + ) + + device = self._device(mocker) + + rules.run_device(device, stop_event) + + mock_publish_log.assert_called_once() + mock_client.loop_stop.assert_called_once() + mock_client.disconnect.assert_called_once() + + def test_subscribes_to_mqtt_toggle_topics_on_connect(self, mocker): + stop_event = threading.Event() + mock_client = mocker.Mock() + mock_build_client = mocker.patch( + 'iotserver.apps.device.utils.rules._build_mqtt_client', + return_value=mock_client, + ) + mocker.patch('iotserver.apps.device.utils.rules.requests.get') + mocker.patch( + 'iotserver.apps.device.utils.rules._publish_status', + side_effect=lambda *a, **k: stop_event.set(), + ) + + device = self._device( + mocker, + pins=[ + { + 'identifier': 'mqtt-toggle', + 'interval': 1, + 'rule': { + 'action': 'mqtt_toggle', + 'input': {'topic': 'iot-devices/device-1/toggle'}, + }, + } + ], + ) + + rules.run_device(device, stop_event) + + _, on_connect, _on_message = mock_build_client.call_args[0] + on_connect(mock_client, None, None, None, None) + + mock_client.subscribe.assert_called_once_with('iot-devices/device-1/toggle') From 6c9fd1ebe1a38cc27871e82376c0efd28371df31 Mon Sep 17 00:00:00 2001 From: Jonathan Bydendyk Date: Thu, 24 Sep 2026 08:18:22 +0200 Subject: [PATCH 12/12] Added systemd files for running the rules engine management command --- entrypoint.sh | 1 + systemd/iot.rulesengine.service | 12 ++++++++++++ 2 files changed, 13 insertions(+) create mode 100644 systemd/iot.rulesengine.service diff --git a/entrypoint.sh b/entrypoint.sh index 34dff81..b94a9db 100755 --- a/entrypoint.sh +++ b/entrypoint.sh @@ -3,4 +3,5 @@ ./manage.py migrate --noinput ./manage.py collectstatic --noinput ./manage.py mqtt & +./manage.py rules & gunicorn iotserver.wsgi:application -w 2 -b :8000 --reload diff --git a/systemd/iot.rulesengine.service b/systemd/iot.rulesengine.service new file mode 100644 index 0000000..91af082 --- /dev/null +++ b/systemd/iot.rulesengine.service @@ -0,0 +1,12 @@ +[Unit] +Description=rules engine daemon +After=network.target + +[Service] +User=pi +Group=pi +WorkingDirectory=/home/pi/iotserver +ExecStart=/home/pi/.local/bin/poetry run /home/pi/iotserver/manage.py rules + +[Install] +WantedBy=multi-user.target