Skip to content

Repository files navigation

Water Budget System

Adaptive household water monitoring — tracks tap sessions via MQTT, runs a weekly ML optimizer, and pushes updated alert thresholds back to the hardware.


Project structure

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

Quick start

cp .env.example .env          # fill in MQTT_BROKER_HOST etc.
pip install -r requirements.txt
python run.py                 # starts API on :8000 + scheduler

For development with auto-reload:

uvicorn api.main:app --reload --port 8000

Run the tests:

pytest tests/ -v

MQTT protocol

Device → Server (raw FlowRing telemetry)

The 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.

Server → Device (firmware command)

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


REST API

Create a household

POST /households
{
  "id": "hh_001",
  "name": "The Smiths",
  "weekly_budget_s": 72000,
  "num_residents": 4,
  "has_garden": true,
  "has_children": true
}

Register a tap (triggers cold-start + MQTT publish)

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

Ingest a session (REST fallback — MQTT is preferred)

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 current thresholds

GET /households/hh_001/thresholds

Get weekly stats

GET /households/hh_001/stats

Get reports

GET /households/hh_001/reports

Force an optimizer run (testing / admin)

POST /households/hh_001/run_optimizer

Algorithm phases

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

Forecast weights

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.

Reduction weights

weight_i = share_i × flexibility_i × (0.5 + compressibility_i)

Flexibility defaults: shower 0.8, kitchen 0.7, bathroom 0.9, garden 1.5

Damping

new = 0.8×old + 0.2×target   capped at ±20% of old

Environment variables

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


Extending

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.

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages