Carousell new-listing monitor: NocoDB archive + Telegram alerts
This commit is contained in:
@@ -0,0 +1,9 @@
|
||||
.git
|
||||
.gitignore
|
||||
.env
|
||||
.env.example
|
||||
__pycache__
|
||||
*.pyc
|
||||
*.md
|
||||
SECRETS.md
|
||||
docker-compose.yml
|
||||
@@ -0,0 +1,14 @@
|
||||
# Copy to .env and fill in. All values are required.
|
||||
|
||||
# NocoDB (container DNS when on bridge_hoelee; LAN IP from a desktop):
|
||||
NOCODB_URL=http://nocodb:10380
|
||||
NOCODB_BASE_ID=poqw1zjw3hnsk37
|
||||
NOCODB_TOKEN=
|
||||
|
||||
# Telegram alerts (@HoeleeAgentBot):
|
||||
TELEGRAM_BOT_TOKEN=
|
||||
TELEGRAM_CHAT_ID=5648309582
|
||||
|
||||
# Optional tuning:
|
||||
TICK_SECONDS=60
|
||||
HEALTH_STALE_SECONDS=600
|
||||
@@ -0,0 +1,4 @@
|
||||
.env
|
||||
__pycache__/
|
||||
*.pyc
|
||||
.DS_Store
|
||||
@@ -0,0 +1,47 @@
|
||||
# AGENTS.md
|
||||
|
||||
Project: Carousell new-listing monitor (Python stdlib, Docker, NocoDB, Telegram).
|
||||
|
||||
## What it does
|
||||
|
||||
`monitor.py` polls Carousell search URLs (sort_by=3 = recent), extracts listings
|
||||
from the server-rendered `<script type="application/json">` Redux state
|
||||
(`SearchListing.listingCards`), dedupes by `product_url` (param-less), archives to a
|
||||
NocoDB base, and alerts Telegram `"<title>: N new listings"`. Runs 24/7 as a Docker
|
||||
container on DSM (network `bridge_hoelee`, reaches NocoDB at `http://nocodb:10380`).
|
||||
|
||||
## Iron rules
|
||||
|
||||
- Secrets NEVER in code or compose — only env vars / `SECRETS.md` (private repo).
|
||||
Operational knobs (`enabled` / `notify` / `check_interval_minutes`) live in the
|
||||
NocoDB **Settings** table, adjustable from the UI without redeploy.
|
||||
- Dedupe key is `product_url` (`https://www.carousell.com.my/p/<id>/`), not the raw
|
||||
listing id and never the query-string URL.
|
||||
- First run per watch seeds the archive with **no** Telegram alert (`last_checked_at`
|
||||
null == unseeded).
|
||||
- Docker HEALTHCHECK: container is unhealthy if the last tick is >10 min old or the
|
||||
last run had a failure (failed extract / 403 / rate-limit).
|
||||
- Image thumbnail: `image` Attachment field stores the remote URL (NocoDB hotlinks
|
||||
it — media.karousell.com is GCS-backed, `Access-Control-Allow-Origin: *`). The raw
|
||||
URL is also kept in `image_url`.
|
||||
|
||||
## Verified facts (2026-09-02)
|
||||
|
||||
- Carousell search page: 1.6 MB HTML, state blob ~1.27 MB, ~49 `listingCards` per
|
||||
load. Card fields: `listingID`, `title`, `price` (e.g. "RM85"), `thumbnailURL`,
|
||||
`seller.username`, `aboveFold[time_created].timestampContent.seconds.low`,
|
||||
`belowFold` paragraphs where `paragraph[1]` = condition (Like new/Brand new/...).
|
||||
- `https://www.carousell.com.my/p/<id>/` 301-redirects to the canonical slug URL.
|
||||
- NocoDB instance: DSM `http://192.168.137.2:10380` (IP drifts; container DNS
|
||||
`nocodb:10380` on bridge_hoelee). Base `Carousell` = `poqw1zjw3hnsk37`.
|
||||
Workspace token in SECRETS.md. v2 API: records use column *titles* as JSON keys;
|
||||
Attachment field accepts a JSON string `[{"path","mimetype","title"}]` and keeps the
|
||||
remote URL (does not re-host).
|
||||
|
||||
## Build / deploy
|
||||
|
||||
```bash
|
||||
docker build -t hoelee/carousell-monitor:latest .
|
||||
docker push hoelee/carousell-monitor:latest
|
||||
# then Portainer stack from /volume1/docker/carousell-monitor/docker-compose.yml
|
||||
```
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
FROM python:3.11-alpine
|
||||
|
||||
WORKDIR /app
|
||||
COPY monitor.py healthcheck.py /app/
|
||||
RUN mkdir -p /data
|
||||
|
||||
VOLUME ["/data"]
|
||||
|
||||
HEALTHCHECK --interval=60s --timeout=15s --start-period=120s --retries=3 \
|
||||
CMD python /app/healthcheck.py
|
||||
|
||||
CMD ["python", "-u", "/app/monitor.py"]
|
||||
@@ -0,0 +1,51 @@
|
||||
# carousell-monitor
|
||||
|
||||
Watches Carousell search pages (sorted by *recent*) for new listings, archives every
|
||||
listing to a NocoDB base (with image URL + thumbnail), and alerts Telegram.
|
||||
|
||||
## How it works
|
||||
|
||||
- `monitor.py` runs in a Docker container on DSM, self-bootstrapping its NocoDB
|
||||
schema (`Listings` + `Settings` tables) and looping forever.
|
||||
- Every `TICK_SECONDS` it reads the watch list from the **Settings** table and polls
|
||||
each enabled watch's URL on its own `check_interval_minutes`.
|
||||
- Dedupe key = `product_url` (param-less listing URL). First run per watch = seed
|
||||
archive only (no Telegram). After that, new listings are archived and alerted as
|
||||
`"<title>: N new listings"`.
|
||||
- The container marks itself **unhealthy** (Docker healthcheck) if a tick fails to
|
||||
extract / gets rate-limited / crashes.
|
||||
|
||||
## Schema
|
||||
|
||||
**Listings** — `product_url` (unique), `title`, `price` (numeric), `condition`
|
||||
(SingleSelect), `image_url`, `image` (Attachment → thumbnail), `seller_name`,
|
||||
`seller_url`, `search_title`, `search_url`, `listed_at`, `first_seen_at`.
|
||||
|
||||
**Settings** — `title`, `url`, `enabled`, `notify`, `check_interval_minutes`,
|
||||
`last_checked_at`. Add/remove watches here from the NocoDB UI; no redeploy needed.
|
||||
|
||||
## Run
|
||||
|
||||
```bash
|
||||
# local (against LAN NocoDB)
|
||||
NOCODB_URL=http://192.168.137.2:10380 \
|
||||
NOCODB_TOKEN=... TELEGRAM_BOT_TOKEN=... TELEGRAM_CHAT_ID=... \
|
||||
python monitor.py
|
||||
|
||||
# docker
|
||||
docker build -t hoelee/carousell-monitor:latest .
|
||||
docker compose up -d
|
||||
```
|
||||
|
||||
## Deploy (DSM via Portainer)
|
||||
|
||||
Image pushed to Docker Hub `hoelee/carousell-monitor:latest`; the compose at
|
||||
`/volume1/docker/carousell-monitor` is deployed as a Portainer stack with the
|
||||
secrets passed as stack environment variables.
|
||||
|
||||
## Files
|
||||
|
||||
- `monitor.py` — main loop, schema bootstrap, fetch/parse, NocoDB IO, Telegram.
|
||||
- `healthcheck.py` — Docker HEALTHCHECK probe (`/data/health.json`).
|
||||
- `Dockerfile`, `docker-compose.yml`, `.env.example`.
|
||||
- `AGENTS.md` — AI-agent entry. `SECRETS.md` — credentials (private repo).
|
||||
@@ -0,0 +1,30 @@
|
||||
# Carousell new-listing monitor -> NocoDB archive + Telegram alerts.
|
||||
# Deployed via Portainer on DSM. Secrets come from the stack's environment
|
||||
# (see .env.example); operational knobs (enabled / notify / interval) live in
|
||||
# the NocoDB "Settings" table.
|
||||
|
||||
services:
|
||||
carousell-monitor:
|
||||
image: hoelee/carousell-monitor:latest
|
||||
container_name: carousell-monitor
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
NOCODB_URL: ${NOCODB_URL:-http://nocodb:10380}
|
||||
NOCODB_TOKEN: ${NOCODB_TOKEN}
|
||||
NOCODB_BASE_ID: ${NOCODB_BASE_ID:-poqw1zjw3hnsk37}
|
||||
TELEGRAM_BOT_TOKEN: ${TELEGRAM_BOT_TOKEN}
|
||||
TELEGRAM_CHAT_ID: ${TELEGRAM_CHAT_ID}
|
||||
TICK_SECONDS: ${TICK_SECONDS:-60}
|
||||
HEALTH_STALE_SECONDS: ${HEALTH_STALE_SECONDS:-600}
|
||||
TZ: Asia/Kuala_Lumpur
|
||||
volumes:
|
||||
- carousell-data:/data
|
||||
networks:
|
||||
- bridge_hoelee
|
||||
|
||||
volumes:
|
||||
carousell-data:
|
||||
|
||||
networks:
|
||||
bridge_hoelee:
|
||||
external: true
|
||||
@@ -0,0 +1,19 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Docker HEALTHCHECK: healthy iff a tick ran recently AND the last run was ok."""
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
|
||||
DATA_DIR = os.environ.get("DATA_DIR", "/data")
|
||||
STALE = int(os.environ.get("HEALTH_STALE_SECONDS", "600"))
|
||||
|
||||
try:
|
||||
with open(os.path.join(DATA_DIR, "health.json"), encoding="utf-8") as f:
|
||||
h = json.load(f)
|
||||
age = time.time() - int(h.get("last_run_epoch", 0))
|
||||
if age <= STALE and h.get("ok") is True:
|
||||
sys.exit(0)
|
||||
except Exception:
|
||||
pass
|
||||
sys.exit(1)
|
||||
+418
@@ -0,0 +1,418 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Carousell new-listing monitor -> NocoDB archive + Telegram alerts.
|
||||
|
||||
Self-bootstrapping: on startup it ensures the NocoDB schema (Listings + Settings
|
||||
tables) exists, then loops forever. Each tick it reads the watch list from the
|
||||
Settings table, polls each enabled watch's Carousell search URL on its own
|
||||
interval, dedupes by product_url (param-less listing URL), archives every listing
|
||||
to NocoDB (title, numeric price, condition, image URL + thumbnail attachment,
|
||||
seller, link, timestamps), and — after the first seed — alerts Telegram with
|
||||
"<title>: N new listing(s)".
|
||||
|
||||
Secrets via env vars; operational knobs (enabled / notify / interval) live in the
|
||||
NocoDB Settings table so they're adjustable from the UI without a redeploy.
|
||||
|
||||
Writes /data/health.json every tick; the Docker HEALTHCHECK flags the container
|
||||
unhealthy on any failed extract / rate-limit / crash.
|
||||
|
||||
Stdlib only.
|
||||
"""
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import sys
|
||||
import time
|
||||
import urllib.error
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Config (env)
|
||||
# --------------------------------------------------------------------------- #
|
||||
NOCODB_URL = os.environ.get("NOCODB_URL", "http://nocodb:10380").rstrip("/")
|
||||
NOCODB_TOKEN = os.environ.get("NOCODB_TOKEN", "")
|
||||
NOCODB_BASE_ID = os.environ.get("NOCODB_BASE_ID", "poqw1zjw3hnsk37")
|
||||
TELEGRAM_BOT_TOKEN = os.environ.get("TELEGRAM_BOT_TOKEN", "")
|
||||
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"))
|
||||
DEFAULT_INTERVAL_MIN = int(os.environ.get("DEFAULT_INTERVAL_MIN", "5"))
|
||||
|
||||
UA = ("Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 "
|
||||
"(KHTML, like Gecko) Chrome/120.0 Safari/537.36")
|
||||
|
||||
CONDITIONS = ["Brand new", "Like new", "Lightly used", "Well used", "Heavily used"]
|
||||
PRODUCT_URL_TMPL = "https://www.carousell.com.my/p/{listing_id}/"
|
||||
SELLER_URL_TMPL = "https://www.carousell.com.my/u/{username}/"
|
||||
|
||||
# Column definitions: table title -> list of (title, uidt)
|
||||
LISTINGS_COLS = [
|
||||
("product_url", "URL"),
|
||||
("title", "SingleLineText"),
|
||||
("price", "Decimal"),
|
||||
("condition", "SingleSelect"), # options set via 2-pass below
|
||||
("image_url", "URL"),
|
||||
("image", "Attachment"),
|
||||
("seller_name", "SingleLineText"),
|
||||
("seller_url", "URL"),
|
||||
("search_title", "SingleLineText"),
|
||||
("search_url", "URL"),
|
||||
("listed_at", "DateTime"),
|
||||
("first_seen_at", "DateTime"),
|
||||
]
|
||||
SETTINGS_COLS = [
|
||||
("title", "SingleLineText"),
|
||||
("url", "URL"),
|
||||
("enabled", "Checkbox"),
|
||||
("notify", "Checkbox"),
|
||||
("check_interval_minutes", "Number"),
|
||||
("last_checked_at", "DateTime"),
|
||||
]
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# HTTP helpers
|
||||
# --------------------------------------------------------------------------- #
|
||||
def _http(method, url, body=None, headers=None, timeout=30):
|
||||
h = {"User-Agent": UA}
|
||||
if headers:
|
||||
h.update(headers)
|
||||
data = None
|
||||
if body is not None:
|
||||
data = json.dumps(body).encode("utf-8")
|
||||
h["Content-Type"] = "application/json"
|
||||
req = urllib.request.Request(url, data=data, method=method, headers=h)
|
||||
with urllib.request.urlopen(req, timeout=timeout) as r:
|
||||
raw = r.read().decode("utf-8", errors="ignore")
|
||||
return r.status, raw
|
||||
|
||||
|
||||
def _json(status_raw):
|
||||
status, raw = status_raw
|
||||
return status, (json.loads(raw) if raw else None)
|
||||
|
||||
|
||||
def nc(method, path, body=None):
|
||||
return _json(_http(method, f"{NOCODB_URL}{path}", body=body,
|
||||
headers={"xc-token": NOCODB_TOKEN}))
|
||||
|
||||
|
||||
def tg(method, payload):
|
||||
return _json(_http(
|
||||
method, f"https://api.telegram.org/bot{TELEGRAM_BOT_TOKEN}/{method}",
|
||||
body=payload))
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Schema bootstrap (idempotent)
|
||||
# --------------------------------------------------------------------------- #
|
||||
def _ensure_table(title, cols):
|
||||
st, tables = nc("GET", f"/api/v2/meta/bases/{NOCODB_BASE_ID}/tables")
|
||||
if st != 200:
|
||||
raise RuntimeError(f"list tables failed: {tables}")
|
||||
tid = None
|
||||
for t in tables.get("list", []):
|
||||
if t.get("title") == title:
|
||||
tid = t["id"]
|
||||
if tid is None:
|
||||
# NocoDB requires >=1 column inline on create; seed with the first.
|
||||
st, out = nc("POST", f"/api/v2/meta/bases/{NOCODB_BASE_ID}/tables",
|
||||
{"title": title, "table_name": title, "type": "table",
|
||||
"columns": [{"title": cols[0][0], "uidt": cols[0][1]}]})
|
||||
if st != 200:
|
||||
raise RuntimeError(f"create table {title} failed: {out}")
|
||||
tid = out["id"]
|
||||
|
||||
def _columns_by_title():
|
||||
st, meta = nc("GET", f"/api/v2/meta/tables/{tid}")
|
||||
return {c.get("title"): c for c in meta.get("columns", [])}
|
||||
|
||||
existing = _columns_by_title()
|
||||
for (ctitle, uidt) in cols:
|
||||
if ctitle in existing:
|
||||
continue
|
||||
st, out = nc("POST", f"/api/v2/meta/tables/{tid}/columns",
|
||||
{"title": ctitle, "uidt": uidt})
|
||||
if st != 200:
|
||||
raise RuntimeError(f"add column {ctitle} failed: {out}")
|
||||
|
||||
# SingleSelect options. Re-read meta for authoritative column ids: POST
|
||||
# /columns returns the *table* object, not the column id.
|
||||
existing = _columns_by_title()
|
||||
for (ctitle, uidt) in cols:
|
||||
if uidt != "SingleSelect" or ctitle not in existing:
|
||||
continue
|
||||
col = existing[ctitle]
|
||||
opts = col.get("colOptions", {}).get("options") or []
|
||||
if len(opts) >= len(CONDITIONS):
|
||||
continue
|
||||
opt_list = [{"title": c, "order": i + 1, "color": None}
|
||||
for i, c in enumerate(CONDITIONS)]
|
||||
st, out = nc("PATCH", f"/api/v2/meta/columns/{col['id']}",
|
||||
{"colOptions": {"options": opt_list},
|
||||
"dtxp": ",".join(CONDITIONS)})
|
||||
if st != 200:
|
||||
raise RuntimeError(f"patch SingleSelect {ctitle} failed: {out}")
|
||||
|
||||
return tid
|
||||
|
||||
|
||||
def bootstrap():
|
||||
listings_tid = _ensure_table("Listings", LISTINGS_COLS)
|
||||
settings_tid = _ensure_table("Settings", SETTINGS_COLS)
|
||||
return listings_tid, settings_tid
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Carousell extraction
|
||||
# --------------------------------------------------------------------------- #
|
||||
def fetch_listings(search_url):
|
||||
st, html = _http("GET", search_url)
|
||||
if st != 200:
|
||||
raise RuntimeError(f"carousell fetch HTTP {st}")
|
||||
blobs = re.findall(r'<script type="application/json">(.*?)</script>',
|
||||
html, re.S)
|
||||
if not blobs:
|
||||
raise RuntimeError("no application/json state found (blocked/ratelimited?)")
|
||||
state = json.loads(max(blobs, key=len))
|
||||
cards = state["SearchListing"]["listingCards"]
|
||||
out = []
|
||||
for c in cards:
|
||||
try:
|
||||
lid = int(c["listingID"])
|
||||
ts = None
|
||||
for item in c.get("aboveFold", []):
|
||||
if item.get("component") == "time_created":
|
||||
ts = item["timestampContent"]["seconds"]["low"]
|
||||
break
|
||||
bf = c.get("belowFold", [])
|
||||
title = next((i["stringContent"] for i in bf
|
||||
if i.get("component") == "header_1"), "")
|
||||
price_raw = next((i["stringContent"] for i in bf
|
||||
if i.get("component") == "header_2"), "")
|
||||
paras = [i.get("stringContent", "") for i in bf
|
||||
if i.get("component") == "paragraph"]
|
||||
cond = paras[1].strip() if len(paras) > 1 else ""
|
||||
if cond not in CONDITIONS:
|
||||
cond = ""
|
||||
thumb = c.get("thumbnailURL", "")
|
||||
seller = (c.get("seller") or {}).get("username", "")
|
||||
out.append({
|
||||
"listing_id": lid,
|
||||
"title": title,
|
||||
"price": price_raw,
|
||||
"condition": cond,
|
||||
"thumbnail": thumb,
|
||||
"seller": seller,
|
||||
"ts": ts,
|
||||
})
|
||||
except (KeyError, TypeError, ValueError):
|
||||
continue
|
||||
return out
|
||||
|
||||
|
||||
def parse_price(s):
|
||||
if not s:
|
||||
return None
|
||||
t = re.sub(r"[^0-9.]", "", s)
|
||||
if not t:
|
||||
return None
|
||||
try:
|
||||
return round(float(t), 2)
|
||||
except ValueError:
|
||||
return None
|
||||
|
||||
|
||||
def mimetype_for(url):
|
||||
p = url.lower()
|
||||
if ".png" in p:
|
||||
return "image/png"
|
||||
if ".webp" in p:
|
||||
return "image/webp"
|
||||
if ".gif" in p:
|
||||
return "image/gif"
|
||||
return "image/jpeg"
|
||||
|
||||
|
||||
def iso_now():
|
||||
# UTC: NocoDB parses naive datetimes as UTC.
|
||||
return time.strftime("%Y-%m-%d %H:%M:%S", time.gmtime())
|
||||
|
||||
|
||||
def iso_from_epoch(epoch):
|
||||
if not epoch:
|
||||
return None
|
||||
return time.strftime("%Y-%m-%d %H:%M:%S", time.gmtime(epoch))
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# NocoDB row IO
|
||||
# --------------------------------------------------------------------------- #
|
||||
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
|
||||
return seen
|
||||
|
||||
|
||||
def load_watches(settings_tid):
|
||||
st, j = nc("GET", f"/api/v2/tables/{settings_tid}/records?limit=1000")
|
||||
if st != 200:
|
||||
raise RuntimeError(f"load watches failed: {j}")
|
||||
watches = []
|
||||
for r in j.get("list", []):
|
||||
if r.get("enabled"):
|
||||
watches.append(r)
|
||||
return watches
|
||||
|
||||
|
||||
def insert_listings(listings_tid, rows):
|
||||
if not rows:
|
||||
return
|
||||
st, out = nc("POST", f"/api/v2/tables/{listings_tid}/records", rows)
|
||||
if st != 200:
|
||||
raise RuntimeError(f"insert listings failed: {out}")
|
||||
|
||||
|
||||
def update_checked(settings_tid, watch_id):
|
||||
nc("PATCH", f"/api/v2/tables/{settings_tid}/records",
|
||||
[{"Id": watch_id, "last_checked_at": iso_now()}])
|
||||
|
||||
|
||||
def build_row(l, watch):
|
||||
row = {
|
||||
"product_url": PRODUCT_URL_TMPL.format(listing_id=l["listing_id"]),
|
||||
"title": l["title"],
|
||||
"image_url": l["thumbnail"],
|
||||
"image": json.dumps([{"path": l["thumbnail"],
|
||||
"mimetype": mimetype_for(l["thumbnail"]),
|
||||
"title": f"{l['listing_id']}.jpg"}]),
|
||||
"seller_name": l["seller"],
|
||||
"seller_url": SELLER_URL_TMPL.format(username=l["seller"]),
|
||||
"search_title": watch.get("title", ""),
|
||||
"search_url": watch.get("url", ""),
|
||||
"first_seen_at": iso_now(),
|
||||
}
|
||||
price = parse_price(l["price"])
|
||||
if price is not None:
|
||||
row["price"] = price
|
||||
if l["condition"]:
|
||||
row["condition"] = l["condition"]
|
||||
listed = iso_from_epoch(l["ts"])
|
||||
if listed:
|
||||
row["listed_at"] = listed
|
||||
return row
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Telegram
|
||||
# --------------------------------------------------------------------------- #
|
||||
def send_telegram(text):
|
||||
if not TELEGRAM_BOT_TOKEN or not TELEGRAM_CHAT_ID:
|
||||
return
|
||||
try:
|
||||
tg("sendMessage", {"chat_id": TELEGRAM_CHAT_ID, "text": text})
|
||||
except urllib.error.HTTPError as e:
|
||||
sys.stderr.write(f"telegram send failed: {e.code} {e.read()[:200]}\n")
|
||||
|
||||
|
||||
def plural(n, word):
|
||||
return f"{n} {word}{'' if n == 1 else 's'}"
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Health
|
||||
# --------------------------------------------------------------------------- #
|
||||
def write_health(ok, error, extra=None):
|
||||
os.makedirs(DATA_DIR, exist_ok=True)
|
||||
h = {"last_run_epoch": int(time.time()), "ok": ok, "error": error or ""}
|
||||
if extra:
|
||||
h.update(extra)
|
||||
tmp = os.path.join(DATA_DIR, "health.json.tmp")
|
||||
with open(tmp, "w", encoding="utf-8") as f:
|
||||
json.dump(h, f)
|
||||
os.replace(tmp, os.path.join(DATA_DIR, "health.json"))
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Main loop
|
||||
# --------------------------------------------------------------------------- #
|
||||
def run_tick(listings_tid, settings_tid, seen, last_run):
|
||||
failures = []
|
||||
new_total = 0
|
||||
watches = load_watches(settings_tid)
|
||||
now = time.time()
|
||||
|
||||
for w in watches:
|
||||
wid = w.get("Id")
|
||||
interval_min = w.get("check_interval_minutes") or DEFAULT_INTERVAL_MIN
|
||||
interval_sec = max(int(interval_min), 1) * 60
|
||||
if wid in last_run and (now - last_run[wid]) < interval_sec:
|
||||
continue
|
||||
|
||||
try:
|
||||
listings = fetch_listings(w.get("url", ""))
|
||||
except Exception as e:
|
||||
failures.append(f"{w.get('title')}: {e}")
|
||||
# still advance so a hard-failing watch doesn't hammer every tick
|
||||
last_run[wid] = now
|
||||
continue
|
||||
|
||||
first_seed = not w.get("last_checked_at")
|
||||
fresh = [l for l in listings
|
||||
if PRODUCT_URL_TMPL.format(listing_id=l["listing_id"]) not in seen]
|
||||
|
||||
rows = [build_row(l, w) for l in fresh]
|
||||
if rows:
|
||||
insert_listings(listings_tid, rows)
|
||||
for l in fresh:
|
||||
seen.add(PRODUCT_URL_TMPL.format(listing_id=l["listing_id"]))
|
||||
|
||||
if fresh and not first_seed and w.get("notify", True):
|
||||
send_telegram(f"{w.get('title')}: {plural(len(fresh), 'new listing')}")
|
||||
|
||||
new_total += len(fresh)
|
||||
update_checked(settings_tid, wid)
|
||||
last_run[wid] = now
|
||||
|
||||
ok = len(failures) == 0
|
||||
return ok, ("; ".join(failures) if failures else ""), {
|
||||
"watch_count": len(watches), "new_this_tick": new_total}
|
||||
|
||||
|
||||
def main():
|
||||
if not NOCODB_TOKEN:
|
||||
sys.stderr.write("NOCODB_TOKEN not set\n")
|
||||
write_health(False, "NOCODB_TOKEN not set")
|
||||
sys.exit(2)
|
||||
|
||||
listings_tid, settings_tid = bootstrap()
|
||||
seen = load_seen(listings_tid)
|
||||
last_run = {}
|
||||
|
||||
sys.stderr.write(f"ready: listings={listings_tid} settings={settings_tid} "
|
||||
f"seen={len(seen)}\n")
|
||||
|
||||
while True:
|
||||
try:
|
||||
ok, err, extra = run_tick(listings_tid, settings_tid, seen, last_run)
|
||||
except Exception as e:
|
||||
ok, err, extra = False, f"tick error: {e}", {}
|
||||
write_health(ok, err, extra)
|
||||
time.sleep(TICK_SECONDS)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user