Dokumentation

Tracing

Jeder Hook eines Aufrufs erhält dieselbe PredictContext.run_id, sodass ein Tracer die Start-, End- und Fehlerereignisse (und alle Spans, die er öffnet) korrelieren kann, ohne eine eigene Buchführung zu führen.

run_id

  • Ein uuid4().hex-String, einmal pro öffentlichem Aufruf erzeugt (predict_batch, system_one, Router.predict, ONNXAgent.system_one). Router.predict_batch erzeugt stattdessen eine pro Anfrage, die run_id, die ein Router.predict-Aufruf für diese Anfrage gehabt hätte.
  • Gemeinsam genutzt von jedem Hook dieses Aufrufs, einschließlich on_error und on_predict_end.
  • Nicht global und nicht persistent: Sie kennzeichnet einen Aufruf innerhalb des Prozesses. Schreibe sie in deine Logs und ausgehenden Payloads, um über Systeme hinweg zu korrelieren.
  • Pro Aufruf verschieden, sodass zwei Aufrufe nie kollidieren.
call A:  run_id=1a2b...  start ──┐
                                 ├─ end
                          error ─┘
call B:  run_id=9f8e...  start ─── end

Ein minimaler Tracer

Führe eine Zuordnung von run_id zum Span, den du beim Start geöffnet hast, und schließe ihn am Ende oder bei einem Fehler.

import time

class Tracer:
    def __init__(self):
        self.spans = {}

    def on_predict_start(self, ctx):
        self.spans[ctx.run_id] = {
            "model": ctx.model,
            "started_at": ctx.started_at,
        }

    def on_predict_end(self, ctx):
        span = self.spans.pop(ctx.run_id, None)
        if span is None:
            return
        emit_span(
            name="laya.predict",
            run_id=ctx.run_id,
            model=ctx.model,
            duration_ms=ctx.elapsed_ms,
            usage=ctx.usage,
            ok=ctx.error is None,
        )

    def on_error(self, ctx):
        # on_error runs before on_predict_end; leaving the span for on_predict_end is fine,
        # or close it here if you prefer.
        pass

agent = laya.load("convaiinnovations/laya", hooks=[Tracer()])

Weil on_predict_end immer läuft, ist es der natürliche Ort, um einen Span zu schließen, und es kann ctx.error auf dem Fehlerpfad sehen.

Span-Lebensdauer

on_predict_start ──► open span (run_id, model, started_at)
      │
      ├─ inference
      │
on_error ──► record ctx.error on the span
      │
on_predict_end ──► close span (elapsed_ms, usage, ok)

Fehler-Spans

on_predict_end läuft auf dem Fehlerpfad mit gesetztem ctx.error, sodass eine einzelne Schließstelle beides abdeckt:

def on_predict_end(self, ctx):
    span = self.spans.pop(ctx.run_id, None)
    if span is None:
        return
    if ctx.error is not None:
        span["status"] = "error"
        span["error_type"] = type(ctx.error).__name__
        span["error_message"] = str(ctx.error)
    span["duration_ms"] = ctx.elapsed_ms
    emit(span)

Wenn du nur on_error implementierst, denke daran, dass es vor on_predict_end ausgelöst wird; schließe den Span nicht in beiden, sonst zählst du doppelt.

OpenTelemetry

Das Beispiel examples/hooks/otel.py zeichnet Zähler und ein Histogramm auf. Für echte Spans steuere die OTel-API aus dem Tracer. Hooks sind synchron, also verwende den synchronen Exporter (oder stelle in eine Queue ein und exportiere aus einem Worker).

from opentelemetry import trace

tracer = trace.get_tracer("laya")

class OTelHooks:
    def __init__(self):
        self.spans = {}

    def on_predict_start(self, ctx):
        span = tracer.start_span("laya.predict", attributes={"laya.run_id": ctx.run_id, "laya.model": ctx.model})
        self.spans[ctx.run_id] = span

    def on_predict_end(self, ctx):
        span = self.spans.pop(ctx.run_id, None)
        if span is None:
            return
        if ctx.usage:
            span.set_attribute("laya.input_tokens", ctx.usage["input_tokens"])
        if ctx.error is not None:
            span.record_exception(ctx.error)
            span.set_status(trace.Status(trace.StatusCode.ERROR))
        span.end()

laya.load("convaiinnovations/laya", hooks=[OTelHooks()], hooks_raise=False)

Setze hooks_raise=False, damit ein Ausfall des Tracers nie eine Anfrage scheitern lässt.

Verteilte Propagation

run_id ist ein einfacher String, also nimm sie in alles auf, was den Prozess verlässt: Log-Zeilen, die an einen Audit-Dienst gesendete Payload oder einen HTTP-Header, wenn eine Entscheidung einen nachgelagerten Aufruf auslöst.

def audit(ctx):
    requests.post(
        "https://audit.example/decisions",
        json=record(ctx),
        headers={"X-Laya-Run-Id": ctx.run_id},
        timeout=2,
    )

Denke daran, dass ein blockierender Aufruf wie dieser den aufrufenden Thread blockiert; stelle ihn stattdessen in eine Queue, wenn du nebenläufig bedienst. Siehe Muster.

Verschachtelte Aufrufe

Ein Hook, der predict erneut aufruft, startet einen neuen Aufruf mit einer neuen run_id. Eltern- und Kindaufruf sind unabhängig, es sei denn, du verknüpfst sie selbst. Fange die ID des Elternaufrufs ein und gib sie weiter:

def enrich(ctx):
    for state in ctx.states:                     # a hook sees every state of the call
        child = enricher.predict(state, EXTRA_QUESTIONS)
        record_child_span(parent_run_id=ctx.run_id, child_run_id=child.get("run_id"))

Eine run_id des Elternaufrufs, ein Kindaufruf pro Zustand. Ein Body, der ctx.states[0] liest, verknüpft die erste Entscheidung eines Batchs und verwirft den Rest stillschweigend.

Schütze dich vor Rekursion (siehe Antimuster); die einfachste Absicherung ist ein separater enricher-Agent ohne Hooks.

Siehe auch