Automatyczne reagowanie na zdarzenia z czujników
Jeśli temperatura > próg → alert → wykonanie działania: sterowanie pętlami IoT przez agenta.
Automatyczne reagowanie na zdarzenia z czujników to bezpłatna lekcja AI Agents na CoddyKit. To lekcja 3 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej AI Agents, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs AI Agents zawiera 4 lekcji w sumie.
Automatyczne reakcje agentów na zdarzenia
Gdy czujnik przekroczy próg, agent musi zareagować automatycznie, bez udziału człowieka. Najważniejsze wyzwania to określenie, jakie działanie należy podjąć, zapewnienie, że to samo zdarzenie nie wywoła zduplikowanych działań, oraz przestrzeganie okresu wyciszenia, aby agent nie zalewał urządzeń wykonawczych poleceniami.
Definiowanie zasad działań
Zasada działania mapuje warunki czujnika na reakcje agenta. Zasady należy definiować deklaratywnie, aby można je było łatwo odczytywać i modyfikować bez zmieniania kodu logiki. Każda zasada zawiera warunek, priorytet i co najmniej jedno działanie.
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']}")
Kolejka działań
Kolejka działań oddziela wykrywanie zdarzeń od wykonywania działań. Zdarzenia są umieszczane w kolejce, a worker pobiera je i wykonuje. Zapobiega to blokowaniu pętli odbierającej MQTT i umożliwia ponawianie prób, jeśli działanie zakończy się niepowodzeniem.
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())
Deduplikacja zdarzeń
Bez deduplikacji temperatura utrzymująca się powyżej 38°C przez 10 minut przy jednym odczycie na sekundę wygeneruje 600 identycznych zdarzeń. Deduplikacja gwarantuje, że ta sama kombinacja (topic, condition, action) zostanie uruchomiona tylko raz dla danego zdarzenia, a jej stan zostanie zresetowany po ustąpieniu warunku.
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'))Okres wyciszenia
Nawet po ustąpieniu i ponownym wystąpieniu zdarzenia okres wyciszenia zapobiega jego szybkiemu ponownemu wywołaniu. Należy połączyć deduplikację (uruchomienie raz, gdy warunek jest spełniony) z okresem wyciszenia (odczekanie N minut po ustąpieniu warunku przed zezwoleniem na ponowne wywołanie tego samego alertu).
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'))Decyzja dotycząca działania wspomagana przez LLM
W złożonych sytuacjach — przy wielu jednoczesnych alertach, sprzecznych zasadach lub nietypowych kombinacjach odczytów — należy przekazać decyzję LLM. LLM otrzymuje pełny kontekst dotyczący czujników i rekomenduje uporządkowany według priorytetów plan działań.
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)Wykonywanie działań przez MQTT
Działania są wykonywane przez publikowanie komunikatów z poleceniami w tematach MQTT właściwych dla poszczególnych urządzeń. Ładunek polecenia jest zgodny ze standardowym schematem: nazwa działania, parametry, identyfikator żądania służący do potwierdzenia oraz TTL (polecenie wygasa, jeśli urządzenie pozostaje zbyt długo offline).
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())
Potwierdzanie działań
Urządzenia powinny potwierdzać odebrane polecenia, publikując wiadomości w temacie potwierdzeń. Agent subskrybuje tematy potwierdzeń i może ponowić próbę, jeśli w określonym czasie oczekiwania nie otrzyma potwierdzenia.
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?')Pełny potok zdarzeń z czujników
Połączenie wszystkich komponentów: odbiór MQTT → dopasowanie zasady → sprawdzenie deduplikacji i okresu wyciszenia → umieszczenie działań w kolejce → worker wykonuje je przez publikację MQTT → śledzenie potwierdzeń. Ta architektura obsługuje tysiące zdarzeń z czujników na minutę bez blokowania.
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'])Ścieżka eskalacji
Niektóre sytuacje wymagają eskalacji do człowieka: powtarzające się niepowodzenia potwierdzania polecenia, utrzymujące się stany krytyczne lub wiele sprzecznych zasad. Należy zdefiniować ścieżkę eskalacji, która wysyła powiadomienie push lub tworzy zgłoszenie.
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')Testowanie potoku zdarzeń
Przed wdrożeniem potoku zdarzeń na produkcji należy napisać testy automatyczne, które symulują zdarzenia z czujników i sprawdzają, czy kolejkowane są właściwe działania. Należy przetestować każdą zasadę osobno, działanie deduplikacji oraz wygaśnięcie okresu wyciszenia.
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)Sprawdzenie wiedzy
Jaki jest główny cel deduplikacji zdarzeń w potoku zdarzeń z czujników?
Podsumowanie: automatyczne reakcje na zdarzenia z czujników
Doskonale! Poznali Państwo:
- Zasady działań: deklaratywne mapowanie warunków na działania z uwzględnieniem priorytetów
- Kolejka działań: oddzielenie wykrywania od wykonywania; wątek roboczy przetwarza działania
- Deduplikacja: uruchomienie raz dla każdego zdarzenia, a nie raz dla każdego odczytu
- Okres wyciszenia: zapobieganie natychmiastowemu ponownemu wywołaniu po ustąpieniu warunku
- Eskalacja do LLM: przekazywanie złożonych sytuacji z wieloma zasadami do analizy LLM
- Śledzenie potwierdzeń: wykrywanie niepotwierdzonych poleceń i ponawianie prób
Następny temat: wdrażanie lekkich agentów na urządzeniach brzegowych, takich jak Raspberry Pi.
Często zadawane pytania
Czy lekcja „Automatyczne reagowanie na zdarzenia z czujników” jest bezpłatna?
Tak — pełny tekst „Automatyczne reagowanie na zdarzenia z czujników” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu AI Agents, przejdź na CoddyKit PRO. Kurs AI Agents zawiera 4 lekcji w sumie.
Co nauczysz się w „Automatyczne reagowanie na zdarzenia z czujników”?
Jeśli temperatura > próg → alert → wykonanie działania: sterowanie pętlami IoT przez agenta. Ćwiczysz AI Agents z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć AI Agents?
Nie wymagamy żadnego doświadczenia. AI Agents w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 3 z 4.
Ile czasu zajmuje lekcja „Automatyczne reagowanie na zdarzenia z czujników”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji AI Agents?
Tak. Każda lekcja AI Agents zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Protokół MQTT do integracji agentów
- Przetwarzanie danych szeregów czasowych przez agentów
- Automatyczne reagowanie na zdarzenia z czujników
- Wdrażanie lekkich agentów na brzegu sieci