Chen Yuan

我的 AI Agent 只有一個工作:寫文章、發佈到平台,然後清理。十次裡面有三次,它會在管線中途崩潰、留下孤立的草稿、丟失暫存檔,然後又重新嘗試發布同一篇文章。結果就是重複草稿、半失敗的貼文——最糟的情況是同一個 feed 裡出現兩篇一模一樣的文章。

我需要讓管線具備「崩潰安全」的特性,不是「最終一致性」那種,而是「原子性」崩潰安全。以下就是實現這個目標的狀態機模式。

解決方案

我把管線拆成多個階段,每個階段只負責單一職責。run ID 把所有東西串在一起,manifest 檔案追蹤哪些暫存檔屬於哪次執行。狀態機只有在當前階段的閘門全部綠燈時才會前進。

from __future__ import annotations
import json
from pathlib import Path
from enum import Enum
from datetime import datetime
from typing import Optional

class Phase(Enum):
    LOAD = "loading"
    SELECT = "topic_selection"
    DEDUP = "topic_dedup"
    WRITE = "article_generation"
    GATE = "de_ai_gate"
    VALIDATE = "pre_publish_validate"
    PUBLISH = "publish"
    VERIFY = "post_publish_verify"
    CLEAN = "cleanup"

class RunState:
    def __init__(self, run_id: str, manifest_path: Path):
        self.run_id = run_id
        self.manifest_path = manifest_path
        self.current_phase = None
        self.phase_results = {}

    def save(self):
        data = {
            "run_id": self.run_id,
            "current_phase": self.current_phase.value if self.current_phase else None,
            "phase_results": self.phase_results,
            "updated_at": datetime.now().isoformat(),
        }
        self.manifest_path.write_text(json.dumps(data, ensure_ascii=False, indent=2))

    @classmethod
    def load(cls, manifest_path: Path):
        if not manifest_path.exists():
            return None
        data = json.loads(manifest_path.read_text(encoding="utf-8"))
        state = cls(data["run_id"], manifest_path)
        state.current_phase = Phase(data["current_phase"]) if data["current_phase"] else None
        state.phase_results = data["phase_results"]
        return state

Enter fullscreen mode Exit fullscreen mode

狀態機本身是一個簡單的迴圈。每個階段都是回傳 True(前進)或拋出例外(停止)的函式。在前進之前,狀態會先寫到磁碟。

class Pipeline:
    def __init__(self, state: RunState, temp_files: list[Path]):
        self.state = state
        self.temp_files = temp_files
        self.guard = Guard(temp_files)

    def run(self, phases: list[Phase]):
        for phase in phases:
            print(f"[{self.state.run_id}] Phase: {phase.value}")
            self.state.current_phase = phase
            self.state.save()

            handler = getattr(self, f"_phase_{phase.value}")
            handler()

            self.state.phase_results[phase.value] = True
            self.state.save()

    def _phase_cleanup(self):
        for f in self.temp_files:
            if f.exists():
                f.unlink()
                print(f"  cleaned: {f.name}")
        self.state.manifest_path.unlink(missing_ok=True)

Enter fullscreen mode Exit fullscreen mode

Guard 類別介於管線與檔案系統之間。它清楚知道哪些檔案是暫存的發布產物,並在技能一致性檢查中忽略它們。

class Guard:
    def __init__(self, allowed_temps: list[Path]):
        self.allowed = {p.resolve() for p in allowed_temps}

    def is_temp_artifact(self, path: Path) -> bool:
        return path.resolve() in self.allowed

    def verify_no_untracked_writes(self, workspace: Path):
        for f in workspace.iterdir():
            if f.is_file() and not self.is_temp_artifact(f):
                raise RuntimeError(
                    f"Untracked file: {f.name}. "
                    "Add to temp manifest or refactor."
                )

Enter fullscreen mode Exit fullscreen mode

為什麼有效

三個設計決策讓這個模式具備崩潰安全。

第一,磁碟優先狀態。狀態機在每個階段處理函式的前後都會寫入磁碟。如果 Agent 在 _phase_validate() 期間崩潰,下一次執行會讀取 manifest.json,看到 current_phase: "pre_publish_validate",並知道要從哪裡恢復。不再有「我剛剛做到哪裡?」的困惑。

第二,暫存檔清單。管線啟動前,會先註冊它將建立的所有檔案。Guard 類別會強制執行:如果寫入的檔案不在清單中,管線就會停止。這讓我抓到兩次問題——一次是我在 Python 中用了 write_file 而非 open(),另一次是日誌函式庫在工作目錄產生 .log 檔。

第三,冪等階段處理函式。每個階段函式在再次執行前,會先檢查工作是否已完成。_phase_dedup() 會讀取現有文章清單,只有在尚未處理時才寫入新資料。這表示崩潰後重播 Phase 2 是安全的——不會重複檢查或重複寫入。

def _phase_dedup(self):
    manifest = self.state.manifest_path
    dedup_flag = manifest.parent / f"{self.state.run_id}_dedup_done"
    if dedup_flag.exists():
        print("  dedup already done, skipping")
        return
    # ... actual dedup logic ...
    dedup_flag.write_text("done")

Enter fullscreen mode Exit fullscreen mode

注意事項

Windows 檔案鎖定不同。在 Linux 上,fcntl.flock() 運作良好;在 Windows 上則需要 msvcrt.locking()。當兩次管線執行意外重疊並嘗試寫入同一個 manifest 檔時,我就遇到了這個問題。解決方法是跨平台鎖定:

import sys

def safe_write(path: Path, data: str):
    if sys.platform == "win32":
        import msvcrt
        with open(path, "w", encoding="utf-8") as f:
            msvcrt.locking(f.fileno(), msvcrt.LK_LOCK, 1)
            f.write(data)
            msvcrt.locking(f.fileno(), msvcrt.LK_UNLCK, 1)
    else:
        import fcntl
        with open(path, "w", encoding="utf-8") as f:
            fcntl.flock(f.fileno(), fcntl.LOCK_EX)
            f.write(data)
            fcntl.flock(f.fileno(), fcntl.LOCK_UN)

Enter fullscreen mode Exit fullscreen mode

manifest 檔本身也是暫存檔。如果你把它納入暫存檔清單,Guard 就會放行。但 _phase_cleanup 處理函式會最後才刪除 manifest——在所有其他暫存檔都清除之後——因此清理期間發生崩潰也不會遺失狀態記錄。

狀態機的好壞取決於階段邊界。我一開始把階段切得太細(10 步驟的管線用了 14 個階段)。狀態檔不斷被覆寫,I/O 開銷讓每次執行多花約 2 秒。後來我把相關步驟合併:de_ai_gate + validate 變成一個階段,publish + verify 也合併成另一個。最後降到 7 個階段,沒有明顯的開銷。

那你呢?

我看過有人用 SQLite 儲存管線狀態、用 Redis 處理分散式管線,甚至把 JSON 檔放在 S3 上。我的方法刻意保持最小——單一 JSON 檔、一個 guard 類別,以及一個迴圈。它之所以有效,是因為管線是單執行緒且在單一機器上執行。

你如何處理崩潰安全的自動化?是使用工作流引擎(Temporal、Prefect),還是自己實作狀態機?很想聽聽你的經驗。