「一個真正實用的 AI Agent,不該被動等待使用者手動貼入資料,而要能主動對接雲端儲存,實現知識庫的自動化增量同步。」
在先前幾天,我們完成了 AI Agent 的大腦推理、記憶與背景學習機制。從今天開始我們正式進入 資料管線自動化與周邊服務整合 的章節!
今天我們將深入探討 drive_sync.py 如何透過 Google Drive API (Service Account) 監控雲端資料夾,自動辨識新增或修改的 Google Docs、.txt 與 .md 檔案,並透過 FastAPI 的 /chat 介面搭配 /learn 特殊指令 進行動態文字切分與 Incremental Ingestion。
在 Self-hosted RHEL 伺服器上,為了讓背景腳本能自動排程(Cron Job)同步雲端文件,我們採用了 Google Service Account(服務帳號)機制:
1. Headless Auth:無須開啟瀏覽器進行 OAuth2 使用者授權,透過 /opt/angelina/config/service-account.json 憑證檔直接取得僅讀(drive.readonly)存取 Token。
2. 輕量化 API 對接:使用 google.oauth2.service_account 搭配 httpx.AsyncClient 非同步呼叫 Google Drive v3 REST API。
3. 狀態鎖定與去重:透過記錄檔案的 modifiedTime(修改時間),確保只有新增或更動過的文件才會被下載與注入,避免重複消耗 Token。
Python
async def download_file_content(access_token: str, file_id: str, mime_type: str) -> str | None:
"""從 Google Drive 下載或導出檔案內容"""
headers = {"Authorization": f"Bearer {access_token}"}
async with httpx.AsyncClient() as client:
if mime_type == "application/vnd.google-apps.document":
# Google Docs 導出為純文字
url = f"{DRIVE_API_BASE}/files/{file_id}/export"
params = {"mimeType": "text/plain"}
response = await client.get(url, params=params, headers=headers)
response.raise_for_status()
return response.text
elif mime_type in ("text/plain", "text/markdown"):
# .txt/.md 檔案直接下載
url = f"{DRIVE_API_BASE}/files/{file_id}"
params = {"alt": "media"}
response = await client.get(url, params=params, headers=headers)
response.raise_for_status()
return response.text
else:
logger.info(f"Skipping unsupported mime type: {mime_type}")
return None
Python
MAX_LEARN_CHUNK = 1500
async def sync_to_knowledge(text: str, filename: str) -> None:
"""將文本切分並透過 /learn 指令寫入 Angelina 知識庫"""
chunks = []
# 智慧邊界切分算法 (優先搜尋 \n 或空格)
while len(text) > MAX_LEARN_CHUNK:
split_pos = text.rfind("\n", 0, MAX_LEARN_CHUNK)
if split_pos == -1:
split_pos = text.rfind(" ", 0, MAX_LEARN_CHUNK)
if split_pos == -1:
split_pos = MAX_LEARN_CHUNK
chunks.append(text[:split_pos])
text = text[split_pos:].lstrip()
if text:
chunks.append(text)
async with httpx.AsyncClient(timeout=60.0) as client:
for i, chunk in enumerate(chunks, 1):
part_label = f" (part {i}/{len(chunks)})" if len(chunks) > 1 else ""
# 格式化為 Angelina 內建指令格式
message = f"/learn [Drive: {filename}{part_label}] {chunk}"
response = await client.post(
f"{ANGELINA_API_URL}/chat",
json={"message": message, "session_id": "drive-sync", "language": "zh-TW"},
)
response.raise_for_status()
logger.info(f"Ingested chunk {i}/{len(chunks)} for '{filename}'")
Python
async def run_sync() -> dict:
summary = {"synced": 0, "skipped": 0, "errors": 0}
access_token = get_drive_credentials()
files = await list_files(access_token)
sync_state = load_sync_state()
for file_info in files:
file_id, filename = file_info["id"], file_info["name"]
modified_time = file_info["modifiedTime"]
# 比對修改時間,未變更則跳過
if file_id in sync_state and sync_state[file_id] == modified_time:
summary["skipped"] += 1
continue
content = await download_file_content(access_token, file_id, file_info["mimeType"])
if content:
await sync_to_knowledge(content, filename)
sync_state[file_id] = modified_time
summary["synced"] += 1
save_sync_state(sync_state)
return summary
透過 drive_sync.py 的實作,我們完成了資料輸入端(Data Ingestion Pipeline)的自動化:
明天(Day 09)我們將進入 daily_analysis.py,探討如何抓取 TWSE 證交所 API 與 Yahoo Finance 美股數據,結合 Gemini 生成每日市場分析報告並推送至 Telegram 機器人!
明日預告:【Day 09】每日自動化盤後分析:TWSE/美股數據爬取、Gemini 分析與 Telegram 機器人推播 (daily_analysis.py)