Skip to content

fix: async-safe embeddings and resilient drain_writes - #5702

Open
MatthiasHowellYopp wants to merge 16 commits into
crewAIInc:mainfrom
MatthiasHowellYopp:feat/valkey-3-embedding-safety
Open

MatthiasHowellYopp wants to merge 16 commits into
crewAIInc:mainfrom
MatthiasHowellYopp:feat/valkey-3-embedding-safety

Conversation

@MatthiasHowellYopp

@MatthiasHowellYopp MatthiasHowellYopp commented May 4, 2026 •

Copy link
Copy Markdown
Contributor

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

    • More robust embedding normalization (bytes, numeric arrays, empty values handled consistently).
    • Embedding calls are safe in async contexts and preserve positions for skipped/empty inputs.
    • Background memory saves now use per-save timeouts with progress logging and no indefinite blocking (defaults to 60s).
  • Tests

    • Added comprehensive tests covering embedding validation, sync/async embedding behavior, and timeout/save handling.

Review Change Stack


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-backed ValkeyCache, shared cache_config URL parsing, and wiring into A2A task cancellation (polling when Valkey is used) and upload cache (cache_type="valkey" with a pluggable backend and lazy crewai import so crewai-files still loads without crewai).

Security / isolation: Agent-card fetch caching is pinned to in-process SimpleMemoryCache so pickled objects are never read from operator-controlled Redis/Valkey, and the global aiocache default alias is no longer reconfigured by agent-card code (fixes cancel-flag loss across instances).

Memory & embeddings: MemoryRecord (and encoding-flow state) normalize bytes embeddings to list[float] for Valkey-backed storage. embed_texts() runs embedders on a thread pool in async contexts, with optional 30s timeout and raise_on_timeout=True on 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/contents JSON.

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.

@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch 3 times, most recently from d425472 to c1a3805 Compare May 7, 2026 20:00
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from c1a3805 to 6df71c1 Compare May 11, 2026 19:53
@coderabbitai

coderabbitai Bot commented May 11, 2026 •

Copy link
Copy Markdown

Note

Reviews paused

It 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 reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

The 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.

Changes

Embedding Normalization and Async-Safe Embedding

Layer / File(s) Summary
Embedding validation and normalization
lib/crewai/src/crewai/memory/types.py, lib/crewai/src/crewai/memory/encoding_flow.py
MemoryRecord.validate_embedding converts embedding inputs into list[float] | None, handling bytes (via NumPy) and empty bytes as None. ItemState field validator normalizes MemoryRecord.embedding entries in similar_records and result_record before parsing.
Async-safe embedding execution
lib/crewai/src/crewai/memory/types.py
Module-level _EMBED_POOL dispatches batch embeddings to a thread pool when an asyncio loop is running; embed_texts detects the loop and submits work with a 30s timeout, or calls the embedder directly when no loop is active. Empty strings are skipped and positions preserved.
Timeout-aware write draining
lib/crewai/src/crewai/memory/unified_memory.py
Memory.drain_writes(timeout_per_save=60.0) waits per-save with timeouts, logging progress and failures, and counts exceptions instead of blocking indefinitely or raising.
Tests for embedding safety
lib/crewai/tests/memory/test_embedding_safety.py
Tests verify embedding normalization (bytes/NumPy/int to float list), empty string handling, embed_texts sync/async behavior with position preservation, and NumPy result conversion.

🎯 3 (Moderate) | ⏱️ ~25 minutes

🐰 Embeddings now flow safe,
Async loops get thread relief,
Timeouts won't block long—
bytes become lists of light,
and writes drain with grace! 🌿

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 39.13% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title directly captures the main changes: async-safe embeddings and resilient drain_writes timeout handling, matching the PR's core objectives.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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_writes timeout resilience is undermined by close() shutdown/close sequencing.

