Skip to content

fix: encode type and handle missing blob in Redis checkpoint writes #57

New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 24 additions & 20 deletions langgraph/checkpoint/redis/aio.py
Original file line number Diff line number Diff line change
Expand Up @@ -1009,29 +1009,33 @@ async def _aload_pending_sends(
"""
# Query checkpoint_writes for parent checkpoint's TASKS channel
parent_writes_query = FilterQuery(
filter_expression=(Tag("thread_id") == to_storage_safe_id(thread_id))
& (Tag("checkpoint_ns") == to_storage_safe_str(checkpoint_ns))
& (Tag("checkpoint_id") == to_storage_safe_id(parent_checkpoint_id))
& (Tag("channel") == TASKS),
return_fields=["type", "blob", "task_path", "task_id", "idx"],
num_results=100, # Adjust as needed
)
parent_writes_results = await self.checkpoint_writes_index.search(
parent_writes_query
filter_expression=(
(Tag("thread_id") == to_storage_safe_id(thread_id))
& (Tag("checkpoint_ns") == checkpoint_ns)
& (Tag("checkpoint_id") == to_storage_safe_id(parent_checkpoint_id))
& (Tag("channel") == TASKS)
),
return_fields=["type", "$.blob", "task_path", "task_id", "idx"],
num_results=100,
)

# Sort results by task_path, task_id, idx (matching Postgres implementation)
sorted_writes = sorted(
parent_writes_results.docs,
key=lambda x: (
getattr(x, "task_path", ""),
getattr(x, "task_id", ""),
getattr(x, "idx", 0),
res = await self.checkpoint_writes_index.search(parent_writes_query)

# Sort results for deterministic order
docs = sorted(
res.docs,
key=lambda d: (
getattr(d, "task_path", ""),
getattr(d, "task_id", ""),
getattr(d, "idx", 0),
),
)

# Extract type and blob pairs
return [(doc.type, doc.blob) for doc in sorted_writes]

# Convert to expected format
return [
(d.type.encode(), blob)
for d in docs
if (blob := getattr(d, "$.blob", getattr(d, "blob", None))) is not None
]

async def _aload_pending_writes(
self,
Expand Down