Skip to content

IoT & Telemetry Engine

The IoT & Telemetry Engine forms the data ingestion backbone of Planovi. It handles high-throughput telemetry from physical smart meters, photovoltaic (PV) inverters, and battery storage units, while reconciling third-party DSO grid data.


Functions Matrix

Function NameInbound TriggerCore ResponsibilityPrimary Data Stores
ingest-telemetryCron (every 60s) / WebhookPolls InCharge API, decodes converter packets, writes to telemetry_rawtelemetry_raw, Deno KV
refresh-live-stateDB Trigger / Post-ingestComputes instantaneous cooperative totals into single-row cachecoop_live_state
cleanup-telemetryNightly Cron (0 2 * * *)Prunes raw records older than 30 daystelemetry_raw
process-tauron-csvUser Upload / S3 TriggerParses Tauron DSO 15-minute interval billing CSVstauron_readings, members
process-alertsEvent / Automated CronDetects inverter fault codes, phase imbalance, and offline unitssystem_alerts, push_notifications

1. ingest-telemetry

Endpoint Specification

  • Path: /functions/v1/ingest-telemetry
  • Method: POST
  • Auth: Service Role Key or Authorized Cron Worker

Architecture & Data Flow

flowchart TD
    Cron[Scheduled Cron Runner] -->|POST| Fn[ingest-telemetry]
    Fn --> CacheCheck{Check Deno KV}
    CacheCheck -->|Valid Token| API[InCharge Hardware API]
    CacheCheck -->|Expired| Auth[InCharge Auth Endpoint]
    Auth -->|Store Token + TTL| CacheCheck
    API -->|Raw Device Payloads| Decoder["Decoder & Normalizer"]
    Decoder --> DB[Insert into telemetry_raw]
    Decoder --> LiveTrigger[Invoke refresh-live-state]

Ingestion Logic & Deno KV Cache

  • Token Caching: Uses Deno.openKv() keys ["incharge", "token"] and ["incharge", "expiry"].
  • Proactive Refresh: Refreshes the token 5 minutes prior to its 24-hour expiration window.
  • Fail-Safe Mechanism: If external IoT servers return 5xx errors or fail to respond within 8,000ms, the function records an incident in system_alerts and exits cleanly without crashing upstream crons.

Payload Schema (Normalized Telemetry)

{
"device_id": "inv_huawei_00941",
"coop_id": "b1b2c3d4-e5f6-7a8b-9c0d-1e2f3a4b5c6d",
"timestamp": "2026-09-16T10:00:00Z",
"metrics": {
"solar_power_kw": 14.85,
"grid_consumption_kw": 2.10,
"grid_feed_in_kw": 12.75,
"battery_charge_soc_pct": 88.5,
"battery_power_kw": -3.2,
"voltage_l1": 230.4,
"voltage_l2": 231.1,
"voltage_l3": 229.8,
"frequency_hz": 50.01
}
}

2. refresh-live-state

Endpoint Specification

  • Path: /functions/v1/refresh-live-state
  • Method: POST
  • Auth: Internal Service Role

Operational Mechanism

Computes the current aggregate state of an energy cooperative by summing the latest metric for all active devices belonging to coop_id.

-- Conceptual aggregation executed in function
INSERT INTO coop_live_state (coop_id, total_generation_kw, total_consumption_kw, battery_soc_avg, last_updated)
VALUES ($1, $2, $3, $4, NOW())
ON CONFLICT (coop_id) DO UPDATE SET
total_generation_kw = EXCLUDED.total_generation_kw,
total_consumption_kw = EXCLUDED.total_consumption_kw,
battery_soc_avg = EXCLUDED.battery_soc_avg,
last_updated = EXCLUDED.last_updated;

3. process-tauron-csv

Endpoint Specification

  • Path: /functions/v1/process-tauron-csv
  • Method: POST (Multipart form or Storage Event)
  • Auth: Authenticated Admin / Technician

Processing Pipeline

  1. File Decoding: Streams the Polish DSO CSV file (Windows-1250 / UTF-8 encoding).
  2. Header Normalization: Identifies column mappings for PPE identification codes, timestamp format (YYYY-MM-DD HH:mm), active energy consumed (P+ [kWh]), and active energy injected (P- [kWh]).
  3. Cooperative Member Reconciliation: Matches PPE codes against registered meter numbers in the members table.
  4. Batch Upsert: Inserts 15-minute interval values into tauron_readings using batch chunks of 500 records to prevent memory overflow.

4. cleanup-telemetry & process-alerts

Retention Policy

  • cleanup-telemetry: Executes daily at 02:00 UTC. Deletes rows in telemetry_raw where created_at < NOW() - INTERVAL '30 days' after verifying that aggregate-daily has completed.
  • process-alerts: Inspects recent readings for critical thresholds:
    • Over-voltage conditions (> 253V sustained for > 10 min — Polish grid compliance standard).
    • Inverter communication dropout (> 15 min without packets).
    • Dispatches notifications via Expo Push and records entries into system_alerts.