0Pricing
AI Agents · Lezione

Risposta automatizzata agli eventi dei sensori

Se temperatura > soglia → avviso → attivazione: cicli di controllo IoT gestiti dall’agente.

Risposta automatizzata agli eventi dei sensori è una lezione AI Agents gratuita su CoddyKit. Questa è la lezione 3 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento AI Agents, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso AI Agents include 4 lezioni in totale.

Risposte automatizzate degli agenti guidate dagli eventi

Quando un sensore supera una soglia, l'agente deve rispondere automaticamente senza intervento umano. Le difficoltà principali sono decidere quale azione intraprendere, garantire che lo stesso evento non attivi azioni duplicate e rispettare un periodo di cooldown, così da non sovraccaricare gli attuatori con comandi.

Definizione dei criteri per le azioni

Un criterio per le azioni associa le condizioni dei sensori alle risposte dell'agente. Definisca i criteri in modo dichiarativo, così che siano facili da leggere e modificare senza intervenire sul codice della logica. Ogni criterio comprende una condizione, una priorità e una o più azioni.

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']}")

Coda delle azioni

Una coda delle azioni separa il rilevamento degli eventi dall'esecuzione delle azioni. Gli eventi vengono inseriti nella coda; un worker li preleva e li esegue. In questo modo si evita di bloccare il ciclo di ricezione MQTT e sono possibili nuovi tentativi se un'azione non riesce.

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())

Deduplicazione degli eventi

Senza deduplicazione, una temperatura che rimane sopra i 38 °C per 10 minuti con una lettura al secondo genera 600 eventi identici. La deduplicazione garantisce che la stessa combinazione (topic, condizione, azione) venga attivata una sola volta per evento e venga reimpostata quando la condizione si annulla.

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'))

Periodo di cool-down

Anche dopo che un evento si annulla e si riattiva, un periodo di cool-down impedisce nuove attivazioni ravvicinate. Combini la deduplicazione (attivazione una sola volta finché la condizione permane) con il cool-down (attesa di N minuti dopo l'annullamento della condizione prima di consentire nuovamente lo stesso avviso).

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'))

Decisione sulle azioni assistita dall'LLM

Per situazioni complesse — più avvisi simultanei, criteri in conflitto o combinazioni insolite di letture — deleghi la decisione all'LLM. L'LLM riceve il contesto completo dei sensori e raccomanda un piano d'azione con priorità.

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)

Esecuzione delle azioni tramite MQTT

Le azioni vengono eseguite pubblicando messaggi di comando su topic MQTT specifici per il dispositivo. Il payload del comando segue uno schema standard: nome dell'azione, parametri, ID della richiesta per la conferma e TTL (il comando scade se il dispositivo rimane offline troppo a lungo).

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())

Conferma di ricezione delle azioni

I dispositivi devono confermare la ricezione dei comandi pubblicando su un topic di conferma. L'agente si sottoscrive ai topic di conferma e può riprovare se non riceve alcuna conferma entro il periodo di timeout.

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?')

Pipeline completa degli eventi dei sensori

Assemblaggio di tutti i componenti: ricezione MQTT → corrispondenza con i criteri → controllo di deduplicazione e cool-down → inserimento delle azioni in coda → esecuzione da parte del worker tramite pubblicazione MQTT → monitoraggio delle conferme. Questa architettura gestisce migliaia di eventi dei sensori al minuto senza blocchi.

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'])

Percorso di escalation

Alcune situazioni richiedono un'escalation verso una persona: ripetuti mancati riconoscimenti di un comando, condizioni critiche persistenti o più criteri in conflitto. Definisca un percorso di escalation che invii una notifica push o crei un ticket.

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')

Test della pipeline degli eventi

Prima di distribuire la pipeline degli eventi in produzione, scriva test automatizzati che simulino gli eventi dei sensori e verifichino che vengano accodate le azioni corrette. Verifichi ogni criterio singolarmente, il comportamento della deduplicazione e la scadenza del cool-down.

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)

Verifica delle conoscenze

Qual è lo scopo principale della deduplicazione degli eventi in una pipeline di eventi dei sensori?

Riepilogo: risposta automatizzata agli eventi dei sensori

Ottimo! Ha imparato:

  • Criteri per le azioni: associazione dichiarativa tra condizioni e azioni con priorità
  • Coda delle azioni: separazione tra rilevamento ed esecuzione; un thread worker elabora le azioni
  • Deduplicazione: attivazione una volta per evento, non una volta per lettura
  • Cool-down: impedisce una nuova attivazione immediata dopo l'annullamento di una condizione
  • Escalation all'LLM: situazioni complesse con più criteri delegate al ragionamento dell'LLM
  • Monitoraggio delle conferme: rilevamento e nuovo tentativo per i comandi senza conferma

Prossimo argomento: distribuzione di agenti leggeri su dispositivi edge come Raspberry Pi.

Domande Frequenti

La lezione «Risposta automatizzata agli eventi dei sensori» è gratuita?

Sì — il testo completo di «Risposta automatizzata agli eventi dei sensori» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso AI Agents, passa a CoddyKit PRO. Il corso AI Agents include 4 lezioni in totale.

Cosa imparerò in «Risposta automatizzata agli eventi dei sensori»?

Se temperatura > soglia → avviso → attivazione: cicli di controllo IoT gestiti dall’agente. Eserciti AI Agents con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.

Ho bisogno di esperienza per iniziare AI Agents?

Non è richiesta alcuna esperienza precedente. AI Agents su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 3 di 4.

Quanto tempo richiede la lezione «Risposta automatizzata agli eventi dei sensori»?

La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.

Posso scrivere ed eseguire codice in questa lezione AI Agents?

Sì. Ogni lezione AI Agents include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.

Tutte le lezioni di questo corso

  1. Protocollo MQTT per l’integrazione degli agenti
  2. Elaborazione dei dati di serie temporali negli agenti
  3. Risposta automatizzata agli eventi dei sensori
  4. Distribuzione edge di agenti leggeri
← Torna a AI Agents