Carousell new-listing monitor: NocoDB archive + Telegram alerts
This commit is contained in:
+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