-
Notifications
You must be signed in to change notification settings - Fork 1.5k
feat(turns): opt-in async write-behind durability mode for the turn-state row store (#6263 Step 3) #6298
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
feat(turns): opt-in async write-behind durability mode for the turn-state row store (#6263 Step 3) #6298
Changes from all commits
c6b39ba
7d5be70
036c658
cf8f390
44ba506
66ac5f7
06ee5ad
d218e32
9c91db4
c12c7b1
d3810e0
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -218,6 +218,104 @@ async fn stub_gateway_send_cancels_recovery_required_and_releases_conversation() | |
| runtime.shutdown().await.unwrap(); | ||
| } | ||
|
|
||
| /// Minimal completing model gateway: every model call returns a plain assistant | ||
| /// reply, so a turn reaches `TurnStatus::Completed` without needing a real LLM. | ||
| #[cfg(feature = "inmemory-turn-state")] | ||
| #[derive(Default)] | ||
| struct AlwaysReplyGateway; | ||
|
|
||
| #[cfg(feature = "inmemory-turn-state")] | ||
| #[async_trait] | ||
| impl HostManagedModelGateway for AlwaysReplyGateway { | ||
| async fn stream_model( | ||
| &self, | ||
| _request: HostManagedModelRequest, | ||
| ) -> Result<HostManagedModelResponse, HostManagedModelError> { | ||
| Ok(HostManagedModelResponse::assistant_reply( | ||
| "done".to_string(), | ||
| )) | ||
| } | ||
|
|
||
| async fn stream_model_with_capabilities( | ||
| &self, | ||
| _request: HostManagedModelRequest, | ||
| _capabilities: Arc<dyn LoopCapabilityPort>, | ||
| ) -> Result<HostManagedModelResponse, HostManagedModelError> { | ||
| Ok(HostManagedModelResponse::assistant_reply( | ||
| "done".to_string(), | ||
| )) | ||
| } | ||
| } | ||
|
|
||
| /// #6263 Step 4 — production-flip wiring at the composition seam. With | ||
| /// `inmemory-turn-state` on, `build_reborn_runtime` composes the durable turn-state | ||
| /// ROW store (`factory.rs`), replacing the former in-memory authority + | ||
| /// block-persistence snapshot. This drives a real turn end to end over that store | ||
| /// (submit → claim → terminal, through the production runtime), then gracefully | ||
| /// `shutdown()`s — which routes through `RebornRuntime::shutdown → | ||
| /// FilesystemTurnStateStoreKind::drain`. The profile ships the row store at the | ||
| /// `WriteThrough` default (WriteBehind is blocked on the row store's non-cache-aware | ||
| /// query paths — see the factory arm), so the shutdown drain is a no-op here; the | ||
| /// test locks that composing the flipped store, serving a real turn over it, and | ||
| /// draining on shutdown all succeed without error/hang/panic. | ||
| /// | ||
| /// Deeper durability is pinned one tier down, over the raw store where | ||
| /// scope/backend are controlled precisely: terminal/gate-park recovery across a | ||
| /// store reopen and the drain-flushes-the-tail contract in | ||
| /// `ironclaw_turns::row_store_crash_consistency` (incl. | ||
| /// `write_behind_drain_flushes_the_async_tail_for_graceful_restart`), and the | ||
| /// block-persistence→row migration in | ||
| /// `filesystem_turn_state_contract::filesystem_turn_state_row_store_migrates_block_persistence_gate_park_snapshot`. | ||
| #[cfg(feature = "inmemory-turn-state")] | ||
| #[tokio::test] | ||
| async fn inmemory_turn_state_row_store_serves_turn_and_drains_on_shutdown() { | ||
| let _guard = runtime_composition_test_guard().await; | ||
| let root = tempfile::tempdir().unwrap(); | ||
| let input = RebornRuntimeInput::from_services( | ||
| RebornBuildInput::local_dev("wb-durable-owner", root.path().join("local-dev")) | ||
| .with_runtime_policy(local_dev_runtime_policy()), | ||
| ) | ||
| .with_identity(RebornRuntimeIdentity { | ||
| tenant_id: "wb-durable-tenant".to_string(), | ||
| agent_id: "wb-durable-agent".to_string(), | ||
| source_binding_id: "wb-durable-source".to_string(), | ||
| reply_target_binding_id: "wb-durable-reply".to_string(), | ||
| }) | ||
| .with_runner_settings( | ||
| TurnRunnerSettings::default() | ||
| .set_heartbeat_interval(Duration::from_secs(60)) | ||
| .set_poll_interval(Duration::from_secs(60)), | ||
| ) | ||
| .with_model_gateway_override(Arc::new(AlwaysReplyGateway)); | ||
|
|
||
| // Compose the durable row store via the production build path and drive a real | ||
| // turn to Completed over it: proves the flipped store serves the full | ||
| // submit → claim → terminal transition set through the production runtime. | ||
| let runtime = build_reborn_runtime(input).await.unwrap(); | ||
| let conversation = runtime.new_conversation().await.unwrap(); | ||
| let reply = tokio::time::timeout( | ||
| SEND_USER_MESSAGE_TIMEOUT, | ||
| runtime.send_user_message(&conversation, "durable please"), | ||
| ) | ||
| .await | ||
| .unwrap() | ||
| .unwrap(); | ||
| assert_eq!( | ||
| reply.status, | ||
| TurnStatus::Completed, | ||
| "turn must complete over the WriteBehind store, got {:?} ({:?})", | ||
| reply.status, | ||
| reply.failure_category | ||
| ); | ||
|
|
||
| // Graceful shutdown drains the WriteBehind tail through | ||
| // `FilesystemTurnStateStoreKind::drain`; a broken drain wiring surfaces here. | ||
|
Comment on lines
+303
to
+312
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value Assertion message and comment say "WriteBehind" but this store ships The test's own doc comment states the profile ships the row store at the 🤖 Prompt for AI Agents |
||
| runtime | ||
| .shutdown() | ||
| .await | ||
| .expect("graceful shutdown drains the WriteBehind tail without error"); | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn send_user_message_with_cancellation_cancels_submitted_run() { | ||
| let _guard = runtime_composition_test_guard().await; | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win
warn!on the shutdown-drain path contradicts the REPL/TUI logging rule.RebornRuntime::shutdownis REPL/TUI-reachable, and this same impl block already chosedebug!overwarn!for exactly this reason (wait_for_terminal_or_gate, ~Line 2745: "debug!notwarn!per the logging rule — this runtime is REPL/TUI-reachable"). A degraded-drain diagnostic is internal, not intentionally-rendered user status, so preferdebug!here for consistency.As per path instructions: "REPL/TUI logging: info!/warn! corrupt the terminal UI — internal diagnostics use debug!".
Proposed change
📝 Committable suggestion
🤖 Prompt for AI Agents
Source: Path instructions