0Pricing
AI Agents · Урок

Автоматическая реакция на события датчиков

Если температура > порога → оповещение → приведение в действие: управляемые агентом циклы управления 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 — локальная установка не требуется.

Все уроки этого курса

  1. Протокол MQTT для интеграции агентов
  2. Обработка временных рядов агентами
  3. Автоматическая реакция на события датчиков
  4. Развёртывание лёгких агентов на периферии
← Назад к AI Agents