パターンとアンチパターン
パターンとアンチパターン
フックは小さな継ぎ目で、うまく使うのも下手に使うのも簡単です。このページは本番で持ちこたえる形と、噛みついてくる形を集めます。
パターン
監査ログ
すべての意思決定を、それを再構成するのに十分な情報とともに記録します。state、質問、答え、モデル、ルーティング判断、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)
フックは呼び出しごとに 1 回発火し、predict_batch 呼び出しはすべての state をその中に運ぶので、レコードは意思決定ごとに書かれます。ctx.states と ctx.results はインデックスで対応します。ctx.usage と ctx.elapsed_ms は呼び出し全体の合計で、各 result は自分の usage を運びます。
ログ行の喪失がリクエストを失敗させてはいけないなら寛容にしてください:hooks_raise=False。監査証跡がコンプライアンス要件なら厳格にしてください。
PII の秘匿化
秘匿化はトークン化の前に on_predict_start で行わなければなりません。そうでなければモデルはすでにデータを見ています。
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)
フックは呼び出しごとに 1 回発火するので、predict_batch では 1 つの state をキーにするだけでは足りません。ctx.skip() は呼び出しが返したであろうすべての結果を置き換えます。並行して提供するときはロックでキャッシュを保護してください。Router ではキャッシュされたペイロードにも routing キーが付くので、戻り値の形は変わりません。
同じ組は hooks= 引数を通じて単一の LangChain ノードでも動きます。これは、そのエージェントの他のすべての呼び出し元が見るものを変えずに、グラフの 1 つのホットなステップをキャッシュする方法です。
メトリクス
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)
state の形に対してテストしてください。ctx.states[0] だけを読むガードは、単一の呼び出しをブロックし、predict_batch 呼び出しが残りのすべての state をフォワードパスへ通すのを許してしまいます。
信頼度ゲート
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 を通して変更してください。それは呼び出しの state ごとに 1 つの dict を持ちます。最初のものだけをゲートすると、他のすべての低信頼度の答えを注釈なしで出荷してしまいます。
ルーティングの上書き
on_route は ctx.decision を置き換え、ある種のトラフィックに対してチェックポイントを固定できます。
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 はチェックポイントを観察します。ウォームアップのログ、メモリの計上、追い出しのアラートに使ってください。それらは 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()])
マルチテナントのコンテキスト
テナント 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 やデバッグフックを付けます。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() を呼んで、あるテストが次のテストへフックを漏らさないようにします。
トークン予算の形成
start フックは 1 回の呼び出しのトークン予算を上げられます。たとえば質問の選択肢が多く、既定の head 予算ではラベルが崩壊してしまうときです。フックが助けるか黙って呼び出しを悪化させるかを決める 4 つの詳細があります。
- start フックの
ctx.head_max_lenはその呼び出しの予算を置き換えます。その前に有効なのは、呼び出し元自身の呼び出しごとの値か、ctx.agent.cfgのチェックポイント既定です —— だからそれと比較してください。素の数値を書くと、呼び出し元がすでに設定した予算を下げてしまいます。 - 1 回の呼び出しは運ぶすべての質問に答えるので、たまたま最初に来たものではなく最も広いものでサイズを決めてください。
- 選択肢が head に収まらなくなると、
laya/common.pyは各々にmax(4, (head_max_len - 16) // k)トークンを与えます。したがって16 + 4 * kはちょうどその下限に落ちます。すべてのラベルが他のものと共有するトークンまで切り詰められ、これがフックが避けようとした崩壊そのものです。16 + 8 * kならそれらを区別可能に保てます。 - state は
max_len - head_max_len - 8トークンを得るので、広げた head は一緒にmax_lenも広げないと state が窓を失います。
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 config に触れないので、並行する呼び出しは影響を受けません。同じつまみは呼び出しごとにも使えます:agent.system_one(state, questions, head_max_len=512, max_len=1024)。
広げるのは無料ではありません。長い窓は大きなテンソルを意味し、チェックポイントは 512(laya)と 1,024 トークンで学習されています。それを超えるなら、予算を引き伸ばすよりも 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 は失敗ごとに 1 回警告しますが、すべて自分で捕まえるフックは本当の問題を隠します。
# bad: no one will ever know the audit trail stopped
def audit(ctx):
try:
ship(record(ctx))
except Exception:
pass
フックが任意なら hooks_raise=False に任せ、警告を見てください。そうでなければ、例外を投げさせてください。
コンテキストの保持
ctx をリストに追加するフックは、state・質問・結果・エージェントの全体を生かし続けます。
# 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 の時点で、モデルはすでに state をトークン化しています。on_predict_start で秘匿化してください。
呼び出しごとのフックでの質問ごとのロジック
呼び出しごとに 1 つの PredictContext があり、1 回のフォワードパスがすべての質問に答えます。質問ごとのイベントはありません。答えは on_predict_end の内部で反復し、state も反復してください。バッチでは、1 つのコンテキストが呼び出しのすべての state を運びます。
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= の素の callable
hooks= はフックオブジェクトを取ります。素の callable はどのイベント向けかを述べないので拒否されます。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 フックで結果が存在すると仮定する
失敗経路では、start フックが設定しない限り ctx.results は None です。必ず確認してください。
def audit(ctx):
if ctx.results is None:
log_failure(ctx.run_id, ctx.error)
return
log_success(ctx.run_id, ctx.results)
順序に依存するフック
別のフックからの変更を読むフックは、順序が固定されていなければ脆くなります。インストールされたフックはリスト順で走り、次に便宜 callable です。結合があれば文書化するか、結合したフックを 1 つのオブジェクトにまとめてください。