fix: async-safe embeddings and resilient drain_writes - #5702
MatthiasHowellYopp wants to merge 16 commits into
Conversation
d425472 to
c1a3805
Compare
c1a3805 to
6df71c1
Compare
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughThe PR normalizes embeddings (including bytes and NumPy arrays), makes batch embedding calls safe when an asyncio loop is active by offloading to a thread pool with a 30s timeout, and replaces indefinite background-write blocking with per-save timeouts and logging. Tests covering embedding normalization and async-safe embedding behavior were added. ChangesEmbedding Normalization and Async-Safe Embedding
🎯 3 (Moderate) | ⏱️ ~25 minutes
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
lib/crewai/src/crewai/memory/unified_memory.py (1)
323-383:⚠️ Potential issue | 🟠 Major | 🏗️ Heavy lift
drain_writestimeout resilience is undermined byclose()shutdown/close sequencing.After a per-save timeout, writes can still be in-flight, but
close()still doesstorage.close()andshutdown(wait=True). That can (1) race writers against a closed storage backend and (2) still block crew shutdown indefinitely.Suggested direction (bounded shutdown path)
def close(self) -> None: """Drain pending saves, flush storage, and shut down the background thread pool.""" self.drain_writes() - if hasattr(self._storage, "close"): - self._storage.close() - self._save_pool.shutdown(wait=True) + with self._pending_lock: + has_inflight = any(not f.done() for f in self._pending_saves) + + if has_inflight: + _logger.warning( + "[CLOSE] In-flight saves remain after drain timeout; " + "skipping blocking shutdown path to avoid hanging crew completion." + ) + self._save_pool.shutdown(wait=False, cancel_futures=True) + return + + self._save_pool.shutdown(wait=True, cancel_futures=True) + if hasattr(self._storage, "close"): + self._storage.close()🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@lib/crewai/src/crewai/memory/unified_memory.py` around lines 323 - 383, The close() sequence currently races with in-flight background saves because it calls storage.close() and then shuts down the pool unconditionally; change close() to first call drain_writes(...) with a bounded total timeout (e.g., timeout_per_save * max(1, len(pending)) or a configurable overall timeout) to allow background saves to finish, then stop blocking shutdown: after that attempt a non-blocking shutdown of the save pool (use _save_pool.shutdown(wait=False) or otherwise avoid indefinite blocking), attempt to cancel any remaining futures in self._pending_saves, and only call _storage.close() after pending futures were awaited/cancelled (or after the non-blocking shutdown) to avoid racing writers against a closed storage backend; refer to drain_writes, close, _save_pool, _storage, and self._pending_saves when making the changes.
🧹 Nitpick comments (1)
lib/crewai/src/crewai/memory/unified_memory.py (1)
365-367: ⚡ Quick winInclude traceback on save failure logs for actionable diagnostics.
The error log only prints
str(e). Addexc_info=Trueso failures in drain mode are debuggable.Small logging improvement
- _logger.error( - "[DRAIN_WRITES] Save %d/%d failed: %s", i + 1, len(pending), e - ) + _logger.error( + "[DRAIN_WRITES] Save %d/%d failed: %s", + i + 1, + len(pending), + e, + exc_info=True, + )🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@lib/crewai/src/crewai/memory/unified_memory.py` around lines 365 - 367, The error log in UnifiedMemory's drain routine uses _logger.error("[DRAIN_WRITES] Save %d/%d failed: %s", i + 1, len(pending), e) which only prints e; update the call in unified_memory.py (inside the drain/save loop in the UnifiedMemory class/function) to pass exc_info=True to the logger so the full traceback is included (i.e., call _logger.error(..., exc_info=True)) ensuring failures during drain mode produce actionable stack traces.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@lib/crewai/src/crewai/memory/encoding_flow.py`:
- Around line 71-94: The validator ensure_embedding_is_list should stop
re-implementing bytes→list[float] coercion and instead delegate to
MemoryRecord's canonical normalization (MemoryRecord.validate_embedding) for
both single records and lists; update ensure_embedding_is_list to iterate over
similar_records when it's a list and replace each record.embedding with the
result of MemoryRecord.validate_embedding (or call the shared normalization
utility) rather than mutating raw bytes in-place, and for the single-record path
call the same MemoryRecord normalization so empty bytes behavior (b"" → None)
remains consistent with MemoryRecord.
---
Outside diff comments:
In `@lib/crewai/src/crewai/memory/unified_memory.py`:
- Around line 323-383: The close() sequence currently races with in-flight
background saves because it calls storage.close() and then shuts down the pool
unconditionally; change close() to first call drain_writes(...) with a bounded
total timeout (e.g., timeout_per_save * max(1, len(pending)) or a configurable
overall timeout) to allow background saves to finish, then stop blocking
shutdown: after that attempt a non-blocking shutdown of the save pool (use
_save_pool.shutdown(wait=False) or otherwise avoid indefinite blocking), attempt
to cancel any remaining futures in self._pending_saves, and only call
_storage.close() after pending futures were awaited/cancelled (or after the
non-blocking shutdown) to avoid racing writers against a closed storage backend;
refer to drain_writes, close, _save_pool, _storage, and self._pending_saves when
making the changes.
---
Nitpick comments:
In `@lib/crewai/src/crewai/memory/unified_memory.py`:
- Around line 365-367: The error log in UnifiedMemory's drain routine uses
_logger.error("[DRAIN_WRITES] Save %d/%d failed: %s", i + 1, len(pending), e)
which only prints e; update the call in unified_memory.py (inside the drain/save
loop in the UnifiedMemory class/function) to pass exc_info=True to the logger so
the full traceback is included (i.e., call _logger.error(..., exc_info=True))
ensuring failures during drain mode produce actionable stack traces.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro Plus
Run ID: e68de661-3cc4-4216-96e5-f6486226094e
📒 Files selected for processing (4)
lib/crewai/src/crewai/memory/encoding_flow.pylib/crewai/src/crewai/memory/types.pylib/crewai/src/crewai/memory/unified_memory.pylib/crewai/tests/memory/test_embedding_safety.py
6df71c1 to
a4332f6
Compare
There was a problem hiding this comment.
🧹 Nitpick comments (1)
lib/crewai/src/crewai/memory/types.py (1)
316-340: 💤 Low valueConsider applying the same async-context handling as
embed_texts.
embed_textsoffloads to the thread pool when called from an async context (lines 386-404), butembed_textcalls the embedder directly. Ifembed_textis ever called from an async context, it will block the event loop.Since the comment acknowledges this is deferred tech debt, consider at minimum adding
embed_textas a wrapper that delegates toembed_textsfor consistency:♻️ Optional refactor
def embed_text(embedder: Any, text: str) -> list[float]: ... if not text or not text.strip(): return [] - - # Just call the embedder directly - the blocking issue needs to be fixed - # at a higher level (making Memory.recall() async) - result = embedder([text]) - - if not result: - return [] - first = result[0] - if hasattr(first, "tolist"): - return list(first.tolist()) - if isinstance(first, list): - return [float(x) for x in first] - return list(first) + results = embed_texts(embedder, [text]) + return results[0] if results else []🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@lib/crewai/src/crewai/memory/types.py` around lines 316 - 340, embed_text currently calls the embedder directly and can block the event loop unlike embed_texts which offloads to a thread pool; update embed_text to delegate to embed_texts for consistent async-aware behavior by calling embed_texts(embedder, [text]) (or its synchronous wrapper that runs embed_texts via the same ThreadPoolExecutor path used in embed_texts), then extract the first embedding, handle empty results, and convert types the same way as the current embed_text logic (preserve checks for .tolist and list). Also update the docstring to note that embed_text delegates to embed_texts to avoid blocking in async contexts.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Nitpick comments:
In `@lib/crewai/src/crewai/memory/types.py`:
- Around line 316-340: embed_text currently calls the embedder directly and can
block the event loop unlike embed_texts which offloads to a thread pool; update
embed_text to delegate to embed_texts for consistent async-aware behavior by
calling embed_texts(embedder, [text]) (or its synchronous wrapper that runs
embed_texts via the same ThreadPoolExecutor path used in embed_texts), then
extract the first embedding, handle empty results, and convert types the same
way as the current embed_text logic (preserve checks for .tolist and list). Also
update the docstring to note that embed_text delegates to embed_texts to avoid
blocking in async contexts.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro Plus
Run ID: df324f71-c9e3-457c-b693-918b702a38e6
📒 Files selected for processing (4)
lib/crewai/src/crewai/memory/encoding_flow.pylib/crewai/src/crewai/memory/types.pylib/crewai/src/crewai/memory/unified_memory.pylib/crewai/tests/memory/test_embedding_safety.py
✅ Files skipped from review due to trivial changes (1)
- lib/crewai/src/crewai/memory/unified_memory.py
🚧 Files skipped from review as they are similar to previous changes (2)
- lib/crewai/src/crewai/memory/encoding_flow.py
- lib/crewai/tests/memory/test_embedding_safety.py
8b8a533 to
bbbe31e
Compare
There was a problem hiding this comment.
♻️ Duplicate comments (1)
lib/crewai/src/crewai/memory/types.py (1)
343-360:⚠️ Potential issue | 🟠 Major | 🏗️ Heavy liftThis path still blocks the event loop.
Line 393 waits with
Future.result(timeout=30)on the caller thread. Ifembed_texts()is reached from a coroutine, that caller thread is the event-loop thread, so the loop stays frozen until the embedder finishes or times out. Timed-out workers also keep running, so two hung calls can still exhaust_EMBED_POOLand turn later async reads into the same empty-embedding fallback. Please move the async case to an actual async helper (await loop.run_in_executor(...)/asyncio.to_thread(...)+asyncio.wait_for(...)) and keep this function sync-only.In Python asyncio, if synchronous code running on the event-loop thread calls concurrent.futures.Future.result(timeout=30) on a ThreadPoolExecutor future, does that block the event loop until completion or timeout? If the timeout expires, does the worker thread continue running, and what is the recommended pattern for offloading blocking work from async code?Also applies to: 384-399
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@lib/crewai/src/crewai/memory/types.py` around lines 343 - 360, The current embed_texts function offloads to _EMBED_POOL but blocks the event-loop by calling Future.result(timeout=30) on the event-loop thread; instead make embed_texts strictly synchronous and move the async-path into a new async helper (e.g., async_embed_texts) that uses await loop.run_in_executor(...) or asyncio.to_thread(...) combined with asyncio.wait_for(..., timeout=30) so the event loop is not blocked and timed-out worker threads continue running without freezing the loop; update callers so coroutines call async_embed_texts while sync callers keep calling embed_texts, and reference _EMBED_POOL, embed_texts, and the new async_embed_texts (or chosen helper name) when making the change.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Duplicate comments:
In `@lib/crewai/src/crewai/memory/types.py`:
- Around line 343-360: The current embed_texts function offloads to _EMBED_POOL
but blocks the event-loop by calling Future.result(timeout=30) on the event-loop
thread; instead make embed_texts strictly synchronous and move the async-path
into a new async helper (e.g., async_embed_texts) that uses await
loop.run_in_executor(...) or asyncio.to_thread(...) combined with
asyncio.wait_for(..., timeout=30) so the event loop is not blocked and timed-out
worker threads continue running without freezing the loop; update callers so
coroutines call async_embed_texts while sync callers keep calling embed_texts,
and reference _EMBED_POOL, embed_texts, and the new async_embed_texts (or chosen
helper name) when making the change.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro Plus
Run ID: 2eb45bad-c2d6-4985-9379-a115365a1142
📒 Files selected for processing (4)
lib/crewai/src/crewai/memory/encoding_flow.pylib/crewai/src/crewai/memory/types.pylib/crewai/src/crewai/memory/unified_memory.pylib/crewai/tests/memory/test_embedding_safety.py
🚧 Files skipped from review as they are similar to previous changes (3)
- lib/crewai/src/crewai/memory/unified_memory.py
- lib/crewai/src/crewai/memory/encoding_flow.py
- lib/crewai/tests/memory/test_embedding_safety.py
f9a7f08 to
2ddeaaf
Compare
|
@greysonlalonde hoping I can get a review on this. |
d2af721 to
e17a45d
Compare
c5dcef6 to
0683317
Compare
0683317 to
dbe8086
Compare
6342c7c to
6d3f769
Compare
Extract duplicated Redis URL parsing into a shared cache_config utility. Introduce ValkeyCache as a lightweight async key/value cache using valkey-glide. Wire it into A2A task handling, agent card caching, and file upload caching. Part 1/4 of Valkey storage implementation.
Sets CLIENT SETNAME to 'crewai_valkey' so connections are identifiable in CLIENT LIST and monitoring tools (Valkey Admin, CloudWatch). Signed-off-by: Matthias Howell <matthias.howell@improving.com>
Set client_info_tag="crewai" on the Glide client so the connection reports lib-name GlidePy(crewai) via CLIENT INFO, letting operators attribute Valkey usage to CrewAI. Bumps the valkey-glide floor to >=2.5.2, where client_info_tag was introduced. client_name (CLIENT SETNAME) is unchanged; the two are independent fields. Metadata only, no behavioural change. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
- crewai-files no longer hard-imports crewai at module load. parse_cache_url moves to a lazy import inside the Valkey backend path, so importing crewai_files without crewai installed works again (the default in-memory cache never needs it). Preserves the crewai[file-processing] -> crewai-files direction. Adds regression tests. - agent_card: pin the AgentCard PickleSerializer @cached to an in-process SimpleMemoryCache. Previously it used the default alias, which A2A wires to VALKEY_URL/REDIS_URL — pickling cache values to a network Redis/Valkey is a code-execution surface on cache read if that store is writable. Pickle now never deserializes from an operator-controlled backend. - ValkeyCache/ValkeyCacheBackend/parse_cache_url now carry use_tls (from rediss/valkeys schemes) so cache connections honor TLS instead of silently opening a plaintext connection to a managed endpoint. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
6d3f769 to
90be3a4
Compare
The agent-card pickle cache is pinned to its own SimpleMemoryCache via the @cached decorator, so it no longer needs the process-wide aiocache default alias configured. The lazy _ensure_cache_configured() helper still called caches.set_config(get_aiocache_config()) on every fetch, which replaced the default alias that the in-memory A2A cancel path (poll_for_cancel_aiocache -> caches.get('default')) depends on. A running poller and a later cancel() could then hold different SimpleMemoryCache instances and never see the cancel flag. Remove the helper, its call, and the now-unused imports. Add regression tests asserting the helper is gone and the default alias is stable. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
90be3a4 to
e7d1089
Compare
_ensure_task_cache built ValkeyCache without use_tls, so a VALKEY_URL with a TLS scheme (valkeys:// or rediss://) opened a plaintext connection and A2A cancel/cancel-watch failed against a TLS-only Valkey. parse_cache_url already extracts use_tls from the scheme; thread it through. Lives on the base branch where the A2A task cache is introduced so it applies to the whole stack. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
e7d1089 to
a35c329
Compare
Setting VALKEY_URL routes A2A task cancellation through ValkeyCache, which imports glide at module load. Without the optional valkey extra installed, that raised a bare ImportError from every cancellable A2A task. Catch it in _ensure_task_cache and re-raise with an actionable message pointing at 'pip install crewai[valkey]' (or unsetting VALKEY_URL). This is the fail-loud option. A graceful fallback to the aiocache Redis path is deliberately left for later, pending the broader third-party provider/config decision. Adds a regression test. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
a35c329 to
5032fa0
Compare
ValkeyCache cached a single GlideClient and asyncio.Lock on the instance. UploadCache's sync methods each call asyncio.run(), which creates then closes a fresh loop, so the second sync get/set reused a client and lock bound to a closed loop and failed. Track the loop the client was created on and drop the stale client/lock when _get_client runs on a different loop. Adds a regression test that drives two asyncio.run() calls. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
5032fa0 to
3d39f2f
Compare
- ValkeyCache._get_client dropped the stale GlideClient on an event-loop change without closing it, leaking a connection and reader task per asyncio.run() cycle (UploadCache sync path). Best-effort close the old client before rebinding. - get_aiocache_config never forwarded use_tls to aiocache.RedisCache, so a rediss:// / valkeys:// REDIS_URL opened a plaintext connection. Pass ssl=True when the URL scheme is TLS. Adds regression tests for both. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
Add None/empty input handling to RecallMemoryTool and RememberTool. Filter empty strings, convert string inputs to lists, and return descriptive error messages. Update tool descriptions in en.json to include explicit parameter examples. Part 2/4 of Valkey storage implementation.
Drop min_length=1 and widen queries/contents to list[str] | str | None on RecallMemorySchema/RememberSchema. min_length=1 rejected empty lists at schema validation, raising a Pydantic error before _run could return its guidance string — so ToolUsage retried the same bad payload. The schema now mirrors what _run accepts (empty lists, bare strings), letting _run own the validation and return a useful message. Adds schema-level tests. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
Add bytes→float validators on MemoryRecord and ItemState to handle Valkey returning embeddings as raw bytes. Make embed_texts() safe when called from an async context by using a thread pool. Improve drain_writes() with per-save timeouts and error logging instead of raising on failure. Part 3/4 of Valkey storage implementation.
…embed - unified_memory: _background_encode_batch caught every RuntimeError and returned [], silently dropping embedder/LLM init and other encode failures. Now only the executor-shutdown RuntimeError is swallowed; any other emits MemorySaveFailedEvent and propagates so memories don't disappear silently. - unified_memory: restore Memory default llm to gpt-5.4-mini (an unintended revert to gpt-4o-mini diverged from the rest of CrewAI); fix the matching model name in the init-failure error message. - types.embed_texts: detect the running loop before invoking the embedder so an embedder RuntimeError can no longer be misread as "no running loop" and retried on the event-loop thread (double-invocation). - Add regression test for generic background RuntimeError. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
…h remote on clear - _background_encode_batch no longer emits MemorySaveFailedEvent itself; it re-raises so the future carries the error and _on_save_done emits the event exactly once. Previously both paths fired, double-reporting each failure. - drain_writes now cancels and untracks a save that exceeds timeout_per_save. It was left in _pending_saves, so every later recall() waited another full timeout and close()'s shutdown(wait=True) could block indefinitely. - UploadCache.aclear now flushes the whole backend namespace via a new backend clear() (ValkeyCache.clear scans+deletes by prefix; AiocacheBackend delegates to cache.clear). It previously deleted only locally-tracked keys, leaving remote entries from other processes/runs behind. Tests updated for single-emit; added drain-timeout untrack test. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
3d39f2f to
c39830f
Compare
embed_texts returned empty vectors on its 30s timeout, so EncodingFlow treated timed-out items as new and persisted them with a missing/zero embedding - unsearchable memories that pollute the index with no signal that encoding failed. Add raise_on_timeout (default False to preserve recall's tolerate-a-miss behavior); the save-path callers in EncodingFlow (batch_embed and the update path) pass raise_on_timeout=True so a timed-out batch raises TimeoutError and fails the save instead. Adds regression tests for both timeout modes. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes using high effort and found 1 potential issue.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 9df13fb. Configure here.
raise_on_timeout only fired in the async-context branch (thread pool + .result(timeout=30)). EncodingFlow's batch_embed/execute_plans run via asyncio.to_thread, so embed_texts is invoked with no running loop and took the direct embedder(...) branch with no timeout — the fail-loud save protection never triggered on exactly the paths that pass raise_on_timeout=True. Route the no-loop branch through the pool with the 30s timeout when raise_on_timeout is set. Adds a no-event-loop timeout regression test. Signed-off-by: MatthiasHowellYopp <matthias.howell@improving.com>

Description:
Part 3/4 of adding Valkey as a storage backend for CrewAI. This PR makes the embedding and memory persistence paths robust enough to work with async storage backends like Valkey.
What changed:
types.py — Added a field_validator on MemoryRecord.embedding that converts bytes to list[float] via numpy. Valkey stores vectors as raw bytes, so this ensures embeddings are always in the expected format regardless of storage backend. Also added a thread pool to embed_texts() so it doesn't block the event loop when called from an async context.
encoding_flow.py — Added a matching field_validator on ItemState.similar_records and result_record to handle the same bytes→float conversion during the consolidation flow.
unified_memory.py — drain_writes() now accepts a timeout_per_save parameter (default 60s) and logs warnings on timeout or failure instead of raising. This prevents a single slow or failed save from blocking crew completion. Added structured debug/warning/error logging throughout the drain cycle.
Testing:
test_embedding_safety.py (15 tests) — Covers bytes→float conversion, empty bytes, numpy arrays, int-to-float coercion, sync and async embed_texts behavior, empty/whitespace input handling.
Summary by CodeRabbit
Bug Fixes
Tests
Note
High Risk
Changes distributed A2A cancellation and upload caching, adds Valkey connectivity, and alters memory persistence/embedding failure modes—areas that affect multi-process correctness and data durability.
Overview
Introduces an optional Valkey distributed cache path (
VALKEY_URL,pip install 'crewai[valkey]') via a new JSON-backedValkeyCache, sharedcache_configURL parsing, and wiring into A2A task cancellation (polling when Valkey is used) and upload cache (cache_type="valkey"with a pluggable backend and lazycrewaiimport socrewai-filesstill loads withoutcrewai).Security / isolation: Agent-card fetch caching is pinned to in-process
SimpleMemoryCacheso pickled objects are never read from operator-controlled Redis/Valkey, and the global aiocachedefaultalias is no longer reconfigured by agent-card code (fixes cancel-flag loss across instances).Memory & embeddings:
MemoryRecord(and encoding-flow state) normalize bytes embeddings tolist[float]for Valkey-backed storage.embed_texts()runs embedders on a thread pool in async contexts, with optional 30s timeout andraise_on_timeout=Trueon save paths so timeouts fail instead of storing empty vectors.drain_writes()gets per-save timeouts, logging, and untracks stuck futures; non-shutdown save failures propagate and emit a single failure event.Agent tooling: Recall/save memory tools accept flexible inputs with clearer errors; i18n prompts require explicit
queries/contentsJSON.Tests cover Valkey cache, cache config, embedding safety, drain timeouts, A2A Valkey-missing-extra, upload-cache lazy import, and agent-card cache stability.
Reviewed by Cursor Bugbot for commit 2c7de49. Bugbot is set up for automated code reviews on this repo. Configure here.