Adaptive household water monitoring — tracks tap sessions via MQTT, runs a weekly ML optimizer, and pushes updated alert thresholds back to the hardware.
water_budget/
├── api/
│ └── main.py FastAPI app — all REST endpoints
├── db/
│ ├── models.py SQLAlchemy ORM (Household, Tap, TapSession, Threshold, WeeklyReport)
│ └── database.py Engine + session factory
├── ml/
│ └── optimizer.py The algorithm — cold start, forecasting, tuner, anomaly detection
├── mqtt/
│ └── client.py paho-mqtt wrapper — publish thresholds, receive sessions
├── scheduler/
│ └── jobs.py APScheduler — Sunday 23:00 optimizer, anomaly checks
├── tests/
│ └── test_optimizer.py 16 unit tests for the ML layer
├── run.py Entry point
├── requirements.txt
└── .env.example
cp .env.example .env # fill in MQTT_BROKER_HOST etc.
pip install -r requirements.txt
python run.py # starts API on :8000 + schedulerFor development with auto-reload:
uvicorn api.main:app --reload --port 8000Run the tests:
pytest tests/ -vThe OTA firmware publishes a raw telemetry message per device on the public MQTT broker:
Topic: D2S/SA/V1/{device_id}
Payload:
"ZZ:0213145456/+22|0-Count:3;1-water_ml:142"The Python bridge converts this into an internal session-like record using the tap's configured flow rate when available.
The firmware branch listens for command/configuration messages on:
Topic: S2D/SA/V1/{device_id}/cfg
Payload:
"CFG|2-5;3-12;7-300;8-120|END"The branch currently uses config indices such as 1-minPressMs, 2-lowFlowLpm, 3-highFlowLpm, 7-publishIntervalSec, and 8-timezoneOffset.
From Python, the writer now looks like this:
mqtt_client.publish_cfg("flowring003", [(2, 5), (3, 12), (7, 300), (8, 120)])To request a readback of the current device settings:
mqtt_client.publish_cfg_readback("flowring003")There is no direct per-tap threshold topic in this firmware branch, so the old water/.../threshold publish path is not compatible with the current device protocol.
Topic: water/{household_id}/{tap_id}/alert
Payload:
{
"type": "leak",
"duration_s": 3600,
"message": "Tap tap_kitchen ran for 60 min between midnight and 6 AM. This could be a dripping tap or stuck valve.",
"fired_at": "2024-05-06T03:15:00"
}Types: "leak" (45+ min at night) or "forgotten" (3× median duration).
POST /households
{
"id": "hh_001",
"name": "The Smiths",
"weekly_budget_s": 72000,
"num_residents": 4,
"has_garden": true,
"has_children": true
}
POST /households/hh_001/taps
{
"id": "tap_shower_main",
"tap_type": "shower",
"label": "Main bathroom shower",
"flow_rate_lpm": 9.0
}
tap_type options: shower | kitchen | bathroom | garden | other
POST /households/hh_001/sessions
{
"tap_id": "tap_shower_main",
"started_at": "2024-05-06T07:14:00",
"ended_at": "2024-05-06T07:19:15",
"duration_s": 315
}
GET /households/hh_001/thresholds
GET /households/hh_001/stats
GET /households/hh_001/reports
POST /households/hh_001/run_optimizer
| Week | Phase | What happens |
|---|---|---|
| 0 | Cold start | Defaults from household profile; all thresholds in ghost mode |
| 1–2 | Observation | Ghost mode: sessions recorded, no alerts; hardware learns |
| End of wk 2 | Personalisation | Thresholds set to p75 of observed sessions; alerts enabled |
| 3+ | Weekly tuner | Every Sunday: forecast → gap → weighted cuts → damped publish |
| Month 2+ | Long-term | Seasonal priors, weekday/weekend splits, drift detection |
n_forecast = 0.5×(w-1) + 0.3×(w-2) + 0.2×(w-3)
Switches to 0.7 / 0.2 / 0.1 when usage diverges >30% for two consecutive weeks.
weight_i = share_i × flexibility_i × (0.5 + compressibility_i)
Flexibility defaults: shower 0.8, kitchen 0.7, bathroom 0.9, garden 1.5
new = 0.8×old + 0.2×target capped at ±20% of old
| Variable | Default | Description |
|---|---|---|
DATABASE_URL |
sqlite:///./water_budget.db |
SQLAlchemy connection string |
MQTT_BROKER_HOST |
localhost |
MQTT broker hostname |
MQTT_BROKER_PORT |
1883 |
MQTT broker port |
MQTT_USERNAME |
(empty) | MQTT username (if auth enabled) |
MQTT_PASSWORD |
(empty) | MQTT password |
MQTT_CLIENT_ID |
water_budget_server |
MQTT client identifier |
For PostgreSQL: DATABASE_URL=postgresql://user:pass@host/water_budget
Add a new tap type: add an entry to TAP_DEFAULTS in ml/optimizer.py
with (threshold_s, floor_s, ceiling_s, openings_per_week, flexibility).
Change the schedule: edit the cron trigger in scheduler/jobs.py.
Add litre-based budgets: multiply each tap's duration_s by its
flow_rate_lpm / 60 — the algorithm shape is identical, only units change.
Weekday/weekend splits: extend TapStats with is_weekend: bool and
maintain two threshold sets per tap once 4+ weeks of data exist.