Skip to content

a2a_completion_notifier — customization

Customization guide for the bite: sender-auth modes, the receiver lifecycle, configuration, and the full API surface.

New here? Start with the main overview. For internals & troubleshooting see the troubleshooting & internals page.

Receiver packaging

The receiver runs as a standalone FastAPI process (create_receiver_app), started from your own entry point (uvicorn or equivalent) and managed by your orchestrator or compose setup, always-on next to the supervisor.

Sender authentication modes (A2A_VERIFY_MODE)

The receiver's sender-verification posture is selected by the supervisor via A2A_VERIFY_MODE on ReceiverSettings (default dev). In every mode the routing callback token is mandatory; the notification is rejected with 400 if metadata.token cannot be unsealed. The mode only governs sender-identity verification:

Mode Callback-token secret RSA signing key + JWKS Behaviour
dev ✅ required ❌ not needed Trusted private network; accepts signed or unsigned. Logs a2a_dev_ignores_sender_jwt.
verify ✅ required ✅ required (JWKS URL + issuer) Verifies signed notifications (failing closed to 401 when a signature is bad), tolerates unsigned (logged). Good mid-migration from dev.
strict ✅ required ✅ required Rejects unsigned (401); verifies every signed notification. Requires A2A_SUBAGENT_JWKS_URL + A2A_SUBAGENT_ISSUER.

dev and verify log a warning if you configure sender-auth settings they will ignore (or if verify has no JWKS/issuer to verify against).

Receiver lifecycle

Notification request lifecycle, in order:

  1. Unseal the callback token: the routing invariant, enforced in every mode.
  2. Sender verification, dispatched by A2A_VERIFY_MODE.
  3. Claim the dedup key exactly once (JWT jti, else the payload id).
  4. Accept only terminal states (completed / failed / ...).
  5. taskId inside the JWT must match the payload id (signed only).
  6. ACK 202 (BackgroundTasks), then Store write + wake.

Wake-up input (DEFAULT_WAKE_INPUT / wake_input)

After the receiver writes the completion to the supervisor's Store it asks the Agent Server for a wake-up run (NotificationMailbox.trigger_if_idle) so the drained notice reaches a model call. By default the wake run carries the real user message DEFAULT_WAKE_INPUT rather than an empty input:

  • On an idle (non-interrupted) thread, a run started with input=None is a silent no-op in some langgraph-api versions (no node executes, so the before_model drain never fires and the supervisor never wakes).
  • The minimal wake message makes the graph run. The MailboxDrainMiddleware injects the completion notice itself. The model does see the wake message as a user turn, and only to make the run execute; the notice injected by the drain is what the supervisor should act on, and the supervisor prompt that defines that behavior (drain step) keys off the injected notice, not this wake text. Overriding wake_input with deployment-specific text is safe.
from langshark_bites.a2a_completion_notifier.receiver import create_receiver_app
from langshark_bites.a2a_completion_notifier.settings import ReceiverSettings

# Default: every wake run carries DEFAULT_WAKE_INPUT (a minimal wake message).
app = create_receiver_app(settings=ReceiverSettings.from_env())

# Opt-out: resume/interrupt flows use a truly input-less wake run (original behavior).
app = create_receiver_app(settings=ReceiverSettings.from_env(), wake_input=None)

# Custom wake input: any graph input shape the supervisor accepts.
app = create_receiver_app(
    settings=ReceiverSettings.from_env(),
    wake_input={"messages": [{"role": "user", "content": "check for completions"}]},
)

The same wake_input parameter is available on NotificationMailbox directly.

Requirements & env vars

Install the bite (the [a2a-notifier] extra pulls fastapi, uvicorn, a2a-sdk):

uv add "langshark-bites[a2a-notifier]"

Key A2A_* env vars (resolved by from_env(); default shown where relevant):

Env var Purpose Default
A2A_SUPERVISOR_URL Supervisor Agent Server base URL http://localhost:8000
A2A_SUPERVISOR_API_KEY Supervisor API key –
A2A_CALLBACK_TOKEN_SECRET Secret for minting/unsealing callback tokens –
A2A_RECEIVER_URL Public receiver base URL (used as JWT aud) http://localhost:8001
A2A_RECEIVER_HOST / A2A_RECEIVER_PORT Receiver bind address 0.0.0.0 / 8001
A2A_VERIFY_MODE Sender-auth posture (dev / verify / strict) dev
A2A_SUBAGENT_JWKS_URL Subagent JWKS endpoint (verify/strict) –
A2A_SUBAGENT_ISSUER Pin the sender iss claim (verify/strict) –
A2A_SUBAGENT_URL Subagent Agent Server base (full-result fetch) –
A2A_SUBAGENT_API_KEY Bearer token for the subagent deployment –
A2A_REDIS_URL Redis URL used for cross-replica jti dedup (via RedisJtiStore) – (in-memory default)

Delivery-failure status: on a Store/wake failure the receiver logs a2a_forward_failed / a2a_forward_wake_failed at error level (surface, not silent) but does not park the notification anywhere for later inspection. Treat those log events as the delivery-failure signal.

Emitter-side env (subagent deployment): A2A_PRIVATE_KEY_PEM (or A2A_PRIVATE_KEY_FILE), A2A_KID, A2A_EMITTER_ISSUER, A2A_EMITTER_AUDIENCE, and the retry policy (A2A_TIMEOUT_SECONDS, A2A_MAX_RETRIES, A2A_RETRY_BACKOFF_SECONDS).

HTTP / error surface

What the receiver can return, and what to check when it does.

Response Meaning What to check
202 {"status":"accepted"} Notification accepted; Store write + wake run in background –
202 {"status":"duplicate"} jti already claimed (at-least-once redelivery) A2A_REDIS_URL wiring; jti TTL vs sender retry window
400 Callback token invalid/expired, or non-terminal state A2A_CALLBACK_TOKEN_SECRET mismatch; token TTL; status.state
401 Bad/missing sender JWT, unknown kid, stale iat JWKS URL, iss/aud pinning, A2A_VERIFY_MODE, iat_staleness_seconds
502 JWKS/key fetch failed Subagent JWKS endpoint reachable; response is valid JSON

Error cases inside the bite log structlog events (see the troubleshooting page for the full table and the alert-worthy log events).

API reference

Grouped by role; each section lists the public symbols of the modules that implement it.

Emitter — subagent server side

Sent by the subagent deployment: the middleware that fires the webhook, the RS256 signer, and the HTTP delivery client.

langshark_bites.a2a_completion_notifier.middleware

A2A push-notification emitter middleware for subagent graphs.

Why this exists

LangChain does not implement A2A push notifications: the Agent Server supports message/send, message/stream, tasks/get, tasks/cancel on the inbound side, but *TaskPushNotificationConfig and SubscribeToTask return -32601 (not implemented). When the supervisor dispatches work to a subagent on a separate server, nothing on the subagent server emits a completion webhook. This middleware provides that missing server-side push behaviour on the subagent graph's stack.

It is middleware, not a tool, by design:

  • aafter_agent fires once per invocation, unconditionally, after the agent completes -- exactly when a completion webhook must go out.
  • awrap_model_call catches exceptions so a failed run also notifies, then re-raises. A tool call could do neither (it depends on the model choosing to call it, and runs inside the turn).
How the dispatch config gets here

The supervisor passes the PushNotificationConfig (webhook url + opaque callback token) at dispatch time via config.configurable. Because middleware hooks do not receive the RunnableConfig, the graph is built by a dynamic graph factory that reads it and constructs this middleware -- the same pattern the reference async-deep-agents completion_notifier uses for parent_thread_id. See build_a2a_notifier_from_config.

The emitter signs the outgoing JWT with the subagent deployment's RS256 key (A2ASigner) and POSTs with retry via PushClient. Delivery failures (the webhook itself is down/errors) are logged and swallowed -- a failed webhook must never take down the subagent's own run. A failure to construct or sign the notification is different: it raises :class:PushEmissionError so the missing push is not silently lost (on the awrap_model_call failed-run path, the model's original error is still re-raised).

A2APushNotifierMiddleware
A2APushNotifierMiddleware(
    push_config,
    signer=None,
    push_client=None,
    task_id_resolver=None,
)

Bases: AgentMiddleware

Emits an A2A completion notification at terminal run states.

Attach to the subagent graph's middleware stack. Add it last so every other middleware has finished before the notification fires.

Parameters:

Name Type Description Default
push_config
PushNotificationConfig

Webhook target + opaque callback token for this run.

required
signer
A2ASigner | None

RS256 signer for the subagent deployment.

None
push_client
PushClient | None

HTTP client that POSTs the signed notification.

None
task_id_resolver
Callable[[Any], str | None] | None

Optional callable mapping the hook runtime to the A2A task id. Defaults to runtime.execution_info.thread_id (then run_id) so the id is fetchable by the supervisor's full-result primitive (the async-agent "thread id == task id" convention).

None
aafter_agent async
aafter_agent(state, runtime)

Fire the 'completed' notification after the agent finishes.

abefore_agent async
abefore_agent(state, runtime)

Resolve and cache the A2A task id for this run.

Fires before the first model call, so awrap_model_call can emit a failed notification even though it does not receive a runtime.

awrap_model_call async
awrap_model_call(request, handler)

Fire a 'failed' notification on exception, then re-raise.

Uses the task id resolved by abefore_agent. If the notification itself fails to build, that PushEmissionError is logged but must not mask the model exception being propagated.

PushEmissionError

Bases: RuntimeError

The completion notification could not be built or signed.

Distinct from delivery failure: PushClient.send already reports its own retry/failure surface. This surfaces a failure to construct the notification (task id missing, signer failed) instead of silently swallowing it -- a non-notified completion is a correctness gap, not a best-effort no-op.

PushNotificationConfig
PushNotificationConfig(
    url,
    token,
    authentication=None,
    mode="dev",
    context_id=None,
)

A2A TaskPushNotificationConfig equivalent for one dispatch.

Mirrors the A2A spec field names so a supervisor-side builder can produce this from a dict without an import cycle.

Attributes:

Name Type Description
url

Webhook URL the subagent server must POST the notification to (the supervisor's receiver /a2a/notifications).

token

Opaque callback token, echoed back uninterpreted in the notification metadata.

authentication

Optional auth mapping (e.g. {"schemes": ["Bearer"]}).

from_dict classmethod
from_dict(data)

Build from a raw mapping (as placed under a2a_push_config).

PushNotificationConfigError

Bases: ValueError

The push config in config.configurable is missing or invalid.

build_a2a_notifier_from_config
build_a2a_notifier_from_config(
    config,
    signer=None,
    push_client=None,
    task_id_resolver=None,
)

Factory for dynamic subagent graphs: read config, build middleware.

Use inside a graph factory so the middleware is constructed per run with the dispatch's webhook config::

def make_graph(config):
    notifier = build_a2a_notifier_from_config(
        config, signer=signer, push_client=client,
    )
    agent = create_agent(model=..., tools=..., middleware=[notifier])
    return agent

Parameters:

Name Type Description Default
config
Mapping[str, Any] | None

The runnable config (containing configurable).

required
signer
A2ASigner | None

Optional deployment signer. When omitted, notifications are sent unsigned -- accepted by dev and (for unsigned senders) verify receivers, rejected by strict receivers.

None
push_client
PushClient | None

Delivery client.

None
task_id_resolver
Callable[[Any], str | None] | None

Override for the A2A task id resolver. The default resolves to execution_info.thread_id (then run_id) so the id is fetchable by the supervisor's full-result primitive.

None

Raises:

Type Description
PushNotificationConfigError

If no a2a_push_config is present, or a strict-mode dispatch arrives without a signer.

extract_push_config
extract_push_config(config)

Read config.configurable["a2a_push_config"], if present.

Returns:

Type Description
PushNotificationConfig | None

A parsed PushNotificationConfig, or None when the runnable

PushNotificationConfig | None

config has no a2a_push_config (the run was not dispatched as a

PushNotificationConfig | None

subagent task, so there is nothing to notify).

langshark_bites.a2a_completion_notifier.signer

A2A push notification JWT signing (RS256) and JWKS export.

Why this exists

The A2A streaming/async spec requires the sending server to authenticate itself to the client's webhook with a signed JWT: claims iss, aud, iat, exp, jti, taskId and a kid header, with public keys published at a JWKS endpoint. LangChain does not implement this push side, so this module provides the signing half for the EmitterSettings.

This is the subagent-deployment side. The signing private key never leaves the subagent deployment; jwks() produces the public key set for the deployment's /.well-known/jwks.json.

Usage
from langshark_bites.a2a_completion_notifier.signer import A2ASigner

signer = A2ASigner(
    private_key_pem=pem, kid="subagent-1",
    issuer="https://subagents.example.com", audience="https://receiver.example.com",
)
jwt_token = signer.sign(task_id="task_123")
jwks_doc = signer.jwks()   # serve this at your JWKS endpoint
A2ASigner dataclass
A2ASigner(
    private_key_pem,
    kid,
    issuer,
    audience,
    jti_ttl_seconds=300,
)

Signs A2A push notification JWTs for a subagent deployment.

Attributes:

Name Type Description
private_key_pem str

RSA private key (PEM, PKCS#8 or PKCS#1) used to sign outgoing notifications.

kid str

Key ID published in the JWT header and the JWKS set. The receiver fetches the matching public key by this value.

issuer str

The iss claim -- the subagent deployment's URL.

audience str

The aud claim -- the receiver's webhook URL. This must match RECEIVER_URL on the supervisor side.

jti_ttl_seconds int

Lifetime of each signed JWT (exp - iat).

jwks
jwks()

Export the public key set for this signer.

Returns:

Type Description
dict[str, Any]

A JWKS document ({"keys": [...]}) containing this signer's

dict[str, Any]

RSA public key keyed by kid. Serve it at your deployment's

dict[str, Any]

JWKS endpoint (e.g. /.well-known/jwks.json).

sign
sign(task_id, *, now=None)

Sign a JWT authenticating the sender of one A2A notification.

Parameters:

Name Type Description Default
task_id str

The A2A task whose terminal event this JWT authenticates. Used for the taskId claim and to derive a unique jti.

required
now int | None

Override the current time (tests).

None

Returns:

Type Description
str

A signed RS256 JWT with a unique jti scoped to task_id.

langshark_bites.a2a_completion_notifier.push_client

HTTP client that POSTs A2A push notifications with retry/backoff.

Why this exists

The emitter middleware produces a signed notification and must deliver it at-least-once. The A2A spec recommends 10-30s webhook timeouts with retry and backoff, so the POST path needs its own retry policy independent of the surrounding SDK client.

This module is deliberately thin: it takes an already-signed Bearer token and a JSON-serializable payload and POSTs them to the registered webhook, returning a bool so middleware can log-and-continue rather than raise into the agent loop.

PushClient
PushClient(*, http=None, settings=None)

POSTs signed A2A notifications to a webhook URL with backoff.

Parameters:

Name Type Description Default
http
AsyncClient | None

Optional client to reuse. When omitted, one is created and owned by this instance (closed by :meth:aclose).

None
settings
PushClientSettings | None

Retry/backoff settings. Defaults to PushClientSettings().

None
aclose async
aclose()

Close the HTTP client if this instance created it.

send async
send(url, *, bearer=None, payload)

Deliver one notification; retry on failure with backoff.

Parameters:

Name Type Description Default
url str

The webhook URL from PushNotificationConfig['url'].

required
bearer str | None

The signed sender JWT (goes in the Authorization header). None sends the notification unsigned -- the dev / verify-unsigned receiver path.

None
payload dict[str, Any]

JSON-serializable notification body.

required

Returns:

Type Description
bool

True if any attempt got an HTTP 2xx; False if all attempts

bool

failed. Never raises -- callers log-and-continue.

PushClientSettings dataclass
PushClientSettings(
    timeout_seconds=30.0,
    max_retries=3,
    retry_backoff_seconds=1.0,
)

Retry policy for A2A webhook delivery.

Attributes:

Name Type Description
timeout_seconds float

Per-attempt HTTP timeout.

max_retries int

Extra attempts after the first POST.

retry_backoff_seconds float

Base of the exponential backoff (base * 2**attempt sleeps between attempts).

Supervisor — dispatch

Minted by the supervisor when it dispatches work to a subagent.

langshark_bites.a2a_completion_notifier.push_config

Supervisor-side push-config helpers: mint the token and build the config.

Why this exists

The supervisor must, at dispatch time, mint an opaque callback token and hand it to the subagent alongside the webhook URL, so the subagent's emission middleware can notify the receiver::

config.configurable["a2a_push_config"] = {
    "url": f"{RECEIVER_URL}/a2a/notifications",
    "token": <minted callback token>,
    "authentication": {"schemes": ["Bearer"]},
}

build_push_config wraps that exact construction so a supervisor graph does not need to know the wire format. The token is signed with the supervisor-side secret and carries thread_id/assistant_id/ dispatch_id; the subagent echoes it back uninterpreted and the receiver unseals it to route.

Usage
from langshark_bites.a2a_completion_notifier.push_config import (
    build_push_config,
)

config = build_push_config(
    thread_id="thread_1",
    assistant_id="supervisor_asst",
    dispatch_id="dispatch_1",
    receiver_url="https://receiver.example",
    callback_token_secret="supervisor-secret",
)
# pass config into the subagent run:
#   await lg.runs.create(thread_id=..., assistant_id=...,
#                        input=..., config=config)
build_push_config

Build the config.configurable push config for one subagent.

Parameters:

Name Type Description Default
thread_id
str

Supervisor thread the completion must be routed to.

required
assistant_id
str

Supervisor assistant to wake when the subagent finishes.

required
dispatch_id
str

Correlation id for the whole fan-out dispatch.

required
receiver_url
str

Public base URL of the supervisor's receiver (the /a2a/notifications path is appended automatically).

required
callback_token_secret
str

CALLBACK_TOKEN_SECRET (supervisor-side only; never sent to the subagent).

required
callback_token_ttl_seconds
int

Lifetime of the callback token.

86400
mode
A2AVerifyMode

The supervisor's A2A_VERIFY_MODE at dispatch time. Stamped into the push config so the emitter can fail fast when a strict dispatch demands signing but the deployment has no key.

DEV
now
int | None

Override the current time (tests).

None

Returns:

Type Description
dict[str, Any]

A dict suitable for config.configurable on the subagent run. The

dict[str, Any]

subagent's build_a2a_notifier_from_config reads it and wires the

dict[str, Any]

emitter middleware automatically.

Receiver — inbound webhook

Terminates the A2A POST, verifies the sender, and deduplicates.

langshark_bites.a2a_completion_notifier.receiver

Supervisor-side A2A notification receiver (FastAPI webhook endpoint).

Why this exists

A LangGraph Agent Server has no generic inbound webhook endpoint -- the A2A push notification has nowhere to land. This receiver is the supervisor- side component that terminates POST /a2a/notifications, authenticates the sender, deduplicates, unseals the callback token, and delivers the completion into the supervisor via the mailbox (NotificationMailbox).

It also exposes POST /a2a/result -- the sister primitive to the async- agent framework's poll-based check_async_task. Because the notifier already proved the task is terminal, the supervisor can fetch the full result of a completed task with a single on-demand call (no status gate): the receiver forwards the task id to the configured subagent deployment and returns its thread state.

Deployment rule (from the design notes): the receiver lives next to the supervisor, not the subagents. It holds the supervisor's API key and the CALLBACK_TOKEN_SECRET; neither ever crosses to the subagent side.

The handler logic is also exported as framework-neutral functions (process_notification / process_task_result / process_health) so alternative hosts can reuse the same webhook logic without depending on the FastAPI app itself.

Notification request lifecycle (in this exact order):

  1. Unseal the callback token -- the routing invariant, enforced in every A2A_MODE (dev / verify / strict).
  2. Sender verification, dispatched by mode: dev accepts signed and unsigned unverified; verify verifies signed notifications (failing closed) while accepting unsigned ones; strict rejects unsigned and always verifies.
  3. Claim the dedup key exactly once (JWT jti when signed, else the payload id).
  4. Accept only terminal states (completed / failed / ...).
  5. taskId inside the JWT must match the payload id (signed only).
  6. ACK with 202 (BackgroundTasks), then mailbox.write + wake.
  7. The sender iss is persisted with the mailbox notice for later full-result routing.
Usage
from langshark_bites.a2a_completion_notifier.receiver import create_receiver_app
from langshark_bites.a2a_completion_notifier.settings import ReceiverSettings

app = create_receiver_app(settings=ReceiverSettings.from_env())
# uvicorn main:app --port 8001
ReceiverError
ReceiverError(status_code, detail)

Bases: Exception

A hard receiver rejection that maps to an HTTP error status.

Framework-neutral: the FastAPI wrapper translates it to HTTPException.

create_receiver_app
create_receiver_app(
    *,
    settings,
    lg=None,
    jti_store=None,
    jwks_client=None,
    mailbox=None,
    wake_input=DEFAULT_WAKE_INPUT,
    result_settings=None,
    result_http=None,
)

Build the receiver FastAPI application.

Every dependency can be injected for tests; defaults are wired lazily when omitted:

  • lg -- langgraph_sdk.get_client(url, api_key) for the supervisor Agent Server.
  • jti_store -- InMemoryJtiStore. In production pass a RedisJtiStore so cross-replica dedup is atomic.
  • jwks_client -- fetches from settings.subagent_jwks_url.
  • mailbox -- NotificationMailbox over lg.
  • result_settings -- defaults to ResultFetchSettings from settings.subagent_url / subagent_api_key; used by the POST /a2a/result sister primitive.
  • result_http -- the outbound HTTP client for full-result fetches.

Parameters:

Name Type Description Default
settings
ReceiverSettings

Supervisor-side receiver configuration.

required
lg
Any | None

LangGraph SDK client (supervisor). Injected for tests.

None
jti_store
JtiStore | None

Single-use jti claim store.

None
jwks_client
JWKSClient | None

Source of the subagent deployment's public keys.

None
mailbox
NotificationMailbox | None

Delivery sink for accepted notifications.

None
wake_input
Any | None

Input for the wake-up run created after a mailbox write. Defaults to NotificationMailbox.DEFAULT_WAKE_INPUT; pass None for a resume-style input-less run.

DEFAULT_WAKE_INPUT
result_settings
ResultFetchSettings | None

Fetch configuration for the subagent deployment.

None
result_http
AsyncClient | None

HTTP client used to fetch completed task results.

None

Returns:

Type Description
FastAPI

A configured FastAPI application.

process_health async
process_health()

Shared liveness body for the FastAPI route.

process_notification async
process_notification(
    payload,
    *,
    settings,
    jti_store,
    jwks_client,
    mailbox,
    authorization=None,
)

Validate one A2A notification and stage its mailbox delivery.

This is the framework-neutral core of POST /a2a/notifications, used by the standalone FastAPI receiver. Full lifecycle: unseal the callback token (routing invariant, all modes) -> mode-dispatched sender verification -> claim the dedup key exactly once -> accept only terminal states -> taskId cross-check (signed only).

Returns:

Type Description
dict[str, Any]

(body, forward): the response body and an optional async

Callable[[], Awaitable[None]] | None

callable that must run after the 202 ACK has left the process

tuple[dict[str, Any], Callable[[], Awaitable[None]] | None]

(writes the mailbox, then wakes an idle supervisor). forward is

tuple[dict[str, Any], Callable[[], Awaitable[None]] | None]

None for duplicate/ignored outcomes.

Raises:

Type Description
ReceiverError

A hard rejection already mapped to an HTTP status (400 for routing/taskId defects, 401/502 for sender-auth).

process_task_result async
process_task_result(body, *, result_settings, result_http)

Fetch a completed subagent task's full result (sister primitive).

Framework-neutral core of POST /a2a/result, used by the FastAPI receiver.

Returns:

Type Description
dict[str, Any]

{"task_id": ..., "status": "ok", "result": {...state...}}.

Raises:

Type Description
ReceiverError

400 for a missing task_id, 502 when the result fetch fails.

langshark_bites.a2a_completion_notifier.auth

Verify the A2A sender's JWT at the receiver webhook (JWKS / RS256).

Why this exists

The receiver terminates the A2A webhook POST. It must confirm the notification really comes from the subagent deployment before routing it into the supervisor. Per the A2A spec's asymmetric flow, the subagent signs its JWT with a private key and publishes public keys at a JWKS endpoint; the receiver fetches the key indicated by the kid header, verifies the signature, and pins aud + iss.

Design notes
  • JWKSClient caches fetched keys by kid and re-fetches only on an unknown kid (key rotation support without a long-lived cache).
  • verify_sender_jwt raises typed SenderAuthError subclasses so the FastAPI layer can map each to a 401 without catching bare exceptions.
  • The iat staleness window is checked separately from JWT exp: PyJWT validates exp but not freshness, so a 24h exp with a 10-minute-old iat must still be rejected.
Usage
from langshark_bites.a2a_completion_notifier.auth import (
    JWKSClient, verify_sender_jwt,
)

jwks = JWKSClient("https://subagents.example.com/.well-known/jwks.json")
claims = await verify_sender_jwt(
    "Bearer eyJ...", jwks_client=jwks,
    audience=settings.receiver_url, issuer=settings.subagent_issuer,
)
ExpiredTokenError

Bases: SenderAuthError

The JWT exp has passed.

InvalidSignatureError

Bases: SenderAuthError

The JWT failed signature/claim verification (PyJWT error surface).

JWKSClient
JWKSClient(jwks_url, *, http=None)

Fetches and caches JWKS public keys by kid.

Keys are cached after the first fetch. When a previously-unseen kid appears (key rotation), the endpoint is re-queried; the new key is cached for the rest of the process lifetime.

aclose async
aclose()

Close the underlying HTTP client if this instance owns it.

get_key async
get_key(kid)

Return the JWK for kid, or None if not found.

MissingBearerTokenError

Bases: SenderAuthError

The Authorization header is absent or not a Bearer token.

MissingJtiError

Bases: SenderAuthError

The JWT lacks the jti claim required for deduplication.

SenderAuthError

Bases: Exception

Base class for sender-authentication failures (maps to HTTP 401).

StaleTokenError

Bases: SenderAuthError

The JWT iat is older than the configured staleness window.

UnknownKidError

Bases: SenderAuthError

The JWT kid does not match any key in the JWKS set.

verify_sender_jwt async

Verify the A2A sender's JWT and return its validated claims.

Parameters:

Name Type Description Default
authorization
str | None

The raw Authorization header value.

required
jwks_client
JWKSClient

Source of the subagent deployment's public keys.

required
audience
str

The expected aud claim (the receiver's public URL).

required
issuer
str

The expected iss claim (the subagent deployment).

required
iat_staleness_seconds
int

Max age of the iat claim; older tokens are rejected as stale redeliveries.

300

Returns:

Type Description
dict[str, Any]

The validated JWT claims (contains jti, taskId, iss,

dict[str, Any]

aud, iat, exp).

Raises:

Type Description
MissingBearerTokenError

No Bearer token in the header.

UnknownKidError

No JWKS key matches the JWT's kid.

InvalidSignatureError

Signature or required-claim verification failed.

ExpiredTokenError

The JWT exp has passed.

StaleTokenError

The JWT iat is older than the staleness window.

MissingJtiError

The JWT lacks the jti claim.

langshark_bites.a2a_completion_notifier.idempotency

Idempotency for incoming A2A notifications (jti single-use claim).

Why this exists

A2A delivery is at-least-once: the subagent server retries its webhook POST with backoff until it gets an ACK. If the receiver accepted a notification but the ACK was lost, the sender retries and the receiver sees the same notification twice. Without dedup, a redelivered completion injects a second completion message into the supervisor's history -- which the supervisor may read as two completions for one task (the re-dispatch loop).

The A2A JWT carries a unique jti. This module claims it exactly once per receiver lifetime window, so retries collapse to a no-op.

Why Redis and not the Store

jti claims are receiver infrastructure state, not conversation state: they need a hard TTL, cross-replica atomicity, and sub-millisecond latency on the ACK path. SET NX EX gives all three in one round trip. A LangGraph Store has no native key expiry and check-then-write is two round trips -- wrong tool for this job.

The in-memory store is a single-process fallback for tests and prototypes; it is NOT safe across receiver replicas.

Usage
from langshark_bites.a2a_completion_notifier.idempotency import RedisJtiStore

store = RedisJtiStore(redis_client, ttl_seconds=900)
is_new = await store.claim("jti-abc")   # True once, then False
InMemoryJtiStore
InMemoryJtiStore(*, ttl_seconds=900)

Single-process jti store for tests and prototypes.

Keeps an expiry map instead of a growing set so long-lived processes do not accumulate stale keys. Not safe across receiver replicas -- use RedisJtiStore in production.

claim async
claim(jti)

Attempt to claim jti exactly once within this process.

Parameters:

Name Type Description Default
jti str

The unique JWT ID from the incoming notification.

required

Returns:

Type Description
bool

True if this is the first time jti was seen; False on a

bool

redelivery/duplicate.

JtiStore

Bases: Protocol

A single-use claim store for A2A notification jti values.

claim async
claim(jti)

Return True the first time jti is claimed, False afterwards.

RedisJtiStore
RedisJtiStore(redis, *, ttl_seconds=900)

Redis-backed jti claim store using atomic SET NX EX.

Attributes:

Name Type Description
ttl_seconds

How long a claimed jti is remembered. Must exceed the sender's total retry window (spec: 10s HTTP timeout + retry budget); 900s (15 min) covers most retry policies.

claim async
claim(jti)

Attempt to claim jti exactly once.

SET key 1 NX EX ttl returns a truthy value only if the key did not previously exist -- check-and-write is a single atomic op, so two receiver replicas cannot both claim the same jti.

Parameters:

Name Type Description Default
jti str

The unique JWT ID from the incoming notification.

required

Returns:

Type Description
bool

True if this is the first time jti was seen; False on a

bool

redelivery/duplicate.

Receiver — delivery (Store write + drain)

Writes completions into the supervisor's Store and injects them into the next model call.

langshark_bites.a2a_completion_notifier.mailbox

Mailbox delivery: write completions to the supervisor's Store + wake it.

Why this exists

This is the supervisor-side delivery half. The receiver has already verified and deduplicated a notification and learned the parent CallbackToken. NotificationMailbox then:

  1. Writes the completion into the LangGraph Store under ("notifications", thread_id) keyed by task_id (write first);
  2. Checks whether a run is pending or running on that thread, and if not, creates a trigger run carrying :data:DEFAULT_WAKE_INPUT -- a minimal user message -- so the supervisor wakes and its before_model middleware drains the mailbox. A run with no input is a silent no-op on an idle (non-interrupted) thread in some Agent Server versions, so the wake must always carry real content (see :data:DEFAULT_WAKE_INPUT).

Write-first/check-second ordering makes the check-then-create race resolve to a harmless empty extra run instead of a lost notification.

MailboxWakeError

Bases: Exception

The wake-up trigger run failed to be created.

Distinct from :class:MailboxWriteError: the mailbox may have been written successfully but the supervisor never woke, so the completion sits undelivered. Callers (e.g. the receiver's post-ACK _forward) must surface this, not swallow it.

MailboxWriteError

Bases: Exception

The Store write failed irrecoverably.

NotificationMailbox
NotificationMailbox(
    lg,
    *,
    store_namespace_prefix=(NOTIFICATION_NAMESPACE_PREFIX,),
    wake_input=DEFAULT_WAKE_INPUT,
    wake_multitask_strategy="enqueue",
)

Writes completion notifications to the supervisor's Store.

Parameters:

Name Type Description Default
lg
Any

A LangGraph SDK client for the supervisor Agent Server (langgraph_sdk.get_client()) exposing store and runs.

required
store_namespace_prefix
tuple[str, ...]

Leading tuple of a Store namespace. The mailbox appends thread_id to form the full namespace.

(NOTIFICATION_NAMESPACE_PREFIX,)
wake_input
Any | None

The trigger-run input. Defaults to :data:DEFAULT_WAKE_INPUT (a minimal nudge message) so the wake run actually executes on an idle thread instead of silently no-op'ing. Pass None for a truly input-less run (original behavior, e.g. interrupt/resume flows); the drain middleware supplies the content either way.

DEFAULT_WAKE_INPUT
wake_multitask_strategy
str

Strategy when a run is already queued. "enqueue" is the only option that does not drop work.

'enqueue'
deliver async
deliver(
    parent, task_id, state, summary=None, *, extra=None
)

Write one completion notice into the thread's mailbox.

Parameters:

Name Type Description Default
parent CallbackToken

Unsealed callback token giving the target thread.

required
task_id str

A2A task id; doubles as the Store key, so a redelivery overwrites instead of appending.

required
state str

Terminal state (completed / failed / ...).

required
summary str | None

Optional short result summary for the drain.

None
extra dict[str, Any] | None

Optional extra fields merged into the stored value. The receiver passes the sender iss here so the full-result primitive can route back to the deployment.

None

Raises:

Type Description
MailboxWriteError

If the Store write fails.

trigger_if_idle async
trigger_if_idle(parent)

Create a wake-up run only if no run is pending/running.

Covers the queue-sweep middle state (pending/running) so a swept run cycling between states is not double-woken.

Parameters:

Name Type Description Default
parent CallbackToken

Unsealed callback token giving the target thread.

required

Returns:

Type Description
bool

True if a wake-up run was created; False if a run was already

bool

pending/running (no wake needed) or the liveness check could not

bool

determine the run state.

Raises:

Type Description
MailboxWakeError

The wake-up run create call failed. The mailbox may already hold the completion, but the supervisor never woke -- callers must surface this rather than swallow it.

Notes

The wake run always carries self._wake_input (default :data:DEFAULT_WAKE_INPUT). Construct with wake_input=None to emit a resume-style input-less run instead.

langshark_bites.a2a_completion_notifier.drain

Supervisor-side mailbox drain middleware (before_model hook).

Why this exists

The receiver writes completion notices into the supervisor's LangGraph Store under ("notifications", thread_id). The supervisor graph must consume that mailbox: at the start of each model call, inject everything that has arrived into the model's context. That is the "drain" half of the mailbox + drain design (see design note §04).

MailboxDrainMiddleware is an AgentMiddleware that:

  1. resolves the current thread_id from runtime.execution_info;
  2. reads the mailbox namespace (("notifications", thread_id));
  3. drops the drained items from the Store (drain is idempotent: the completion text is already committed into the returned state update, so a superstep retry re-drains an already-empty mailbox instead of double-injecting);
  4. returns a single HumanMessage carrying all pending notices, so the model sees the complete picture in one turn (coalescing is inherent).

Attach it to the supervisor graph's middleware stack. It is a no-op when there is no Store, no thread_id, or nothing pending.

Dev-server parity

langgraph dev (the in-memory dev runtime) does not populate runtime.store on the middleware runtime, so this middleware falls back to a langgraph-sdk StoreClient built from A2A_SUPERVISOR_URL and drains the same mailbox over HTTP. The platform / Agent Server (Postgres) runtime supplies runtime.store directly and is used unchanged.

MailboxDrainError

Bases: Exception

The mailbox search failed; pending completions could not be drained.

Raised (not swallowed) so a Store outage is visible in the supervisor's error path instead of silently dropping every pending completion.

MailboxDrainMiddleware
MailboxDrainMiddleware(
    *,
    namespace_prefix=MAILBOX_NAMESPACE_PREFIX,
    render_notice=None,
    max_notices=50,
    store=None,
)

Bases: AgentMiddleware

Drains pending A2A completions into each supervisor model call.

Works with both the runtime Store handle (langgraph up / platform Agent Server) and a langgraph-sdk HTTP StoreClient (e.g. langgraph dev where runtime.store is not populated): pass the store explicitly, or let it build one from A2A_SUPERVISOR_URL (ReceiverSettings.supervisor_url) when the runtime handle is absent.

Parameters:

Name Type Description Default
namespace_prefix
tuple[str, ...]

Store namespace prefix. The thread id is appended to form the full mailbox namespace -- must match the receiver's mailbox configuration (default ("notifications",)).

MAILBOX_NAMESPACE_PREFIX
render_notice
RenderNotice | None

Optional callable mapping (task_id, value) to a prompt string. Defaults to :func:render_notice.

None
max_notices
int

Cap on notices drained per model call.

50
store
StoreLike | None

Optional store handle override. When omitted, the runtime's runtime.store is used; if that is also absent, a langgraph-sdk StoreClient is built from A2A_SUPERVISOR_URL (falling back to http://localhost:8000).

None
abefore_model async
abefore_model(state, runtime)

Notify (no-op) if nothing is pending, else return a state update.

Raises:

Type Description
MailboxDrainError

The mailbox asearch failed. A Store outage must not silently drop pending completions.

render_notice
render_notice(task_id, value)

Render one stored completion notice into a prompt-friendly block.

Matches the [SUBAGENT COMPLETION NOTICE] convention so the model can recognise injected completions.

Parameters:

Name Type Description Default
task_id
str

The A2A task id (Store key).

required
value
dict[str, Any]

The stored notice dict (task_id, state, summary...).

required

Returns:

Type Description
str

A short prompt block for the supervisor's model.

Supervisor — full-result fetch

Fetches a completed subagent task's full result on demand (the sister primitive).

langshark_bites.a2a_completion_notifier.result_fetch

Fetch the full result of a completed subagent task (fetch-only primitive).

Why this exists

The completion notifier already proves a subagent run reached a terminal state -- the webhook arrived, was authenticated and deduplicated, and the mailbox delivered it. The supervisor therefore never needs to poll for completion. The async-agent framework's check_async_task couples two concerns (status + result) and assumes the caller knows neither; this module is the sister primitive that assumes the caller knows the first and only gets the second:

  • check_async_task(task_id) -- status round-trip, then maybe a fetch.
  • fetch_task_result(task_id) -- a single on-demand fetch, no status gate.

The full response is read from the subagent deployment's Agent Protocol thread state (GET /threads/{task_id}/state). That works because the emitter resolves the notification task_id to the subagent's thread_id (see middleware._default_task_id), not its per-run run_id -- the same \"thread ID == task ID\" convention the async-agent framework uses.

Because the receiver pins a single subagent_issuer / subagent_jwks_url today, the deployment is also configured, not discovered: subagent_url is the Agent Server base URL whose /threads/ namespace holds the task.

ResultFetchSettings dataclass
ResultFetchSettings(
    subagent_url="",
    subagent_api_key="",
    timeout_seconds=15.0,
)

Configuration for fetching a completed subagent's full result.

Attributes:

Name Type Description
subagent_url str

Agent Server base URL of the subagent deployment. The thread-state endpoint is <subagent_url>/threads/<task_id>/state.

subagent_api_key str

Optional bearer token for that deployment.

timeout_seconds float

Per-attempt HTTP timeout.

TaskResultError

Bases: Exception

The full-result fetch failed (unconfigured, unreachable, or non-2xx).

fetch_task_result async
fetch_task_result(task_id, *, settings, http)

Fetch the full thread state for a completed subagent task.

One round trip, no status check: the completion notifier established that the task is terminal, so the caller is expected to already know it is done.

Parameters:

Name Type Description Default
task_id
str

The notification's task_id (the subagent's thread id).

required
settings
ResultFetchSettings

Fetch configuration (subagent URL + optional API key).

required
http
AsyncClient

HTTP client used for the outbound request.

required

Returns:

Type Description
dict[str, Any]

The subagent's thread state (Agent Protocol GET /threads/{id}/state

dict[str, Any]

JSON body).

Raises:

Type Description
TaskResultError

If no subagent_url is configured, the request fails, or the subagent returns a non-2xx response.

Wire format and routing tokens (shared)

The notification body contract and the opaque callback token that both sides use.

langshark_bites.a2a_completion_notifier.payload

Wire format for A2A completion notifications (emitter -> receiver).

Why this exists

Both ends of the notifier must agree on the notification body. The emitter (subagent side) builds it; the receiver (supervisor side) parses it. This module is the single home for that contract so either side can change without drifting.

Completion profile (A2A Task, sparse)

Webhook bodies are A2A Task objects (official a2a-sdk / proto JSON), populated only as far as a terminal completion needs::

{
    "id": "<task id>",
    "contextId": "<optional A2A context / supervisor thread>",
    "status": {"state": "TASK_STATE_COMPLETED"},
    "metadata": {
        "token": "<opaque callback token>",
        "summary": "...optional extension..."
    }
}

metadata.token remains the supervisor-minted opaque callback token (extension for routing). Core identity is the A2A Task fields (id / contextId / status.state).

Legacy bite dicts (status.state as completed and top-level summary) are still accepted by the extract helpers for one compatibility window.

NotificationPayloadError

Bases: ValueError

Webhook body is not a valid (sparse) A2A Task notification.

build_notification
build_notification(
    task_id,
    state,
    token,
    *,
    context_id=None,
    summary=None,
    artifacts=None,
)

Build a sparse A2A Task body for the receiver webhook.

Parameters:

Name Type Description Default
task_id
str

A2A task id (Task.id).

required
state
str

App label (completed / failed / ...) or A2A enum name.

required
token
str

Opaque callback token from PushNotificationConfig (extension).

required
context_id
str | None

Optional A2A contextId (e.g. supervisor thread id).

None
summary
str | None

Optional short summary stored under metadata.summary.

None
artifacts
list[dict[str, Any]] | None

Optional A2A artifact dicts (advanced; usually omitted).

None

Returns:

Type Description
dict[str, Any]

JSON-serializable A2A Task dict (proto JSON / camelCase).

Raises:

Type Description
NotificationPayloadError

Unknown state or invalid Task after build.

extract_callback_token
extract_callback_token(payload)

Return the opaque callback token from metadata.token.

extract_context_id
extract_context_id(payload)

Return Task.contextId when present.

extract_summary
extract_summary(payload)

Return optional summary from metadata.summary or legacy top-level.

extract_task_id
extract_task_id(payload)

Return Task.id from a (normalized) notification body.

extract_terminal_state
extract_terminal_state(payload)

Return the app terminal state label, or None if non-terminal/unknown.

Accepts A2A enum names and legacy short labels.

normalize_notification_payload
normalize_notification_payload(payload)

Return an A2A Task dict, accepting legacy bite bodies for one release.

Legacy shape::

{"id", "status": {"state": "completed"}, "metadata": {"token"}, "summary"?}

Raises:

Type Description
NotificationPayloadError

Cannot interpret as Task or legacy completion.

parse_a2a_task
parse_a2a_task(payload)

Parse and validate payload as an A2A Task.

Raises:

Type Description
NotificationPayloadError

Not a valid Task message.

task_to_dict
task_to_dict(task)

Serialize an A2A Task to proto JSON (camelCase).

langshark_bites.a2a_completion_notifier.tokens

Opaque callback tokens (HS256 JWT) for the completion notifier.

Why this exists

When the supervisor dispatches work to a subagent it must later tell the receiver which supervisor thread + assistant a completion belongs to -- without leaking those details to the subagent server. The token is minted at dispatch time, travels as the opaque PushNotificationConfig.token, and is echoed back inside the A2A notification payload's metadata. The receiver unseals it to learn the routing target.

Because the token is signed (HS256) with a secret only the supervisor deployment holds, a subagent cannot read it, cannot forge a different routing, and cannot redirect a notification to another thread. The exp claim bounds the token's lifetime so a stale subagent cannot inject into a thread days later.

Usage
from langshark_bites.a2a_completion_notifier.tokens import (
    mint_callback_token,
    unseal_callback_token,
)

token = mint_callback_token(
    thread_id="t1", assistant_id="asst_1", dispatch_id="d1",
    secret="supervisor-secret",
)
parent = unseal_callback_token(token, "supervisor-secret")
assert parent.thread_id == "t1"
CallbackToken dataclass
CallbackToken(thread_id, assistant_id, dispatch_id, exp)

Decoded routing payload carried inside an opaque callback token.

Attributes:

Name Type Description
thread_id str

Supervisor thread the completion must be routed to.

assistant_id str

Supervisor assistant to create the wake-up run against.

dispatch_id str

Correlation id for the whole fan-out dispatch.

exp int

Unix timestamp when the token expires.

CallbackTokenError

Bases: Exception

Base error for callback-token minting/unsealing (receiver side).

CallbackTokenExpiredError

Bases: CallbackTokenError

The token's exp has passed; it can no longer be used for routing.

CallbackTokenInvalidError

Bases: CallbackTokenError

The token is malformed, invalidly signed, or missing required claims.

mint_callback_token
mint_callback_token(
    thread_id,
    assistant_id,
    dispatch_id,
    secret,
    *,
    ttl_seconds=86400,
    now=None,
)

Mint a signed, opaque callback token (HS256).

Parameters:

Name Type Description Default
thread_id
str

Supervisor thread id to encode.

required
assistant_id
str

Supervisor assistant id to encode.

required
dispatch_id
str

Correlation id for the dispatch.

required
secret
str

CALLBACK_TOKEN_SECRET -- never sent to the subagent side.

required
ttl_seconds
int

Lifetime of the token. Default 24h.

86400
now
int | None

Override the current time (tests).

None

Returns:

Type Description
str

The signed JWT string, safe to hand to the subagent as an opaque

str

value in PushNotificationConfig.token.

unseal_callback_token
unseal_callback_token(token, secret)

Verify and decode a callback token signed with secret.

Parameters:

Name Type Description Default
token
str

The opaque JWT echoed back by the subagent.

required
secret
str

The same CALLBACK_TOKEN_SECRET used at mint time.

required

Returns:

Type Description
CallbackToken

The decoded routing payload.

Raises:

Type Description
CallbackTokenExpiredError

Token exp has passed.

CallbackTokenInvalidError

Bad signature, malformed token, or a missing required claim.

Settings and security modes

Both deployments' configuration and the A2A_VERIFY_MODE ladder.

langshark_bites.a2a_completion_notifier.settings

Runtime configuration for the a2a_completion_notifier bite.

Why this exists

The completion notifier spans two deployments that must not share credentials, so configuration is split into two dataclasses:

  • ReceiverSettings -- the supervisor deployment. Owns the callback token secret, the supervisor Agent Server URL + API key, the push webhook URL it places in PushNotificationConfig, and (for fetching full results) the subagent Agent Server URL + API key.
  • EmitterSettings -- the subagent deployment. Owns the RS256 private key + kid used to sign the outgoing A2A push JWT, and the retry/backoff policy for the webhook POST.

Environment variable conventions mirror api_rate_limiter: uppercase, _-separated, prefixed with A2A_, resolved by from_env() classmethods.

Usage
from langshark_bites.a2a_completion_notifier.settings import ReceiverSettings

receiver = ReceiverSettings.from_env()
A2AVerifyMode

Bases: StrEnum

Sender-verification mode for the receiver, set by the supervisor.

Ordered by strictness; DEV < VERIFY < STRICT.

  • dev -- trusted private network. Only the callback token is checked; sender JWTs are accepted unverified.
  • verify -- verify-if-signed. A signed notification is verified and fails closed (401). An unsigned one is accepted (logged), so a fleet mid-migration tolerates not-yet-signed senders.
  • strict -- signing is mandatory. Unsigned -> 401; signed but unverifiable -> 401; sender-auth config is a startup requirement.
EmitterSettings dataclass
EmitterSettings(
    private_key_pem="",
    kid="subagent",
    issuer="",
    audience="",
    jti_ttl_seconds=300,
    timeout_seconds=30.0,
    max_retries=3,
    retry_backoff_seconds=1.0,
)

Subagent-side emitter (push-notification signer) configuration.

The per-dispatch push config (webhook url + opaque callback token) is not here -- it arrives at runtime via the supervisor's PushNotificationConfig in config.configurable. This dataclass only carries the deployment-level signing material and retry policy.

from_env classmethod
from_env()

Read settings from A2A_* environment variables.

private_key_pem may be supplied inline (A2A_PRIVATE_KEY_PEM) or as a path on disk (A2A_PRIVATE_KEY_FILE); the file wins.

ReceiverSettings dataclass
ReceiverSettings(
    supervisor_url="http://localhost:8000",
    supervisor_api_key="",
    callback_token_secret="",
    receiver_url="http://localhost:8001",
    subagent_issuer="",
    subagent_jwks_url="",
    subagent_url="",
    subagent_api_key="",
    receiver_port=8001,
    redis_url="",
    jti_ttl_seconds=900,
    iat_staleness_seconds=300,
    callback_token_ttl_seconds=86400,
    mode=DEV,
)

Supervisor-side receiver configuration.

A2A_VERIFY_MODE selects the sender-verification posture (dev / verify / strict); see :class:A2AVerifyMode. redis_url enables Redis-backed jti deduplication (atomic across receiver replicas). The dead-letter tool is read-only today -- nothing parks notifications in the queue yet (the receiver logs a2a_forward_failed / a2a_forward_wake_failed on undelivered deliveries instead of parking them). All other fields except the numeric/timeouts are read from A2A_* env vars with the defaults shown.

ensure_valid
ensure_valid()

Fail fast on sender-auth misconfiguration for the chosen mode.

Raises:

Type Description
ValueError

strict without sender-auth config, or verify with partially populated sender-auth config.

from_env classmethod
from_env()

Read settings from A2A_* environment variables.

Env mapping:

============================= ================================== Env var Field ============================= ================================== A2A_SUPERVISOR_URL supervisor_url A2A_SUPERVISOR_API_KEY supervisor_api_key A2A_CALLBACK_TOKEN_SECRET callback_token_secret A2A_RECEIVER_URL receiver_url A2A_SUBAGENT_ISSUER subagent_issuer A2A_SUBAGENT_JWKS_URL subagent_jwks_url A2A_SUBAGENT_URL subagent_url A2A_SUBAGENT_API_KEY subagent_api_key A2A_RECEIVER_PORT receiver_port A2A_REDIS_URL redis_url A2A_JTI_TTL_SECONDS jti_ttl_seconds A2A_IAT_STALENESS_SECONDS iat_staleness_seconds A2A_CALLBACK_TOKEN_TTL callback_token_ttl_seconds A2A_VERIFY_MODE mode ============================= ==================================

Supervisor consumption

Build an MCP server's tools as reconnectable native LangChain tools for the supervisor graph. Uses a short-lived discovery scope so tools reconnect per invocation (never hold the adapter across runs).

langshark_bites.a2a_completion_notifier.supervisor_tools

Supervisor-side helper exposing an MCP server's tools to LangChain.

Why this exists

A supervisor deployment consumes an MCP server's tools as native LangChain tools via langchain.mcp.MCPAdapter -- but the adapter must be used only for discovery, not held across runs.

as_langchain_tool builds tools that "open the client themselves whether or not a connection is already held elsewhere": each invocation does its own async with client: (fastmcp reentrancy counter 0→1→0). If the adapter's client is left connected across runs (nesting_counter stuck at 1), the second tool call crashes with fastmcp's "Internal error: nesting counter should be 0 when starting new session, got 1".

This helper is the canonical, regression-tested way to build the supervisor tools: open the adapter for a short discovery scope, return the tools, close the adapter. Each tool reconnects per invocation.

Blockbuster-safe pre-warm

mcp.client.session.call_tool lazily imports jsonschema (via validate_tool_result) on the first tool call, and that import builds jsonschema_specifications' schema registry with a synchronous directory scan. Inside langgraph dev's ASGI event loop, blockbuster intercepts that scan and fails the run ("Blocking call to ScandirIterator.next").

The scan must therefore run at graph-module import time, before the event loop (and with it blockbuster) is active -- the same pattern langgraph_api itself uses for ddtrace (see langgraph_api.graph). Call prewarm_supervisor_tools() at the top of the graph module that will host the supervisor tools (import it at module scope, not inside a coroutine):

# src/.../supervisor.py
from langshark_bites.a2a_completion_notifier.supervisor_tools import (
    prewarm_supervisor_tools,
)
prewarm_supervisor_tools()

Do NOT call it from inside create_supervisor_tools or any other coroutine: under a running loop the synchronous scan trips blockbuster regardless of which thread the loop lives in.

create_supervisor_tools async
create_supervisor_tools(target, *, cache_mode='use')

Discover an MCP server's tools as reconnectable LangChain tools.

Opens langchain.mcp.MCPAdapter for a short discovery scope and closes it before returning. Each returned tool holds the client and reconnects per invocation (the as_langchain_tool contract), so consecutive calls -- including to the SAME tool on the SAME server -- cannot trip fastmcp's reentrancy-nesting crash.

The mcp SDK's lazy jsonschema import must already be warmed BEFORE any coroutine runs: call prewarm_supervisor_tools() at module import (see the module docstring). This function deliberately does not import jsonschema itself -- inside a running loop the synchronous scan would trip blockbuster.

Parameters:

Name Type Description Default
target
Any

MCP target accepted by langchain.mcp.MCPAdapter -- a URL string (http(s)://.../mcp) or an mcpServers dict for stdio / streamable-http.

required
cache_mode
Literal['use', 'refresh', 'bypass']

How discovery reads the client-side response cache (SEP-2549). Defaults to "use" so a configured cache is honored, matching MCPAdapter.list_tools.

'use'

Returns:

Type Description
list[Any]

The server's tools as LangChain BaseTool objects.

Raises:

Type Description
ImportError

If langchain.mcp is not installed (add the langchain[mcp] extra).

prewarm_supervisor_tools
prewarm_supervisor_tools()

Import everything the MCP tool path touches lazily, at module import.

langgraph dev's blockbuster rejects blocking file IO inside any running loop ("Blocking call to ..."). Several lazy imports on the MCP discovery / tool-call path would otherwise do a blocking read at the worst time:

  • mcp.client.session.validate_tool_result lazily imports jsonschema, whose jsonschema_specifications registry scans a directory;
  • langchain.mcp (first import) pulls fastmcp + pydantic_settings, which read .env;
  • fastmcp's HTTP transport lazily imports httpx2/httpcore2 at connect time, whose importlib.metadata read pulls a package METADATA file.

Call this at the top level of the graph module that hosts the supervisor's MCP tools -- before the event loop / blockbuster exist -- so every one of those imports is a no-op cache hit during discovery and per-call reconnects. Exactly like langgraph_api.graph eagerly imports ddtrace. Never call it from inside a coroutine (a running loop makes the scan fatal).