文档导航

模式与反模式

模式与反模式

钩子是一条小小的扩展缝,用好用坏都很容易。这一页收集在生产里站得住的形状,以及会咬人的那些。

模式

审计日志

记录每一个决策,并记下足以重建它的东西:状态、问题、答案、模型、路由决策、usage 和延迟。

import json

def audit(ctx):
    for state, result in zip(ctx.states, ctx.results or []):
        json.dump({
            "run_id": ctx.run_id,
            "model": ctx.model,
            "state": state,
            "routing": result.get("routing"),
            "answers": result["answers"],
            "usage": result.get("usage"),
            "call_usage": ctx.usage,
            "call_elapsed_ms": round(ctx.elapsed_ms or 0.0, 3),
        }, sys.stdout)
        sys.stdout.write("\n")

laya.load("convaiinnovations/laya", on_predict_end=audit)

一个钩子每次调用触发一次,而一次 predict_batch 调用把每个状态都装在里面,所以记录是按决策写 的:ctx.states 和 ctx.results 按下标对齐。ctx.usage 和 ctx.elapsed_ms 是整次调用的合计; 每个结果带它自己的 usage。

如果丢一行日志绝不能导致请求失败,就让它宽松:hooks_raise=False。如果审计轨迹是一项合规要求, 就让它严格。

PII 脱敏

脱敏必须发生在 on_predict_start 里、tokenization 之前,否则模型已经见过那些数据了。

import re
EMAIL = re.compile(r"\b[\w.+-]+@[\w-]+\.[\w.-]+\b")

def redact(ctx):
    ctx.states = [
        EMAIL.sub("[email]", s) if isinstance(s, str) else s
        for s in ctx.states
    ]

laya.load("convaiinnovations/laya", on_predict_start=redact)

脱敏钩子是一条策略钩子:保持 hooks_raise=True,因为一个静默坏掉的脱敏器就是一次数据泄漏。

缓存

一个 start 钩子查缓存并调用 ctx.skip(...);一个 end 钩子填充它。命中时前向传播被跳过。

import hashlib, json

CACHE = {}

def key(ctx, index):
    # Not sort_keys=True: criteria order is positional, so two orders are two questions,
    # and the checkpoint and token budget change the answer too.
    payload = json.dumps([ctx.states[index], ctx.questions, ctx.model,
                          ctx.max_len, ctx.head_max_len], default=str)
    return hashlib.sha256(payload.encode()).hexdigest()

def read(ctx):
    hits = [CACHE.get(key(ctx, i)) for i in range(len(ctx.states))]
    if all(hit is not None for hit in hits):
        ctx.skip(hits)   # one per state: skip replaces the whole call

def write(ctx):
    for i, result in enumerate(ctx.results or []):
        CACHE[key(ctx, i)] = result

laya.load("convaiinnovations/laya", on_predict_start=read, on_predict_end=write)

钩子每次调用触发一次,所以在 predict_batch 上只按一个状态做 key 是不够的:ctx.skip() 会替换 这次调用本会返回的每一个结果。并发服务时用一把锁守住缓存。在 Router 上,被缓存的载荷仍然会拿到 一个 routing 键,所以返回形状不变。

同一对钩子也能通过单个 LangChain 节点的 hooks= 参数用在它上面,这是在不改变 该 agent 其他每个调用方所见内容的前提下,缓存图里某一步热点的办法。

指标

从 ctx.model、ctx.usage 和 ctx.elapsed_ms 得到计数器和直方图。让它宽松。

COUNTS, LATENCIES = {}, []

def metrics(ctx):
    COUNTS[ctx.model] = COUNTS.get(ctx.model, 0) + 1
    if ctx.elapsed_ms is not None:
        LATENCIES.append(ctx.elapsed_ms)

laya.load("convaiinnovations/laya", on_predict_end=metrics, hooks_raise=False)

防护栏

一条策略钩子抛出以拦下一个请求。hooks_raise=True(默认)让拦截抵达调用方;on_error 和 on_predict_end 仍然运行,于是审计轨迹记下它。

class Blocked(Exception):
    pass

def guard(ctx):
    if any("ssn" in str(state).lower() for state in ctx.states):
        raise Blocked("possible PII in state")

laya.load("convaiinnovations/laya", on_predict_start=guard)

拿状态形状去测它。一个只读 ctx.states[0] 的防护栏会拦下一次单调用,却让一次 predict_batch 调用把剩下的每一个状态都送进前向传播。

置信度门控

一个 end 钩子把低置信度的答案改写成安全的回退值,或者为下游逻辑给它加标注。这是一次结果改动,不是 一次拒绝。

def gate(ctx):
    for result in ctx.results or []:
        answer = result["answers"].get("dept")
        if answer and answer["confidence"] < 0.6:
            answer["choice"] = "human-review"
            answer["gated"] = True

