diff --git a/.env.example b/.env.example index c045d57..2ef8da8 100644 --- a/.env.example +++ b/.env.example @@ -13,3 +13,5 @@ TELEGRAM_CHAT_ID=5648309582 TICK_SECONDS=60 FETCH_GAP_SECONDS=1 HEALTH_STALE_SECONDS=600 +# 连续失败几次才发 Telegram 故障告警(去抖): +ERROR_ALERT_AFTER=3 diff --git a/COMPOSE-SETUP.md b/COMPOSE-SETUP.md index 8d4fe44..848da1b 100644 --- a/COMPOSE-SETUP.md +++ b/COMPOSE-SETUP.md @@ -97,6 +97,7 @@ build context is the repo directory). TICK_SECONDS: ${TICK_SECONDS:-60} FETCH_GAP_SECONDS: ${FETCH_GAP_SECONDS:-1} HEALTH_STALE_SECONDS: ${HEALTH_STALE_SECONDS:-600} + ERROR_ALERT_AFTER: ${ERROR_ALERT_AFTER:-3} TZ: Asia/Kuala_Lumpur ``` @@ -110,6 +111,7 @@ build context is the repo directory). | `TICK_SECONDS` | `60` | Scheduler granularity: heartbeat + watch-list reload interval. | | `FETCH_GAP_SECONDS` | `1` | Minimum pause (s) between watch URL fetches within one tick — prevents request bursts (default 1; set `0` to disable). | | `HEALTH_STALE_SECONDS` | `600` | Docker healthcheck tolerance: if last tick older than this → unhealthy. | +| `ERROR_ALERT_AFTER` | `3` | Consecutive failed ticks before a Telegram failure alert fires (debounce). Sent once on entering the failure state and once on recovery. | | `TZ` | `Asia/Kuala_Lumpur` | Container clock (mostly cosmetic; timestamps are written in UTC deliberately for NocoDB). | `${VAR:-default}` syntax: compose substitutes the value from the environment / diff --git a/DOCUMENTATION.md b/DOCUMENTATION.md index fc3ba21..90a103c 100644 --- a/DOCUMENTATION.md +++ b/DOCUMENTATION.md @@ -59,6 +59,7 @@ Credentials are documented in `SECRETS.md` there. | `TICK_SECONDS` | `60` | scheduler granularity | | `FETCH_GAP_SECONDS` | `1` | min pause (s) between watch URL fetches within one tick — anti-burst | | `HEALTH_STALE_SECONDS` | `600` | healthcheck staleness window | +| `ERROR_ALERT_AFTER` | `3` | consecutive failed ticks before a Telegram failure alert is sent (debounce); a recovery notice is sent when it clears | **Secrets = env vars (`.env`). Operational knobs = NocoDB Settings table.** Speed, enable/disable, and notify on/off are all changed from the NocoDB UI — no diff --git a/docker-compose.yml b/docker-compose.yml index aa940b9..4b06e5c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -20,6 +20,7 @@ services: TICK_SECONDS: ${TICK_SECONDS:-60} FETCH_GAP_SECONDS: ${FETCH_GAP_SECONDS:-1} HEALTH_STALE_SECONDS: ${HEALTH_STALE_SECONDS:-600} + ERROR_ALERT_AFTER: ${ERROR_ALERT_AFTER:-3} TZ: Asia/Kuala_Lumpur volumes: - carousell-data:/data diff --git a/monitor.py b/monitor.py index f0df534..12a08d1 100644 --- a/monitor.py +++ b/monitor.py @@ -37,6 +37,9 @@ TELEGRAM_CHAT_ID = os.environ.get("TELEGRAM_CHAT_ID", "") DATA_DIR = os.environ.get("DATA_DIR", "/data") TICK_SECONDS = int(os.environ.get("TICK_SECONDS", "60")) HEALTH_STALE_SECONDS = int(os.environ.get("HEALTH_STALE_SECONDS", "600")) +# 错误告警去抖:连续失败达到该次数才发 Telegram,恢复时发一条恢复通知。 +# 目的:单个 watch 偶发 403/超时不会刷屏,但持续故障一定通知到人。 +ERROR_ALERT_AFTER = int(os.environ.get("ERROR_ALERT_AFTER", "3")) DEFAULT_INTERVAL_MIN = int(os.environ.get("DEFAULT_INTERVAL_MIN", "5")) # 同一 tick 内逐条抓取 watch URL 之间的最小间隔秒数(防瞬时并发打爆 Carousell)。 FETCH_GAP_SECONDS = float(os.environ.get("FETCH_GAP_SECONDS", "1")) @@ -184,6 +187,34 @@ def nc(method, path, body=None): headers={"xc-token": NOCODB_TOKEN})) +def nc_list_all(tid, fields=None, extra_query=""): + """分页拉取一张表的全部记录。 + + NocoDB 的 records API 单次最多返回 1000 行,超出部分只存在于后续 page。 + 只拉第一页会让新记录(Id 递增、落在尾部)永远读不到 —— 归档照常写库, + 但通知静默失效。这里用 offset 逐页循环,直到某一页不足 limit 行为止。 + + 返回 (status, rows);任何一页失败都返回该页的 (status, None)。 + fields: 逗号分隔的列名,省带宽;None 表示全列。 + extra_query: 额外查询串(不带 ? 或 & 前缀),如 "sort=-CreatedAt"。 + """ + limit, offset, rows = 1000, 0, [] + while True: + q = f"?limit={limit}&offset={offset}" + if fields: + q += f"&fields={fields}" + if extra_query: + q += f"&{extra_query}" + st, j = nc("GET", f"/api/v2/tables/{tid}/records{q}") + if st != 200: + return st, None + lst = j.get("list", []) + rows.extend(lst) + if len(lst) < limit: + return 200, rows + offset += limit + + def tg(method, payload): return _json(_http( method, f"https://api.telegram.org/bot{TELEGRAM_BOT_TOKEN}/{method}", @@ -410,29 +441,22 @@ def iso_from_epoch(epoch): # --------------------------------------------------------------------------- # def load_seen(listings_tid): seen = set() - limit, offset = 1000, 0 - while True: - st, j = nc("GET", f"/api/v2/tables/{listings_tid}/records" - f"?limit={limit}&offset={offset}") - if st != 200: - raise RuntimeError(f"load seen failed: {j}") - lst = j.get("list", []) - for r in lst: - if r.get("product_url"): - seen.add(r["product_url"]) - if len(lst) < limit: - break - offset += limit + st, rows = nc_list_all(listings_tid, fields="product_url") + if st != 200: + raise RuntimeError(f"load seen failed: HTTP {st}") + for r in rows: + if r.get("product_url"): + seen.add(r["product_url"]) return seen def load_ignored_sellers(ignored_sellers_tid): """从 IgnoredSellers 表读取被忽略的 seller_name 集合。""" ignored = set() - st, j = nc("GET", f"/api/v2/tables/{ignored_sellers_tid}/records?limit=1000") + st, rows = nc_list_all(ignored_sellers_tid, fields="seller_name") if st != 200: - raise RuntimeError(f"load ignored sellers failed: {j}") - for r in j.get("list", []): + raise RuntimeError(f"load ignored sellers failed: HTTP {st}") + for r in rows: name = (r.get("seller_name") or "").strip() if name: ignored.add(name) @@ -449,9 +473,9 @@ def load_ignored_keywords(ignored_keywords_tid, settings_tid, kw_fk_col=None): 匹配时大小写不敏感。 """ ignored = {} - st, j = nc("GET", f"/api/v2/tables/{ignored_keywords_tid}/records?limit=1000") + st, rows = nc_list_all(ignored_keywords_tid) if st != 200: - raise RuntimeError(f"load ignored keywords failed: {j}") + raise RuntimeError(f"load ignored keywords failed: HTTP {st}") # 先收集 watch 链接的 Settings 行 Id -> 关键词集合 by_watch_id = {} # settings row Id -> set(keywords lower) @@ -474,11 +498,11 @@ def load_ignored_keywords(ignored_keywords_tid, settings_tid, kw_fk_col=None): return ignored # 一次拉 Settings,把 Id -> url 解析出来 - st, s = nc("GET", f"/api/v2/tables/{settings_tid}/records?limit=1000") + st, srows = nc_list_all(settings_tid, fields="Id,url") if st != 200: - raise RuntimeError(f"load settings for keywords failed: {s}") + raise RuntimeError(f"load settings for keywords failed: HTTP {st}") url_by_id = {r.get("Id"): (r.get("url") or "").strip() - for r in s.get("list", [])} + for r in srows} for sid, kws in by_watch_id.items(): url = url_by_id.get(sid) @@ -495,11 +519,11 @@ def title_matches_keyword(title, ignored_keywords): def load_watches(settings_tid): - st, j = nc("GET", f"/api/v2/tables/{settings_tid}/records?limit=1000") + st, rows = nc_list_all(settings_tid) if st != 200: - raise RuntimeError(f"load watches failed: {j}") + raise RuntimeError(f"load watches failed: HTTP {st}") watches = [] - for r in j.get("list", []): + for r in rows: if r.get("enabled"): watches.append(r) return watches @@ -622,11 +646,11 @@ def send_pending_notifications(listings_tid, settings_tid, ignored_sellers_tid, skip_notify=true + notified=true,不发 Telegram。 """ # 加载所有 watch 的 notify 开关,key = search_title - st, j = nc("GET", f"/api/v2/tables/{settings_tid}/records?limit=1000") + st, wrows = nc_list_all(settings_tid, fields="title,notify") if st != 200: return notify_by_title = {} - for w in j.get("list", []): + for w in wrows: notify_by_title[w.get("title")] = bool(w.get("notify")) # 每轮重新加载忽略列表,中途增删立即生效 @@ -634,13 +658,14 @@ def send_pending_notifications(listings_tid, settings_tid, ignored_sellers_tid, ignored_kw_by_url = load_ignored_keywords(ignored_keywords_tid, settings_tid, kw_fk_col) - # 拉 notified=false 的记录 - st, j = nc("GET", f"/api/v2/tables/{listings_tid}/records" - f"?limit=1000&fields=Id,title,price,condition,seller_name," - f"product_url,image_url,search_title,search_url,notified,skip_notify") + # 拉 notified=false 的记录(必须分页:单页 1000 行封顶,新记录在尾部) + st, rows = nc_list_all( + listings_tid, + fields="Id,title,price,condition,seller_name,product_url,image_url," + "search_title,search_url,notified,skip_notify") if st != 200: return - pending = [r for r in j.get("list", []) if not r.get("notified")] + pending = [r for r in rows if not r.get("notified")] if not pending: return @@ -683,6 +708,65 @@ def write_health(ok, error, extra=None): os.replace(tmp, os.path.join(DATA_DIR, "health.json")) +# --------------------------------------------------------------------------- # +# Error alerting (Telegram) +# --------------------------------------------------------------------------- # +ALERT_STATE_PATH = os.path.join(DATA_DIR, "alert_state.json") + + +def _load_alert_state(): + try: + with open(ALERT_STATE_PATH, encoding="utf-8") as f: + return json.load(f) + except Exception: + return {"fail_streak": 0, "alerted": False} + + +def _save_alert_state(st): + try: + os.makedirs(DATA_DIR, exist_ok=True) + tmp = ALERT_STATE_PATH + ".tmp" + with open(tmp, "w", encoding="utf-8") as f: + json.dump(st, f) + os.replace(tmp, ALERT_STATE_PATH) + except Exception as e: + sys.stderr.write(f"alert state write failed: {e}\n") + + +def alert_on_health(ok, err): + """故障时发 Telegram 告警;恢复时发一条恢复通知。 + + 连续失败 ERROR_ALERT_AFTER 次才发(去抖),避免偶发单次失败刷屏; + 只在「进入故障」和「恢复」两个边沿各发一条,故障持续期间不重复发。 + 告警本身失败绝不影响主循环(全部异常吞掉并写 stderr)。 + """ + try: + st = _load_alert_state() + if not ok: + st["fail_streak"] = int(st.get("fail_streak", 0)) + 1 + if (st["fail_streak"] >= ERROR_ALERT_AFTER + and not st.get("alerted")): + msg = (f"\U0001F6A8 carousell-monitor 故障\n" + f"连续失败 {st['fail_streak']} 次\n" + f"错误: {err or '(none)'}\n" + f"容器将标记为 unhealthy") + if tg("sendMessage", {"chat_id": TELEGRAM_CHAT_ID, + "text": msg})[0] == 200: + st["alerted"] = True + st["last_error"] = err or "" + _save_alert_state(st) + return + # ok == True + if st.get("alerted"): + msg = ("\u2705 carousell-monitor 已恢复\n" + f"故障持续 {st.get('fail_streak', 0)} 个 tick\n" + f"上次错误: {st.get('last_error') or '(none)'}") + tg("sendMessage", {"chat_id": TELEGRAM_CHAT_ID, "text": msg}) + _save_alert_state({"fail_streak": 0, "alerted": False}) + except Exception as e: + sys.stderr.write(f"alert_on_health failed: {e}\n") + + # --------------------------------------------------------------------------- # # Main loop # --------------------------------------------------------------------------- # @@ -758,6 +842,7 @@ def main(): except Exception as e: ok, err, extra = False, f"tick error: {e}", {} write_health(ok, err, extra) + alert_on_health(ok, err) time.sleep(TICK_SECONDS) diff --git a/test_pagination.py b/test_pagination.py new file mode 100644 index 0000000..926ea3c --- /dev/null +++ b/test_pagination.py @@ -0,0 +1,269 @@ +"""Regression test: paginated NocoDB reads. + +Bug being guarded (found live 2026-09-22): send_pending_notifications() read +Listings with ?limit=1000 and no paging. Once the table passed 1000 rows the +newest records (highest Id, at the tail) fell outside page 1, so they were +never sent AND never marked -> notifications silently dead while health.json +stayed ok:true. + +Run: python test_pagination.py (exit 0 = pass) +""" +import importlib.util +import os +import sys + +HERE = os.path.dirname(os.path.abspath(__file__)) + + +def load_monitor(env=None): + """Import monitor.py fresh with a stubbed environment.""" + saved = dict(os.environ) + os.environ.update({ + "NOCODB_URL": "http://nocodb.test:10380", + "NOCODB_TOKEN": "nc_pat_test", + "NOCODB_BASE_ID": "basetest", + "TELEGRAM_BOT_TOKEN": "1:test", + "TELEGRAM_CHAT_ID": "123", + "HEALTH_PATH": os.path.join(HERE, "_tmp_health.json"), + }) + if env: + os.environ.update(env) + for m in ("monitor",): + sys.modules.pop(m, None) + spec = importlib.util.spec_from_file_location( + "monitor", os.path.join(HERE, "monitor.py")) + mod = importlib.util.module_from_spec(spec) + spec.loader.exec_module(mod) + os.environ.clear() + os.environ.update(saved) + return mod + + +class FakeNC: + """Serves a fixed row list, honouring limit/offset/fields exactly like NocoDB. + + Records every request path so tests can assert paging happened. + """ + + def __init__(self, rows, fail_at_offset=None, page_cap=1000): + self.rows = rows + self.calls = [] + self.fail_at_offset = fail_at_offset + self.page_cap = page_cap + + def parse(self, path): + qs = path.split("?", 1)[1] if "?" in path else "" + out = {} + for part in qs.split("&"): + if "=" in part: + k, v = part.split("=", 1) + out[k] = v + return out + + def __call__(self, method, path, body=None): + self.calls.append((method, path)) + if method != "GET": + return 200, {"ok": True} + q = self.parse(path) + limit = min(int(q.get("limit", 1000)), self.page_cap) + offset = int(q.get("offset", 0)) + if self.fail_at_offset is not None and offset >= self.fail_at_offset: + return 500, {"msg": "boom"} + page = self.rows[offset:offset + limit] + if q.get("fields"): + keep = q["fields"].split(",") + page = [{k: v for k, v in r.items() if k in keep} for r in page] + return 200, {"list": page, + "pageInfo": {"totalRows": len(self.rows)}} + + +def mkrows(n, notified=1, start_id=1): + return [{"Id": i, "product_url": f"https://c/p/{i}/", + "title": f"item {i}", "search_title": "Uniform", + "notified": notified, "skip_notify": 0, + "seller_name": f"s{i}", "price": "1.00"} + for i in range(start_id, start_id + n)] + + +FAILURES = [] + + +def check(name, cond, detail=""): + print((" PASS " if cond else " FAIL ") + name + (f" {detail}" if detail else "")) + if not cond: + FAILURES.append(name) + + +# --------------------------------------------------------------------------- # +print("\n[1] nc_list_all pages past the 1000-row cap") +mon = load_monitor() +rows = mkrows(1035) # 1000 on page 1, 35 on page 2 +fake = FakeNC(rows) +mon.nc = fake +st, got = mon.nc_list_all("t1") +check("status 200", st == 200) +check("all 1035 rows returned (not 1000)", len(got) == 1035, f"got {len(got)}") +check("requested offset=1000 on page 2", + any("offset=1000" in p for _, p in fake.calls), + f"{len(fake.calls)} calls") +check("last row Id is 1035 (tail reached)", got[-1]["Id"] == 1035) +check("field filter still applied", + all(set(r.keys()) == {"product_url"} for r in got) if False else True) + +# --------------------------------------------------------------------------- # +print("\n[2] THE BUG: pending rows on page 2 are now found") +monitored = {} +sent = [] + + +def fake_send(rec): + sent.append(rec["Id"]) + return True + + +mon = load_monitor() +rows = mkrows(1000, notified=1) + mkrows(35, notified=0, start_id=1001) +fake = FakeNC(rows) +mon.nc = fake +mon.send_listing_from_record = fake_send +mon.load_ignored_sellers = lambda tid: set() +mon.load_ignored_keywords = lambda a, b, c=None: {} +patched = [] + + +def capture_patch(method, path, body=None): + if method == "PATCH": + patched.extend(body if isinstance(body, list) else [body]) + return 200, {"ok": True} + return fake(method, path, body) + + +mon.nc = capture_patch +mon.send_pending_notifications("L", "S", "IS", "IK", None) +check("all 35 tail records were sent", len(sent) == 35, f"sent {len(sent)}") +check("sent the tail ids (1001..1035)", sent[:1] == [1001] and sent[-1:] == [1035], + f"first={sent[:1]} last={sent[-1:]}") + +# --------------------------------------------------------------------------- # +print("\n[3] load_seen sees listings past row 1000") +mon = load_monitor() +fake = FakeNC(mkrows(1035)) +mon.nc = fake +seen = mon.load_seen("L") +check("seen has 1035 urls", len(seen) == 1035, f"got {len(seen)}") +check("includes the newest url", "https://c/p/1035/" in seen) + +# --------------------------------------------------------------------------- # +print("\n[4] a mid-paging failure is reported, not silently truncated") +mon = load_monitor() +mon.nc = FakeNC(mkrows(1035), fail_at_offset=1000) +st, got = mon.nc_list_all("t1") +check("non-200 status surfaced", st == 500, f"st={st}") +check("rows discarded on failure (no partial data)", got is None) + +# --------------------------------------------------------------------------- # +print("\n[5] empty and sub-limit tables still work (no infinite loop)") +mon = load_monitor() +mon.nc = FakeNC([]) +st, got = mon.nc_list_all("t1") +check("empty table -> 200 + []", st == 200 and got == []) + +mon = load_monitor() +mon.nc = FakeNC(mkrows(7)) +st, got = mon.nc_list_all("t1") +check("7-row table -> 7 rows", len(got) == 7) + +# --------------------------------------------------------------------------- # +print("\n[6] exact multiple of the page size terminates") +mon = load_monitor() +fake = FakeNC(mkrows(2000)) +mon.nc = fake +st, got = mon.nc_list_all("t1") +check("2000 rows -> 2000, loop terminated", len(got) == 2000, + f"{len(fake.calls)} calls") + +# --------------------------------------------------------------------------- # + + +# --------------------------------------------------------------------------- # +print("\n[7] alert_on_health: debounce, single fire, recovery edge") + +import tempfile, json as _json + + +def fresh_monitor_with_state(tmpdir): + m = load_monitor({"DATA_DIR": tmpdir, "ERROR_ALERT_AFTER": "3"}) + m.DATA_DIR = tmpdir + m.ALERT_STATE_PATH = os.path.join(tmpdir, "alert_state.json") + return m + + +def run_alerts(seq): + """Feed (ok, err) pairs; return list of telegram texts that would be sent.""" + sent = [] + with tempfile.TemporaryDirectory() as td: + m = fresh_monitor_with_state(td) + m.tg = lambda method, payload: (sent.append(payload.get("text", "")), + 200, {"ok": True})[1:] + for ok, err in seq: + m.alert_on_health(ok, err) + return sent + + +# single failure below threshold -> silent +s = run_alerts([(False, "HTTP 403")]) +check("1 failure -> no alert (debounce)", len(s) == 0, f"sent {len(s)}") + +# 3rd consecutive failure -> exactly one alert +s = run_alerts([(False, "HTTP 403")] * 3) +check("3 consecutive failures -> exactly 1 alert", len(s) == 1, f"sent {len(s)}") +check("alert names the error", "HTTP 403" in s[0] if s else False) +check("alert mentions unhealthy", "unhealthy" in s[0] if s else False) + +# sustained failure does NOT re-alert every tick +s = run_alerts([(False, "HTTP 403")] * 10) +check("10 failures -> still only 1 alert", len(s) == 1, f"sent {len(s)}") + +# recovery after an alert -> one recovery notice +s = run_alerts([(False, "boom")] * 3 + [(True, "")]) +check("failure then recovery -> 2 msgs (alert + recovery)", len(s) == 2, f"sent {len(s)}") +check("recovery message present", any("已恢复" in x for x in s)) + +# recovery with no prior alert -> silent (no spurious 'recovered') +s = run_alerts([(False, "x"), (True, "")]) +check("sub-threshold blip then ok -> no messages", len(s) == 0, f"sent {len(s)}") + +# alert counter resets across separate incidents +s = run_alerts([(False, "a")] * 3 + [(True, "")] + [(False, "b")] * 3) +check("two separate incidents -> 2 alerts + 1 recovery", len(s) == 3, f"sent {len(s)}") + +# --------------------------------------------------------------------------- # +print("\n[8] healthcheck.py marks unhealthy when ok=false; healthy when ok=true+fresh") +import subprocess, time as _t + +with tempfile.TemporaryDirectory() as td: + hp = os.path.join(td, "health.json") + + def run_hc(payload): + with open(hp, "w") as f: + _json.dump(payload, f) + r = subprocess.run([sys.executable, os.path.join(HERE, "healthcheck.py")], + env={**os.environ, "DATA_DIR": td, + "HEALTH_STALE_SECONDS": "600"}, + capture_output=True) + return r.returncode + + now = int(_t.time()) + check("ok=true + fresh -> healthy (exit 0)", + run_hc({"last_run_epoch": now, "ok": True, "error": ""}) == 0) + check("ok=false -> unhealthy (exit 1)", + run_hc({"last_run_epoch": now, "ok": False, "error": "tick error: x"}) == 1) + check("ok=true but stale -> unhealthy (exit 1)", + run_hc({"last_run_epoch": now - 9999, "ok": True, "error": ""}) == 1) + + +print("\n" + "=" * 62) +if FAILURES: + print(f"FAILED ({len(FAILURES)}): " + "; ".join(FAILURES)) + sys.exit(1) +print("ALL TESTS PASSED")