Repository navigation
feat(memory): schedule STMO jobs after conversation persistence (NoOp writer — follow-up for Oracle wiring) - #1346
spalimpaaces-star wants to merge 5 commits into
Conversation
|
Warning You have reached your daily quota limit. Please wait up to 24 hours and I will start processing your requests again! |
📝 WalkthroughWalkthroughThis PR extends the conversation memory system with Short-Term Memory (STM) configuration and enqueuing logic. It adds STM enablement flags and condenser model IDs to context and header objects, implements deterministic user-turn counting, and conditionally creates STMO enqueue tasks during conversation persistence (triggered at turn 4, then every 3 turns thereafter). Changes
Sequence Diagram(s)sequenceDiagram
participant Client
participant Route as Route Handler
participant Context as Memory Context
participant Persister as persist_conversation_items
participant Writer as ConversationMemoryWriter
participant DB as Database
Client->>Route: ResponsesRequest (with STM enabled)
Route->>Route: count_conversation_user_turns(input)
Route->>Context: Store turn_count in ResponsesPayloadState
Route->>Persister: Call with memory_writer, memory_context, turn_count
Persister->>Persister: Link input/output items to conversation
Persister->>Persister: Check: stm_enabled && eligible_turn?<br/>(turn 4, then every 3)
alt Turn is Eligible for STMO Enqueue
Persister->>Persister: Compute target_item_end<br/>Build memory_config JSON<br/>(condenser_model, last_index, target_item_end)
Persister->>Writer: Create NewConversationMemory<br/>(type: Stmo, status: Ready)
Writer->>DB: INSERT memory record
DB-->>Writer: Success/Failure
Writer-->>Persister: Result (warn on failure)
else Turn Not Eligible
Persister->>Persister: Skip STMO enqueue
end
Persister-->>Route: Persistence complete
Route-->>Client: Response
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~27 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 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 |
|
Hi @spalimpaaces-star, the DCO sign-off check has failed. All commits must include a To fix existing commits: # Sign off the last N commits (replace N with the number of unsigned commits)
git rebase HEAD~N --signoff
git push --force-with-leaseTo sign off future commits automatically:
|
…cutionContext x-conversation-memory-config was parsed twice per request — once in middleware to build MemoryExecutionContext, and again in route_responses for the inject_memory_context no-op stub. Remove the second parse and the stub entirely. Signed-off-by: Saikiran Palimpati <saikiran.palimpati@oracle.com> Signed-off-by: spalimpaaces-star <saikiran.palimpati@oracle.com>
… into persist_conversation_items Signed-off-by: saikiranpalimpati <34260562+saikiranpalimpati@users.noreply.github.com> Signed-off-by: spalimpaaces-star <saikiran.palimpati@oracle.com>
…n flow Signed-off-by: saikiranpalimpati <34260562+saikiranpalimpati@users.noreply.github.com> Signed-off-by: spalimpaaces-star <saikiran.palimpati@oracle.com>
a3b0237 to
20ed351
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: a3b0237ffb
ℹ️ 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".
Previously enqueue_stmo_if_needed derived the user turn count from original_body.input at persistence time, which only contains the current request payload. STMO trigger boundaries (4/7/10/...) were never reached for multi-turn conversations. Fix: count_conversation_user_turns runs over request_body.input after load_input_history assembles the full conversation context. The count is threaded through ResponsesPayloadState and StorageHandles to both the streaming and non-streaming persistence paths. Signed-off-by: spalimpaaces-star <saikiran.palimpati@oracle.com>
stm_enabled was set directly from the request header without checking runtime.enabled, allowing STMO scheduling to be activated even when the master memory feature gate is off. Apply the same runtime guard used by store_ltm and recall. Signed-off-by: spalimpaaces-star <saikiran.palimpati@oracle.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 1c910fe89d
ℹ️ 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".
| let conversation_user_turn_count = if ctx.memory_execution_context.stm_enabled { | ||
| Some(super::history::count_conversation_user_turns( | ||
| &request_body.input, | ||
| )) |
There was a problem hiding this comment.
Use full conversation length for STMO turn count
conversation_user_turn_count is derived from request_body.input, but in conversation-id flows that input is assembled by load_input_history, which still limits history to MAX_CONVERSATION_HISTORY_ITEMS = 100 (history.rs) before this count is taken. Fresh evidence: the loader fetches a capped slice, so after long conversations the count stops reflecting total user turns, and STMO boundary checks (4/7/10/...) will fire at the wrong times or not at all. This directly affects production threads once they exceed the history cap.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 3
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
model_gateway/src/routers/grpc/common/responses/utils.rs (1)
154-171: 🧹 Nitpick | 🔵 TrivialLGTM — double-gated STMO inhibition on the gRPC path.
Passing
NoOpConversationMemoryWriter, a defaultMemoryExecutionContext(stm_enabled=false), andNonefor the turn count ensures the gRPC persistence path never enqueues STMO jobs, matching the PR's "NoOp writer — follow-up for Oracle wiring" scope.Minor nit (optional):
Arc::new(NoOpConversationMemoryWriter::new())allocates on every persisted response. If this path is hot, consider aOnceLock<Arc<dyn ConversationMemoryWriter>>or storing the Arc on a shared components struct to avoid per-call allocation. Not required for this PR.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/routers/grpc/common/responses/utils.rs` around lines 154 - 171, The current code calls Arc::new(NoOpConversationMemoryWriter::new()) each time before calling persist_conversation_items which may allocate on every hot path; change this to a shared, lazily-initialized Arc so the NoOpConversationMemoryWriter is created once and reused (e.g., store Arc<dyn ConversationMemoryWriter> in a static OnceLock or on the shared components struct) and pass that shared Arc into persist_conversation_items instead of creating a new one per call.model_gateway/src/routers/common/header_utils.rs (1)
496-511: 🧹 Nitpick | 🔵 TrivialAdd STM assertions to this test to lock in the independence of STM and LTM gates.
The input header already carries
"short_term_memory":{"enabled":true,"condenser_model_id":"cond-1"}, but no assertion covers the STM fields. Per the newfrom_http_headerslogic, STM should remain enabled regardless of LTM state — adding the assertions below exercises that contract and guards against an accidental future coupling.✏️ Proposed test additions
let view = MemoryHeaderView::from_http_headers(&headers); assert_eq!(view.policy, None); assert_eq!(view.subject_id, None); assert_eq!(view.embedding_model, None); assert_eq!(view.extraction_model, None); + // STM is independent of the LTM enabled flag. + assert!(view.stm_enabled); + assert_eq!(view.stm_condenser_model_id.as_deref(), Some("cond-1")); }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/routers/common/header_utils.rs` around lines 496 - 511, In test_memory_header_view_defaults_when_ltm_disabled, add assertions after creating view = MemoryHeaderView::from_http_headers(&headers) to lock the STM/LTM independence: assert that the short-term memory flag on MemoryHeaderView is enabled (e.g., view.short_term_enabled or view.short_term_memory_enabled is true) and that the condenser model id is preserved (e.g., view.condenser_model_id or view.condenser_model equals "cond-1"), so the test explicitly verifies STM remains active and its condenser model is parsed even when LTM is disabled.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/memory/context.rs`:
- Around line 89-91: The MemoryExecutionContext currently sets
stm_condenser_model_id unconditionally which can leave a model id present when
stm_enabled is false; update the constructor/initializer that builds
MemoryExecutionContext so stm_condenser_model_id is only assigned when both
headers.stm_condenser_model_id.is_some() and runtime.enabled are true (otherwise
set it to None) to keep it consistent with stm_enabled, and add a unit test
similar to stm_enabled_gated_off_when_runtime_disabled that asserts
ctx.stm_condenser_model_id.is_none() when runtime.enabled is false but the
header provided a condenser id; refer to MemoryExecutionContext, stm_enabled,
stm_condenser_model_id, headers and runtime to locate affected code.
In `@model_gateway/src/routers/common/persistence_utils.rs`:
- Around line 605-793: Tests cover STMO boundaries and config shapes but miss
the stm_enabled=false short-circuit; add a tokio::test that constructs a
MemoryExecutionContext with stm_enabled: false (keep other fields like
stm_condenser_model_id as needed), use a RecordingConversationMemoryWriter and
call enqueue_stmo_if_needed with a turn that would normally enqueue (e.g.,
Some(4)), then assert that writer.rows remains empty; reference the existing
RecordingConversationMemoryWriter, MemoryExecutionContext, and
enqueue_stmo_if_needed to implement the test consistent with the other test
patterns.
- Around line 540-603: The STMO enqueue behavior is undocumented and currently
sets scope_id = None which can produce duplicate jobs on retries; update the
code in enqueue_stmo_if_needed to include an explicit inline comment documenting
the chosen semantic (either "best-effort / duplicates allowed; consumer must be
idempotent" OR "at-most-once / deduplicate by scope_id") and follow the same
pattern used in enqueue_conversation_memory_rows (PR `#1343`): if you decide
duplicates are acceptable, leave scope_id = None and state that clearly next to
NewConversationMemory creation and the create_memory call; if you decide to
enforce uniqueness, compute and set a deterministic scope_id (e.g., based on
conversation_id + user_turns + response_id or target_item_end) and document that
producers/writers must use scope_id to dedupe when implementing Postgres/Oracle
writers; ensure the comment references enqueue_stmo_if_needed,
NewConversationMemory.scope_id, and ConversationMemoryWriter.create_memory so
future writers know how to implement deduplication.
---
Outside diff comments:
In `@model_gateway/src/routers/common/header_utils.rs`:
- Around line 496-511: In test_memory_header_view_defaults_when_ltm_disabled,
add assertions after creating view =
MemoryHeaderView::from_http_headers(&headers) to lock the STM/LTM independence:
assert that the short-term memory flag on MemoryHeaderView is enabled (e.g.,
view.short_term_enabled or view.short_term_memory_enabled is true) and that the
condenser model id is preserved (e.g., view.condenser_model_id or
view.condenser_model equals "cond-1"), so the test explicitly verifies STM
remains active and its condenser model is parsed even when LTM is disabled.
In `@model_gateway/src/routers/grpc/common/responses/utils.rs`:
- Around line 154-171: The current code calls
Arc::new(NoOpConversationMemoryWriter::new()) each time before calling
persist_conversation_items which may allocate on every hot path; change this to
a shared, lazily-initialized Arc so the NoOpConversationMemoryWriter is created
once and reused (e.g., store Arc<dyn ConversationMemoryWriter> in a static
OnceLock or on the shared components struct) and pass that shared Arc into
persist_conversation_items instead of creating a new one per call.
🪄 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: ASSERTIVE
Plan: Pro
Run ID: 41805c2b-8b75-420e-8174-7da0b7b92356
📒 Files selected for processing (9)
model_gateway/src/memory/context.rsmodel_gateway/src/routers/common/header_utils.rsmodel_gateway/src/routers/common/persistence_utils.rsmodel_gateway/src/routers/grpc/common/responses/utils.rsmodel_gateway/src/routers/openai/context.rsmodel_gateway/src/routers/openai/responses/history.rsmodel_gateway/src/routers/openai/responses/non_streaming.rsmodel_gateway/src/routers/openai/responses/route.rsmodel_gateway/src/routers/openai/responses/streaming.rs
| stm_enabled: headers.stm_enabled && runtime.enabled, | ||
| stm_condenser_model_id: headers.stm_condenser_model_id.clone(), | ||
| } |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Minor inconsistency: stm_condenser_model_id is not gated by runtime.enabled.
When headers.stm_enabled=true but runtime.enabled=false, the resulting MemoryExecutionContext has stm_enabled=false yet stm_condenser_model_id=Some(...). This is fine as long as every downstream consumer gates on stm_enabled (not on stm_condenser_model_id.is_some()). To remove the footgun entirely and keep the two STM fields consistent, consider clearing the model id when the gate trips:
♻️ Proposed diff
- stm_enabled: headers.stm_enabled && runtime.enabled,
- stm_condenser_model_id: headers.stm_condenser_model_id.clone(),
+ stm_enabled: headers.stm_enabled && runtime.enabled,
+ stm_condenser_model_id: if headers.stm_enabled && runtime.enabled {
+ headers.stm_condenser_model_id.clone()
+ } else {
+ None
+ },Adding a test asserting ctx.stm_condenser_model_id.is_none() when runtime is disabled but the header set a condenser would lock this in alongside the existing stm_enabled_gated_off_when_runtime_disabled test.
📝 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.
| stm_enabled: headers.stm_enabled && runtime.enabled, | |
| stm_condenser_model_id: headers.stm_condenser_model_id.clone(), | |
| } | |
| stm_enabled: headers.stm_enabled && runtime.enabled, | |
| stm_condenser_model_id: if headers.stm_enabled && runtime.enabled { | |
| headers.stm_condenser_model_id.clone() | |
| } else { | |
| None | |
| }, | |
| } |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/memory/context.rs` around lines 89 - 91, The
MemoryExecutionContext currently sets stm_condenser_model_id unconditionally
which can leave a model id present when stm_enabled is false; update the
constructor/initializer that builds MemoryExecutionContext so
stm_condenser_model_id is only assigned when both
headers.stm_condenser_model_id.is_some() and runtime.enabled are true (otherwise
set it to None) to keep it consistent with stm_enabled, and add a unit test
similar to stm_enabled_gated_off_when_runtime_disabled that asserts
ctx.stm_condenser_model_id.is_none() when runtime.enabled is false but the
header provided a condenser id; refer to MemoryExecutionContext, stm_enabled,
stm_condenser_model_id, headers and runtime to locate affected code.
| async fn enqueue_stmo_if_needed( | ||
| conversation_memory_writer: &Arc<dyn ConversationMemoryWriter>, | ||
| memory_execution_context: &MemoryExecutionContext, | ||
| conversation_user_turn_count: Option<usize>, | ||
| conversation_id: &ConversationId, | ||
| response_id: &ResponseId, | ||
| output_items: &[Value], | ||
| input_item_count: usize, | ||
| ) { | ||
| if !memory_execution_context.stm_enabled { | ||
| return; | ||
| } | ||
|
|
||
| let Some(user_turns) = conversation_user_turn_count else { | ||
| return; | ||
| }; | ||
|
|
||
| if !should_enqueue_stmo_for_current_turn(user_turns) { | ||
| return; | ||
| } | ||
|
|
||
| let target_item_end = input_item_count + output_items.len(); | ||
| // STMO worker config semantics: | ||
| // - `last_index`: latest observed user-turn count at enqueue time. | ||
| // - `target_item_end`: exclusive end index for items included in this run. | ||
| let mut job_config = Map::new(); | ||
| if let Some(condenser_model) = memory_execution_context.stm_condenser_model_id.as_deref() { | ||
| job_config.insert( | ||
| STMO_CFG_KEY_CONDENSER_MODEL.to_string(), | ||
| Value::String(condenser_model.to_string()), | ||
| ); | ||
| } | ||
| job_config.insert(STMO_CFG_KEY_LAST_INDEX.to_string(), json!(user_turns)); | ||
| job_config.insert( | ||
| STMO_CFG_KEY_TARGET_ITEM_END.to_string(), | ||
| json!(target_item_end), | ||
| ); | ||
| let job_config = Value::Object(job_config).to_string(); | ||
|
|
||
| let row = NewConversationMemory { | ||
| conversation_id: conversation_id.clone(), | ||
| conversation_version: None, | ||
| response_id: Some(response_id.clone()), | ||
| memory_type: ConversationMemoryType::Stmo, | ||
| status: ConversationMemoryStatus::Ready, | ||
| attempt: 0, | ||
| owner_id: None, | ||
| next_run_at: Utc::now(), | ||
| lease_until: None, | ||
| content: None, | ||
| memory_config: Some(job_config), | ||
| scope_id: None, | ||
| error_msg: None, | ||
| }; | ||
|
|
||
| if let Err(err) = conversation_memory_writer.create_memory(row).await { | ||
| warn!( | ||
| conversation_id = %conversation_id.0, | ||
| response_id = %response_id.0, | ||
| error = %err, | ||
| "Failed to enqueue STMO job (best-effort; request flow continues)" | ||
| ); | ||
| } | ||
| } |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Does the writer (or its Oracle impl) enforce uniqueness on scope_id / some key?
rg -nP -C8 'create_memory\s*\(|scope_id'
ast-grep --pattern $'impl ConversationMemoryWriter for $_ {
$$$
}'Repository: lightseekorg/smg
Length of output: 14198
🏁 Script executed:
# Check Postgres and Oracle create_memory implementations
fd -e rs '\.(postgres|oracle)\.rs' | xargs rg -l 'create_memory|conversation_memory'Repository: lightseekorg/smg
Length of output: 903
🏁 Script executed:
# Look for STMO worker implementation to understand duplicate handling
fd -e rs | xargs rg -l 'stmo|STMO' | head -20Repository: lightseekorg/smg
Length of output: 147
🏁 Script executed:
# Check for uniqueness constraints in schema or implementation
rg -n 'UNIQUE|unique|PRIMARY' crates/data_connector/src/schema.rs | grep -i 'conversation_memory\|scope'Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Find Postgres and Oracle data connector implementations
fd -e rs 'postgres|oracle' crates/data_connector/src/Repository: lightseekorg/smg
Length of output: 229
🏁 Script executed:
# Search for STMO worker implementation
rg -n 'stmo_worker|StmoWorker|STMO.*worker' --type rsRepository: lightseekorg/smg
Length of output: 86
🏁 Script executed:
# Check if scope_id is used for uniqueness enforcement anywhere
rg -n 'scope_id' crates/data_connector/src/postgres.rs crates/data_connector/src/oracle.rs 2>/dev/null | head -30Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Look at schema definition for conversation_memory table
rg -A 20 'conversation_memory.*=>' crates/data_connector/src/schema.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Search for create_memory in postgres and oracle files
rg 'create_memory' crates/data_connector/src/postgres.rs crates/data_connector/src/oracle.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Look for STMO worker
rg -n 'struct.*Stmo|impl.*Stmo' --max-count 20Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check conversation_memory table schema in postgres migrations
rg -A 30 'conversation_memory' crates/data_connector/src/postgres_migrations.rs | head -60Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check for idempotency handling in STMO
rg -n 'idempotent|dedupe|duplicate' crates/workers/ 2>/dev/null || rg -n 'idempotent|dedupe|duplicate' --max-count 10Repository: lightseekorg/smg
Length of output: 16525
🏁 Script executed:
# Get postgres create_memory implementation
rg -A 30 'impl ConversationMemoryWriter' crates/data_connector/src/postgres.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Get oracle create_memory implementation
rg -A 30 'impl ConversationMemoryWriter' crates/data_connector/src/oracle.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check if there's a unique constraint on scope_id in migrations
rg -n 'scope_id|UNIQUE|unique' crates/data_connector/src/postgres_migrations.rs | head -40Repository: lightseekorg/smg
Length of output: 173
🏁 Script executed:
# Search for STMO worker or condensing logic
find . -name '*.rs' | xargs rg -l 'condense|Condense' | grep -v test | head -10Repository: lightseekorg/smg
Length of output: 204
🏁 Script executed:
# Search for ConversationMemoryWriter implementations across the codebase
rg 'impl.*ConversationMemoryWriter' --type rsRepository: lightseekorg/smg
Length of output: 86
🏁 Script executed:
# Check if there's any INSERT logic for conversation_memory in postgres/oracle
rg -n 'INSERT.*conversation_memory|conversation_memory.*INSERT' crates/data_connector/src/Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Look at the data_connector factory to see which writers are used
rg -A 20 'fn.*conversation_memory_writer|create_conversation_memory' crates/data_connector/src/Repository: lightseekorg/smg
Length of output: 1759
🏁 Script executed:
# Check for any dedup or uniqueness logic in the conversation memory handling
rg -B 5 -A 10 'conversation_memory' crates/data_connector/src/postgres.rs | head -60Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Try without -type flag
rg 'impl.*ConversationMemoryWriter'Repository: lightseekorg/smg
Length of output: 632
🏁 Script executed:
# List all files in data_connector src
ls -la crates/data_connector/src/Repository: lightseekorg/smg
Length of output: 1281
🏁 Script executed:
# Check if Postgres/Oracle writers exist in the factory or elsewhere
rg 'PostgresConversationMemoryWriter|OracleConversationMemoryWriter'Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Look at how conversation_memory rows are actually inserted
rg -n 'NewConversationMemory' crates/data_connector/src/postgres.rs crates/data_connector/src/oracle.rs 2>/dev/nullRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Search for conversation_memory table handling in postgres.rs
head -100 crates/data_connector/src/postgres.rsRepository: lightseekorg/smg
Length of output: 3681
🏁 Script executed:
# Look for the actual conversation_memory table structure
rg -B 5 -A 15 'conversation_memory' crates/data_connector/src/postgres.rs | head -100Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check factory to see which writer is actually instantiated
cat crates/data_connector/src/factory.rs | head -100Repository: lightseekorg/smg
Length of output: 4229
🏁 Script executed:
# Look for STMO-related code and understand the worker flow
rg -n 'STMO\|stmo' model_gateway/src/routers/common/persistence_utils.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Get the complete factory to see the Oracle/Postgres backend handling
grep -A 50 'HistoryBackend::Postgres =>' crates/data_connector/src/factory.rs | head -60Repository: lightseekorg/smg
Length of output: 1833
🏁 Script executed:
# Look at how enqueue_conversation_memory_rows is called and documented
rg -B 5 -A 15 'enqueue_conversation_memory_rows' model_gateway/src/routers/common/persistence_utils.rs | head -80Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check where STMO enqueue is called from
rg -B 10 'enqueue_stmo_if_needed' model_gateway/src/routers/common/persistence_utils.rsRepository: lightseekorg/smg
Length of output: 3068
🏁 Script executed:
# Look at the complete context of both enqueuing functions
cat model_gateway/src/routers/common/persistence_utils.rs | head -300 | tail -200Repository: lightseekorg/smg
Length of output: 8351
🏁 Script executed:
# Look at enqueue_conversation_memory_rows to understand the documented pattern
rg -B 10 -A 40 'fn enqueue_conversation_memory_rows' model_gateway/src/routers/common/persistence_utils.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check if there's any documentation or inline comment about LTM/OnDemand semantics
rg -B 5 -A 20 'LTM.*OnDemand|best-effort' model_gateway/src/routers/common/persistence_utils.rsRepository: lightseekorg/smg
Length of output: 1660
🏁 Script executed:
# Check if Postgres/Oracle backends create a real memory writer
rg -A 5 'create_postgres_storage\|create_oracle_storage' crates/data_connector/src/factory.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check what create_postgres_storage returns
rg -A 30 'fn create_postgres_storage' crates/data_connector/src/factory.rsRepository: lightseekorg/smg
Length of output: 1493
🏁 Script executed:
# Check what create_oracle_storage returns
rg -A 30 'fn create_oracle_storage' crates/data_connector/src/factory.rsRepository: lightseekorg/smg
Length of output: 1499
🏁 Script executed:
# Search for all ConversationMemoryWriter implementations to confirm which ones exist
rg 'ConversationMemoryWriter' crates/data_connector/src/ | grep -v test | grep -v '//'Repository: lightseekorg/smg
Length of output: 2070
Document STMO enqueue as best-effort, matching LTM/OnDemand semantics.
enqueue_stmo_if_needed currently sets scope_id = None and logs failures as best-effort. When a persistent ConversationMemoryWriter is implemented for Postgres/Oracle (currently only NoOp backends exist), duplicate STMO jobs can be enqueued on client/router retry at the same user_turns boundary, causing the worker to condense twice. Following the pattern established in PR #1343 for enqueue_conversation_memory_rows, clarify whether STMO enqueue should accept at-most-once semantics (duplicates allowed, idempotent consumer) or require dedup via scope_id. Document the chosen semantics inline so the future Postgres/Oracle writer implementations know whether to enforce uniqueness.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/routers/common/persistence_utils.rs` around lines 540 -
603, The STMO enqueue behavior is undocumented and currently sets scope_id =
None which can produce duplicate jobs on retries; update the code in
enqueue_stmo_if_needed to include an explicit inline comment documenting the
chosen semantic (either "best-effort / duplicates allowed; consumer must be
idempotent" OR "at-most-once / deduplicate by scope_id") and follow the same
pattern used in enqueue_conversation_memory_rows (PR `#1343`): if you decide
duplicates are acceptable, leave scope_id = None and state that clearly next to
NewConversationMemory creation and the create_memory call; if you decide to
enforce uniqueness, compute and set a deterministic scope_id (e.g., based on
conversation_id + user_turns + response_id or target_item_end) and document that
producers/writers must use scope_id to dedupe when implementing Postgres/Oracle
writers; ensure the comment references enqueue_stmo_if_needed,
NewConversationMemory.scope_id, and ConversationMemoryWriter.create_memory so
future writers know how to implement deduplication.
| #[cfg(test)] | ||
| mod tests { | ||
| use std::sync::{Arc, Mutex}; | ||
|
|
||
| use serde_json::Value; | ||
| use smg_data_connector::{ConversationMemoryId, ConversationMemoryResult}; | ||
|
|
||
| use super::*; | ||
|
|
||
| struct RecordingConversationMemoryWriter { | ||
| rows: Mutex<Vec<NewConversationMemory>>, | ||
| } | ||
|
|
||
| impl RecordingConversationMemoryWriter { | ||
| fn new() -> Self { | ||
| Self { | ||
| rows: Mutex::new(Vec::new()), | ||
| } | ||
| } | ||
| } | ||
|
|
||
| #[async_trait::async_trait] | ||
| impl ConversationMemoryWriter for RecordingConversationMemoryWriter { | ||
| async fn create_memory( | ||
| &self, | ||
| input: NewConversationMemory, | ||
| ) -> ConversationMemoryResult<ConversationMemoryId> { | ||
| self.rows.lock().expect("rows mutex poisoned").push(input); | ||
| Ok(ConversationMemoryId::from("mem_test")) | ||
| } | ||
| } | ||
|
|
||
| #[test] | ||
| fn stmo_turn_boundary_matches_expected_sequence() { | ||
| let cases = [ | ||
| (1, false), | ||
| (2, false), | ||
| (3, false), | ||
| (4, true), | ||
| (5, false), | ||
| (6, false), | ||
| (7, true), | ||
| (8, false), | ||
| (9, false), | ||
| (10, true), | ||
| (11, false), | ||
| (12, false), | ||
| (13, true), | ||
| ]; | ||
|
|
||
| for (turn, expected) in cases { | ||
| assert_eq!( | ||
| should_enqueue_stmo_for_current_turn(turn), | ||
| expected, | ||
| "turn={turn}" | ||
| ); | ||
| } | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn enqueue_stmo_if_needed_enqueues_expected_row_on_boundary() { | ||
| let writer = Arc::new(RecordingConversationMemoryWriter::new()); | ||
| let writer_dyn: Arc<dyn ConversationMemoryWriter> = writer.clone(); | ||
|
|
||
| let memory_execution_context = MemoryExecutionContext { | ||
| stm_enabled: true, | ||
| stm_condenser_model_id: Some("condense-1".to_string()), | ||
| ..MemoryExecutionContext::default() | ||
| }; | ||
|
|
||
| let output_items = vec![json!({ "type": "message", "role": "assistant" })]; | ||
| let conversation_id = ConversationId::from("conv_test"); | ||
| let response_id = ResponseId::from("resp_test"); | ||
|
|
||
| enqueue_stmo_if_needed( | ||
| &writer_dyn, | ||
| &memory_execution_context, | ||
| Some(4), | ||
| &conversation_id, | ||
| &response_id, | ||
| &output_items, | ||
| 4, | ||
| ) | ||
| .await; | ||
|
|
||
| let rows = writer.rows.lock().expect("rows mutex poisoned"); | ||
| assert_eq!(rows.len(), 1, "should enqueue exactly one STMO row"); | ||
|
|
||
| let row = &rows[0]; | ||
| assert_eq!(row.conversation_id, conversation_id); | ||
| assert_eq!(row.response_id, Some(response_id)); | ||
| assert_eq!(row.memory_type, ConversationMemoryType::Stmo); | ||
| assert_eq!(row.status, ConversationMemoryStatus::Ready); | ||
|
|
||
| let config = row | ||
| .memory_config | ||
| .as_deref() | ||
| .expect("memory_config should be set"); | ||
| let config_json: Value = | ||
| serde_json::from_str(config).expect("memory_config must be valid JSON"); | ||
|
|
||
| assert_eq!( | ||
| config_json | ||
| .get(STMO_CFG_KEY_CONDENSER_MODEL) | ||
| .and_then(Value::as_str), | ||
| Some("condense-1") | ||
| ); | ||
| assert_eq!( | ||
| config_json | ||
| .get(STMO_CFG_KEY_LAST_INDEX) | ||
| .and_then(Value::as_u64), | ||
| Some(4) | ||
| ); | ||
| assert_eq!( | ||
| config_json | ||
| .get(STMO_CFG_KEY_TARGET_ITEM_END) | ||
| .and_then(Value::as_u64), | ||
| Some(5) | ||
| ); | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn enqueue_stmo_if_needed_skips_when_turn_count_missing() { | ||
| let writer = Arc::new(RecordingConversationMemoryWriter::new()); | ||
| let writer_dyn: Arc<dyn ConversationMemoryWriter> = writer.clone(); | ||
|
|
||
| let memory_execution_context = MemoryExecutionContext { | ||
| stm_enabled: true, | ||
| stm_condenser_model_id: Some("condense-1".to_string()), | ||
| ..MemoryExecutionContext::default() | ||
| }; | ||
|
|
||
| enqueue_stmo_if_needed( | ||
| &writer_dyn, | ||
| &memory_execution_context, | ||
| None, | ||
| &ConversationId::from("conv_test"), | ||
| &ResponseId::from("resp_test"), | ||
| &[], | ||
| 0, | ||
| ) | ||
| .await; | ||
|
|
||
| let rows = writer.rows.lock().expect("rows mutex poisoned"); | ||
| assert!( | ||
| rows.is_empty(), | ||
| "should not enqueue when turn count is absent" | ||
| ); | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn enqueue_stmo_if_needed_enqueues_without_condenser_model() { | ||
| let writer = Arc::new(RecordingConversationMemoryWriter::new()); | ||
| let writer_dyn: Arc<dyn ConversationMemoryWriter> = writer.clone(); | ||
|
|
||
| let memory_execution_context = MemoryExecutionContext { | ||
| stm_enabled: true, | ||
| stm_condenser_model_id: None, | ||
| ..MemoryExecutionContext::default() | ||
| }; | ||
|
|
||
| enqueue_stmo_if_needed( | ||
| &writer_dyn, | ||
| &memory_execution_context, | ||
| Some(4), | ||
| &ConversationId::from("conv_test"), | ||
| &ResponseId::from("resp_test"), | ||
| &[json!({ "type": "message", "role": "assistant" })], | ||
| 4, | ||
| ) | ||
| .await; | ||
|
|
||
| let rows = writer.rows.lock().expect("rows mutex poisoned"); | ||
| assert_eq!(rows.len(), 1, "should enqueue without condenser model"); | ||
|
|
||
| let config = rows[0] | ||
| .memory_config | ||
| .as_deref() | ||
| .expect("memory_config should be set"); | ||
| let config_json: Value = | ||
| serde_json::from_str(config).expect("memory_config must be valid JSON"); | ||
| assert!(config_json.get(STMO_CFG_KEY_CONDENSER_MODEL).is_none()); | ||
| assert_eq!( | ||
| config_json | ||
| .get(STMO_CFG_KEY_LAST_INDEX) | ||
| .and_then(Value::as_u64), | ||
| Some(4) | ||
| ); | ||
| } |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Good test coverage for boundary + config shape.
The boundary table (1..=13) pins the 4/7/10/13 contract; the condenser-present / condenser-absent / missing-turn-count cases cover the enqueue decision tree well. One small gap: no test for the stm_enabled: false short-circuit — worth adding since the gate is the primary runtime switch that determines whether STMO ever fires.
♻️ Suggested additional test
+ #[tokio::test]
+ async fn enqueue_stmo_if_needed_skips_when_stm_disabled() {
+ let writer = Arc::new(RecordingConversationMemoryWriter::new());
+ let writer_dyn: Arc<dyn ConversationMemoryWriter> = writer.clone();
+
+ let memory_execution_context = MemoryExecutionContext {
+ stm_enabled: false,
+ stm_condenser_model_id: Some("condense-1".to_string()),
+ ..MemoryExecutionContext::default()
+ };
+
+ enqueue_stmo_if_needed(
+ &writer_dyn,
+ &memory_execution_context,
+ Some(4),
+ &ConversationId::from("conv_test"),
+ &ResponseId::from("resp_test"),
+ &[json!({ "type": "message", "role": "assistant" })],
+ 4,
+ )
+ .await;
+
+ assert!(writer.rows.lock().expect("rows mutex poisoned").is_empty());
+ }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/routers/common/persistence_utils.rs` around lines 605 -
793, Tests cover STMO boundaries and config shapes but miss the
stm_enabled=false short-circuit; add a tokio::test that constructs a
MemoryExecutionContext with stm_enabled: false (keep other fields like
stm_condenser_model_id as needed), use a RecordingConversationMemoryWriter and
call enqueue_stmo_if_needed with a turn that would normally enqueue (e.g.,
Some(4)), then assert that writer.rows remains empty; reference the existing
RecordingConversationMemoryWriter, MemoryExecutionContext, and
enqueue_stmo_if_needed to implement the test consistent with the other test
patterns.
Description
Problem
Solution
stm_enabledandstm_condenser_model_idfrom thex-conversation-memory-configheader throughMemoryHeaderViewandMemoryExecutionContextConversationMemoryWriterandMemoryExecutionContextintopersist_conversation_items; gRPC path uses a NoOp writerenqueue_stmo_if_needed— fires at turns 4, 7, 10, 13... matching the upstream worker trigger contractstm_enabledby the runtime memory switch — header alone is not sufficient to activate STMO schedulingChanges
model_gateway/src/memory/context.rsstm_enabledandstm_condenser_model_idtoMemoryExecutionContextstm_enabledbyruntime.enabled— same pattern as LTMstore_ltmandrecallmodel_gateway/src/routers/common/header_utils.rsstm_enabledandstm_condenser_model_idfields toMemoryHeaderViewx-conversation-memory-configJSON headermodel_gateway/src/routers/openai/responses/route.rsx-conversation-memory-configload_input_history, computeconversation_user_turn_count(gated onstm_enabled) and store inResponsesPayloadStatemodel_gateway/src/routers/openai/context.rsconversation_user_turn_count: Option<usize>toResponsesPayloadStateandStorageHandlesinto_streaming_contextmodel_gateway/src/routers/openai/responses/history.rscount_conversation_user_turns(input: &ResponseInput) -> usize— counts from typedResponseInputvariants, handles bothTextandItemsmodel_gateway/src/routers/common/persistence_utils.rsconversation_memory_writer,memory_execution_context, andconversation_user_turn_countparams topersist_conversation_itemsenqueue_stmo_if_needed— fires at turns 4, 7, 10, 13...model_gateway/src/routers/openai/responses/non_streaming.rsconversation_memory_writer,memory_execution_context, andconversation_user_turn_counttopersist_conversation_itemsmodel_gateway/src/routers/openai/responses/streaming.rsconversation_memory_writer,memory_execution_context, andconversation_user_turn_countat bothpersist_conversation_itemscall sitesmodel_gateway/src/routers/grpc/common/responses/utils.rsNoOpConversationMemoryWriterandMemoryExecutionContext::default()— gRPC path never has STMO enabledTest Plan
condenser_model,last_index,target_item_endin job configChecklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit