energy_price.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302
  1. """The electricity price energy is costed at, and the cost of energy over time (#1251).
  2. The price is either the fixed ``energy_cost_per_kwh`` setting or, with
  3. ``energy_price_source = "homeassistant"``, a Home Assistant sensor. A sensor
  4. reading is written back into ``energy_cost_per_kwh``, so that setting always
  5. holds the last price known: the Settings page shows it, and it is what is used
  6. whenever Home Assistant can't be reached.
  7. Energy is costed at the price of the hour it was used. Every hourly energy
  8. snapshot carries the price at that moment, and the energy a plug used between
  9. two snapshots is costed at the price of the first. Each snapshot also keeps a
  10. running total of the energy and cost to date, so the cost of any span is the
  11. difference of two rows, like the energy itself. A print or a date range
  12. then costs the sum of its hours, instead of all of it at whatever the price
  13. happens to be when it ends, or today.
  14. """
  15. from __future__ import annotations
  16. import logging
  17. from dataclasses import dataclass
  18. from datetime import datetime
  19. import httpx
  20. from sqlalchemy import select
  21. from sqlalchemy.ext.asyncio import AsyncSession
  22. from backend.app.models.smart_plug_energy_snapshot import SmartPlugEnergySnapshot
  23. from backend.app.utils.local_time import to_naive_utc
  24. logger = logging.getLogger(__name__)
  25. DEFAULT_PRICE = 0.15
  26. _HA_TIMEOUT = 5.0
  27. PRICE_SOURCE_FIXED = "fixed"
  28. PRICE_SOURCE_HOMEASSISTANT = "homeassistant"
  29. # A sensor's unit, after the currency, scaled to per kWh. Day-ahead market
  30. # sensors (Nord Pool, ENTSO-E) are often per MWh.
  31. _ENERGY_UNIT_SCALE = {"kwh": 1.0, "mwh": 0.001, "wh": 1000.0}
  32. async def stored_price(db: AsyncSession) -> float:
  33. """The ``energy_cost_per_kwh`` setting: the fixed price, or the last one read."""
  34. from backend.app.api.routes.settings import get_setting
  35. raw = await get_setting(db, "energy_cost_per_kwh")
  36. try:
  37. return float(raw) if raw else DEFAULT_PRICE
  38. except (TypeError, ValueError):
  39. return DEFAULT_PRICE
  40. def price_from_state(state: dict | None) -> float | None:
  41. """A per-kWh price from a Home Assistant state, or None when it has none.
  42. Negative prices are paid out in some markets; they are costed as free,
  43. as a cost below zero is not something the rest of Bambuddy expects.
  44. """
  45. if not state:
  46. return None
  47. try:
  48. value = float(state.get("state"))
  49. except (TypeError, ValueError):
  50. return None
  51. unit = str((state.get("attributes") or {}).get("unit_of_measurement") or "")
  52. per = unit.rsplit("/", 1)[-1].strip().lower() if "/" in unit else ""
  53. value *= _ENERGY_UNIT_SCALE.get(per, 1.0)
  54. if value != value or value in (float("inf"), float("-inf")):
  55. return None
  56. return max(0.0, value)
  57. async def current_price(db: AsyncSession, *, remember: bool = False) -> float:
  58. """The electricity price now.
  59. With Home Assistant as the source, reads the sensor; falls back to the last
  60. price remembered when it can't be read. ``remember`` writes a new reading
  61. into ``energy_cost_per_kwh`` (the caller commits). Only the hourly loop and
  62. the settings save do, so the print paths never write a setting in the same
  63. transaction as the energy reading they exist to save.
  64. """
  65. from backend.app.api.routes.settings import get_homeassistant_settings, get_setting, set_setting
  66. fallback = await stored_price(db)
  67. if (await get_setting(db, "energy_price_source") or PRICE_SOURCE_FIXED) != PRICE_SOURCE_HOMEASSISTANT:
  68. return fallback
  69. entity_id = (await get_setting(db, "energy_price_ha_entity") or "").strip()
  70. if not entity_id:
  71. return fallback
  72. ha = await get_homeassistant_settings(db)
  73. if not ha["ha_enabled"] or not ha["ha_url"] or not ha["ha_token"]:
  74. return fallback
  75. price = price_from_state(await _read_state(ha["ha_url"], ha["ha_token"], entity_id))
  76. if price is None:
  77. logger.info("Electricity price sensor %s has no usable value; keeping %s", entity_id, fallback)
  78. return fallback
  79. if remember and price != fallback:
  80. await set_setting(db, "energy_cost_per_kwh", str(price))
  81. return price
  82. async def _read_state(url: str, token: str, entity_id: str) -> dict | None:
  83. """One entity's state, or None. Its own short timeout: this runs at print
  84. start and end, which an unreachable Home Assistant must not hold up."""
  85. try:
  86. async with httpx.AsyncClient(timeout=_HA_TIMEOUT) as client:
  87. response = await client.get(
  88. f"{url.rstrip('/')}/api/states/{entity_id}",
  89. headers={"Authorization": f"Bearer {token}"},
  90. )
  91. response.raise_for_status()
  92. return response.json()
  93. except Exception as e:
  94. logger.debug("Could not read electricity price from %s: %s", entity_id, e)
  95. return None
  96. def cost_of_readings(readings: list[tuple[float, float | None]], fallback_price: float) -> tuple[float, float]:
  97. """``(kwh, cost)`` over consecutive ``(lifetime_kwh, price)`` readings.
  98. Each step is costed at the price of the reading it starts from. A step
  99. where the counter went backwards (a device reset) counts as nothing.
  100. """
  101. kwh = 0.0
  102. cost = 0.0
  103. for (start, price), (end, _) in zip(readings, readings[1:], strict=False):
  104. delta = end - start
  105. if delta <= 0:
  106. continue
  107. kwh += delta
  108. cost += delta * (price if price is not None else fallback_price)
  109. return kwh, cost
  110. def cost_at_average_price(kwh: float, costed_kwh: float, costed: float, fallback_price: float) -> float:
  111. """Cost ``kwh`` at the average price of the energy that could be costed.
  112. The energy figure shown and the energy summed hour by hour are the same
  113. unless a counter was reset along the way; this keeps cost and energy
  114. consistent either way.
  115. """
  116. if costed_kwh > 0:
  117. return kwh * (costed / costed_kwh)
  118. return kwh * fallback_price
  119. async def print_energy_cost(
  120. db: AsyncSession,
  121. *,
  122. plug_id: int | None,
  123. start_at: datetime | None,
  124. start_kwh: float,
  125. start_price: float | None,
  126. end_kwh: float,
  127. end_plug_id: int,
  128. price_now: float,
  129. ) -> float:
  130. """What the energy a print used cost, hour by hour.
  131. Falls back to the price now for a print whose start predates #1251, or
  132. whose end reading came from a different plug than its start.
  133. """
  134. energy_used = end_kwh - start_kwh
  135. if start_at is None or plug_id is None or plug_id != end_plug_id:
  136. return energy_used * price_now
  137. result = await db.execute(
  138. select(SmartPlugEnergySnapshot.lifetime_kwh, SmartPlugEnergySnapshot.price_per_kwh)
  139. .where(
  140. SmartPlugEnergySnapshot.plug_id == plug_id,
  141. SmartPlugEnergySnapshot.recorded_at > to_naive_utc(start_at),
  142. )
  143. .order_by(SmartPlugEnergySnapshot.recorded_at, SmartPlugEnergySnapshot.id)
  144. )
  145. readings: list[tuple[float, float | None]] = [(start_kwh, start_price)]
  146. readings.extend((row[0], row[1]) for row in result.all())
  147. readings.append((end_kwh, None))
  148. costed_kwh, costed = cost_of_readings(readings, price_now)
  149. return cost_at_average_price(energy_used, costed_kwh, costed, price_now)
  150. @dataclass
  151. class SnapshotCost:
  152. """Energy and cost between two points of the snapshot history, summed over plugs.
  153. ``kwh``/``cost`` come from the running totals and are what is costed;
  154. ``last_lifetime_kwh`` is the sum of each plug's latest lifetime counter,
  155. which only the all-time figure needs.
  156. """
  157. kwh: float = 0.0
  158. cost: float = 0.0
  159. last_lifetime_kwh: float = 0.0
  160. def _running_totals(row: SmartPlugEnergySnapshot, fallback_price: float) -> tuple[float, float]:
  161. """A snapshot's energy and cost to date. A row without them (none should be
  162. left after the upgrade backfill) counts its whole counter at its price."""
  163. kwh = row.kwh_to_date if row.kwh_to_date is not None else row.lifetime_kwh
  164. if row.cost_to_date is not None:
  165. return kwh, row.cost_to_date
  166. price = row.price_per_kwh if row.price_per_kwh is not None else fallback_price
  167. return kwh, row.lifetime_kwh * price
  168. async def new_snapshot(
  169. db: AsyncSession, *, plug_id: int, recorded_at: datetime, lifetime_kwh: float, price: float
  170. ) -> SmartPlugEnergySnapshot:
  171. """The next snapshot row for a plug, its running totals carried forward.
  172. The energy since the previous snapshot is added at the previous snapshot's
  173. price, so each hour is costed at the price that held during it. A counter
  174. that went backwards (a device reset) adds nothing. The first snapshot of a
  175. plug counts what its counter had already reached at the price then.
  176. """
  177. s = SmartPlugEnergySnapshot
  178. prev = (
  179. await db.execute(select(s).where(s.plug_id == plug_id).order_by(s.recorded_at.desc(), s.id.desc()).limit(1))
  180. ).scalar_one_or_none()
  181. if prev is None:
  182. kwh_to_date, cost_to_date = lifetime_kwh, lifetime_kwh * price
  183. else:
  184. prev_kwh, prev_cost = _running_totals(prev, price)
  185. prev_price = prev.price_per_kwh if prev.price_per_kwh is not None else price
  186. step = max(0.0, lifetime_kwh - prev.lifetime_kwh)
  187. kwh_to_date, cost_to_date = prev_kwh + step, prev_cost + step * prev_price
  188. return s(
  189. plug_id=plug_id,
  190. recorded_at=recorded_at,
  191. lifetime_kwh=lifetime_kwh,
  192. price_per_kwh=price,
  193. kwh_to_date=kwh_to_date,
  194. cost_to_date=cost_to_date,
  195. )
  196. async def _snapshot_at_or_before(
  197. db: AsyncSession, plug_id: int, moment: datetime | None
  198. ) -> SmartPlugEnergySnapshot | None:
  199. s = SmartPlugEnergySnapshot
  200. query = select(s).where(s.plug_id == plug_id)
  201. if moment is not None:
  202. query = query.where(s.recorded_at <= moment)
  203. return (await db.execute(query.order_by(s.recorded_at.desc(), s.id.desc()).limit(1))).scalar_one_or_none()
  204. async def snapshot_cost(
  205. db: AsyncSession,
  206. *,
  207. fallback_price: float,
  208. dt_from: datetime | None = None,
  209. dt_to: datetime | None = None,
  210. ) -> SnapshotCost:
  211. """Energy and cost over the same span ``_sum_snapshot_deltas`` measures.
  212. Per plug: from the last snapshot at or before ``dt_from`` (or its first
  213. snapshot) to the last at or before ``dt_to``, as the difference of their
  214. running totals. Without ``dt_from``, from zero: the whole history. Two
  215. indexed lookups per plug, however long the history is.
  216. """
  217. from backend.app.models.smart_plug import SmartPlug
  218. s = SmartPlugEnergySnapshot
  219. dt_from = to_naive_utc(dt_from)
  220. dt_to = to_naive_utc(dt_to)
  221. result = SnapshotCost()
  222. for plug_id in (await db.execute(select(SmartPlug.id))).scalars().all():
  223. end = await _snapshot_at_or_before(db, plug_id, dt_to)
  224. if end is None:
  225. continue
  226. end_kwh, end_cost = _running_totals(end, fallback_price)
  227. result.last_lifetime_kwh += end.lifetime_kwh
  228. base_kwh = base_cost = 0.0
  229. if dt_from is not None:
  230. base = await _snapshot_at_or_before(db, plug_id, dt_from)
  231. if base is None:
  232. base = (
  233. await db.execute(select(s).where(s.plug_id == plug_id).order_by(s.recorded_at, s.id).limit(1))
  234. ).scalar_one_or_none()
  235. if base is not None:
  236. base_kwh, base_cost = _running_totals(base, fallback_price)
  237. kwh, cost = end_kwh - base_kwh, end_cost - base_cost
  238. # Negative only across a reset in rows from before the upgrade, whose
  239. # totals are the raw counter; the energy figure clamps those to 0 too.
  240. if kwh > 0 and cost >= 0:
  241. result.kwh += kwh
  242. result.cost += cost
  243. return result
  244. async def all_time_cost(db: AsyncSession, live_total_kwh: float, price_now: float) -> float:
  245. """What the plugs' lifetime energy cost, hour by hour where it can be.
  246. Energy since each plug's latest snapshot, and every plug with no snapshots
  247. (an MQTT plug), is costed at the price now.
  248. """
  249. hist = await snapshot_cost(db, fallback_price=price_now)
  250. tail = max(0.0, live_total_kwh - hist.last_lifetime_kwh)
  251. return cost_at_average_price(live_total_kwh, hist.kwh + tail, hist.cost + tail * price_now, price_now)