fix(data-connector): harden storage backends against data corruption and races - #505
Conversation
- Remove unused `dashmap` dependency from Cargo.toml - Fix grammatically broken error messages in PostgresConfig::validate() - Optimize hex string generation using std::fmt::Write with pre-allocated buffer - Fix double-serialization bug in postgres.rs for tool_calls and metadata fields - Add missing responses_safety_idx index in postgres schema initialization - Rename build_response_from_now → build_response_from_row for clarity - Reduce unnecessary cloning in postgres.rs by destructuring and using references - Extract build_response_from_map helper in redis.rs to eliminate ~50 lines of duplication - Consolidate parse_conversation_metadata into common.rs, used by all 3 backends - Fix redis storage struct visibility from pub to pub(super) for consistency - Add 106 unit tests across all modules (config, core, common, noop, memory, factory) - Add README.md with architecture, configuration, and usage documentation Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
…and races Fix 5 production-readiness issues identified during code review: 1. postgres.rs: Delegate parse_metadata to shared crate::common::parse_conversation_metadata helper, aligning with Oracle and Redis backends that already use it. The inline Postgres version had subtly different whitespace/null handling. 2. memory.rs: Consolidate get_response_chain from two lock acquisitions into a single read-lock scope. The previous code collected IDs under one lock, dropped it, then re-acquired to fetch responses — allowing concurrent writers to delete chain entries between locks and silently lose data. 3. redis.rs: Add build_item_from_map helper that returns proper errors for corrupted content (serde failures) and missing/invalid created_at timestamps. Previously, malformed data was silently replaced with Value::Null or Utc::now(), masking data corruption in production. Also harden build_response_from_map with the same created_at error propagation. 4. redis.rs: Replace sequential HGET + DEL in delete_response with a Redis pipeline so the safety_identifier read and key deletion happen atomically. Eliminates a race where the identifier could change between the read and the delete. 5. redis.rs: Fix cursor pagination in list_items to use inclusive score bounds with post-filtering by item_id, matching the composite (added_at, item_id) cursor semantics of Postgres and Oracle. The previous exclusive-bound approach silently skipped items sharing the same millisecond timestamp as the cursor. All 106 existing tests pass. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
📝 WalkthroughWalkthroughCentralizes metadata parsing, refactors hex ID generation, tightens Redis storage visibility, simplifies DB call-sites and ownership handling, makes in-memory response-chain traversal single-pass under one lock, removes a dashmap workspace flag, and adds extensive unit tests plus a new crate README. Changes
Sequence Diagram(s)(omitted) Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches
🧪 Generate unit tests (beta)
Comment |
Summary of ChangesHello @slin1237, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request addresses several critical production-readiness issues across the Highlights
Changelog
Activity
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for Github and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 47431735fd
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
There was a problem hiding this comment.
Code Review
The pull request significantly hardens the storage backends by addressing race conditions, improving error propagation, and aligning pagination behavior across different databases. Key improvements include consolidating lock scopes in the memory backend to prevent TOCTOU races and extracting shared metadata parsing logic. However, the Redis implementation of cursor-based pagination tie-breaking is currently incorrect and could lead to duplicate items in result sets. Additionally, the Redis deletion logic, while improved with a pipeline, is not yet atomic and remains susceptible to rare race conditions. Addressing these points will ensure the backends are truly production-ready.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@data_connector/src/redis.rs`:
- Around line 403-467: The current post-filter only compares IDs and assumes
tied-score items are at the result boundary; instead modify the retrieval to
fetch (score, id) pairs (use the zrange/zrevrangebyscore variant that returns
withscores — e.g., zrangebyscore_withscores_limit /
zrevrangebyscore_withscores_limit) into item_ids (rename to items_with_scores),
then implement composite-key filtering using cursor_score and cursor_id: for
SortOrder::Asc keep items where (score > cursor_score) OR (score == cursor_score
AND id > cursor_id); for SortOrder::Desc keep items where (score < cursor_score)
OR (score == cursor_score AND id < cursor_id); finally take params.limit and
collect IDs, preserving the existing over-fetch logic and error handling around
zscore, cursor_score, cursor_id, and fetch_limit.
…and pipeline atomicity 1. redis.rs list_items: Replace `filter(|id| id != c_id)` with `skip_while(|id| id != c_id).skip(1)` for cursor post-filtering. Redis returns same-score members in lexicographic order, so the previous filter incorrectly included same-score predecessors that belonged to the previous page, causing duplicate items across pages. 2. redis.rs delete_response: Add `.atomic()` to the Redis pipeline so HGET + DEL are wrapped in MULTI/EXEC. A plain pipeline only batches commands without transactional guarantees — other clients can still interleave between the two operations. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
e04b889 to
f890d07
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@data_connector/src/redis.rs`:
- Around line 250-295: The function build_item_from_map currently defaults
item_type to an empty string which hides corrupted/missing data; change the
extraction to fail fast: replace the item_type =
map.get("item_type").cloned().unwrap_or_default() line with logic that returns
Err(ConversationItemStorageError::StorageError(...)) when "item_type" is missing
or when its value is the empty string, e.g. use map.get("item_type").map(|s|
s.clone()).filter(|s| !s.is_empty()).ok_or_else(||
ConversationItemStorageError::StorageError(format!("item {fallback_id} missing
or empty item_type"))) so the function returns an error for missing/empty
item_type instead of silently accepting "".
---
Duplicate comments:
In `@data_connector/src/redis.rs`:
- Around line 405-459: The current post-filtering using zscore +
zrangebyscore_limit can drop items if the cursor id is missing or the same-score
group exceeds the overfetch; instead obtain the cursor position with
ZRANK/ZREVRANK (use the same SortOrder branch where you currently call zscore
and the item fetch), compute a start index = rank + 1, then fetch via
ZRANGE/ZREVRANGE with start and (limit + padding) to ensure you include
subsequent items; finally trim the returned Vec to params.limit without relying
on skip_while/skip(1). Update code paths that reference cursor_score/cursor_id,
zscore, zrangebyscore_limit, zrevrangebyscore_limit and replace the post-filter
block with rank-based start/limit logic while still falling back to the existing
range-by-score behavior if ZRANK fails or cursor is not present.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f890d07e74
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
…item_from_map Fail fast when item_type is missing or empty instead of silently defaulting to "". This aligns with the error-propagation approach already used for content and created_at in the same function. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@data_connector/src/redis.rs`:
- Around line 427-459: The current +32 over-fetch (fetch_limit) can still fail
when many members share the same score; fix by fetching member scores and doing
exact (score, id) cursor comparison instead of relying on a fixed margin. Change
the calls to the Redis fetch functions (currently zrangebyscore_limit /
zrevrangebyscore_limit) to versions that return WITHSCORES (or use a helper that
returns Vec<(String, f64)>), then in the post-filter use the tuple (score, id)
to skip until you hit the exact (cursor_score, cursor_id) pair and drop only
that one, finally take params.limit results; alternatively, if WITHSCORES
variants are unavailable, increase fetch_limit conservatively (e.g., to
params.limit + some larger bound) and document the limitation. Ensure you update
the item_ids handling (cursor_score, cursor_id, fetch_limit, and the
post-filter) to operate on scored items.
---
Duplicate comments:
In `@data_connector/src/redis.rs`:
- Line 262: The current deserialization silently defaults item_type via let
item_type = map.get("item_type").cloned().unwrap_or_default(); — change this to
propagate an error when "item_type" is missing or empty so corrupted data isn't
accepted; locate the deserialization code that sets item_type (the variable
named item_type) and replace the unwrap_or_default logic with a check that
returns a Result::Err (or appropriate error variant) when map.get("item_type")
is None or the string is empty, ensuring callers (and make_item_id usage)
receive the error instead of an empty prefix.
| // Over-fetch to handle same-score ties that need filtering | ||
| let fetch_limit = if cursor_score.is_some() { | ||
| // Fetch extra to compensate for items we'll filter out at the cursor boundary | ||
| (params.limit + 32) as isize | ||
| } else { | ||
| params.limit as isize | ||
| }; | ||
|
|
||
| let item_ids: Vec<String> = match params.order { | ||
| SortOrder::Asc => { | ||
| // ZRANGEBYSCORE key min max LIMIT offset count | ||
| conn.zrangebyscore_limit(&key, min, max, 0, params.limit as isize) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))? | ||
| } | ||
| SortOrder::Desc => { | ||
| // ZREVRANGEBYSCORE key max min LIMIT offset count | ||
| conn.zrevrangebyscore_limit(&key, max, min, 0, params.limit as isize) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))? | ||
| } | ||
| SortOrder::Asc => conn | ||
| .zrangebyscore_limit(&key, min, max, 0, fetch_limit) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))?, | ||
| SortOrder::Desc => conn | ||
| .zrevrangebyscore_limit(&key, max, min, 0, fetch_limit) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))?, | ||
| }; | ||
|
|
||
| // Post-filter: skip past the cursor item and all same-score predecessors. | ||
| // Redis returns same-score members in lexicographic order (ASC) or | ||
| // reverse-lex (DESC), so `skip_while` advances past items that appeared | ||
| // on the previous page, then `skip(1)` drops the cursor item itself. | ||
| let item_ids: Vec<String> = if let (Some(_), Some(ref c_id)) = (cursor_score, &cursor_id) { | ||
| item_ids | ||
| .into_iter() | ||
| .skip_while(|id| id != c_id) | ||
| .skip(1) | ||
| .take(params.limit) | ||
| .collect() | ||
| } else { | ||
| item_ids.into_iter().take(params.limit).collect() | ||
| }; |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
The +32 over-fetch margin may be insufficient when many items share the same score.
The skip_while(|id| id != cursor_id).skip(1) logic correctly implements cursor pagination by relying on Redis's documented lexicographic ordering of same-score members. However, if more than 32 items share the cursor's millisecond timestamp, the post-filter will skip past the entire over-fetched buffer and return fewer than params.limit items.
For typical timestamps with low collision rates, this works fine. In high-cardinality scenarios (e.g., bulk inserts with the same timestamp), consider increasing the margin or fetching WITHSCORES to filter by exact (score, id) comparison, which avoids relying on arbitrary margin sizing.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@data_connector/src/redis.rs` around lines 427 - 459, The current +32
over-fetch (fetch_limit) can still fail when many members share the same score;
fix by fetching member scores and doing exact (score, id) cursor comparison
instead of relying on a fixed margin. Change the calls to the Redis fetch
functions (currently zrangebyscore_limit / zrevrangebyscore_limit) to versions
that return WITHSCORES (or use a helper that returns Vec<(String, f64)>), then
in the post-filter use the tuple (score, id) to skip until you hit the exact
(cursor_score, cursor_id) pair and drop only that one, finally take params.limit
results; alternatively, if WITHSCORES variants are unavailable, increase
fetch_limit conservatively (e.g., to params.limit + some larger bound) and
document the limitation. Ensure you update the item_ids handling (cursor_score,
cursor_id, fetch_limit, and the post-filter) to operate on scored items.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: df2c2eb0b3
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| .skip_while(|id| id != c_id) | ||
| .skip(1) |
There was a problem hiding this comment.
Guard missing cursor before skipping Redis page entries
In list_items, the cursor post-filter skip_while(|id| id != c_id).skip(1) assumes the cursor member is present in the fetched item_ids. If the cursor item is deleted/expired after ZSCORE but before ZRANGEBYSCORE, c_id is absent and skip_while consumes the entire iterator, so the API returns an empty page even though later items still exist. This creates pagination gaps under concurrent deletes/retention expiry; only apply this skip logic when the cursor ID is actually found in the fetched window.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@data_connector/src/redis.rs`:
- Around line 250-279: build_item_from_map currently treats a missing "content"
key as Value::Null which hides corrupted/missing data; change the logic in
build_item_from_map so that when map.get("content") is None it returns
Err(ConversationItemStorageError::StorageError(...)) (similar to how missing
"item_type" is handled), using the fallback_id to construct a clear error
message; keep the existing serde_json::from_str error mapping for
present-but-invalid JSON.
- Around line 411-467: The post-filtering currently assumes the cursor id
(cursor_id / params.after) is present in the fetched item_ids, which races if
the cursor member expired and causes empty pages; update the logic after
fetching item_ids (from zrangebyscore_limit / zrevrangebyscore_limit) to check
whether c_id is actually present in the returned Vec: if c_id is found, keep the
existing skip_while(|id| id != c_id).skip(1).take(params.limit).collect()
behavior, but if c_id is not found, fall back to score-only pagination by simply
taking the first params.limit items from the fetched set (i.e.,
item_ids.into_iter().take(params.limit).collect()), thus preserving pagination
when the cursor member is missing.
| /// Parse a Redis hash map into a `ConversationItem`, returning errors for | ||
| /// corrupted data instead of silently substituting defaults. | ||
| fn build_item_from_map( | ||
| map: &std::collections::HashMap<String, String>, | ||
| fallback_id: &str, | ||
| ) -> Result<ConversationItem, ConversationItemStorageError> { | ||
| let id = ConversationItemId( | ||
| map.get("id") | ||
| .cloned() | ||
| .unwrap_or_else(|| fallback_id.to_string()), | ||
| ); | ||
| let response_id = map.get("response_id").cloned(); | ||
| let item_type = map | ||
| .get("item_type") | ||
| .filter(|s| !s.is_empty()) | ||
| .cloned() | ||
| .ok_or_else(|| { | ||
| ConversationItemStorageError::StorageError(format!( | ||
| "item {fallback_id} missing item_type" | ||
| )) | ||
| })?; | ||
| let role = map.get("role").cloned(); | ||
| let status = map.get("status").cloned(); | ||
|
|
||
| let content = match map.get("content") { | ||
| Some(s) => { | ||
| serde_json::from_str(s).map_err(ConversationItemStorageError::SerializationError)? | ||
| } | ||
| None => Value::Null, | ||
| }; |
There was a problem hiding this comment.
Treat missing content as corruption (don’t default to Null).
The helper’s doc comment says it surfaces corrupted data, but a missing content key silently becomes Value::Null, masking corruption and making it indistinguishable from an explicit null. Consider failing fast like item_type/created_at.
🔧 Proposed fix
- let content = match map.get("content") {
- Some(s) => {
- serde_json::from_str(s).map_err(ConversationItemStorageError::SerializationError)?
- }
- None => Value::Null,
- };
+ let content_str = map.get("content").ok_or_else(|| {
+ ConversationItemStorageError::StorageError(format!(
+ "item {fallback_id} missing content"
+ ))
+ })?;
+ let content =
+ serde_json::from_str(content_str).map_err(ConversationItemStorageError::SerializationError)?;📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| /// Parse a Redis hash map into a `ConversationItem`, returning errors for | |
| /// corrupted data instead of silently substituting defaults. | |
| fn build_item_from_map( | |
| map: &std::collections::HashMap<String, String>, | |
| fallback_id: &str, | |
| ) -> Result<ConversationItem, ConversationItemStorageError> { | |
| let id = ConversationItemId( | |
| map.get("id") | |
| .cloned() | |
| .unwrap_or_else(|| fallback_id.to_string()), | |
| ); | |
| let response_id = map.get("response_id").cloned(); | |
| let item_type = map | |
| .get("item_type") | |
| .filter(|s| !s.is_empty()) | |
| .cloned() | |
| .ok_or_else(|| { | |
| ConversationItemStorageError::StorageError(format!( | |
| "item {fallback_id} missing item_type" | |
| )) | |
| })?; | |
| let role = map.get("role").cloned(); | |
| let status = map.get("status").cloned(); | |
| let content = match map.get("content") { | |
| Some(s) => { | |
| serde_json::from_str(s).map_err(ConversationItemStorageError::SerializationError)? | |
| } | |
| None => Value::Null, | |
| }; | |
| /// Parse a Redis hash map into a `ConversationItem`, returning errors for | |
| /// corrupted data instead of silently substituting defaults. | |
| fn build_item_from_map( | |
| map: &std::collections::HashMap<String, String>, | |
| fallback_id: &str, | |
| ) -> Result<ConversationItem, ConversationItemStorageError> { | |
| let id = ConversationItemId( | |
| map.get("id") | |
| .cloned() | |
| .unwrap_or_else(|| fallback_id.to_string()), | |
| ); | |
| let response_id = map.get("response_id").cloned(); | |
| let item_type = map | |
| .get("item_type") | |
| .filter(|s| !s.is_empty()) | |
| .cloned() | |
| .ok_or_else(|| { | |
| ConversationItemStorageError::StorageError(format!( | |
| "item {fallback_id} missing item_type" | |
| )) | |
| })?; | |
| let role = map.get("role").cloned(); | |
| let status = map.get("status").cloned(); | |
| let content_str = map.get("content").ok_or_else(|| { | |
| ConversationItemStorageError::StorageError(format!( | |
| "item {fallback_id} missing content" | |
| )) | |
| })?; | |
| let content = | |
| serde_json::from_str(content_str).map_err(ConversationItemStorageError::SerializationError)?; |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@data_connector/src/redis.rs` around lines 250 - 279, build_item_from_map
currently treats a missing "content" key as Value::Null which hides
corrupted/missing data; change the logic in build_item_from_map so that when
map.get("content") is None it returns
Err(ConversationItemStorageError::StorageError(...)) (similar to how missing
"item_type" is handled), using the fallback_id to construct a clear error
message; keep the existing serde_json::from_str error mapping for
present-but-invalid JSON.
| let mut min = "-inf".to_string(); | ||
| let mut max = "+inf".to_string(); | ||
| // Track cursor score + id for post-filtering same-millisecond ties, | ||
| // matching the composite (added_at, item_id) cursor of Postgres/Oracle. | ||
| let mut cursor_score: Option<f64> = None; | ||
| let mut cursor_id: Option<String> = None; | ||
|
|
||
| if let Some(after_id) = ¶ms.after { | ||
| let score: Option<f64> = conn | ||
| .zscore(&key, after_id) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))?; | ||
| if let Some(s) = score { | ||
| cursor_score = Some(s); | ||
| cursor_id = Some(after_id.clone()); | ||
| // Use inclusive bound so we can post-filter ties by item_id. | ||
| // Over-fetch slightly to account for items at the cursor's score. | ||
| match params.order { | ||
| SortOrder::Asc => min = format!("({s}"), | ||
| SortOrder::Desc => max = format!("({s}"), | ||
| SortOrder::Asc => min = s.to_string(), | ||
| SortOrder::Desc => max = s.to_string(), | ||
| } | ||
| } | ||
| } | ||
|
|
||
| // Over-fetch to handle same-score ties that need filtering | ||
| let fetch_limit = if cursor_score.is_some() { | ||
| // Fetch extra to compensate for items we'll filter out at the cursor boundary | ||
| (params.limit + 32) as isize | ||
| } else { | ||
| params.limit as isize | ||
| }; | ||
|
|
||
| let item_ids: Vec<String> = match params.order { | ||
| SortOrder::Asc => { | ||
| // ZRANGEBYSCORE key min max LIMIT offset count | ||
| conn.zrangebyscore_limit(&key, min, max, 0, params.limit as isize) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))? | ||
| } | ||
| SortOrder::Desc => { | ||
| // ZREVRANGEBYSCORE key max min LIMIT offset count | ||
| conn.zrevrangebyscore_limit(&key, max, min, 0, params.limit as isize) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))? | ||
| } | ||
| SortOrder::Asc => conn | ||
| .zrangebyscore_limit(&key, min, max, 0, fetch_limit) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))?, | ||
| SortOrder::Desc => conn | ||
| .zrevrangebyscore_limit(&key, max, min, 0, fetch_limit) | ||
| .await | ||
| .map_err(|e| ConversationItemStorageError::StorageError(e.to_string()))?, | ||
| }; | ||
|
|
||
| // Post-filter: skip past the cursor item and all same-score predecessors. | ||
| // Redis returns same-score members in lexicographic order (ASC) or | ||
| // reverse-lex (DESC), so `skip_while` advances past items that appeared | ||
| // on the previous page, then `skip(1)` drops the cursor item itself. | ||
| let item_ids: Vec<String> = if let (Some(_), Some(ref c_id)) = (cursor_score, &cursor_id) { | ||
| item_ids | ||
| .into_iter() | ||
| .skip_while(|id| id != c_id) | ||
| .skip(1) | ||
| .take(params.limit) | ||
| .collect() | ||
| } else { | ||
| item_ids.into_iter().take(params.limit).collect() | ||
| }; |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
# First, let's examine the full context around the cursor pagination logic
cat -n data_connector/src/redis.rs | sed -n '380,480p'Repository: lightseekorg/smg
Length of output: 4652
🏁 Script executed:
# Check for any item expiry/TTL logic in the codebase
rg -i "expire|ttl|delete|remove" data_connector/src/redis.rs -A 2 -B 2Repository: lightseekorg/smg
Length of output: 2053
🏁 Script executed:
# Look for test cases related to cursor pagination
fd -e rs -e toml | xargs rg -l "cursor|pagination|zrange" | head -5Repository: lightseekorg/smg
Length of output: 198
🏁 Script executed:
# Check if there are any guards or error handling for missing cursor scenarios elsewhere
rg "skip_while|position|cursor" data_connector/src/redis.rs -B 3 -A 3Repository: lightseekorg/smg
Length of output: 2080
Guard against missing cursor IDs to avoid empty pages during pagination.
The cursor entry can expire or be deleted between the ZSCORE lookup (line 419) and the ZRANGEBYSCORE query (lines 444–451). If this occurs, the skip_while(|id| id != c_id).skip(1) logic will iterate through the entire result set without finding the cursor ID, exhaust the iterator, and return an empty page even when later items exist. The over-fetch buffer of +32 items compensates for same-score ties but does not address this race condition.
Detect when the cursor is missing and fall back to score-only pagination to ensure continuous pagination across item lifecycles.
Proposed mitigation
- let item_ids: Vec<String> = if let (Some(_), Some(ref c_id)) = (cursor_score, &cursor_id) {
- item_ids
- .into_iter()
- .skip_while(|id| id != c_id)
- .skip(1)
- .take(params.limit)
- .collect()
- } else {
- item_ids.into_iter().take(params.limit).collect()
- };
+ let item_ids: Vec<String> = if let (Some(_), Some(ref c_id)) = (cursor_score, &cursor_id) {
+ if let Some(pos) = item_ids.iter().position(|id| id == c_id) {
+ item_ids
+ .into_iter()
+ .skip(pos + 1)
+ .take(params.limit)
+ .collect()
+ } else {
+ // Cursor vanished between ZSCORE and ZRANGEBYSCORE; fall back to score-only paging.
+ item_ids.into_iter().take(params.limit).collect()
+ }
+ } else {
+ item_ids.into_iter().take(params.limit).collect()
+ };🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@data_connector/src/redis.rs` around lines 411 - 467, The post-filtering
currently assumes the cursor id (cursor_id / params.after) is present in the
fetched item_ids, which races if the cursor member expired and causes empty
pages; update the logic after fetching item_ids (from zrangebyscore_limit /
zrevrangebyscore_limit) to check whether c_id is actually present in the
returned Vec: if c_id is found, keep the existing skip_while(|id| id !=
c_id).skip(1).take(params.limit).collect() behavior, but if c_id is not found,
fall back to score-only pagination by simply taking the first params.limit items
from the fetched set (i.e., item_ids.into_iter().take(params.limit).collect()),
thus preserving pagination when the cursor member is missing.
Summary
Fixes 5 production-readiness issues in the
data_connectorstorage backends (Postgres, Memory, Redis) identified during code review of #504 (commit363bbe25).Refs:
363bbe25(refactor(data-connector): comprehensive code quality improvements)What changed
data_connector/src/postgres.rsparse_metadatawith delegation to the sharedcrate::common::parse_conversation_metadatahelper, matching Oracle and Redis backends. The Postgres-specific version had subtly different whitespace/null handling.data_connector/src/memory.rsget_response_chainfrom two separate lock acquisitions into a single read-lock scope. The original code collected chain IDs under one lock, dropped it, then re-acquired to fetch responses — allowing concurrent writers to delete entries between locks and silently lose chain data.data_connector/src/redis.rsbuild_item_from_maphelper with proper error propagation for corruptedcontent(serde failures) and missing/invalidcreated_attimestamps. Previously, malformed data was silently replaced withValue::NullorUtc::now(), masking data corruption. Also hardenedbuild_response_from_mapwith the samecreated_aterror handling.HGET+DELindelete_responsewith a Redis pipeline for atomic read+delete. Eliminates a TOCTOU race where thesafety_identifiercould change between the read and the delete.list_itemsto use inclusive score bounds with post-filtering byitem_id, matching the composite(added_at, item_id)cursor semantics of Postgres and Oracle. The exclusive-bound approach silently skipped items sharing the same millisecond timestamp as the cursor.Why
The original refactor commit introduced several patterns that could cause silent data loss or corruption in production:
Utc::now()andValue::Nullin Redis masks corrupted recordsHow
build_item_from_mapthat propagatesSerializationErrorandStorageErrorinstead of substituting defaultsbuild_response_from_mapnow errors on missing/invalidcreated_atredis::pipe()batchesHGET+DELinto one round-tripZRANGEBYSCORE/ZREVRANGEBYSCOREbounds with over-fetch (+32) and post-filter to skip the cursor item, thentake(limit)Test plan
cargo test -p data-connectorcargo check -p data-connectorSummary by CodeRabbit
Documentation
Improvements
Bug Fixes
Tests
Chores