laya.load("convaiinnovations/laya", on_predict_end=gate)

通过 ctx.results 来改动,它为这次调用的每个状态保存一个 dict:只门控第一个,就会把其他每一个 低置信度的答案不加标注地发出去。

路由覆盖

on_route 可以替换 ctx.decision,为某一类流量钉住一个 checkpoint。

from laya.router import RouteDecision

def pin(ctx):
    if "refund" in str(ctx.states[0]).lower():
        ctx.decision = RouteDecision(
            model="typed-decisions",
            repo="convaiinnovations/laya/typed-decisions",
            reason="refund workflow",
            detection=None,
            workflow=None,
        )

Router(hooks=[pin])

模型生命周期

on_load 和 on_evict 观察 checkpoint。用它们做预热日志、内存核算或驱逐告警。它们运行在 Router 锁之外,所以钩子可以回调进 Router。

class Lifecycle:
    def on_load(self, ctx):
        print("loaded", ctx.model)

    def on_evict(self, ctx):
        print("evicted", ctx.model)

Router(hooks=[Lifecycle()])

多租户上下文

通过在钩子闭包里捕获它,或者从一个 context-local 里读它,来把租户 id 串起来。不要在钩子对象上存逐 请求的状态,除非上锁。

def make_audit(tenant):
    def audit(ctx):
        ship(tenant, ctx.run_id, ctx.results)
    return audit

agent = laya.load("convaiinnovations/laya", on_predict_end=make_audit("acme"))

组合

几种不同种类的钩子自然地组合;已安装的钩子按顺序先运行。

agent = laya.load(
    "convaiinnovations/laya",
    hooks=[Metrics(), Guardrail()],     # metrics first, then policy
    on_predict_start=redact,            # convenience callables appended after hooks
    hooks_raise=True,                   # policy failures are fatal
)

让这个顺序刻意且记录在案,因为更晚的钩子会看到更早那个的改动。

限定作用域的埋点

只给需要它的代码挂上一个 tracer 或调试钩子,而不是重建整个 agent。hooks_installed 在退出时恢复 之前的列表,即便代码块抛异常也一样。

with agent.hooks_installed(DebugDump()):
    agent.system_one(state, questions)   # DebugDump only here

add_hook/remove_hook 不用代码块也能做同样的事,适合一个活得和进程一样久的 tracer。

进程级埋点

一个每个决策都该看到的 tracer 或指标钩子可以注册一次,而不是传给每个 Agent 和 Router。默认值 在实例钩子和逐调用钩子之前运行。

from laya import BaseHook, hooks

class Metrics(BaseHook):
    def on_predict_end(self, ctx):
        record(ctx.model, ctx.elapsed_ms)

hooks.set_default_hooks(hooks=[Metrics()])

这是全局状态,所以要有意识地限定它:在启动时设置一次,并在测试里 clear_default_hooks(),这样 一个测试就不会把钩子泄漏进下一个。

token 预算的塑造

