Releases: eggai-tech/EggAI
Releases · eggai-tech/EggAI
Release list
v0.6.0
What's Changed
Added
- Redis
subscribe(..., delete_on_ack=True): every entry the subscription
acks is alsoXDEL'd (sameMULTI), so consumed messages stop occupying
Redis memory instead of waiting forMAXLENtrimming. Covers handler
success (including filtered-out messages), the retry stream, and the
reclaimer's moves to the retry stream / DLQ. A failing handler's entry
stays in the PEL under the defaultNACK_ON_ERROR, and is deleted under
AckPolicy.ACK/ACK_FIRST(which ack it anyway). Opt-in because it
assumes a single consumer group per stream:XDELremoves the entry for
every group.connect()raisesRuntimeErrorif another group already
reads the stream; if one joins later, the group monitor turns deletion off
for that stream (plainXACK) and logs an error.
Fixed
- MCP adapter reads tool schemas from fastmcp's own
Toolinstead of the
protocolTool, so it works on fastmcp 3 and 4 without touching the
fields MCP SDK 2 renamed. Themcpextra now allowsfastmcp>=3,<5. - MCP adapter subscriptions declare
data_type, so handlers receive parsed
request models instead of raw dicts. Covered by a new end-to-end test. setup_tracingtype errors: exporter selection is typed,endpointis
passed explicitly, and the unused provider reassignment is gone.- A2A adapter type errors:
A2AConfigimportsSecurityScheme
unconditionally, and the executor no longer falls back to enqueuing a raw
dict the event queue cannot model. - CI lint and test jobs install all extras, so the tracing, a2a and mcp
modules are type-checked and tested.
Installation
pip install eggai==0.6.0v0.5.0
What's Changed
Migration note for DLQ consumers. The shape of DLQ entries changes in
this release; the SDK API, retry-stream entries, handler-side
_retry_count / _original_message_id, stream key names and per-handler
trimming do not. If you read DLQ entries yourself (XRANGE, re-drive
scripts, monitoring) or read the dict passed to on_dlq:
- read the exhausted retry budget from
_dlq_retries, not_retry_count
(which is now"0"on every DLQ entry); - expect the additional
_dlq_source,_dlq_handler,_dlq_at,
_dlq_reasonkeys; - poison entries are no longer raw bytes: recover the original payload by
base64-decoding_dlq_raw_b64, andon_dlqreceives that decoded dict; - stream entries with no
__data__field (only possible from a non-eggai
producer) are no longer retried at all; they are dead-lettered as poison on
the first reclaim.
The on_dlq(body, msg_id, retry_count) signature and its retry_count
argument are unchanged.
Added
RedisTransport: newdlq_channelsubscribe option. Dead-letters into ONE
stream you name instead of the per-handler{channel}.{handler_suffix}.dlq,
so a single consumer can subscribe to every handler's (or every service's)
failures with a plain@agent.subscribe(channel=dlq, group_start="0").
Accepts aChannelor a topic name (namespaced likeChannel(name)) on
Agent.subscribe/Channel.subscribe; a full key on
RedisTransport.subscribe. Requiresretry_on_idle_msand a non-None
max_retries; rejects the subscribed channel and the handler's retry stream
as targets. The.retrystream stays per-handler (#225) — only the terminal
sink is shared. A shared DLQ is written withoutMAXLENby default
(retry_max_lenkeeps applying to retry streams and per-handler DLQs):
stream-wide trimming would let one writer's cap delete other writers'
unconsumed dead letters.RedisTransport(dlq_max_len=...)is the opt-in
hard ceiling for a shared DLQ. Default behaviour withoutdlq_channelis
unchanged.- Every DLQ entry now carries provenance in its JSON body, alongside the
existing_retry_count/_original_message_id:_dlq_source(the channel
key the handler subscribed to),_dlq_handler(handler suffix / consumer
group),_dlq_at(ISO-8601 UTC) and_dlq_reason("max_retries"or
"poison") and_dlq_retries(how many retries actually ran). Additive:
typed models ignore the extra keys, as they already do for_retry_count.
On an entry dead-lettered a second time (by a DLQ consumer that gave up)
the origin keys_dlq_source/_dlq_handler/_dlq_atkeep the first
failure, while_dlq_reason/_dlq_retriesdescribe the latest hop.
Changed
_retry_countis reset to"0"on the DLQ write (the exhausted budget
moves to_dlq_retries). A DLQ entry is also a first delivery to whatever
consumes the DLQ; carrying the exceeded count over meant a DLQ consumer with
its ownretry_on_idle_msgot zero retries and, with backoff, an escalated
first reclaim. Anything that read_retry_countoff a DLQ entry should read
_dlq_retries; theon_dlqcallback'sretry_countargument is unchanged.- Poison messages (envelopes the reclaimer cannot parse) are no longer copied to
the DLQ verbatim. They are wrapped in a well-formed envelope whose body holds
the_dlq_*fields,_original_message_id, and the original bytes as
_dlq_raw_b64, so a DLQ subscriber always receives a decodable JSON object
instead of raw bytes. Envelope headers are preserved when only the body was
unusable. Re-drive scripts that handled raw poison entries should read
_dlq_raw_b64. Consequentlyon_dlqnow receives that decoded dict for
poison entries too, instead of the raw{b"__data__": bytes}fields. Agent.subscribe/Channel.subscribereject adlq_channelstring that
already starts withEGGAI_NAMESPACE(e.g. a pastedchannel.get_name()),
which would otherwise be namespaced twice and route dead letters to an
unwatched<ns>.<ns>.dlq. Pass the bare topic name or aChannel.RedisTransport.subscriberejectsretry_on_idle_mswithout a consumer
group (handler_id=/group=): the reclaimer works on a group's PEL, so
without one every reclaim cycle just errored.Agent/Channelalways set a
group; this only affects direct transport callers.dlq_channelrejects any key ending in.retry, not just the handler's own
retry stream: every.retrystream is auto-consumed by some handler, so
dead-lettering into one would feed that handler's retry loop. The check runs
before the broker subscriber is registered.
Fixed
PendingReclaimer: a pending entry with a parseable envelope but a non-object
JSON body (e.g. a published list), a non-integer_retry_count, or no
__data__field at all (non-eggai producer) escaped the poison handling. The
first two raised outside the guarded parse and aborted the whole reclaim
cycle afterXCLAIM, head-of-line blocking every later pending entry; the
third ping-ponged between main and retry stream forever with a retry count
that could never grow. All three are now dead-lettered as poison; in
particular, entries without__data__are no longer retried at all.
Installation
pip install eggai==0.5.0v0.4.1
What's Changed
Fixed
BaseMessage.datahas its{}default back, now declared with
validate_default=True. 0.4.0 made the field required to stop typed
subclasses (BaseMessage[Order]) from silently accepting a missing payload;
validating the default achieves the same without breaking
BaseMessage(source=..., type=...)for the untyped base. Typed subclasses
still reject a missing or malformed payload with aValidationError.
Installation
pip install eggai==0.4.1v0.4.0
What's Changed
Added
- Static type checking with mypy: the
eggaipackage is type-checked in CI
(poetry run mypy, config inpyproject.toml);mypyis a new dev dependency. RedisTransport: newgroup_startsubscribe option ("$"default,"0"or
a stream id) chooses where a NEW consumer group starts reading, so a consumer
added to a channel that already carries traffic can pick up the existing
backlog. The transport now creates its consumer groups itself before the
broker starts (#260).
Fixed
- A2A executor: the JSON-serialisation fallback caught the non-existent
json.JSONEncodeError, so a failingjson.dumpsraisedAttributeErrorand
masked the real error. It now catches(TypeError, ValueError). - A2A executor: the "unknown skill" error path enqueued a raw dict instead of a
proper A2A message; a skill registered withdata_type=Noneraised
TypeErrorwhen building the message and now falls back toBaseMessage. - Typed subscriptions: a
data_typewhosetypefield has no default (e.g. raw
BaseMessage) silently dropped every message;subscribe()now raises
ValueErrorasking for a default discriminator. Agent.subscribe(): plugin-prefixed kwargs (e.g.a2a_*) on an agent whose
plugin was never initialised raised an opaqueKeyError; nowValueError.RedisTransport:agent.stop()could hang forever on Python 3.10/3.11 with
redis-py >= 8. redis-py 8 sends every command throughasyncio.wait_for
(socket_timeoutnow defaults to 5s), and on CPython < 3.12wait_forcan
swallow theCancelledErrorwhen the inner await completes concurrently, so
the reclaimer task survived its own cancellation. The reclaimer loop now exits
on a running flag as well as on cancellation.- RedisTransport retries with a non-default
EGGAI_NAMESPACE(#261): the
reclaimer, group monitor,.retryand.dlqstreams were keyed under
eggai.<ns>.<topic>while consumption and publishing used<ns>.<topic>, so
retry_on_idle_ms/max_retriesnever fired outside the default namespace.
The transport no longer adds its owneggai.prefix;Channelalready
namespaces the name exactly once. After upgrading, entries stuck in the PEL
under a custom namespace are retried immediately and dead-lettered once the
retry budget is exhausted. Stale emptyeggai.<ns>.*shadow keys can be
deleted. PendingReclaimer: the retry envelope's header length prefix counted
characters instead of bytes, truncating non-ASCII header values on every
retry delivery (#263).Agent.subscribe()without a channel listened on the literaleggai.channel
instead of<EGGAI_NAMESPACE>.channel, so under a custom namespace a bare
Channel().publish()never reached it (#264).
Changed
- Type-annotation cleanups across transports, channel, agent and hooks to
satisfy mypy (no behaviour change).Transport.subscribe()'s second
parameter is now namedhandlerin the abstract base and all transports. - BREAKING:
BaseMessage.datais now required. The generic base previously
declareddata: TData = Field(default_factory=dict). Pydantic does not
validate defaults, so an envelope missingdatavalidated cleanly on any
typed subclass (BaseMessage[Order]) withdataas a plain{}— the
handler then crashed on first attribute access, and underNACK_ON_ERROR
that single malformed message wedged the subscription. A missing or malformed
payload now fails validation, which typed subscriptions
(wrap_handler_with_filters) already treat as "not ours: skip and ack".
The concreteMessagekeeps its{}default, where the annotation really is
a dict. Migration: construct payload-less envelopes withMessage(or pass
data=explicitly). - Allow ruff 0.16 (
ruff >=0.14.4,<0.17) and exclude Markdown from ruff, which
now formats fenced code blocks by default. RedisTransport:last_idother than">"together with a consumer group
now raisesValueError. That combination never delivered new entries (Redis
returns only the consumer's own pending entries for an explicit id) and
hot-looped XREADGROUP; the earlier note about replaying a backlog with
last_id="0"was wrong. Usegroup_start="0"instead.
Installation
pip install eggai==0.4.0v0.3.4
What's Changed
Fixed
- Forward original args/kwargs to wrapped handler: 0.3.3 changed traced_handler's signature to *args/**kwargs but still
called handler(message) positionally
Installation
pip install eggai==0.3.4v0.3.3
What's Changed
Fixed
- Tracing wrapper handler dispatch:
traced_handlernow accepts*args/**kwargs,
fixingTypeError: got an unexpected keyword argumentfor handlers whose message
parameter isn't namedmessageunder FastStream 0.7's keyword-based dispatch.
Installation
pip install eggai==0.3.3v0.3.2
What's Changed
Fixed
- RedisTransport background clients: connection-resilience kwargs passed to the
transport (socket_timeout,socket_connect_timeout,socket_keepalive,
socket_keepalive_options,health_check_interval,retry_on_timeout,
retry_on_error,max_connections, and thessl_*options) are now forwarded
to both long-lived background clients — the PEL reclaimer and the consumer-group
monitor. Previously each created its client with no socket timeout or keepalive
regardless of the transport's settings, so a silently dropped connection (e.g.
cloud Redis failover or an idle-connection reaper) left its blocking reads hung
indefinitely with no way to recover — and for the monitor, that hang struck
during the very failover it exists to recover from.decode_responsesremains
pinned per client (Falsefor the reclaimer's binary passthrough,Truefor the
monitor's string commands) and cannot be overridden by callers.
Added
- RedisTransport: Exponential retry backoff for SDK-managed retries. New
subscribe()optionsretry_backoff_multiplier(default1.0= the previous
constant cadence),retry_backoff_max_ms(cap on the escalated delay), and
retry_backoff_jitter(spread retries across a worker fleet to avoid a
thundering herd). The PEL reclaimer now treats a failing message as due once it
has been idle forretry_on_idle_ms * (retry_backoff_multiplier ** retry_count)
(capped atretry_backoff_max_ms), so a repeatedly-failing message is retried
progressively less often (e.g. 30s → 60s → 120s …), giving an overloaded
downstream room to recover instead of being hammered on a fixed clock. Backoff
is opt-in and fully backward compatible: the defaultmultiplier=1.0reproduces
the existing fixedretry_on_idle_msspacing exactly.
Installation
pip install eggai==0.3.2v0.3.1
What's Changed
Fixed
- Redis & Kafka transports:
filter_by_message,data_type, and
filter_by_datasubscriptions raisedTypeErrorat subscribe time under
FastStream 0.7, which removed publisher/subscriber-level middlewares. Filtering
and typed-message handling are now applied in EggAI's own handler wrapper
(application code) rather than via FastStream subscriber middlewares — aligning
with FastStream 0.7's removal of subscriber/publisher middlewares. This keeps the
per-subscription filter logic independent of FastStream's middleware API (it is
the same approach the in-memory transport already uses). As part of this,
data_typesubscriptions on Redis/Kafka now deliver the typed model instance
to the handler (matching the in-memory transport and the documented behaviour),
rather than the raw dict. Invalid filter-option combinations now raise
ValueErrorinstead of silently dropping an option:data_typeand
filter_by_messageare mutually exclusive, andfilter_by_datarequires
data_type. These validations are consistent across the Redis, Kafka, and
in-memory transports.
Added
- RedisTransport: New
max_lenandretry_max_lenconstructor options to cap
Redis stream growth via approximate trimming (XADD ... MAXLEN ~).max_len
(defaultNone/unbounded) caps the producer/publish()path; it is opt-in
becauseMAXLENtrims the oldest entries by count regardless of ack state, so a
value belowthroughput × consumer-lagcan drop un-delivered messages.
retry_max_len(default10_000) caps the SDK-managed retry and DLQ streams,
bounding the blast radius of a runaway retry loop. This wires up the previously
documented-but-inertmax_lenknob.
Changed
- RedisTransport: The stream consumer name now defaults to a per-process-unique
value ({handler_id}-{hostname}-{pid}) while the consumer group still defaults
to the stablehandler_id. A fleet of workers running the same handler now shares
one group (Redis load-balances the stream across them) while each worker owns a
distinct slice of the pending-entries list — the competing-consumers pattern. The
auto-created retry-stream subscriber gets the same per-process-unique consumer, so
retried messages load-balance across a worker fleet too. Pass an explicit
consumer=to opt out.
Installation
pip install eggai==0.3.1v0.3.0
What's Changed
Security
- Bump dev dependency
pytestto^9.0.0(9.0.3), withpytest-asyncioto>=1.0,<2(1.4.0) for compatibility, to remediate GHSA-6w46-j5rx-g56g (insecure tmpdir handling). Test-only; no runtime impact.
Changed
- Bump
fastmcpfrom^2.14.0to^3.0.0(resolves to 3.4.0), which pulls inauthlib1.7.2. The MCP adapter now usesTool.to_mcp_tool()to read tool schemas, since fastmcp 3.xlist_tools()returnsFunctionToolobjects that expose schemas viaparameters/output_schemainstead ofinputSchema/outputSchema.
Added
- Distributed tracing via OpenTelemetry: Agents now propagate a shared
trace_idacross every message hop. Opt-in viasetup_tracing(); zero-cost when not configured.
Fixed
- RedisTransport (#225): Retry and DLQ stream keys are now per-handler
({channel}.{handler_suffix}.retry/.dlq) instead of per-channel.
Previously, multiple handlers subscribing to the same channel with different
consumer groups shared a single{channel}.retrystream, so one handler's
reclaimed failure would be redelivered to every handler on that channel
via its auto-created-retrygroup. This caused work amplification (N
replays per nack, where N is the number of consumer groups on the channel)
and visible duplicate side-effects in downstream systems with non-idempotent
handlers. - RedisTransport: A permanently-unparseable ("poison") retry message is now
routed to the DLQ — or dropped with an error log when no DLQ is configured —
instead of being re-queued to the retry stream forever. Its retry count can
never be incremented, so it would previously livelock the per-handler retry
stream. - RedisTransport: The auto-created retry-stream subscriber now always reads
only new entries (last_id=">") regardless of the main subscription's
last_id. Restarting withlast_id="0"to replay the main backlog no longer
also replays the retry stream's history on every restart. - RedisTransport: A
RedisTransportshared across multipleAgents no
longer spawns duplicate consumer loops or orphans its consumer-group monitor
task when each agent callsconnect(). Re-entrantconnect()now starts only
not-yet-running subscribers and reuses the existing monitor. - Tracing: Retry-stream consumption now emits spans with the retry stream as
messaging.destinationinstead of the original channel, so retry attempts can
be observed separately from main-stream processing. - RedisTransport:
subscribe()now rolls back the in-memory reclaimer configs
and stream-subscription entries it registered if retry-stream setup fails, so a
retriedstart()no longer accumulates orphaned registrations. Identical
stream-group subscriptions are also de-duplicated.
Changed
- RedisTransport:
retry_on_idle_msnow rejectsack_policy=AckPolicy.ACK
andAckPolicy.ACK_FIRSTwith aValueError. Those acknowledge messages even
when the handler fails, emptying the PEL the reclaimer scans, which would
silently disable retries and the DLQ. Use the defaultNACK_ON_ERROR.
Migration
- Breaking on-wire change. Existing inflight messages in
{channel}.retry
and{channel}.dlqstreams will not be picked up by upgraded consumers —
they will continue to exist but new retries will land in the per-handler
streams. Drain the legacy streams before deploying, orXADDany pending
entries into the new per-handler keys.
Installation
pip install eggai==0.3.0v0.2.14
What's Changed
Fixed
- RedisTransport: Auto-recover from NOGROUP errors when Redis loses streams
(restart without persistence, failover, memory eviction). A background stream
group monitor periodically ensures consumer groups exist viaXGROUP CREATE
withMKSTREAM, andPendingReclaimerManagernow recreates consumer groups
on NOGROUP errors instead of logging unhandled exceptions.
Changed
- Bump
faststreamfrom 0.6.5 to 0.6.6 — fixesxautoclaimcrash on Redis < 7.0. - Bump
a2a-sdkfrom 0.3.22 to 0.3.24. - Bump
pyasn1from 0.6.1 to 0.6.2 — fixes CVE-2026-23490.
Installation
pip install eggai==0.2.14