Автоматическая реакция на события датчиков
Если температура > порога → оповещение → приведение в действие: управляемые агентом циклы управления IoT.
«Автоматическая реакция на события датчиков» — бесплатный урок AI Agents на CoddyKit. Это урок 3 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AI Agents, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AI Agents содержит 4 уроков всего.
Автоматические ответы агента на события
Когда датчик пересекает пороговое значение, агент должен автоматически отреагировать без участия человека. Основные задачи: решить, какое действие выполнить, не допустить запуска повторных действий одним и тем же событием и соблюдать период охлаждения, чтобы агент не перегружал исполнительные устройства командами.
Определение политик действий
Политика действий сопоставляет условия датчика с ответами агента. Определяйте политики декларативно, чтобы их было легко читать и изменять без правки кода логики. Каждая политика содержит условие, приоритет и одно или несколько действий.
ACTION_POLICIES = [
{
'name': 'HIGH_TEMP_ALERT',
'topic': 'sensors/temperature',
'condition': lambda v: v > 38,
'priority': 'critical',
'actions': ['TURN_ON_COOLING', 'ALERT_MAINTENANCE', 'LOG_EVENT']
},
{
'name': 'HIGH_TEMP_WARNING',
'topic': 'sensors/temperature',
'condition': lambda v: 35 < v <= 38,
'priority': 'warning',
'actions': ['ALERT_MAINTENANCE', 'LOG_EVENT']
},
{
'name': 'LOW_HUMIDITY',
'topic': 'sensors/humidity',
'condition': lambda v: v < 30,
'priority': 'warning',
'actions': ['TURN_ON_HUMIDIFIER', 'LOG_EVENT']
}
]
def match_policies(topic: str, value: float) -> list:
return [
p for p in ACTION_POLICIES
if p['topic'] == topic and p['condition'](value)
]
if __name__ == '__main__':
matches = match_policies('sensors/temperature', 39)
print('Matched policies for temperature=39:')
for p in matches:
print(f" {p['name']} ({p['priority']}): {p['actions']}")
Очередь действий
Очередь действий отделяет обнаружение событий от выполнения действий. События помещаются в очередь, а рабочий поток извлекает их и выполняет. Это не блокирует цикл получения MQTT и позволяет повторять попытки, если действие завершилось ошибкой.
import queue
import threading
from datetime import datetime
action_queue: queue.Queue = queue.Queue(maxsize=500)
def enqueue_action(action_name: str, context: dict, priority: str = 'normal'):
item = {
'action': action_name,
'context': context,
'priority': priority,
'enqueued_at': datetime.utcnow().isoformat()
}
try:
action_queue.put_nowait(item)
print(f'Enqueued: {action_name}')
except queue.Full:
print(f'WARNING: Action queue full, dropping {action_name}')
def action_worker(executor_fn):
"""Run in a background thread, executing actions from the queue."""
while True:
item = action_queue.get()
try:
executor_fn(item['action'], item['context'])
except Exception as e:
print(f'Action failed: {item["action"]} — {e}')
finally:
action_queue.task_done()
# Start worker thread:
# worker_thread = threading.Thread(target=action_worker, args=(execute_action,), daemon=True)
# worker_thread.start()
if __name__ == '__main__':
enqueue_action('TURN_ON_COOLING', {'zone': 'server-room'}, priority='critical')
enqueue_action('LOG_EVENT', {'msg': 'temperature nominal'})
print('Queue size:', action_queue.qsize())
Дедупликация событий
Без дедупликации температура, которая остаётся выше 38 °C в течение 10 минут при одном показании в секунду, создаёт 600 одинаковых событий. Дедупликация гарантирует, что одна и та же комбинация (топик, условие, действие) срабатывает только один раз для каждого события и сбрасывается после исчезновения условия.
class EventDeduplicator:
def __init__(self):
# active_events: (topic, policy_name) -> event_start_time
self._active: dict = {}
def is_new_event(self, topic: str, policy_name: str) -> bool:
key = (topic, policy_name)
return key not in self._active
def mark_active(self, topic: str, policy_name: str):
self._active[(topic, policy_name)] = datetime.utcnow()
def clear_event(self, topic: str, policy_name: str):
key = (topic, policy_name)
if key in self._active:
duration = (datetime.utcnow() - self._active.pop(key)).seconds
print(f'Event cleared: {policy_name} (lasted {duration}s)')
def clear_topic_if_normal(
self, topic: str, value: float, normal_fn
):
if normal_fn(value):
keys = [k for k in self._active if k[0] == topic]
for k in keys:
self.clear_event(k[0], k[1])
dedup = EventDeduplicator()
dedup.mark_active('sensors/temperature', 'HIGH_TEMP_ALERT')
print('New event?', dedup.is_new_event('sensors/temperature', 'HIGH_TEMP_ALERT'))Период охлаждения
Даже после исчезновения и повторного срабатывания события период охлаждения предотвращает слишком частые повторные срабатывания. Объединяйте дедупликацию (срабатывать один раз, пока условие выполняется) с периодом охлаждения (ждать N минут после исчезновения условия, прежде чем снова разрешить то же оповещение).
from datetime import datetime, timedelta
class CoolDownManager:
def __init__(self, cool_down_minutes: int = 15):
self.cool_down = timedelta(minutes=cool_down_minutes)
self._cleared_at: dict = {} # (topic, policy) -> cleared_datetime
def is_in_cool_down(self, topic: str, policy_name: str) -> bool:
key = (topic, policy_name)
cleared_at = self._cleared_at.get(key)
if cleared_at is None:
return False
return datetime.utcnow() - cleared_at < self.cool_down
def record_clear(self, topic: str, policy_name: str):
self._cleared_at[(topic, policy_name)] = datetime.utcnow()
def time_remaining(self, topic: str, policy_name: str) -> int:
key = (topic, policy_name)
cleared_at = self._cleared_at.get(key)
if cleared_at is None:
return 0
elapsed = datetime.utcnow() - cleared_at
remaining = self.cool_down - elapsed
return max(0, int(remaining.total_seconds()))
cooldown = CoolDownManager(cool_down_minutes=15)
cooldown.record_clear('sensors/temperature', 'HIGH_TEMP_ALERT')
print('In cool-down?', cooldown.is_in_cool_down('sensors/temperature', 'HIGH_TEMP_ALERT'))Принятие решений о действиях с помощью LLM
В сложных ситуациях — при нескольких одновременных оповещениях, противоречащих друг другу политиках или необычных сочетаниях показаний — передайте принятие решения LLM. LLM получает полный контекст данных датчиков и рекомендует план действий с приоритетами.
import anthropic
import json
def llm_decide_actions(
sensor_readings: dict,
active_policies: list
) -> list:
client = anthropic.Anthropic(api_key='YOUR_API_KEY')
context = json.dumps({
'readings': sensor_readings,
'triggered_policies': [p['name'] for p in active_policies]
}, indent=2)
prompt = (
f'Current sensor state:\n{context}\n\n'
'Multiple alert policies are active. '
'Recommend an ordered list of actions to take. '
'Consider conflicting effects (e.g., humidifier and cooling may conflict).\n'
'Return JSON: {"recommended_actions": [str], "reasoning": str}'
)
response = client.messages.create(
model='claude-opus-4-5', max_tokens=512,
messages=[{'role': 'user', 'content': prompt}]
)
return json.loads(response.content[0].text)Выполнение действий через MQTT
Действия выполняются публикацией командных сообщений в специализированные MQTT-топики устройств. Полезная нагрузка команды соответствует стандартной схеме: имя действия, параметры, ID запроса для подтверждения и TTL (срок действия команды истекает, если устройство слишком долго находится в автономном режиме).
import json
import uuid
from datetime import datetime, timedelta
ACTION_TOPICS = {
'TURN_ON_COOLING': 'devices/hvac/commands',
'TURN_OFF_COOLING': 'devices/hvac/commands',
'TURN_ON_HUMIDIFIER': 'devices/humidifier/commands',
'ALERT_MAINTENANCE': 'notifications/maintenance',
'LOG_EVENT': 'logs/agent_events'
}
def execute_action(action_name: str, context: dict, mqtt_client) -> str:
topic = ACTION_TOPICS.get(action_name)
if not topic:
print(f'No topic defined for action: {action_name}')
return 'unknown_action'
request_id = str(uuid.uuid4())[:8]
ttl = (datetime.utcnow() + timedelta(minutes=5)).isoformat()
payload = json.dumps({
'action': action_name,
'request_id': request_id,
'context': context,
'ttl': ttl
})
mqtt_client.publish(topic, payload, qos=1)
print(f'Executed {action_name} -> {topic} (req={request_id})')
return request_id
if __name__ == '__main__':
class FakeMQTT:
def publish(self, topic, payload, qos=1):
pass
execute_action('TURN_ON_COOLING', {'zone': 'server-room'}, FakeMQTT())
Подтверждение действий
Устройства должны подтверждать получение команд, публикуя сообщение в топике подтверждений. Агент подписывается на топики подтверждений и может повторить попытку, если не получает подтверждение в течение заданного времени ожидания.
import threading
from collections import defaultdict
class AckTracker:
def __init__(self, timeout_seconds: int = 30):
self.timeout = timeout_seconds
self._pending: dict = {} # request_id -> {'action', 'send_time', 'ack_event'}
def register(self, request_id: str, action_name: str):
event = threading.Event()
self._pending[request_id] = {
'action': action_name,
'send_time': datetime.utcnow(),
'ack_event': event
}
# Schedule timeout check
t = threading.Timer(self.timeout, self._on_timeout, args=[request_id])
t.daemon = True
t.start()
def acknowledge(self, request_id: str):
entry = self._pending.pop(request_id, None)
if entry:
entry['ack_event'].set()
print(f'Ack received for {entry["action"]} (req={request_id})')
def _on_timeout(self, request_id: str):
if request_id in self._pending:
action = self._pending.pop(request_id)['action']
print(f'TIMEOUT: No ack for {action} (req={request_id}) — retry?')Полный конвейер событий датчиков
Все компоненты объединяются в следующую последовательность: получение MQTT → сопоставление с политиками → проверка дедупликации и периода охлаждения → постановка действий в очередь → выполнение рабочим потоком через публикацию MQTT → отслеживание подтверждений. Такая архитектура обрабатывает тысячи событий датчиков в минуту без блокировки.
class IoTAgentPipeline:
def __init__(self, mqtt_client):
self.mqtt = mqtt_client
self.dedup = EventDeduplicator()
self.cooldown = CoolDownManager(cool_down_minutes=15)
self.ack_tracker = AckTracker(timeout_seconds=30)
def on_sensor_message(self, topic: str, value: float):
policies = match_policies(topic, value)
self.dedup.clear_topic_if_normal(
topic, value,
normal_fn=lambda v: v <= 35 # below warning threshold
)
for policy in policies:
name = policy['name']
if not self.dedup.is_new_event(topic, name):
continue # already active, skip
if self.cooldown.is_in_cool_down(topic, name):
print(f'In cool-down: {name}')
continue
self.dedup.mark_active(topic, name)
ctx = {'topic': topic, 'value': value, 'policy': name}
for action in policy['actions']:
enqueue_action(action, ctx, policy['priority'])Путь эскалации
Некоторые ситуации требуют передачи вопроса человеку: повторяющиеся сбои подтверждения команды, длительные критические состояния или несколько противоречащих друг другу политик. Определите путь эскалации, который отправляет уведомление или создаёт заявку.
import requests
def escalate_to_human(
reason: str,
sensor_data: dict,
webhook_url: str = 'https://hooks.slack.com/services/YOUR/SLACK/WEBHOOK'
):
message = {
'text': (
f'*IoT Agent Escalation* \n'
f'Reason: {reason}\n'
f'Sensor data: {sensor_data}\n'
f'Time: {datetime.utcnow().isoformat()}'
)
}
try:
response = requests.post(webhook_url, json=message, timeout=5)
response.raise_for_status()
print(f'Escalation sent: {reason}')
except requests.RequestException as e:
print(f'Escalation failed: {e}')
# Fall back: log to file
with open('escalations.log', 'a') as f:
import json
f.write(json.dumps({'reason': reason, 'data': sensor_data}) + '\n')Проверка конвейера событий
Перед развёртыванием конвейера событий в рабочей среде напишите автоматизированные проверки, которые имитируют события датчиков и проверяют постановку правильных действий в очередь. Проверьте каждую политику отдельно, работу дедупликации и завершение периода охлаждения.
import time
def test_high_temp_policy_fires_once():
dedup = EventDeduplicator()
cooldown = CoolDownManager(cool_down_minutes=0) # disable cooldown for test
pipeline = IoTAgentPipeline(None)
pipeline.dedup = dedup
pipeline.cooldown = cooldown
actions_fired = []
action_queue.queue.clear()
# Fire same event 5 times in a row
for _ in range(5):
pipeline.on_sensor_message('sensors/temperature', 40.0)
# Only 1 set of actions should have been enqueued
actions = list(action_queue.queue)
assert len(actions) > 0, 'At least one action should fire'
print(f'Actions enqueued: {len(actions)} (expected: just 1 event worth)')
return True
result = test_high_temp_policy_fires_once()
print('Test passed:', result)Проверка знаний
Какова основная цель дедупликации событий в конвейере событий датчиков?
Итоги: автоматические ответы на события датчиков
Отлично! Вы узнали:
- Политики действий: декларативное сопоставление условий с действиями с учётом приоритетов
- Очередь действий: отделяет обнаружение от выполнения; рабочий поток обрабатывает действия
- Дедупликация: срабатывание один раз для каждого события, а не для каждого показания
- Период охлаждения: предотвращает немедленное повторное срабатывание после исчезновения условия
- Эскалация к LLM: передача сложных ситуаций с несколькими политиками на рассуждение LLM
- Отслеживание подтверждений: обнаружение команд без подтверждения и повторная попытка
Далее: развёртывание лёгких агентов на периферийных устройствах, например Raspberry Pi.
Часто задаваемые вопросы
Урок «Автоматическая реакция на события датчиков» бесплатный?
Да — полный текст урока «Автоматическая реакция на события датчиков» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Agents, подпишись на CoddyKit PRO. Курс AI Agents содержит 4 уроков всего.
Чему я научусь в уроке «Автоматическая реакция на события датчиков»?
Если температура > порога → оповещение → приведение в действие: управляемые агентом циклы управления IoT. Ты практикуешь AI Agents с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать AI Agents?
Предыдущий опыт не требуется. AI Agents на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 3 из 4.
Сколько времени занимает урок «Автоматическая реакция на события датчиков»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке AI Agents?
Да. Каждый урок AI Agents включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Протокол MQTT для интеграции агентов
- Обработка временных рядов агентами
- Автоматическая реакция на события датчиков
- Развёртывание лёгких агентов на периферии