419 lines
14 KiB
Python
419 lines
14 KiB
Python
#!/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()
|