一个 start 钩子可以为一次调用抬高 token 预算,例如当一个问题的选项很多、默认的 head 预算会让标签 塌掉时。有四个细节决定这个钩子是有帮助,还是悄悄把这次调用变得更糟:

  • 一个 start 钩子的 ctx.head_max_len 替换这次调用的预算。在它之前生效的是调用方自己的逐调用 值,或者 ctx.agent.cfg 里的 checkpoint 默认值 —— 所以要拿它做比较,而写一个光秃秃的数字可能 会拉低调用方已经设好的预算。
  • 一次调用回答它携带的每一个问题,所以按其中最宽的那个定尺寸,而不是按碰巧最先出现的那个。
  • 一旦选项装不下 head,laya/common.py 给每个选项 max(4, (head_max_len - 16) // k) 个 token。 因此 16 + 4 * k 正好落在那条下限上:每个标签仍然被削减到它和其他标签共享的那些 token,这 正是这个钩子写来避免的塌陷。16 + 8 * k 让它们保持可区分。
  • 状态拿到 max_len - head_max_len - 8 个 token,所以拓宽的 head 必须连带拓宽 max_len,否则 状态就丢了自己的窗口。
def widen_for_high_cardinality(ctx):
    k = max((len(q.get("criteria", {}) or {}) for q in ctx.questions.values()), default=0)
    if k < 50:
        return
    cfg = getattr(ctx.agent, "cfg", None) or {}
    head = ctx.head_max_len if ctx.head_max_len is not None else cfg.get("head_max_len", 192)
    window = ctx.max_len if ctx.max_len is not None else cfg.get("max_len", 512)
    need = 16 + 8 * k                          # 8 tokens per label, not the core's floor of 4
    if need > head:                            # only ever widen, never lower
        ctx.head_max_len = need
        ctx.max_len = max(window, need + 8 + 64)   # 8 reserved, then room for the state

agent = laya.load("convaiinnovations/laya", on_predict_start=widen_for_high_cardinality)

这不碰共享的 agent 配置,所以并发的调用不受影响。同样的旋钮也可以逐调用使用: agent.system_one(state, questions, head_max_len=512, max_len=1024)。

拓宽不是免费的:更长的窗口意味着更大的张量,而 checkpoint 是在 512(laya)和 1,024 个 token 上训练的。超过之后,用 predict_shortlist 收窄候选比硬撑预算更好。

反模式

阻塞的工作

钩子运行在调用线程上,而 laya.serve 用一个单独的推理 worker。一个 sleep、等一次网络往返或调用 input() 的钩子会把它后面的每个请求都堵住。

# bad: blocks the whole server
def audit(ctx):
    requests.post("https://slow.example/decisions", json=..., timeout=30)

# better: enqueue, let a background worker ship it
def audit(ctx):
    QUEUE.put_nowait(record(ctx))

如果你非要做慢工作,设置 hooks_concurrent=False 至少让钩子本身不重叠,并让 laya.serve 跑在 一个队列后面。

从 end 钩子抛出以控制流程

on_predict_end 在推理之后运行。在那里抛出会丢掉一个已经算出来的结果,而且在成功路径上还会浮现给 调用方。要拦截,用 start 钩子,在付推理代价之前;要改答案,改写 ctx.results。

没有锁的共享可变状态

同一个钩子实例在很多线程上运行。self.counter += 1 会竞争。

# bad
class Count:
    def __init__(self): self.n = 0
    def on_predict_end(self, ctx): self.n += 1

# good
import threading
class Count:
    def __init__(self):
        self.n = 0
        self._lock = threading.Lock()
    def on_predict_end(self, ctx):
        with self._lock:
            self.n += 1

静默失败

hooks_raise=False 每次失败警告一次,但一个自己捕获一切的钩子把真正的问题藏了起来。

# bad: no one will ever know the audit trail stopped
def audit(ctx):
    try:
        ship(record(ctx))
    except Exception:
        pass

如果一个钩子是可选的,让 hooks_raise=False 处理它,并盯着那些警告。如果不是,让它抛。

留住上下文

一个把 ctx 追加进列表的钩子会让整个状态、问题、结果和 agent 都活着。

# bad: unbounded memory growth
SEEN = []
def audit(ctx):
    SEEN.append(ctx)

# good: keep only what you need
SEEN = []
def audit(ctx):
    SEEN.append((ctx.run_id, ctx.model, ctx.elapsed_ms))

脱敏太晚

到了 on_predict_end,模型已经把状态 tokenize 过了。请在 on_predict_start 里脱敏。

在逐调用钩子里做逐问题逻辑

每次调用只有一个 PredictContext,一次前向传播回答每一个问题。没有逐问题的事件。在 on_predict_end 里遍历答案,也遍历状态:在一个批次上,一个上下文携带这次调用的每个状态。

def flag(ctx):
    for result in ctx.results or []:
        for qid, answer in result["answers"].items():
            if answer.get("confidence", 1.0) < 0.5:
                alert(qid, ctx.run_id)

递归 predict

一个调用 agent.predict/system_one 的钩子会再次运行钩子。没有深度保护的话它会递归下去。

# bad
def enrich(ctx):
    ctx.results = [agent.predict(ctx.states[0], EXTRA_QUESTIONS)]

# good: guard, or use a separate agent with no hooks
def enrich(ctx):
    if getattr(ctx, "_enriched", False):
        return
    ctx._enriched = True
    ctx.results = [enricher.predict(state, EXTRA_QUESTIONS) for state in ctx.states]

hooks= 里的普通可调用对象

hooks= 接受钩子对象;一个裸的可调用对象说不清它是为什么事件,所以被拒绝。用 on_predict_start= / on_predict_end=。

# bad: TypeError
laya.load("convaiinnovations/laya", hooks=[lambda ctx: None])

# good
laya.load("convaiinnovations/laya", on_predict_end=lambda ctx: None)

在 end 钩子里假定结果存在

失败路径上 ctx.results 是 None,除非某个 start 钩子设过它。总要检查。

def audit(ctx):
    if ctx.results is None:
        log_failure(ctx.run_id, ctx.error)
        return
    log_success(ctx.run_id, ctx.results)

依赖顺序的钩子

一个从另一个钩子读取改动的钩子,除非顺序被钉住,否则很脆弱。已安装的钩子按列表顺序运行,然后是 便捷可调用对象;把任何耦合记录下来,或者把耦合的钩子合并成一个对象。

另见

  • 错误:这些反模式背后几个的失败矩阵。
  • 示例:上面这些模式的更完整版本。