Documentação

Tracing

Cada hook de uma chamada recebe o mesmo PredictContext.run_id, por isso um tracer pode correlacionar os eventos de início, fim e erro (e quaisquer spans que abra) sem manter a sua própria contabilidade.

run_id

  • Uma string uuid4().hex, criada uma vez por chamada pública (predict_batch, system_one, Router.predict, ONNXAgent.system_one). O Router.predict_batch cria um por pedido em vez disso, o run_id que teria uma chamada Router.predict para esse pedido.
  • Partilhado por cada hook dessa chamada, incluindo on_error e on_predict_end.
  • Não é global nem persistido: identifica uma chamada dentro do processo. Coloca-o nos teus logs e payloads de saída para correlacionar entre sistemas.
  • Distinto por chamada, por isso duas chamadas nunca colidem.
call A:  run_id=1a2b...  start ──┐
                                 ├─ end
                          error ─┘
call B:  run_id=9f8e...  start ─── end

Um tracer mínimo

Mantém um mapa de run_id para o span que abriste no início, e fecha-o no fim ou no erro.

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()])

Como o on_predict_end corre sempre, é o lugar natural para fechar um span, e consegue ver ctx.error no caminho de falha.

Tempo de vida dos spans

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)

Spans de erro

O on_predict_end corre no caminho de falha com ctx.error definido, por isso um único local de fecho trata de ambos:

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)

Se só implementares on_error, lembra-te de que dispara antes de on_predict_end; não feches o span em ambos, ou vais contar a dobrar.

OpenTelemetry

O exemplo examples/hooks/otel.py registra contadores e um histograma. Para spans reais, conduz a API OTel a partir do tracer. Os hooks são síncronos, por isso usa o exporter síncrono (ou enfileira e exporta a partir de um 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)

Define hooks_raise=False para que uma avaria do tracer nunca falhe um pedido.

Propagação distribuída

run_id é uma string simples, por isso inclui-o em tudo o que sai do processo: linhas de log, o payload enviado para um serviço de auditoria, ou um cabeçalho HTTP se uma decisão desencadear uma chamada a jusante.

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

Lembra-te de que uma chamada bloqueante como esta trava a thread que chama; enfileira-a em vez disso quando serves em simultâneo. Vê padrões.

Chamadas aninhadas

Um hook que chama predict de novo inicia uma nova chamada com um novo run_id. O pai e o filho são independentes a menos que os ligues tu próprio. Captura o id do pai e passa-o adiante:

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"))

Um run_id pai, uma chamada filha por estado. Um corpo que lê ctx.states[0] liga a primeira decisão de um lote e descarta silenciosamente o resto.

Protege-te contra a recursão (vê antipadrões); a proteção mais fácil é um agente enricher separado sem hooks.

Ver também