#!/usr/bin/env python3 """ auto_accept_worker.py — Worker tự động cho mỗi phòng. Chạy mỗi 1-2 phút (cronjob no_agent). Kết nối với Lark Base làm nguồn sự thật (Task Hub). Tự động Accept → Tự động chạy verify → Tự động Done & Cập nhật trạng thái + kết quả lên Lark Base + Telegram. """ import json import os import subprocess import sys from datetime import datetime, timezone from pathlib import Path LIBDIR = Path(__file__).resolve().parent sys.path.insert(0, str(LIBDIR)) import tg_topics DEFAULT_CHAT_ID = "-1003707758328" BASE_TOKEN = "QIx1bOHuzapATFslt4yl8zWTgnt" TABLE_ID = "tblN85Tinj2x1tnf" NLM_CLI = "/opt/ai-os/products/ceo/integrations/notebooklm-mcp-cli/.venv/bin/nlm" DEPARTMENTS = { "it-ai": {"name": "IT & AI", "chat_id": DEFAULT_CHAT_ID, "thread_id": 18}, "r-d": {"name": "R&D", "chat_id": DEFAULT_CHAT_ID, "thread_id": 12}, "strategy-marketing": {"name": "Strategy & Marketing", "chat_id": DEFAULT_CHAT_ID, "thread_id": 14}, "writers": {"name": "Writers", "chat_id": DEFAULT_CHAT_ID, "thread_id": None}, "grill-qa": {"name": "Grill & QA", "chat_id": DEFAULT_CHAT_ID, "thread_id": 20}, } # ========== HELPERS ========== def telegram_send(thread_id, text): if thread_id is None: print("Skip Telegram message: thread_id is None") return True, None token = tg_topics.load_token() payload = { "chat_id": DEFAULT_CHAT_ID, "message_thread_id": int(thread_id), "text": text, "disable_web_page_preview": True, } ok, res = tg_topics.api(token, "sendMessage", payload) return ok, res def run_nlm_describe(notebook_name): """Gọi nlm CLI thật, trả về kết quả hoặc chi tiết lỗi.""" try: env = os.environ.copy() env["PATH"] = "/usr/local/bin:/usr/bin:/bin:" + "/".join(NLM_CLI.split("/")[:-1]) res = subprocess.run( [NLM_CLI, "notebook", "describe", notebook_name], capture_output=True, text=True, timeout=45, env=env ) if res.returncode == 0: return res.stdout.strip() else: err_msg = res.stderr.strip() if res.stderr.strip() else res.stdout.strip() return f"⚠️ CLI error:\n{err_msg[:800]}" except subprocess.TimeoutExpired: return "❌ Timeout after 45s" except Exception as e: return f"❌ Automation error: {str(e)}" def get_lark_tasks_by_status(status): """Lấy danh sách các task từ Lark Base theo trạng thái. Trả về list các dict có record_id + fields.""" list_cmd = [ "lark-cli", "base", "+record-list", "--base-token", BASE_TOKEN, "--table-id", TABLE_ID, "--as", "user", "--format", "json" ] try: res = subprocess.run(list_cmd, capture_output=True, text=True, timeout=30) if res.returncode != 0: print(f"List error: {res.stderr[:200]}") return [] data = json.loads(res.stdout) fields_names = data.get("data", {}).get("fields", []) rows = data.get("data", {}).get("data", []) record_ids = data.get("data", {}).get("record_id_list", []) filtered = [] for idx, row in enumerate(rows): fields = {} for i, val in enumerate(row): if i < len(fields_names): name = fields_names[i] if isinstance(val, list) and len(val) == 1: fields[name] = val[0] else: fields[name] = val if fields.get("Handoff Status") == status: rec = { "record_id": record_ids[idx] if idx < len(record_ids) else None, "task_id": fields.get("Task ID"), "from_dept": fields.get("Phòng gửi"), "to_name": fields.get("Phòng nhận"), "task": fields.get("Task ID"), "input": fields.get("Note"), "output_desc": fields.get("Output mong muốn"), "status": fields.get("Status"), "handoff_status": fields.get("Handoff Status") } filtered.append(rec) return filtered except Exception as e: print(f"Error listing tasks: {e}") return [] def lark_update(task_id, record_id, fields_subset): """Cập nhật record trên Lark Base bằng --record-id.""" cmd = [ "lark-cli", "base", "+record-upsert", "--base-token", BASE_TOKEN, "--table-id", TABLE_ID, "--as", "user", "--record-id", record_id, "--json", json.dumps(fields_subset) ] try: res = subprocess.run(cmd, capture_output=True, text=True, timeout=30) if res.returncode != 0: print(f"Update error for {task_id}: {res.stderr[:200]}") return res.returncode == 0 except Exception as e: print(f"Update error: {e}") return False # ========== MAIN PROCESS ========== def process_task(task): task_id = task.get("task_id") record_id = task.get("record_id") from_dept = task.get("from_dept", "IT & AI") to_name = task.get("to_name", "") task_desc = task.get("task") inp = task.get("input", "") output_desc = task.get("output_desc") # Map back to key to_dept_key = None for k, v in DEPARTMENTS.items(): if v["name"] == to_name: to_dept_key = k break if not to_dept_key: return False target = DEPARTMENTS[to_dept_key] thread_id = target["thread_id"] now = datetime.now(timezone.utc).astimezone().strftime("%Y-%m-%d %H:%M:%S") # --- 1) AUTO-ACCEPT (Lớp 1 & Lớp 2) --- lark_update(task_id, record_id, {"Handoff Status": "Accepted", "Status": "In progress", "Note": f"Accepted automatically at {now}"}) accept_msg = ( f"📋 **XÁC NHẬN NHẬN BÀN GIAO**\n" f"━━━━━━━━━━━━━━━━━━━━━━━\n" f"🆔 **Task:** `{task_id}`\n" f"⏰ **Nhận lúc:** {now}\n" f"🏢 **Từ:** {from_dept} → **Đến:** {to_name}\n" f"📌 **Nội dung:** {task_desc}\n" f"━━━━━━━━━━━━━━━━━━━━━━━\n" f"⚡ Tự động xác nhận và đang xử lý..." ) telegram_send(thread_id, accept_msg) # --- 2) AUTO-RUN VERIFY (Lớp 3) --- notebook = inp if inp else "Báo Cáo Du Lịch" real_result = run_nlm_describe(notebook) # --- 3) AUTO-REPORT (Done & Reported) --- done_msg = ( f"✅ **BÁO CÁO HOÀN TẤT**\n" f"━━━━━━━━━━━━━━━━━━━━━━━\n" f"🆔 **Task:** `{task_id}`\n" f"📊 **Kết quả thực tế từ NotebookLM ({notebook}):**\n" f"{real_result}\n" f"━━━━━━━━━━━━━━━━━━━━━━━\n" f"📤 **Output mong muốn:** {output_desc}\n" f"━━━━━━━━━━━━━━━━━━━━━━━\n" f"✅ **Trạng thái:** Hoàn tất — đã báo cáo tự động." ) telegram_send(thread_id, done_msg) # Cập nhật kết quả lên Lark Base lark_update(task_id, record_id, { "Handoff Status": "Reported", "Status": "Done", "Output thực tế": real_result, "Note": f"Completed and reported automatically at {now}" }) return True def main(): # Quét Lark Base lấy các task ở trạng thái "Handed off" pending = get_lark_tasks_by_status("Handed off") if not pending: print("No pending handoffs found in Lark Base.") return processed_count = 0 for task in pending: success = process_task(task) if success: processed_count += 1 print(f"Successfully processed {processed_count} task(s).") if __name__ == "__main__": main()