GitHub issue #3451: [Reborn] Add direct DB operations for loop checkpoint mappings - #3468
Conversation
There was a problem hiding this comment.
Code Review
This pull request introduces a significant architectural expansion to the IronClaw Reborn system, adding several new crates for conversation management, event projections, runtime policy, and durable event storage. It implements LibSql and PostgreSQL backends for capability leases, run states, and conversation bindings, alongside a new runtime planner and production wiring validation. Feedback highlights several database-related improvements, including the need to replace inefficient wipe-and-reload persistence patterns and full-state memory loading with targeted read/write paths to prevent memory exhaustion. Additionally, the reviewer recommends decomposing the packed owner_key JSON column into individual columns for better queryability and optimizing lease lookups by filtering at the database level. Finally, a documentation discrepancy was noted in the ironclaw_conversations guardrails regarding the inclusion of durable adapters.
| async fn load_state_from_conn( | ||
| conn: &::libsql::Connection, | ||
| ) -> Result<PersistedConversationState, InboundTurnError> { | ||
| let revision = load_revision(conn).await?; | ||
| let mut state = InMemoryState::default(); | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| "SELECT key_payload, user_id FROM reborn_conversation_actor_pairings", | ||
| (), | ||
| ) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| let key: ActorKey = from_json(&row.get::<String>(0).map_err(db_error)?)?; | ||
| let user_id = ironclaw_host_api::UserId::new(row.get::<String>(1).map_err(db_error)?) | ||
| .map_err(|error| InboundTurnError::DurableState { | ||
| reason: error.to_string(), | ||
| })?; | ||
| state.pairings.insert(key, user_id); | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| "SELECT key_payload, payload FROM reborn_conversation_bindings", | ||
| (), | ||
| ) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| let key: BindingKey = from_json(&row.get::<String>(0).map_err(db_error)?)?; | ||
| let binding: BindingRecord = from_json(&row.get::<String>(1).map_err(db_error)?)?; | ||
| state.source_bindings.insert( | ||
| binding.source_binding_ref.as_str().to_string(), | ||
| binding.clone(), | ||
| ); | ||
| state.bindings.insert(key, binding); | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query("SELECT payload FROM reborn_conversation_reply_targets", ()) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| let reply_target: ReplyTargetRecord = from_json(&row.get::<String>(0).map_err(db_error)?)?; | ||
| state.reply_targets.insert( | ||
| reply_target.reply_target_binding_ref.as_str().to_string(), | ||
| reply_target, | ||
| ); | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query("SELECT payload FROM reborn_conversation_threads", ()) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| let thread_key_and_record: (ThreadKey, ThreadRecord) = | ||
| from_json(&row.get::<String>(0).map_err(db_error)?)?; | ||
| state | ||
| .threads | ||
| .insert(thread_key_and_record.0, thread_key_and_record.1); | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| "SELECT tenant_id, thread_id, user_id FROM reborn_conversation_thread_participants", | ||
| (), | ||
| ) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| let tenant_id = ironclaw_host_api::TenantId::new(row.get::<String>(0).map_err(db_error)?) | ||
| .map_err(|error| InboundTurnError::DurableState { | ||
| reason: error.to_string(), | ||
| })?; | ||
| let thread_id = ironclaw_host_api::ThreadId::new(row.get::<String>(1).map_err(db_error)?) | ||
| .map_err(|error| InboundTurnError::DurableState { | ||
| reason: error.to_string(), | ||
| })?; | ||
| let user_id = ironclaw_host_api::UserId::new(row.get::<String>(2).map_err(db_error)?) | ||
| .map_err(|error| InboundTurnError::DurableState { | ||
| reason: error.to_string(), | ||
| })?; | ||
| if let Some(thread) = state | ||
| .threads | ||
| .get_mut(&ThreadKey::new(&tenant_id, &thread_id)) | ||
| { | ||
| thread.participants.insert(user_id); | ||
| } | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| "SELECT key_payload, identity_payload FROM reborn_conversation_external_event_routes", | ||
| (), | ||
| ) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| state.external_event_routes.insert( | ||
| from_json(&row.get::<String>(0).map_err(db_error)?)?, | ||
| from_json(&row.get::<String>(1).map_err(db_error)?)?, | ||
| ); | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| "SELECT payload FROM reborn_conversation_accepted_messages", | ||
| (), | ||
| ) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| let message: ThreadMessageRecord = from_json(&row.get::<String>(0).map_err(db_error)?)?; | ||
| let idempotency_key = MessageIdempotencyKey { | ||
| tenant_id: message.accepted.tenant_id.clone(), | ||
| source_binding_ref: message.accepted.source_binding_ref.as_str().to_string(), | ||
| external_event_id: message.external_event_id.clone(), | ||
| }; | ||
| state | ||
| .message_idempotency | ||
| .insert(idempotency_key, message.accepted.clone()); | ||
| state.messages.push(message); | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| "SELECT key_payload, payload FROM reborn_conversation_message_replays", | ||
| (), | ||
| ) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| state.message_replays.insert( | ||
| from_json(&row.get::<String>(0).map_err(db_error)?)?, | ||
| from_json(&row.get::<String>(1).map_err(db_error)?)?, | ||
| ); | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| "SELECT message_ref, idempotency_key FROM reborn_conversation_submission_keys", | ||
| (), | ||
| ) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| let message_ref = | ||
| ironclaw_turns::AcceptedMessageRef::new(row.get::<String>(0).map_err(db_error)?) | ||
| .map_err(|reason| InboundTurnError::InvalidCanonicalRef { reason })?; | ||
| let key = ironclaw_turns::IdempotencyKey::new(row.get::<String>(1).map_err(db_error)?) | ||
| .map_err(|reason| InboundTurnError::InvalidCanonicalRef { reason })?; | ||
| state.submission_keys.insert(message_ref, key); | ||
| } | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| "SELECT message_ref, payload FROM reborn_conversation_submit_responses", | ||
| (), | ||
| ) | ||
| .await | ||
| .map_err(db_error)?; | ||
| while let Some(row) = rows.next().await.map_err(db_error)? { | ||
| let message_ref = | ||
| ironclaw_turns::AcceptedMessageRef::new(row.get::<String>(0).map_err(db_error)?) | ||
| .map_err(|reason| InboundTurnError::InvalidCanonicalRef { reason })?; | ||
| state.submitted_message_responses.insert( | ||
| message_ref, | ||
| from_json(&row.get::<String>(1).map_err(db_error)?)?, | ||
| ); | ||
| } | ||
|
|
||
| Ok(PersistedConversationState { state, revision }) | ||
| } |
There was a problem hiding this comment.
The load_state function fetches the entire database content into memory. While this aligns with the current snapshot-adapter design pattern, it poses a risk of memory exhaustion as data grows. Per project guidelines, we should defer this optimization for now, but you must document the requirement for targeted read paths as a follow-up task.
References
- If a design pattern like snapshot-adapter necessitates loading full state, defer optimization and document the requirement for targeted read paths as a follow-up task.
| conn.execute(&format!("DELETE FROM {table}"), ()) | ||
| .await | ||
| .map_err(db_error)?; | ||
| } |
There was a problem hiding this comment.
The save_state implementation wipes all tables and re-inserts the entire state. While this is part of the snapshot-adapter pattern, it is inefficient. As per our rules, we should defer the optimization to granular operations for now, but please document the need for a more efficient persistence strategy as a follow-up.
References
- If a design pattern like snapshot-adapter necessitates an inefficient implementation, defer optimization and document the requirement for targeted paths as a follow-up task.
| async fn load_state_from_txn( | ||
| txn: &deadpool_postgres::Transaction<'_>, | ||
| ) -> Result<PersistedConversationState, InboundTurnError> { | ||
| let revision = load_revision(txn).await?; | ||
| let mut state = InMemoryState::default(); | ||
|
|
||
| for row in txn | ||
| .query( | ||
| "SELECT key_payload, user_id FROM reborn_conversation_actor_pairings", | ||
| &[], | ||
| ) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| let key: ActorKey = from_json(row.get::<_, &str>(0))?; | ||
| let user_id = ironclaw_host_api::UserId::new(row.get::<_, String>(1)).map_err(|error| { | ||
| InboundTurnError::DurableState { | ||
| reason: error.to_string(), | ||
| } | ||
| })?; | ||
| state.pairings.insert(key, user_id); | ||
| } | ||
|
|
||
| for row in txn | ||
| .query( | ||
| "SELECT key_payload, payload FROM reborn_conversation_bindings", | ||
| &[], | ||
| ) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| let key: BindingKey = from_json(row.get::<_, &str>(0))?; | ||
| let binding: BindingRecord = from_json(row.get::<_, &str>(1))?; | ||
| state.source_bindings.insert( | ||
| binding.source_binding_ref.as_str().to_string(), | ||
| binding.clone(), | ||
| ); | ||
| state.bindings.insert(key, binding); | ||
| } | ||
|
|
||
| for row in txn | ||
| .query("SELECT payload FROM reborn_conversation_reply_targets", &[]) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| let reply_target: ReplyTargetRecord = from_json(row.get::<_, &str>(0))?; | ||
| state.reply_targets.insert( | ||
| reply_target.reply_target_binding_ref.as_str().to_string(), | ||
| reply_target, | ||
| ); | ||
| } | ||
|
|
||
| for row in txn | ||
| .query("SELECT payload FROM reborn_conversation_threads", &[]) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| let (key, record): (ThreadKey, ThreadRecord) = from_json(row.get::<_, &str>(0))?; | ||
| state.threads.insert(key, record); | ||
| } | ||
|
|
||
| for row in txn | ||
| .query( | ||
| "SELECT tenant_id, thread_id, user_id FROM reborn_conversation_thread_participants", | ||
| &[], | ||
| ) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| let tenant_id = | ||
| ironclaw_host_api::TenantId::new(row.get::<_, String>(0)).map_err(|error| { | ||
| InboundTurnError::DurableState { | ||
| reason: error.to_string(), | ||
| } | ||
| })?; | ||
| let thread_id = | ||
| ironclaw_host_api::ThreadId::new(row.get::<_, String>(1)).map_err(|error| { | ||
| InboundTurnError::DurableState { | ||
| reason: error.to_string(), | ||
| } | ||
| })?; | ||
| let user_id = ironclaw_host_api::UserId::new(row.get::<_, String>(2)).map_err(|error| { | ||
| InboundTurnError::DurableState { | ||
| reason: error.to_string(), | ||
| } | ||
| })?; | ||
| if let Some(thread) = state | ||
| .threads | ||
| .get_mut(&ThreadKey::new(&tenant_id, &thread_id)) | ||
| { | ||
| thread.participants.insert(user_id); | ||
| } | ||
| } | ||
|
|
||
| for row in txn | ||
| .query( | ||
| "SELECT key_payload, identity_payload FROM reborn_conversation_external_event_routes", | ||
| &[], | ||
| ) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| state.external_event_routes.insert( | ||
| from_json(row.get::<_, &str>(0))?, | ||
| from_json(row.get::<_, &str>(1))?, | ||
| ); | ||
| } | ||
|
|
||
| for row in txn | ||
| .query( | ||
| "SELECT payload FROM reborn_conversation_accepted_messages", | ||
| &[], | ||
| ) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| let message: ThreadMessageRecord = from_json(row.get::<_, &str>(0))?; | ||
| let idempotency_key = MessageIdempotencyKey { | ||
| tenant_id: message.accepted.tenant_id.clone(), | ||
| source_binding_ref: message.accepted.source_binding_ref.as_str().to_string(), | ||
| external_event_id: message.external_event_id.clone(), | ||
| }; | ||
| state | ||
| .message_idempotency | ||
| .insert(idempotency_key, message.accepted.clone()); | ||
| state.messages.push(message); | ||
| } | ||
|
|
||
| for row in txn | ||
| .query( | ||
| "SELECT key_payload, payload FROM reborn_conversation_message_replays", | ||
| &[], | ||
| ) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| state.message_replays.insert( | ||
| from_json(row.get::<_, &str>(0))?, | ||
| from_json(row.get::<_, &str>(1))?, | ||
| ); | ||
| } | ||
|
|
||
| for row in txn | ||
| .query( | ||
| "SELECT message_ref, idempotency_key FROM reborn_conversation_submission_keys", | ||
| &[], | ||
| ) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| let message_ref = ironclaw_turns::AcceptedMessageRef::new(row.get::<_, String>(0)) | ||
| .map_err(|reason| InboundTurnError::InvalidCanonicalRef { reason })?; | ||
| let key = ironclaw_turns::IdempotencyKey::new(row.get::<_, String>(1)) | ||
| .map_err(|reason| InboundTurnError::InvalidCanonicalRef { reason })?; | ||
| state.submission_keys.insert(message_ref, key); | ||
| } | ||
|
|
||
| for row in txn | ||
| .query( | ||
| "SELECT message_ref, payload FROM reborn_conversation_submit_responses", | ||
| &[], | ||
| ) | ||
| .await | ||
| .map_err(pg_error)? | ||
| { | ||
| let message_ref = ironclaw_turns::AcceptedMessageRef::new(row.get::<_, String>(0)) | ||
| .map_err(|reason| InboundTurnError::InvalidCanonicalRef { reason })?; | ||
| state | ||
| .submitted_message_responses | ||
| .insert(message_ref, from_json(row.get::<_, &str>(1))?); | ||
| } | ||
|
|
||
| Ok(PersistedConversationState { state, revision }) | ||
| } |
There was a problem hiding this comment.
The load_state implementation for Postgres fetches the entire state into memory. Following the snapshot-adapter pattern, we will defer optimization for now, but please ensure the requirement for targeted read paths is documented as a follow-up task to mitigate OOM risks.
References
- If a design pattern like snapshot-adapter necessitates loading full state, defer optimization and document the requirement for targeted read paths as a follow-up task.
| txn.execute(&format!("DELETE FROM {table}"), &[]) | ||
| .await | ||
| .map_err(pg_error)?; | ||
| } |
There was a problem hiding this comment.
The save_state implementation for Postgres uses the destructive wipe-and-reload pattern. Since this follows the snapshot-adapter design, we will defer refactoring to granular updates, but this requirement must be documented as a follow-up task.
References
- If a design pattern like snapshot-adapter necessitates an inefficient implementation, defer optimization and document the requirement for targeted paths as a follow-up task.
| CREATE TABLE IF NOT EXISTS reborn_capability_lease_records ( | ||
| owner_key TEXT NOT NULL, | ||
| invocation_id TEXT NOT NULL, | ||
| lease_id TEXT NOT NULL, | ||
| status TEXT NOT NULL, | ||
| payload TEXT NOT NULL, | ||
| PRIMARY KEY (owner_key, invocation_id, lease_id) | ||
| ); |
There was a problem hiding this comment.
Packing multiple identifiers into a single JSON string column (owner_key) and using it as a primary key component is a poor database design. It prevents efficient indexing and querying of individual fields (like user_id or project_id) for analytics or cross-invocation management. Consider using individual columns for each identifier in the ResourceScope to allow for better query performance and data integrity.
References
- To prevent key collisions and improve queryability, use structured keys or individual columns instead of string concatenation or packed strings.
| conn: &libsql::Connection, | ||
| scope: &ResourceScope, | ||
| ) -> Result<Vec<CapabilityLease>, CapabilityLeaseError> { | ||
| let mut rows = conn.query("SELECT invocation_id, lease_id, status, payload FROM reborn_capability_lease_records WHERE owner_key = ?1 ORDER BY lease_id", libsql::params![owner_key(scope)?]).await.map_err(db_error)?; |
There was a problem hiding this comment.
The query in libsql_leases_for_scope filters only by owner_key, but the ResourceScope passed in contains a specific invocation_id. Since the table has an invocation_id column and the primary consumer (active_leases_for_context) filters by it in memory, it would be significantly more efficient to perform this filtering at the database level to reduce data transfer and processing overhead.
let mut rows = conn.query("SELECT invocation_id, lease_id, status, payload FROM reborn_capability_lease_records WHERE owner_key = ?1 AND invocation_id = ?2 ORDER BY lease_id", libsql::params![owner_key(scope)?, scope.invocation_id.to_string()]).await.map_err(db_error)?;References
- Use targeted database queries to fetch specific records instead of loading all records and filtering in the application to prevent performance bottlenecks.
| client: &impl deadpool_postgres::GenericClient, | ||
| scope: &ResourceScope, | ||
| ) -> Result<Vec<CapabilityLease>, CapabilityLeaseError> { | ||
| let rows = client.query("SELECT invocation_id, lease_id, status, payload::text FROM reborn_capability_lease_records WHERE owner_key = $1 ORDER BY lease_id", &[&owner_key(scope)?]).await.map_err(db_error)?; |
There was a problem hiding this comment.
Similar to the LibSql implementation, postgres_leases_for_scope should filter by invocation_id in the SQL query to avoid fetching unnecessary rows and filtering them in memory.
let rows = client.query("SELECT invocation_id, lease_id, status, payload::text FROM reborn_capability_lease_records WHERE owner_key = $1 AND invocation_id = $2 ORDER BY lease_id", &[&owner_key(scope)?, &scope.invocation_id.to_string()]).await.map_err(db_error)?;References
- Use targeted database queries to fetch specific records instead of loading all records and filtering in the application to prevent performance bottlenecks.
| - Preserve typed `ironclaw_turns::TurnError` values across this boundary instead of flattening turn failures to strings. | ||
| - Automatic first-contact binding must not trust raw adapter-supplied agent/project scope hints; use explicit linking or a future trusted thread-creation seam for scoped thread selection. | ||
| - Explicit links are idempotent only for the same target thread; never silently retarget an already-bound external conversation to another thread. | ||
| - Keep durable PostgreSQL/libSQL adapters out of this crate until the transcript/thread storage boundary has a scoped implementation plan with parity tests. |
There was a problem hiding this comment.
The addition of durable PostgreSQL/libSQL adapters in this PR directly contradicts the guardrail stating they should be kept out of this crate until a scoped implementation plan exists. If this PR constitutes that plan, the guardrail should be updated to reflect the new status.
References
- Documentation for complex logic and architectural state must precisely match the code implementation and distinguish between current and future states.
ReviewSummary: Replaces wipe-and-reload snapshot persistence for Findings:
Minor: test helpers use Verdict: Request changes — (1) and (2) are blocking. Tests and dual-backend coverage otherwise solid. |
…e LoopCheckpointStore trait
- Add inline safety comments required by no-panics gate for static checkpoint sentinel refs
|
@zmanian addressed in
Verification:
|
zmanian
left a comment
There was a problem hiding this comment.
Re-review: approved
All five blocking findings resolved at a37f794d:
.unwrap()in production — RESOLVED.store.rs:113-120replaces both sites with serde defaults:default_checkpoint_kind()returnsLoopCheckpointKind::BeforeBlockanddefault_checkpoint_state_ref()returns a new infallibleLoopCheckpointStateRef::legacy_unknown().state_refis now a real plumbed field onBlockRunRequest/TurnRunnerOutcome::Blocked/record_checkpointrather than a sentinel.- Non-transactional INSERT+SELECT race — RESOLVED. Both backends do a single atomic UPSERT with a guard. libSQL:
INSERT ... ON CONFLICT(checkpoint_id) DO UPDATE SET checkpoint_id = ... WHERE payload = excluded.payloadreturning row count (rows == 1⇒ success,0⇒Conflict). Postgres mirrors withRETURNING payload::textinsideclient.transaction(). - JSON re-deserialization drift — RESOLVED.
ensure_loop_checkpoint_record_matches_requestcompares persistedscope/turn_id/run_id/checkpoint_idagainst the request and raisesTurnError::Conflicton drift — silent miss replaced by loud error. - libSQL migration robustness — RESOLVED.
libsql_column_existsviaPRAGMA table_infowith identifier whitelist replaces the string match; Postgres usesADD COLUMN IF NOT EXISTS. LegacyNOT NULL DEFAULT ''rows documented indocs/reborn/contracts/turn-persistence.mdwith explicit future-read-path requirements. - Brittle
is_empty()tests — RESOLVED. New contract tests assert exact loop-checkpoint cardinality + absence of state_ref inturn_checkpoints, plus cross-scope/cross-run miss and drift conflict for both backends.
Substantive fixes, not annotations — particularly the state_ref plumbing in #1.
GitHub issue nearai#3451: [Reborn] Add direct DB operations for loop checkpoint mappings
Refs #3451
KB task: KB-002
GitHub issue #3451: [Reborn] Add direct DB operations for loop checkpoint mappings
URL: #3451
Repo: nearai/ironclaw
Labels: risk: medium, scope: db, reborn
Assignees: none
Scope: IronClaw Reborn issue imported from GitHub.
Before coding: read full issue body/comments, check linked PRs/blockers, work from current origin/reborn-integration unless issue says otherwise, keep PR tightly scoped.
Acceptance: satisfy GitHub issue acceptance criteria; include tests/verification evidence; never merge without explicit user approval.
Auto-opened by kb when task reached In Review. Auto-merge remains disabled; do not merge without explicit operator approval.