#!/opt/ai-os/products/ceo/integrations/notebooklm-mcp-cli/.venv/bin/python """ auto_accept_worker.py — Worker tự động cho 5 phòng ban. Status flow: Ready for Handoff → Accepted → Handed off / Blocked Config from centralized schema: config/lark_schema.json """ import json, os, re, subprocess, sys, time from datetime import datetime, timezone from pathlib import Path LIBDIR = Path(__file__).resolve().parent sys.path.insert(0, str(LIBDIR)) sys.path.insert(0, "/opt/ai-os/core/lib") import tg_topics # --- Config from centralized schema --- SCHEMA_PATH = Path("/opt/ai-os/products/ceo/config/lark_schema.json") if SCHEMA_PATH.exists(): with open(SCHEMA_PATH) as f: SCHEMA = json.load(f) BASE_TOKEN = SCHEMA.get("base_token", "QIx1bOHuzapATFslt4yl8zWTgnt") TABLE_ID = SCHEMA.get("table_id", "tblN85Tinj2x1tnf") VALID_STATUSES = SCHEMA.get("status_options", []) VALID_DEPARTMENTS = SCHEMA.get("departments", []) else: BASE_TOKEN = "QIx1bOHuzapATFslt4yl8zWTgnt" TABLE_ID = "tblN85Tinj2x1tnf" VALID_STATUSES = [] VALID_DEPARTMENTS = [] DEFAULT_CHAT_ID = "-1003707758328" DEPARTMENTS = { "policy-lab": {"name": "Policy Lab", "chat_id": DEFAULT_CHAT_ID, "thread_id": 12}, "it-ai": {"name": "IT & AI", "chat_id": DEFAULT_CHAT_ID, "thread_id": 18}, "research-inno": {"name": "Research & Inno.", "chat_id": DEFAULT_CHAT_ID, "thread_id": 16}, "marketing": {"name": "Marketing", "chat_id": DEFAULT_CHAT_ID, "thread_id": 14}, "grill-brainstorm": {"name": "Grill & Brainstorm", "chat_id": DEFAULT_CHAT_ID, "thread_id": 20}, } AUTO_PROCESS_STATUSES = ["Ready for Handoff"] NLM_SRC = "/opt/ai-os/products/ceo/integrations/notebooklm-mcp-cli/src" NLM_BIN = "/opt/ai-os/products/ceo/integrations/notebooklm-mcp-cli/.venv/bin/nlm" # ========== TELEGRAM ========== def tg_send(thread_id, text): t = tg_topics.load_token() ok, _ = tg_topics.api(t, "sendMessage", { "chat_id": DEFAULT_CHAT_ID, "message_thread_id": int(thread_id), "text": text, "disable_web_page_preview": True, }) return ok # ========== NOTEBOOKLM ========== def _reauthenticate(): subprocess.run([NLM_BIN, "login", "--provider", "openclaw", "--cdp-url", "http://127.0.0.1:9222"], capture_output=True, text=True, timeout=30) meta_f = Path.home() / ".notebooklm-mcp-cli" / "profiles" / "default" / "metadata.json" cook_f = Path.home() / ".notebooklm-mcp-cli" / "profiles" / "default" / "cookies.json" dst = Path.home() / ".notebooklm-mcp-cli" / "auth.json" if meta_f.exists() and cook_f.exists(): meta, cookies = json.loads(meta_f.read_text()), json.loads(cook_f.read_text()) dst.write_text(json.dumps({"csrf_token": meta.get("csrf_token",""), "session_id": meta.get("session_id",""), "build_label": meta.get("build_label",""), "cookies": cookies, "extracted_at": time.time()})) print("Auth synced to auth.json") def _call_nlm(fn_name, *args): sys.path.insert(0, NLM_SRC) mod = __import__("notebooklm_tools.mcp.tools.notebooks", fromlist=[fn_name]) fn = getattr(mod, fn_name) result = fn(*args) if result.get("status") == "success": return result err = str(result) if "auth" in err.lower() or "expired" in err.lower(): print("Auth expired, re-authenticating...") _reauthenticate() result = fn(*args) if result.get("status") == "success": return result return None def notebook_list(): result = _call_nlm("notebook_list", 100) return result.get("notebooks", []) if result else [] def notebook_describe(uuid): result = _call_nlm("notebook_describe", uuid) if not result: return None summary = result.get("summary", "") return "\n\n".join(summary) if isinstance(summary, list) else summary def resolve_nb(hint): if re.match(r'^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$', hint): return hint for nb in notebook_list(): if hint.lower() in nb.get("title","").lower(): return nb["id"] return "a0f8646b-90b1-4a65-bddc-07280aa06a99" # ========== LARK ========== def lark_run(args): cmd = ["lark-cli"] + args + ["--as", "user"] r = subprocess.run(cmd, capture_output=True, text=True, timeout=30) if r.returncode != 0: print(f"Lark error: {r.stderr[:200]}") return None return json.loads(r.stdout) def get_tasks(status): """Lấy task từ Lark Base theo status, trả về list dict.""" data = lark_run(["base", "+record-list", "--base-token", BASE_TOKEN, "--table-id", TABLE_ID, "--format", "json"]) if not data: return [] fields_names = data.get("data", {}).get("fields", []) rows = data.get("data", {}).get("data", []) rids = data.get("data", {}).get("record_id_list", []) result = [] for idx, row in enumerate(rows): d = {} for i, v in enumerate(row): if i < len(fields_names): n = fields_names[i] d[n] = v[0] if isinstance(v, list) and len(v)==1 else v if isinstance(d.get("Status"), list) and len(d["Status"])==1: d["Status"] = d["Status"][0] if d.get("Status") == status: result.append({ "record_id": rids[idx] if idx < len(rids) else None, "task_id": d.get("Task ID"), "sender": d.get("Phòng gửi"), "receiver": d.get("Phòng nhận"), "desc": d.get("Task"), "input": d.get("Input",""), "note": d.get("Note",""), "output_desc": d.get("Output mong muốn"), }) return result def lark_update(record_id, fields): return lark_run(["base", "+record-upsert", "--base-token", BASE_TOKEN, "--table-id", TABLE_ID, "--record-id", record_id, "--json", json.dumps(fields)]) # ========== MAIN ========== def process(t): task_id = t.get("task_id") rec_id = t.get("record_id") sender = t.get("sender","IT & AI") receiver = t.get("receiver","") desc = t.get("desc","") inp = t.get("input","") output_desc = t.get("output_desc","") now = datetime.now(timezone.utc).astimezone().strftime("%Y-%m-%d %H:%M:%S") # map dept → thread to_key = None for k, v in DEPARTMENTS.items(): if v["name"] == receiver: to_key = k break if not to_key: print(f"Cannot map: {receiver}") return False thr = DEPARTMENTS[to_key]["thread_id"] # 1) Auto-accept lark_update(rec_id, {"Status": "Ready for Handoff", "Note": f"Handoff sent from {sender} at {now}"}) tg_send(thr, f"📩 **BÀN GIAO MỚI**\n🆔 **{task_id}**\n📌 {desc or '(xem Note trên Lark)'}\n🏢 {sender} → {receiver}\n━━━\n📥 Các bạn dùng:\n💬 `/accept {task_id}` — nhận việc\n💬 `/block {task_id} [lý do]` — nếu kẹt\n💬 `/done {task_id}` — khi xong") return True def main(): print(f"Worker running with: {sys.executable}") pending = [] for s in AUTO_PROCESS_STATUSES: pending += get_tasks(s) if not pending: print("No pending tasks.") return count = sum(1 for t in pending if process(t)) print(f"Processed {count}/{len(pending)} task(s).") if __name__ == "__main__": main()