From 6428c97214cfe5f7edf8e2ec91bdbc93489c080a Mon Sep 17 00:00:00 2001 From: hoelee Date: Wed, 9 Sep 2026 05:31:55 +0800 Subject: [PATCH] Fix Telegram notifications: multipart upload + richer message format MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 通知重构:归档与通知解耦(Listings.notified 列),tick 末尾统一发,1s 间隔,失败自动重试(30s tick) - 每商品一条图文消息:title/price/condition/seller 名/product_url + 高清图 - 图片用 multipart 上传本地字节(原先传 URL 让 Telegram 下载,遇 carousell CDN 不稳定导致间歇 400) - 图片 URL 去 _progressive_thumbnail 后缀(高清原图) - condition 归一化:New→Brand new、Used→Used,加第 6 档;从所有 paragraph 找 condition(修位置漂移) - listed_at 加 active_bump fallback(被顶置商品) - compose extra_hosts 钉 api.telegram.org 到 IPv4 149.154.166.110(容器无 IPv6 时 DNS 只返 AAAA) - bot 换 @carousellFoundBot,chat @MrFullStackDev --- DOCUMENTATION.md | 32 ++++--- docker-compose.yml | 2 + monitor.py | 212 ++++++++++++++++++++++++++++++++++++++++----- 3 files changed, 215 insertions(+), 31 deletions(-) diff --git a/DOCUMENTATION.md b/DOCUMENTATION.md index 36de921..5e4e980 100644 --- a/DOCUMENTATION.md +++ b/DOCUMENTATION.md @@ -15,7 +15,9 @@ Docker container (carousell-monitor, DSM, network bridge_hoelee) │ ├─ parse embedded JSON state → SearchListing.listingCards[] │ ├─ dedupe by product_url (param-less listing URL) │ ├─ INSERT new rows into NocoDB Listings table - │ └─ (after first seed) send Telegram ": N new listings" + │ └─ archive w/ notified flag (first seed = silent; else pending-notify) + ├─ end of each tick ─ send ALL pending (notified=false) listings, 1s apart + │ └─ each notice: photo + title/price/condition/seller/url ├─ writes /data/health.json each tick → Docker HEALTHCHECK └─ reaches: NocoDB http://nocodb:10380 (container DNS, bridge_hoelee) Carousell www.carousell.com.my (public internet) @@ -51,7 +53,7 @@ Credentials are documented in `SECRETS.md` there. | `NOCODB_URL` | `http://nocodb:10380` | container DNS on `bridge_hoelee`; LAN form `http://192.168.137.2:10380` | | `NOCODB_TOKEN` | *(secret)* | workspace-scoped NocoDB PAT | | `NOCODB_BASE_ID` | `poqw1zjw3hnsk37` | base "Carousell" | -| `TELEGRAM_BOT_TOKEN` | *(secret)* | @HoeleeAgentBot | +| `TELEGRAM_BOT_TOKEN` | *(secret)* | @carousellFoundBot | | `TELEGRAM_CHAT_ID` | `5648309582` | alert destination | | `TICK_SECONDS` | `60` | scheduler granularity | | `HEALTH_STALE_SECONDS` | `600` | healthcheck staleness window | @@ -75,15 +77,16 @@ deleting a table is safe; it is recreated on the next start). | `product_url` | URL | **unique dedupe key** — `https://www.carousell.com.my/p/<id>/`, no query params | | `title` | SingleLineText | listing title | | `price` | Decimal | numeric price, "RM" stripped (e.g. `85.00`) — filterable/sortable | -| `condition` | SingleSelect | Brand new / Like new / Lightly used / Well used / Heavily used | -| `image_url` | URL | raw Carousell thumbnail URL (media.karousell.com) | -| `image` | Attachment | thumbnail — NocoDB hotlinks the URL; renders in grid view | +| `condition` | SingleSelect | Brand new / Like new / Lightly used / Well used / Heavily used / Used (归一化自 Carousell 的 New/Used 等写法) | +| `image_url` | URL | 高清图 URL(已去 `_progressive_thumbnail` 后缀) | +| `image` | Attachment | 高清图 — NocoDB hotlinks the URL; renders in grid view | | `seller_name` | SingleLineText | seller username | | `seller_url` | URL | `https://www.carousell.com.my/u/<username>/` | | `search_title` | SingleLineText | which watch found it (denormalized) | | `search_url` | URL | the watch's search URL | -| `listed_at` | DateTime (UTC) | listing's `time_created` on Carousell | +| `listed_at` | DateTime (UTC) | 上架时间(优先 `time_created`,被顶置商品 fallback `active_bump`) | | `first_seen_at` | DateTime (UTC) | when the monitor first captured it | +| `notified` | Checkbox | false = 待发通知;发完/静默归档后置 true(防重复通知) | ### `Settings` (the watch list — you manage this) @@ -106,10 +109,11 @@ deleting a table is safe; it is recreated on the next start). - **Dedupe**: `product_url` is the identity. On startup the container loads every existing `product_url` from NocoDB into an in-memory set; a listing is "new" only if its URL is not in that set. -- **First run per watch** (`last_checked_at` is null): seeds the current ~49 listings - as an archive with **no Telegram alert**. This prevents a 98-message flood on setup. -- **Afterwards**: only genuinely-new listings are inserted **and** alerted, as - `"<title>: N new listings"` (one message per watch that had new items, no details). +- **First run per watch** (`last_checked_at` is null): seeds the current listings as an + archive with `notified=true` (silent, no Telegram). Prevents a flood on setup. +- **Afterwards**: new listings are inserted with `notified=false` (pending). At the end + of each tick the monitor sends every `notified=false` listing **one message each** + (photo + title/price/condition/seller/url), 1 second apart, then sets `notified=true`. - **Failure handling**: a fetch/parse error on any watch marks that tick failed; the container becomes **unhealthy** until the next fully-successful tick. `last_checked_at` is only advanced on success, and a hard-failing watch is rate-limited to one attempt @@ -188,7 +192,7 @@ Health file lives at `/data/health.json` inside the container: | Container | `carousell-monitor` (network `bridge_hoelee`) | | NocoDB base | `Carousell` = `poqw1zjw3hnsk37` (workspace `wal4hatt`) | | Tables | `Listings` + `Settings` (bootstrap finds by title) | -| Telegram | `@HoeleeAgentBot` → chat `5648309582` | +| Telegram | `@carousellFoundBot` → chat `5648309582` (@MrFullStackDev) | --- @@ -198,3 +202,9 @@ Health file lives at `/data/health.json` inside the container: never committed. Credential inventory: `SECRETS.md` in the repo. - The image is private (built on DSM, never pushed to Docker Hub). - NocoDB token is workspace-scoped; regenerate in NocoDB if it leaks and update `.env`. + +--- + +## 11. Changelog + +- **2026-09-08** 通知重构:每商品一条图文消息(title/price/condition/seller/url),归档与通知解耦(`notified` 列 + tick 末尾统一发 + 1s 间隔)。图片改用高清 URL(去 `_progressive_thumbnail`)。condition 归一化(New→Brand new、Used→Used,加第 6 档)。listed_at 加 `active_bump` fallback。修复 Telegram IPv6/DNS 问题(compose `extra_hosts` 钉 IPv4)。bot 换 `@carousellFoundBot`。 diff --git a/docker-compose.yml b/docker-compose.yml index 13217a8..19bc988 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -9,6 +9,8 @@ services: image: carousell-monitor:latest container_name: carousell-monitor restart: unless-stopped + extra_hosts: + - "api.telegram.org:149.154.166.110" environment: NOCODB_URL: ${NOCODB_URL:-http://nocodb:10380} NOCODB_TOKEN: ${NOCODB_TOKEN} diff --git a/monitor.py b/monitor.py index f1d2032..c20ee5d 100644 --- a/monitor.py +++ b/monitor.py @@ -42,10 +42,27 @@ 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"] +CONDITIONS = ["Brand new", "Like new", "Lightly used", "Well used", "Heavily used", "Used"] +# Carousell 不同类目对 condition 的写法不一:New ≈ Brand new;Used 是笼统二手。 +CONDITION_MAP = { + "brand new": "Brand new", + "new": "Brand new", + "like new": "Like new", + "lightly used": "Lightly used", + "well used": "Well used", + "heavily used": "Heavily used", + "used": "Used", +} PRODUCT_URL_TMPL = "https://www.carousell.com.my/p/{listing_id}/" SELLER_URL_TMPL = "https://www.carousell.com.my/u/{username}/" + +def strip_thumbnail_suffix(url): + """去掉 Carousell 缩略图后缀 _progressive_thumbnail,得到高清原图 URL。""" + if not url: + return url + return re.sub(r"_progressive_thumbnail(?=\.\w+$)", "", url) + # Column definitions: table title -> list of (title, uidt) LISTINGS_COLS = [ ("product_url", "URL"), @@ -60,6 +77,7 @@ LISTINGS_COLS = [ ("search_url", "URL"), ("listed_at", "DateTime"), ("first_seen_at", "DateTime"), + ("notified", "Checkbox"), ] SETTINGS_COLS = [ ("title", "SingleLineText"), @@ -88,6 +106,57 @@ def _http(method, url, body=None, headers=None, timeout=30): return r.status, raw +def _download_image(url, timeout=20): + """下载图片到内存 bytes;失败返回 None。""" + try: + req = urllib.request.Request(url, headers={"User-Agent": UA}) + with urllib.request.urlopen(req, timeout=timeout) as r: + return r.read() + except Exception: + return None + + +def _send_photo_multipart(chat_id, image_bytes, filename, caption): + """用 multipart/form-data 上传本地图片字节发 sendPhoto。 + + 逐字段拼装字节(不用 join,避免破坏二进制图片数据)。 + """ + boundary = "----carousellmonitor" + str(int(time.time() * 1000)) + "boundary" + CRLF = b"\r\n" + parts = [] + + def field(name, value): + p = ("--" + boundary).encode("utf-8") + CRLF + p += ("Content-Disposition: form-data; name=\"" + name + "\"").encode("utf-8") + CRLF + p += CRLF + p += value.encode("utf-8") + CRLF + return p + + def file_field(name, filename, data, content_type): + p = ("--" + boundary).encode("utf-8") + CRLF + p += ("Content-Disposition: form-data; name=\"" + name + + "\"; filename=\"" + filename + "\"").encode("utf-8") + CRLF + p += ("Content-Type: " + content_type).encode("utf-8") + CRLF + p += CRLF + p += data + CRLF + return p + + body = b"" + body += field("chat_id", str(chat_id)) + body += field("caption", caption) + body += file_field("photo", filename, image_bytes, "image/jpeg") + body += ("--" + boundary + "--").encode("utf-8") + CRLF + + url = "https://api.telegram.org/bot" + TELEGRAM_BOT_TOKEN + "/sendPhoto" + req = urllib.request.Request(url, data=body, method="POST", headers={ + "User-Agent": UA, + "Content-Type": "multipart/form-data; boundary=" + boundary, + }) + with urllib.request.urlopen(req, timeout=30) 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) @@ -181,22 +250,33 @@ def fetch_listings(search_url): for c in cards: try: lid = int(c["listingID"]) + # 上架时间优先取 time_created;被顶置(bump)的商品只提供 + # active_bump 时间戳,fallback 到它。两者结构相同(timestampContent)。 ts = None - for item in c.get("aboveFold", []): - if item.get("component") == "time_created": - ts = item["timestampContent"]["seconds"]["low"] + for comp in ("time_created", "active_bump"): + for item in c.get("aboveFold", []): + if item.get("component") == comp: + tc = item.get("timestampContent") or {} + sec = tc.get("seconds") or {} + ts = sec.get("low") + break + if ts: 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 + paras = [i.get("stringContent", "").strip() 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", "") + # 从所有 paragraph 里找第一个匹配的 condition 值(忽略 Size: 等噪声) + cond = "" + for p in paras: + key = p.lower() + if key in CONDITION_MAP: + cond = CONDITION_MAP[key] + break + thumb = strip_thumbnail_suffix(c.get("thumbnailURL", "")) seller = (c.get("seller") or {}).get("username", "") out.append({ "listing_id": lid, @@ -291,7 +371,7 @@ def update_checked(settings_tid, watch_id): [{"Id": watch_id, "last_checked_at": iso_now()}]) -def build_row(l, watch): +def build_row(l, watch, notified=False): row = { "product_url": PRODUCT_URL_TMPL.format(listing_id=l["listing_id"]), "title": l["title"], @@ -304,6 +384,7 @@ def build_row(l, watch): "search_title": watch.get("title", ""), "search_url": watch.get("url", ""), "first_seen_at": iso_now(), + "notified": notified, } price = parse_price(l["price"]) if price is not None: @@ -319,17 +400,107 @@ def build_row(l, watch): # --------------------------------------------------------------------------- # # Telegram # --------------------------------------------------------------------------- # -def send_telegram(text): +def _tg_send(method, payload): if not TELEGRAM_BOT_TOKEN or not TELEGRAM_CHAT_ID: - return + return False try: - tg("sendMessage", {"chat_id": TELEGRAM_CHAT_ID, "text": text}) + tg(method, {"chat_id": TELEGRAM_CHAT_ID, **payload}) + return True except urllib.error.HTTPError as e: - sys.stderr.write(f"telegram send failed: {e.code} {e.read()[:200]}\n") + sys.stderr.write(f"telegram {method} failed: {e.code} {e.read()[:200]}\n") + return False -def plural(n, word): - return f"{n} {word}{'' if n == 1 else 's'}" +def send_telegram(text): + _tg_send("sendMessage", {"text": text}) + + +def send_listing_from_record(rec): + """Send one listing notice from a NocoDB record. + + rec fields: title, price, condition, seller_name, product_url, image_url. + Prefer photo; fall back to text-only if the image fails. + """ + title = rec.get("title") or "(no title)" + price = rec.get("price") or "" + condition = rec.get("condition") or "" + seller = rec.get("seller_name") or "" + product_url = rec.get("product_url") or "" + caption_lines = [f"🛒 {title}"] + if price: + caption_lines.append(f"💰 {price}") + if condition: + caption_lines.append(f"📦 {condition}") + if seller: + caption_lines.append(f"👤 {seller}") + if product_url: + caption_lines.append(product_url) + caption = "\n".join(caption_lines) + thumb = rec.get("image_url") or "" + if thumb: + # 方案1:下载图后 multipart 上传(最稳,Telegram 无需访问 carousell CDN) + img_bytes = _download_image(thumb) + if img_bytes: + try: + st, raw = _send_photo_multipart( + TELEGRAM_CHAT_ID, img_bytes, "listing.jpg", caption) + if st == 200: + return True + sys.stderr.write(f"multipart sendPhoto: HTTP {st}\n") + except Exception as e: + sys.stderr.write(f"multipart sendPhoto failed: {type(e).__name__}\n") + # 方案2:退回让 Telegram 直接下载 URL + if _tg_send("sendPhoto", {"photo": thumb, "caption": caption}): + return True + # 方案3:纯文本 + return _tg_send("sendMessage", {"text": caption}) + + +def mark_notified(listings_tid, rec_ids): + """把已发通知的记录 notified 置 true。""" + if not rec_ids: + return + updates = [{"Id": rid, "notified": True} for rid in rec_ids] + nc("PATCH", f"/api/v2/tables/{listings_tid}/records", updates) + + +def send_pending_notifications(listings_tid, settings_tid): + """tick 末尾统一发:查 notified=false 的记录,逐条发(间隔 1s),发完置 true。 + + 仅发「其 watch 仍 notify=true」的记录;watch 已关 notify 的则静默置 true。 + """ + # 加载所有 watch 的 notify 开关,key = search_title + st, j = nc("GET", f"/api/v2/tables/{settings_tid}/records?limit=1000") + if st != 200: + return + notify_by_title = {} + for w in j.get("list", []): + notify_by_title[w.get("title")] = bool(w.get("notify")) + + # 拉 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,notified") + if st != 200: + return + pending = [r for r in j.get("list", []) if not r.get("notified")] + + if not pending: + return + + for rec in pending: + st_title = rec.get("search_title") or "" + should_notify = notify_by_title.get(st_title, True) + if should_notify: + ok = send_listing_from_record(rec) + else: + # 词条关了 notify:静默标记,不发 + ok = True + # 只有发送成功(或无需发)才标记 notified=true; + # 失败则保留 false,下个 tick 末尾自动重试。 + if ok: + mark_notified(listings_tid, [rec["Id"]]) + time.sleep(1) # 最快 1 秒一条,防限流 # --------------------------------------------------------------------------- # @@ -374,19 +545,20 @@ def run_tick(listings_tid, settings_tid, seen, last_run): 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] + # 首次 seed: 静默归档 (notified=True);后续新商品: 待发 (notified=False) + rows = [build_row(l, w, notified=first_seed) 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 + # 归档完成后,统一发送待通知的记录(解耦:归档成功才通知) + send_pending_notifications(listings_tid, settings_tid) + ok = len(failures) == 0 return ok, ("; ".join(failures) if failures else ""), { "watch_count": len(watches), "new_this_tick": new_total}