After a per-save timeout, writes can still be in-flight, but close() still does storage.close() and shutdown(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 win

Include traceback on save failure logs for actionable diagnostics.

The error log only prints str(e). Add exc_info=True so 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

📥 Commits

Reviewing files that changed from the base of the PR and between 63a9e7e and 6df71c1.

📒 Files selected for processing (4)
  • lib/crewai/src/crewai/memory/encoding_flow.py
  • lib/crewai/src/crewai/memory/types.py
  • lib/crewai/src/crewai/memory/unified_memory.py
  • lib/crewai/tests/memory/test_embedding_safety.py

Comment thread lib/crewai/src/crewai/memory/encoding_flow.py
Comment thread lib/crewai/src/crewai/memory/types.py Outdated
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from 6df71c1 to a4332f6 Compare May 11, 2026 20:37

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick comments (1)
lib/crewai/src/crewai/memory/types.py (1)

316-340: 💤 Low value

Consider applying the same async-context handling as embed_texts.

embed_texts offloads to the thread pool when called from an async context (lines 386-404), but embed_text calls the embedder directly. If embed_text is 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_text as a wrapper that delegates to embed_texts for 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

📥 Commits

Reviewing files that changed from the base of the PR and between 6df71c1 and a4332f6.

📒 Files selected for processing (4)
  • lib/crewai/src/crewai/memory/encoding_flow.py
  • lib/crewai/src/crewai/memory/types.py
  • lib/crewai/src/crewai/memory/unified_memory.py
  • lib/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

@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch 3 times, most recently from 8b8a533 to bbbe31e Compare May 13, 2026 14:48

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

♻️ Duplicate comments (1)
lib/crewai/src/crewai/memory/types.py (1)

343-360: ⚠️ Potential issue | 🟠 Major | 🏗️ Heavy lift

This path still blocks the event loop.

Line 393 waits with Future.result(timeout=30) on the caller thread. If embed_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_POOL and 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

📥 Commits

Reviewing files that changed from the base of the PR and between 8b8a533 and bbbe31e.

📒 Files selected for processing (4)
  • lib/crewai/src/crewai/memory/encoding_flow.py
  • lib/crewai/src/crewai/memory/types.py
  • lib/crewai/src/crewai/memory/unified_memory.py
  • lib/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

@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch 2 times, most recently from f9a7f08 to 2ddeaaf Compare May 19, 2026 18:17
@MatthiasHowellYopp

Copy link
Copy Markdown
Contributor Author

@greysonlalonde hoping I can get a review on this.

@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch 8 times, most recently from d2af721 to e17a45d Compare May 26, 2026 13:58
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch 4 times, most recently from c5dcef6 to 0683317 Compare June 2, 2026 15:17
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from 0683317 to dbe8086 Compare June 5, 2026 14:42

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread lib/crewai/src/crewai/a2a/utils/task.py
Comment thread lib/crewai/src/crewai/memory/types.py
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from 6342c7c to 6d3f769 Compare September 29, 2026 15:27

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread lib/crewai/src/crewai/utilities/cache_config.py
Matthias Howell and others added 4 commits September 29, 2026 13:05
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>
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from 6d3f769 to 90be3a4 Compare September 29, 2026 17:05

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread lib/crewai/src/crewai/memory/storage/valkey_cache.py
Comment thread lib/crewai/src/crewai/memory/storage/valkey_cache.py
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>
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from 90be3a4 to e7d1089 Compare September 29, 2026 19:07
Comment thread lib/crewai/src/crewai/a2a/utils/task.py
_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>
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from e7d1089 to a35c329 Compare September 29, 2026 19:14
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>
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from a35c329 to 5032fa0 Compare September 29, 2026 19:20

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread lib/crewai/src/crewai/a2a/utils/task.py
Comment thread lib/crewai/src/crewai/memory/unified_memory.py
Comment thread lib/crewai/src/crewai/tools/memory_tools.py
Comment thread lib/crewai-files/src/crewai_files/cache/upload_cache.py
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>
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from 5032fa0 to 3d39f2f Compare September 29, 2026 19:36
MatthiasHowellYopp and others added 6 commits September 29, 2026 15:55
- 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>
@MatthiasHowellYopp
MatthiasHowellYopp force-pushed the feat/valkey-3-embedding-safety branch from 3d39f2f to c39830f Compare September 29, 2026 20:01

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread lib/crewai/src/crewai/memory/types.py
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>

@cursor cursor Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Cursor Bugbot has reviewed your changes using high effort and found 1 potential issue.

Fix All in Cursor

❌ 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.

Comment thread lib/crewai/src/crewai/memory/types.py
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>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant