From 4472cf68b5f9dd90981fa472d64b8287d3f7ea59 Mon Sep 17 00:00:00 2001 From: Aleksandar Grbic Date: Sun, 16 Aug 2026 12:26:24 +0200 Subject: [PATCH] fix: re-handshake in place so a dropped session costs 0 polls, not 1 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every plug was answering only ~half its polls in production: avg_over_time(tapo_plug_up[3h]) ASUS Ascent GX10 0.503 DGX Spark 0.503 k3s Cluster 0.503 Not flaky hardware — the alternation was built into poll(). A Tapo session here survives exactly one full read; the next one answers 403 Forbidden / Tapo(SessionTimeout). The old code responded by setting up=0, clearing self.device and waiting for the NEXT cycle to handshake. So the loop could only ever go success, expiry, re-handshake, expiry, forever. poll() now retries once in place: on any failure it drops the session, re-handshakes and reads again, and only reports up=0 if that also fails. A recoverable expiry now costs nothing. test_poller.py models the observed plug (session dies after one read) and pins the old behaviour at 50% versus 100% now. Each poller also gets its own ApiClient. Sharing one across plugs let their handshakes race, which is the likeliest reason sessions were being invalidated after a single use in the first place. Three metrics so this class of failure is visible from outside: tapo_plug_last_success_seconds gauges keep their last good values on failure, so nothing looks missing — this is how a consumer measures the true age of a reading tapo_plug_reauth_total recoverable expiries, absorbed silently tapo_plug_poll_failures_total polls that failed even after a re-handshake CI now runs the test on pull requests as well, and only publishes :latest from main. Consumer-side alerting added separately in argocd-app-of-apps (manifests/observability/monitoring/tapo-alerts.yaml): TapoPlugPollsFailing fires below 90% poll success over 30m. --- .github/workflows/docker-build.yml | 17 +++ __pycache__/tapo-exporter.cpython-314.pyc | Bin 0 -> 11005 bytes tapo-exporter.py | 59 ++++++++-- test_poller.py | 127 ++++++++++++++++++++++ 4 files changed, 191 insertions(+), 12 deletions(-) create mode 100644 __pycache__/tapo-exporter.cpython-314.pyc create mode 100644 test_poller.py diff --git a/.github/workflows/docker-build.yml b/.github/workflows/docker-build.yml index e80fc17..4c324f6 100644 --- a/.github/workflows/docker-build.yml +++ b/.github/workflows/docker-build.yml @@ -3,6 +3,8 @@ name: Build and Push on: push: branches: [main] + pull_request: + branches: [main] env: REGISTRY: ghcr.io @@ -11,7 +13,22 @@ env: IMAGE_NAME: programmer-network/tapo-exporter jobs: + # The poller talks to real hardware, so the only thing worth testing is its + # error handling — and that is exactly where the ~50% poll-success bug lived. + test: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: "3.12" + - run: pip install --no-cache-dir prometheus-client==0.21.1 + - run: python3 test_poller.py + build: + needs: test + # Pull requests only validate; publishing :latest happens on main. + if: github.event_name == 'push' runs-on: ubuntu-latest permissions: contents: read diff --git a/__pycache__/tapo-exporter.cpython-314.pyc b/__pycache__/tapo-exporter.cpython-314.pyc new file mode 100644 index 0000000000000000000000000000000000000000..86db3e35cf261a24a5111e5380efe3f63cb140e8 GIT binary patch literal 11005 zcmb_CZEzdMb$h_!`|v?hB={xr#Bb326!k@YP=ZBLq$yFzks_;vGuIL86Hj6_nJGvx26f>~GHv=J|0qit$CLi( z+dUirij+Up?f`H1-tNA)`)>EWw{Pzmm(z|QeX#J-3qP+!=)drTo-CQd!>a~_<`IV! z^gQCIK}w-WZcq#)H!4Pw(+W*;lVT#dSusOy7_Bx_YaFx<*cJPLLvai+3NzqT zoYiR37y!*1DU6h7b3ONgBocRB+1D^iJA$pj{dgDW#JPaCCDXuAq&YEaeLK(pr^L zMrw+YN|-hcCKnuO*N^h=GaX& z?4$$cvxBSIRI`KB)N*y3YHGN8FX9@!DB@9SixKp1gy$w*a!MWNgXd;aU(Y?`MM}eD zBiFL;FAxgDx~xQ#4LONc!`iDQe_v}N2`_S#kH^?RXJ-c+kYclfG9x5pR=60CNs1t` z|2h9E8;?aIGAqyWlETI#iD{YTqhZoS7GxH>NMcB4&x<@eFghA&v)hlxqEq5D+ayHK zvoVRCity92Z@-;ohx{Y{;Hi;5ZrFdKk8RhU1OAbb@snIH%d%uFd?*x4L=`q931K0s zh&(WFl7w>!5fDD0GU^YU3YQa%pE`*er#0k?dW9gD=If^ALNUG@=pO_ZZ0$G8V8=O&; zcuYIX zhRzO^!EAzfRR$CBxspy+oFZ$K6=fF25t`x8L-2E5e+k^H40> z!gd^rO--rAx!QQ_f*=Ji@QNbORUA!75(pI*0U+71#9x4y%92X8FxN+tWic2&Gdou~ zF1GfItSnAPVI^S6M5mP*+$a3(tXiJy5sL<8F&Yv;`ebR$m5yZA2tjKMg)ke7=Jyh! zf;4?8sKmnjrQq2MGjm0K1j7;xW^;CYW{VE9v8XabI+W@iWEOO zLh#nJ1FSz46)%#N$i}8v&>Adm#%7vLB_c#^s3kcv@H!H}R*5NmWUlRm7>1n&i&+v{ zXF#9GGyGZ5Cj1oEIy}fW7LUV(LBq%*%zAr{4(=74;>Acp666fIX99qfVHmu^;wtC~ z3pSc1m(L`Gtyji^A#8oH#d}x>U>=R3=MDS*2J|N2D$y6H8H%GOphipWpgDbpKj1#9 zSrU|l6ot)3hsBWMqa-V!WlXn5_2AKF7k^v)lF9^wd^8GESA<|t+5xS0p3uZ{=abRBckFnsYcKls)<-p z&`P{0sWjHApdVwhBIEZ))V{=GBKAm(4-@4kOI?5stXRj8V&g2Z*Q}~N0Gb^z|tRxMZj?KpfwUQZ1S*h zUJMDkmcSO0X#bXcts+LJu~LgicxV}mvVtVVB=A!Dz&a2m1)X9SqB1KgGA^)L>KnN zRz^f$?1aKWYHQ6lLVXEN76cjzq%mx#LGmVPJ9JcOkg}v&fRr#v(0fK$ z?mh!id#W=S41xZau__D(WxNh*=E))^x=!`xUSnInCLDx@vJb$Yuok3E6%R1DF|v-c zyGKd(hD@?|Ye}89=<*dD4c%d?)`!QT0ICtIV3NiQb%i&EqB(3wR9~iO{74iFi6jmb zUlh~IJk_GR)G){aH&)Fd;!Wk|ubQ=`2V0t!g~*h&7n(`?Fv!Y?Z081?wC!vS2F0id z@==R6h(kjXQDEHxo4cg?CY>s2zf50stZpm4YSTbT`DOZ>4ov{r=|Bf8jIScZJlq3g ztwGg50W>loA5T*tA<*a9=?AjY&9&Tt`NpkI9@nr>@wCO3p)O9@SlhOG5EW78%4@znDd_m9S43J_Q?1!%TM!S6!a2SWTA}P zn{mckvZ{!v)+shB5Dhf!vrC6zjH-o~FGWLQOtoI%VZDR0lX@_oj>9pNs1Mbi=f}Ij zcm}oY1RG>3rg|_U0jK|?LTn9MMO#QXNq(h)+F^e}6!!Lnk{T*m9V0Su|O?@U2F zBh{0?=`iLbw36Qf@CVp2zjL@34qQ2~?5J8}8`EsZGTZU`wl^!@s7QDEmpc9L)cwNu zt}oSj{8ule*wg9q)0ca%4y_eeri<&Bi|cQOZ=bz&Hr?P~YVfBTj=uNBRB<5f4qT?! zyj5v$Gy0lTrQ$8EGE#4A+aAok0k6Nf=;nM z{CChUnohwAax1OUPQ|GTRDo-=$}~d-GJRh>zV|OcTd6iqPi)iM23&e6mqV<>)Q-le zXDJ#5yg9sH0zWaM3Bivm)v-@Q2|Uqqy?-v(eMRj&@H}+_Vs?2Pvk|wa;z(~5naiWj z>f?AS>V`m#-anV;QLodnPe~oS)sLk->P>ppQ&Bf=6+z^TEe@S(?gRqwy+y}9C9gKU ze=dI_uMWM<)9{+Wo)bxdY)Qu;l#9{j-UkDDg4mU=7u z%o#?v+_8Y&!fub6c4at#OtZH8$@3i5cI7A7LL6%xMbWJ10DM19L%;|$@4$lWNY@eD zfll0X*wkk@UaFNuN<;^Y+mT2;CYQ@ZF6PgoRWjq>N~`4Wab1C-|4&Key*KPQK9)H8 zbGmgi=3-~0H z7dioRKTU>uWZ}qM7m40vf-U$wlixcNc**N4V1jp%uTt}mhk=f2>$?~d;tG!M7$dQ1 zX$X*-d*ITd%Ghv2%t%(v5&n!2kyRs{dNd!ymT8fQ!$AZN-%RG2MDtWz#+?AldIj+e zRObdKLUrhTAj+iDIBuu9GblYIr&^$d@y@48*AVd-_@KO}y>$4EIMmmK>LF55GJCaUdahaw>Nmf)!kg;{f{tYC}#i8-b zP~CYR4o1KQ0Fvv#-Nz2jYGLVO-D|!ZzGOk&W%Ig$wRNr=oX+mmmX04jclWt}I-YLW zx8$jPg+8GT?0K{2^`5)@JLguI!6&r3Bc$8>OKtuY=GisIn`Ww(nW`nG>1GLlG}E%gwA>l} z;pw}lSD3wPOhuZhS!QYgd2;s6e@FQ9x)+H5`B5|Sls!ULTjwL|Q@o26rhT=%inyqZdtvCxP_p#IYH8(bwi~va zyKZ;i>b}E&WTc#(A0hyMUUvhZ|B)}iYI)f`R7-uU=$HZhBh`PHzHNch5B$Z0F7#8n zw;QsbA2L9TpC6{6;ui+Tpu_kJr@wj7Y9v=|M&L#?gCToA1wL5dF< zi>dab!WvWOIs!n>>ZI3(;0!*1V?+2kg3ekv#;RW-;qVKTH$(W4uVDp-{z6ippLx#BTyi! z-|kv-6)!};5naOLNL6(wUEKSlOD-;L<&stoSI&N8c5xu(s!v<%m#y^=jL0@Z-M7|% zEaL_Ba*f|feUI|n=zrJ_rFRXK-$B1?wBWLX0vLiu*@S&dvb=u)@YC)pU=_hS+M9V% z`1WDT+Y3X^W17I$kE974|8Hka%59(2ld%A7E^p@g9y&L-skXlw#Uvq>7S9Us&E z#OL0N9WWI5jv$w)07uq)w zs*|+Vzn6fE{N1W^!?z)sNBz@~xtel+yWv(g^Fap=VZx=*{`zpc=J?Ak%-(w(Qa@cW086MvaH$fx$5AuU6z4!C2iTyj*V998qZzjrtn_Fvh*&~v3{u|4Hz zp6^|A7c9(OnOk`A%8QG!l)G(y;I|I<^^0p($AayOEm_E|SV!=K^&8e#3a&bS&y@V_ zak%;%IPo8jdyZc@S6pNFoNqP%OWoDE#kxCnAGmklbIu=Ia*ctHZ7O*sd|kXIe!uD0 zbk`bf`TD?114-w|s^p(@!HZC6fdj_Jsv$tl!yEpBAWJV_D>q9ROy=q%*ynX1_ zp_Hfd1A6Dj4-1fcoRaYid+(t4Sb-TkuvnhpQ#MKG04-g`fLM16oO!2VooP$t8pcln zfQ3gDwn+ zd;AhE;e{j3V(ErTt%`O7k|rCQE&<&u+jn1lyk6GqMG`CAnvG1KL4a z?@I0oY(Tqp^w0*hMtf6h=7ALzl_pE~CHM2m@C(V%+2r{bR|>x}@48=DmMq_y+%uIF zW67ELIzoR%^%@=_RMcn4)OIHi%p}jAOTHjM?Tb`DsXb dict: @@ -56,22 +63,45 @@ def load_plugs() -> dict: class PlugPoller: - """Holds a reused device session per plug; re-handshakes only on error.""" + """Holds a reused device session per plug, re-handshaking in place on error. - def __init__(self, client: ApiClient, name: str, host: str): - self.client = client + Each poller owns its own ApiClient. Sharing one across plugs makes their + handshakes race — the plugs answer 403 Forbidden for every request carrying + a session another plug's handshake has since superseded. + """ + + def __init__(self, user: str, password: str, name: str, host: str): + self.client = ApiClient(user, password) self.name = name self.host = host self.device = None + async def _read(self): + """One full read against the current session, handshaking if needed.""" + if self.device is None: + self.device = await asyncio.wait_for(self.client.p110(self.host), OP_TIMEOUT) + info = await asyncio.wait_for(self.device.get_device_info(), OP_TIMEOUT) + energy = await asyncio.wait_for(self.device.get_energy_usage(), OP_TIMEOUT) + power = await asyncio.wait_for(self.device.get_current_power(), OP_TIMEOUT) + return info, energy, power + async def poll(self): labels = {"plug": self.name} try: - if self.device is None: - self.device = await asyncio.wait_for(self.client.p110(self.host), OP_TIMEOUT) - info = await asyncio.wait_for(self.device.get_device_info(), OP_TIMEOUT) - energy = await asyncio.wait_for(self.device.get_energy_usage(), OP_TIMEOUT) - power = await asyncio.wait_for(self.device.get_current_power(), OP_TIMEOUT) + try: + info, energy, power = await self._read() + except Exception as first: + # A dropped session is the overwhelmingly common failure here and + # it is recoverable immediately — the plug is fine, our token is + # not. Retrying on the NEXT cycle instead of this one is what + # pinned poll success at ~50%: every session survived exactly one + # poll, so the loop alternated success, expiry, re-handshake, + # expiry, forever. Reconnect and retry once, in place. + log.info("session lost for plug %s (%s): %s — re-handshaking", + self.name, self.host, first) + metric_reauths.labels(**labels).inc() + self.device = None + info, energy, power = await self._read() metric_up.labels(**labels).set(1) metric_state.labels(**labels).set(1 if info.device_on else 0) @@ -80,10 +110,16 @@ async def poll(self): metric_on_since.labels(**labels).set(info.on_time) metric_today.labels(**labels).set(energy.today_energy / 1000.0) # Wh -> kWh metric_month.labels(**labels).set(energy.month_energy / 1000.0) # Wh -> kWh + metric_last_ok.labels(**labels).set(time.time()) except Exception as e: + # Both the original attempt and the retry failed, so this is not a + # stale session — the plug is genuinely unreachable, off the Wi-Fi, + # or the credentials are wrong. metric_up.labels(**labels).set(0) + metric_fails.labels(**labels).inc() self.device = None # force a fresh handshake next cycle - log.warning("poll failed for plug %s (%s): %s", self.name, self.host, e) + log.warning("poll failed for plug %s (%s) after re-handshake: %s", + self.name, self.host, e) async def poll_loop(pollers, interval: int): @@ -105,8 +141,7 @@ def main(): sys.exit(1) plugs = load_plugs() - client = ApiClient(user, password) - pollers = [PlugPoller(client, name, host) for name, host in plugs.items()] + pollers = [PlugPoller(user, password, name, host) for name, host in plugs.items()] # Start the HTTP server FIRST (daemon thread) so /metrics — and thus the k8s # readiness probe — is serving immediately, even if the plugs are unreachable. diff --git a/test_poller.py b/test_poller.py new file mode 100644 index 0000000..d958e8d --- /dev/null +++ b/test_poller.py @@ -0,0 +1,127 @@ +#!/usr/bin/env python3 +"""Regression test for the ~50% poll-success bug. + +In the field, every Tapo session survived exactly one poll and then answered +403 Forbidden / Tapo(SessionTimeout). The old poll() reacted by clearing the +device and giving up until the NEXT cycle, so the loop alternated success, +expiry, re-handshake, expiry — pinning `avg_over_time(tapo_plug_up[3h])` at +0.503 on all three plugs while power gauges kept updating from the good half. + +This test models exactly that plug: a session that dies after one full read. +The poller must absorb it and report every poll as successful. + +Run: python3 test_poller.py (needs prometheus-client; no network, no plugs) +""" +import asyncio +import importlib.util +import os +import sys +import types + + +# ── a plug whose session dies after exactly one complete read ──────────────── +class FakeInfo: + device_on, rssi, on_time = True, -55, 1234 + + +class FakeEnergy: + today_energy, month_energy = 612.0, 18000.0 + + +class FakePower: + current_power = 50.0 + + +class FakeDevice: + def __init__(self): + self.spent = False + + async def _guard(self): + if self.spent: + raise RuntimeError("Tapo(SessionTimeout): 403 Forbidden") + + async def get_device_info(self): + await self._guard() + return FakeInfo() + + async def get_energy_usage(self): + await self._guard() + return FakeEnergy() + + async def get_current_power(self): + await self._guard() + self.spent = True # the read that consumes the session + return FakePower() + + +class FakeApiClient: + def __init__(self, user, password): + self.handshakes = 0 + + async def p110(self, host): + self.handshakes += 1 + return FakeDevice() + + +# The exporter is a single script, not a package — stub `tapo` before loading it. +fake_tapo = types.ModuleType("tapo") +fake_tapo.ApiClient = FakeApiClient +sys.modules["tapo"] = fake_tapo + +_here = os.path.dirname(os.path.abspath(__file__)) +_spec = importlib.util.spec_from_file_location( + "tapo_exporter", os.path.join(_here, "tapo-exporter.py")) +exporter = importlib.util.module_from_spec(_spec) +_spec.loader.exec_module(exporter) + +PLUG = "DGX Spark" + + +def gauge(metric): + return metric.labels(plug=PLUG)._value.get() + + +async def test_expiring_session_still_polls_cleanly(): + poller = exporter.PlugPoller("user", "password", PLUG, "192.168.10.217") + + ups = [] + for _ in range(10): + await poller.poll() + ups.append(gauge(exporter.metric_up)) + + rate = sum(ups) / len(ups) + print(f"up per poll : {[int(u) for u in ups]}") + print(f"success rate : {rate:.0%}") + print(f"reauths : {int(gauge(exporter.metric_reauths))}") + print(f"hard failures: {int(gauge(exporter.metric_fails))}") + + assert rate == 1.0, f"expected every poll to succeed, got {rate:.0%}" + assert gauge(exporter.metric_fails) == 0, "a recoverable expiry was reported as a failure" + assert gauge(exporter.metric_power) == 50.0, "power gauge not populated" + assert gauge(exporter.metric_last_ok) > 0, "last-success timestamp not set" + + +async def test_unreachable_plug_is_reported_down(): + """The retry must not paper over a plug that is genuinely gone.""" + class DeadClient(FakeApiClient): + async def p110(self, host): + raise RuntimeError("connection refused") + + poller = exporter.PlugPoller("user", "password", PLUG, "192.168.10.217") + poller.client = DeadClient("user", "password") + before = gauge(exporter.metric_fails) + await poller.poll() + + assert gauge(exporter.metric_up) == 0, "unreachable plug should report up=0" + assert gauge(exporter.metric_fails) == before + 1, "failure counter not incremented" + print("unreachable plug correctly reported down") + + +async def main(): + await test_expiring_session_still_polls_cleanly() + await test_unreachable_plug_is_reported_down() + print("\nOK") + + +if __name__ == "__main__": + asyncio.run(main())