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:
- Unseal the callback token: the routing invariant, enforced in every mode.
- Sender verification, dispatched by
A2A_VERIFY_MODE. - Claim the dedup key exactly once (JWT
jti, else the payloadid). - Accept only terminal states (
completed/failed/ ...). taskIdinside the JWT must match the payloadid(signed only).- 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=Noneis a silent no-op in some langgraph-api versions (no node executes, so thebefore_modeldrain never fires and the supervisor never wakes). - The minimal wake message makes the graph run. The
MailboxDrainMiddlewareinjects 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. Overridingwake_inputwith 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):
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_failedat 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_agentfires once per invocation, unconditionally, after the agent completes -- exactly when a completion webhook must go out.awrap_model_callcatches 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 |
|---|---|---|---|
|
PushNotificationConfig
|
Webhook target + opaque callback token for this run. |
required |
|
A2ASigner | None
|
RS256 signer for the subagent deployment. |
None
|
|
PushClient | None
|
HTTP client that POSTs the signed notification. |
None
|
|
Callable[[Any], str | None] | None
|
Optional callable mapping the hook |
None
|
aafter_agent
async
¶
Fire the 'completed' notification after the agent finishes.
abefore_agent
async
¶
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
¶
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
¶
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 |
|
token |
Opaque callback token, echoed back uninterpreted in the
notification |
|
authentication |
Optional auth mapping
(e.g. |
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 |
|---|---|---|---|
|
Mapping[str, Any] | None
|
The runnable |
required |
|
A2ASigner | None
|
Optional deployment signer. When omitted, notifications are
sent unsigned -- accepted by |
None
|
|
PushClient | None
|
Delivery client. |
None
|
|
Callable[[Any], str | None] | None
|
Override for the A2A task id resolver. The default
resolves to |
None
|
Raises:
| Type | Description |
|---|---|
PushNotificationConfigError
|
If no |
extract_push_config
¶
Read config.configurable["a2a_push_config"], if present.
Returns:
| Type | Description |
|---|---|
PushNotificationConfig | None
|
A parsed |
PushNotificationConfig | None
|
config has no |
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
¶
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 |
audience |
str
|
The |
jti_ttl_seconds |
int
|
Lifetime of each signed JWT ( |
jwks
¶
Export the public key set for this signer.
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
A JWKS document ( |
dict[str, Any]
|
RSA public key keyed by |
dict[str, Any]
|
JWKS endpoint (e.g. |
sign
¶
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 |
required |
now
¶ |
int | None
|
Override the current time (tests). |
None
|
Returns:
| Type | Description |
|---|---|
str
|
A signed RS256 JWT with a unique |
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
¶
POSTs signed A2A notifications to a webhook URL with backoff.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
|
AsyncClient | None
|
Optional client to reuse. When omitted, one is created and
owned by this instance (closed by :meth: |
None
|
|
PushClientSettings | None
|
Retry/backoff settings. Defaults to |
None
|
send
async
¶
Deliver one notification; retry on failure with backoff.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
url
¶ |
str
|
The webhook URL from |
required |
bearer
¶ |
str | None
|
The signed sender JWT (goes in the Authorization header).
|
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
¶
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
( |
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_push_config(
*,
thread_id,
assistant_id,
dispatch_id,
receiver_url,
callback_token_secret,
callback_token_ttl_seconds=86400,
mode=DEV,
now=None,
)
Build the config.configurable push config for one subagent.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
|
str
|
Supervisor thread the completion must be routed to. |
required |
|
str
|
Supervisor assistant to wake when the subagent finishes. |
required |
|
str
|
Correlation id for the whole fan-out dispatch. |
required |
|
str
|
Public base URL of the supervisor's receiver (the
|
required |
|
str
|
|
required |
|
int
|
Lifetime of the callback token. |
86400
|
|
A2AVerifyMode
|
The supervisor's |
DEV
|
|
int | None
|
Override the current time (tests). |
None
|
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
A dict suitable for |
dict[str, Any]
|
subagent's |
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):
- Unseal the callback token -- the routing invariant, enforced in every
A2A_MODE(dev/verify/strict). - Sender verification, dispatched by mode:
devaccepts signed and unsigned unverified;verifyverifies signed notifications (failing closed) while accepting unsigned ones;strictrejects unsigned and always verifies. - Claim the dedup key exactly once (JWT
jtiwhen signed, else the payloadid). - Accept only terminal states (
completed/failed/ ...). taskIdinside the JWT must match the payloadid(signed only).- ACK with
202(BackgroundTasks), then mailbox.write + wake. - The sender
issis 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
¶
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 aRedisJtiStoreso cross-replica dedup is atomic.jwks_client-- fetches fromsettings.subagent_jwks_url.mailbox--NotificationMailboxoverlg.result_settings-- defaults toResultFetchSettingsfromsettings.subagent_url/subagent_api_key; used by thePOST /a2a/resultsister primitive.result_http-- the outbound HTTP client for full-result fetches.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
|
ReceiverSettings
|
Supervisor-side receiver configuration. |
required |
|
Any | None
|
LangGraph SDK client (supervisor). Injected for tests. |
None
|
|
JtiStore | None
|
Single-use |
None
|
|
JWKSClient | None
|
Source of the subagent deployment's public keys. |
None
|
|
NotificationMailbox | None
|
Delivery sink for accepted notifications. |
None
|
|
Any | None
|
Input for the wake-up run created after a mailbox write.
Defaults to |
DEFAULT_WAKE_INPUT
|
|
ResultFetchSettings | None
|
Fetch configuration for the subagent deployment. |
None
|
|
AsyncClient | None
|
HTTP client used to fetch completed task results. |
None
|
Returns:
| Type | Description |
|---|---|
FastAPI
|
A configured FastAPI application. |
process_notification
async
¶
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]
|
|
Callable[[], Awaitable[None]] | None
|
callable that must run after the |
tuple[dict[str, Any], Callable[[], Awaitable[None]] | None]
|
(writes the mailbox, then wakes an idle supervisor). |
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
¶
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]
|
|
Raises:
| Type | Description |
|---|---|
ReceiverError
|
400 for a missing |
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¶
JWKSClientcaches fetched keys bykidand re-fetches only on an unknownkid(key rotation support without a long-lived cache).verify_sender_jwtraises typedSenderAuthErrorsubclasses so the FastAPI layer can map each to a 401 without catching bare exceptions.- The
iatstaleness window is checked separately from JWTexp: PyJWT validatesexpbut not freshness, so a 24hexpwith a 10-minute-oldiatmust 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
¶
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.
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_sender_jwt(
authorization,
*,
jwks_client,
audience,
issuer,
iat_staleness_seconds=300,
)
Verify the A2A sender's JWT and return its validated claims.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
|
str | None
|
The raw |
required |
|
JWKSClient
|
Source of the subagent deployment's public keys. |
required |
|
str
|
The expected |
required |
|
str
|
The expected |
required |
|
int
|
Max age of the |
300
|
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
The validated JWT claims (contains |
dict[str, Any]
|
|
Raises:
| Type | Description |
|---|---|
MissingBearerTokenError
|
No Bearer token in the header. |
UnknownKidError
|
No JWKS key matches the JWT's |
InvalidSignatureError
|
Signature or required-claim verification failed. |
ExpiredTokenError
|
The JWT |
StaleTokenError
|
The JWT |
MissingJtiError
|
The JWT lacks the |
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
¶
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 |
bool
|
redelivery/duplicate. |
JtiStore
¶
Bases: Protocol
A single-use claim store for A2A notification jti values.
RedisJtiStore
¶
Redis-backed jti claim store using atomic SET NX EX.
Attributes:
| Name | Type | Description |
|---|---|---|
ttl_seconds |
How long a claimed |
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 |
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:
- Writes the completion into the LangGraph
Storeunder("notifications", thread_id)keyed bytask_id(write first); - 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 itsbefore_modelmiddleware 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 |
|---|---|---|---|
|
Any
|
A LangGraph SDK client for the supervisor Agent Server
( |
required |
|
tuple[str, ...]
|
Leading tuple of a Store namespace.
The mailbox appends |
(NOTIFICATION_NAMESPACE_PREFIX,)
|
|
Any | None
|
The trigger-run input. Defaults to
:data: |
DEFAULT_WAKE_INPUT
|
|
str
|
Strategy when a run is already queued.
|
'enqueue'
|
deliver
async
¶
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 ( |
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 |
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:
- resolves the current
thread_idfromruntime.execution_info; - reads the mailbox namespace (
("notifications", thread_id)); - 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);
- returns a single
HumanMessagecarrying 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 |
|---|---|---|---|
|
tuple[str, ...]
|
Store namespace prefix. The thread id is appended
to form the full mailbox namespace -- must match the receiver's
mailbox configuration (default |
MAILBOX_NAMESPACE_PREFIX
|
|
RenderNotice | None
|
Optional callable mapping |
None
|
|
int
|
Cap on notices drained per model call. |
50
|
|
StoreLike | None
|
Optional store handle override. When omitted, the runtime's
|
None
|
abefore_model
async
¶
Notify (no-op) if nothing is pending, else return a state update.
Raises:
| Type | Description |
|---|---|
MailboxDrainError
|
The mailbox |
render_notice
¶
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 |
|---|---|---|---|
|
str
|
The A2A task id (Store key). |
required |
|
dict[str, Any]
|
The stored notice dict ( |
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
¶
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_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 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 |
|---|---|---|---|
|
str
|
The notification's |
required |
|
ResultFetchSettings
|
Fetch configuration (subagent URL + optional API key). |
required |
|
AsyncClient
|
HTTP client used for the outbound request. |
required |
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
The subagent's thread state (Agent Protocol |
dict[str, Any]
|
JSON body). |
Raises:
| Type | Description |
|---|---|
TaskResultError
|
If no |
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 a sparse A2A Task body for the receiver webhook.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
|
str
|
A2A task id ( |
required |
|
str
|
App label ( |
required |
|
str
|
Opaque callback token from |
required |
|
str | None
|
Optional A2A |
None
|
|
str | None
|
Optional short summary stored under |
None
|
|
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
¶
Return the opaque callback token from metadata.token.
extract_summary
¶
Return optional summary from metadata.summary or legacy top-level.
extract_terminal_state
¶
Return the app terminal state label, or None if non-terminal/unknown.
Accepts A2A enum names and legacy short labels.
normalize_notification_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 and validate payload as an A2A Task.
Raises:
| Type | Description |
|---|---|
NotificationPayloadError
|
Not a valid Task message. |
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
¶
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 |
|---|---|---|---|
|
str
|
Supervisor thread id to encode. |
required |
|
str
|
Supervisor assistant id to encode. |
required |
|
str
|
Correlation id for the dispatch. |
required |
|
str
|
|
required |
|
int
|
Lifetime of the token. Default 24h. |
86400
|
|
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 |
unseal_callback_token
¶
Verify and decode a callback token signed with secret.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
|
str
|
The opaque JWT echoed back by the subagent. |
required |
|
str
|
The same |
required |
Returns:
| Type | Description |
|---|---|
CallbackToken
|
The decoded routing payload. |
Raises:
| Type | Description |
|---|---|
CallbackTokenExpiredError
|
Token |
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 inPushNotificationConfig, and (for fetching full results) the subagent Agent Server URL + API key.EmitterSettings-- the subagent deployment. Owns the RS256 private key +kidused 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
¶
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
¶
Fail fast on sender-auth misconfiguration for the chosen mode.
Raises:
| Type | Description |
|---|---|
ValueError
|
|
from_env
classmethod
¶
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 |
|---|---|---|---|
|
Any
|
MCP target accepted by |
required |
|
Literal['use', 'refresh', 'bypass']
|
How discovery reads the client-side response cache
(SEP-2549). Defaults to |
'use'
|
Returns:
| Type | Description |
|---|---|
list[Any]
|
The server's tools as LangChain |
Raises:
| Type | Description |
|---|---|
ImportError
|
If |
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_resultlazily importsjsonschema, whosejsonschema_specificationsregistry scans a directory;langchain.mcp(first import) pullsfastmcp+pydantic_settings, which read.env;- fastmcp's HTTP transport lazily imports
httpx2/httpcore2at connect time, whoseimportlib.metadataread 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).