End-to-End Data Flows

Device Hardware → MQTT → Daemon → Carousel → Workers → Back to Device

Segment 1: Device Telemetry Publishing via MQTT

Max Device MQTT Client EMQX Broker :1883 / :8083 Topics: kv/{node_id}/ telemetry kv/{node_id}/ commands
Device publish: Max device collects telemetry (battery %, network %, relay state, heap) every heartbeat interval (2s active, 5m idle) → Publishes via MQTT client on WiFi or LoRa gateway → EMQX broker (:1883 internal, :8083 WebSocket) ingests on topic `kv/{node_id}/telemetry`. Broker holds last-will message if device disconnects.
// Device heartbeat publishes: { "device_id": "KV-1234", "battery_mv": 4200, "rssi": -85, "relay_state": 1, "free_heap": 156000, "uptime_sec": 3600 }

Segment 2: Telegraf Ingest — MQTT to Time-Series Storage

EMQX Subscribe Telegraf Parser & Aggregate HTTP InfluxDB :8086 Mona queries (REST) Bucket: telegraf Measurements: • device_telemetry • device_cost_daily
Telegraf ingest pipeline: Telegraf subscribes to MQTT topics (kv/*/telemetry) → Parses MessagePack/JSON payloads → Aggregates fields (battery_mv → float, relay_state → int) → Writes to InfluxDB (:8086) bucket "telegraf" with 1-minute resolution. Time-series data: 7-day retention (configurable). Mona dashboard queries via REST API for historical charts (battery %, cost trend, queue depth).
// Telegraf writes (InfluxDB line protocol): device_telemetry,device_id=KV-1234,location=lab battery_mv=4200i,rssi=-85i,heap_free=156000i 1715000000000000000 device_telemetry,device_id=KV-5678,location=workshop battery_mv=3900i,rssi=-92i,heap_free=128000i 1715000000000000000

Segment 3: Device Registry Lookup → Carousel Entitlement Check

Device Daemon RX Device Registry Carousel Evaluator Checks: 1. Device in registry? 2. Hardware class? 3. User tier? (starter/pro/enterprise) 4. Daily cost budget? OK? YES Route to Worker NO HTTP 402
Device command routing: Device sends command via MQTT → Daemon receives on holon queue → Daemon looks up device in registry (device_id → IP, hardware_class, tier) → Passes to Carousel evaluator with device context → Carousel checks: device exists? hardware class valid? user tier permits this action? daily cost budget available? → YES: route to worker pool for execution → NO: return HTTP 402 (cost cap exceeded) or 403 (unauthorized). Carousel enforces tier hierarchy: starter limited to GPIO reads, pro gets relay control, enterprise unlimited.

Segment 4: Worker Execution → Cost Tracking → Result Broadcasting

Worker Pool Holon Executor Billing Domain Cost Impact: • Holon execution: $0.005 • Model tier: Haiku -90% • Daily budget: $10.00 • Remaining: $9.87 Publish result via MQTT Device
Execution → cost → broadcast: Worker claims holon from queue → HolonExecutor runs (relay switch, schedule task, fetch sensor) → Cost tracked: execution fee ($0.005) + model tier adjustment (Haiku -90% vs Sonnet) → Billing domain deducts from daily budget + updates tier usage → Result serialized (status + error/success) → Published back to device via MQTT on `kv/{node_id}/response` topic → Device receives within 100ms. User sees state change in app immediately.
// Result published to device: { "request_id": "flynn-claude-KAN-123-abc123-1715000001234-f7a9", "status": "success", "action": "relay_set", "relay_id": 1, "new_state": 1, "timestamp_ms": 1715000005000, "cost_usd": 0.005, "daily_remaining_usd": 9.865 }

Segment 5: Complete Round-Trip (Device → Cloud → Device, 200ms–4s)

T+0ms User sends command via BLE/WiFi T+5ms Device TX via MQTT T+20ms Daemon receives T+30ms Carousel evaluates T+50ms Worker executes T+100ms Result via MQTT Critical Path (Typical): Device Infrastructure
Complete round-trip cycle (~100ms typical, 200ms–4s with concurrent workers):

Device side (5–20ms): User clicks button in BLE app → LocalPersist record → Send command via BLE (5ms) → Device receives → Serializes to MQTT → Publishes on `kv/{id}/commands` (20ms total)

Infrastructure (20–30ms): MQTT broker (:1883) receives → Daemon holon queue router subscribes (20ms) → Daemon processes on port 8001

Cloud evaluation (30–50ms): Carousel evaluator tier check + daily budget (10ms) → Routes to worker pool with semaphore (5ms) → Worker claims from queue (5ms)

Execution (50–100ms): HolonExecutor runs command (varies: relay GPIO instant, API call 100ms+) → Collects telemetry → Billing domain records cost → Publishes result on `kv/{id}/response` (50–100ms)

Device feedback (100–150ms): Device receives result via MQTT → Updates local state → UI reflects change. Total: ~100ms for relay toggle, ~500ms for HTTP calls, ~4s for holon worker delays (concurrent execution limits via semaphore).