From 045a60cdb010edafa6f2ee4045f35435cb2a8d34 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Tue, 7 Jul 2026 18:59:58 +0800 Subject: [PATCH 01/16] fix(reborn): persist chat timeline timestamps --- .../src/compaction_task.rs | 2 + .../src/inbound_turn.rs | 2 + .../tests/reborn_services_contract.rs | 2 + .../src/loop_exit_applier/tests/mod.rs | 4 + .../src/observability/trace_capture.rs | 6 ++ crates/ironclaw_threads/src/contract.rs | 11 +++ .../src/filesystem_service.rs | 24 ++++++ crates/ironclaw_threads/src/in_memory.rs | 44 ++++++++-- .../filesystem_session_thread_contract.rs | 85 +++++++++++++++++++ .../tests/session_thread_contract.rs | 83 ++++++++++++++++++ .../js/pages/chat/lib/history-messages.js | 17 +++- .../pages/chat/lib/history-messages.test.mjs | 32 +++++++ 12 files changed, 302 insertions(+), 10 deletions(-) diff --git a/crates/ironclaw_loop_support/src/compaction_task.rs b/crates/ironclaw_loop_support/src/compaction_task.rs index 8c3abe1368d..077fe84a4ba 100644 --- a/crates/ironclaw_loop_support/src/compaction_task.rs +++ b/crates/ironclaw_loop_support/src/compaction_task.rs @@ -787,6 +787,8 @@ mod tests { sequence: 1, kind, status: MessageStatus::Finalized, + created_at: None, + updated_at: None, actor_id: None, source_binding_id: None, reply_target_binding_id: None, diff --git a/crates/ironclaw_product_workflow/src/inbound_turn.rs b/crates/ironclaw_product_workflow/src/inbound_turn.rs index 03d1cb1c59b..23055748768 100644 --- a/crates/ironclaw_product_workflow/src/inbound_turn.rs +++ b/crates/ironclaw_product_workflow/src/inbound_turn.rs @@ -1163,6 +1163,8 @@ mod tests { sequence: 1, kind: ironclaw_threads::MessageKind::User, status: ironclaw_threads::MessageStatus::Submitted, + created_at: None, + updated_at: None, actor_id: None, source_binding_id: None, reply_target_binding_id: None, diff --git a/crates/ironclaw_product_workflow/tests/reborn_services_contract.rs b/crates/ironclaw_product_workflow/tests/reborn_services_contract.rs index d568568f912..2ecc10f16d4 100644 --- a/crates/ironclaw_product_workflow/tests/reborn_services_contract.rs +++ b/crates/ironclaw_product_workflow/tests/reborn_services_contract.rs @@ -186,6 +186,8 @@ fn fake_thread_history(owner: &WebUiAuthenticatedCaller, thread_id: &str) -> Thr sequence: 1, kind: MessageKind::User, status: MessageStatus::Submitted, + created_at: None, + updated_at: None, actor_id: Some(owner.user_id.as_str().to_string()), source_binding_id: Some("webui-src:test".to_string()), reply_target_binding_id: Some("webui-reply:test".to_string()), diff --git a/crates/ironclaw_reborn/src/loop_exit_applier/tests/mod.rs b/crates/ironclaw_reborn/src/loop_exit_applier/tests/mod.rs index 77114c01c12..d26f4f83b2b 100644 --- a/crates/ironclaw_reborn/src/loop_exit_applier/tests/mod.rs +++ b/crates/ironclaw_reborn/src/loop_exit_applier/tests/mod.rs @@ -843,6 +843,8 @@ async fn thread_checkpoint_evidence_rejects_wrong_run_and_malformed_result_ref_r sequence: next_sequence, kind: MessageKind::ToolResultReference, status: MessageStatus::Finalized, + created_at: None, + updated_at: None, actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -860,6 +862,8 @@ async fn thread_checkpoint_evidence_rejects_wrong_run_and_malformed_result_ref_r sequence: next_sequence + 1, kind: MessageKind::ToolResultReference, status: MessageStatus::Finalized, + created_at: None, + updated_at: None, actor_id: None, source_binding_id: None, reply_target_binding_id: None, diff --git a/crates/ironclaw_reborn_composition/src/observability/trace_capture.rs b/crates/ironclaw_reborn_composition/src/observability/trace_capture.rs index d011458a9aa..0752d736dca 100644 --- a/crates/ironclaw_reborn_composition/src/observability/trace_capture.rs +++ b/crates/ironclaw_reborn_composition/src/observability/trace_capture.rs @@ -106,6 +106,8 @@ fn context_window_to_records(window: ContextWindow) -> Vec sequence: message.sequence, kind: message.kind, status: MessageStatus::Finalized, + created_at: None, + updated_at: None, actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -545,6 +547,8 @@ mod tests { sequence: 0, kind, status, + created_at: None, + updated_at: None, actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -693,6 +697,8 @@ mod tests { sequence: 0, kind: MessageKind::ToolResultReference, status: MessageStatus::Finalized, + created_at: None, + updated_at: None, actor_id: None, source_binding_id: None, reply_target_binding_id: None, diff --git a/crates/ironclaw_threads/src/contract.rs b/crates/ironclaw_threads/src/contract.rs index 361e4dcee37..b07eb7c5681 100644 --- a/crates/ironclaw_threads/src/contract.rs +++ b/crates/ironclaw_threads/src/contract.rs @@ -203,6 +203,15 @@ pub struct ThreadMessageRecord { pub sequence: u64, pub kind: MessageKind, pub status: MessageStatus, + /// When this message row was first persisted. `None` for legacy records + /// written before per-message timestamps existed. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub created_at: Option>, + /// Last time the message row was materially changed. Draft finalization + /// updates this so UI refreshes can show the reply completion time instead + /// of the browser refresh time. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub updated_at: Option>, pub actor_id: Option, pub source_binding_id: Option, pub reply_target_binding_id: Option, @@ -717,6 +726,8 @@ mod tests { }"#; let record: ThreadMessageRecord = serde_json::from_str(json).unwrap(); assert!(record.attachments.is_empty()); + assert_eq!(record.created_at, None); + assert_eq!(record.updated_at, None); assert_eq!(record.content.as_deref(), Some("legacy row")); } } diff --git a/crates/ironclaw_threads/src/filesystem_service.rs b/crates/ironclaw_threads/src/filesystem_service.rs index 4e551d31f90..225a84c0735 100644 --- a/crates/ironclaw_threads/src/filesystem_service.rs +++ b/crates/ironclaw_threads/src/filesystem_service.rs @@ -1435,12 +1435,15 @@ where // transactional backends the thread counter, message, sequence index, // and idempotency record commit together; fallback backends reserve // immediately before the legacy message write. + let now = Utc::now(); let mut message = ThreadMessageRecord { message_id, thread_id: thread_id.clone(), sequence: 0, kind: MessageKind::User, status: MessageStatus::Accepted, + created_at: Some(now), + updated_at: Some(now), actor_id: Some(actor_id.clone()), source_binding_id: source_binding_id.clone(), reply_target_binding_id: reply_target_binding_id.clone(), @@ -1645,12 +1648,15 @@ where let sequence = self .reserve_sequence(&request.scope, &request.thread_id) .await?; + let now = Utc::now(); let message = ThreadMessageRecord { message_id: ThreadMessageId::new(), thread_id: request.thread_id.clone(), sequence, kind: MessageKind::Assistant, status: MessageStatus::Draft, + created_at: Some(now), + updated_at: Some(now), actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -1708,9 +1714,11 @@ where existing.message_id, |message| { ensure_draft(message)?; + let now = Utc::now(); message.status = MessageStatus::Finalized; message.content = Some(content.clone().into_text()); message.attachments = Vec::new(); + message.updated_at = Some(now); Ok(()) }, ) @@ -1724,12 +1732,15 @@ where let sequence = self .reserve_sequence(&request.scope, &request.thread_id) .await?; + let now = Utc::now(); let message = ThreadMessageRecord { message_id: ThreadMessageId::new(), thread_id: request.thread_id.clone(), sequence, kind: MessageKind::Assistant, status: MessageStatus::Finalized, + created_at: Some(now), + updated_at: Some(now), actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -1839,6 +1850,7 @@ where .map_err(SessionThreadError::Serialization)? { message.content = Some(content); + message.updated_at = Some(Utc::now()); } } Ok(()) @@ -1858,12 +1870,15 @@ where let sequence = self .reserve_sequence(&request.scope, &request.thread_id) .await?; + let now = Utc::now(); let message = ThreadMessageRecord { message_id: ThreadMessageId::new(), thread_id: request.thread_id.clone(), sequence, kind: MessageKind::ToolResultReference, status: MessageStatus::Finalized, + created_at: Some(now), + updated_at: Some(now), actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -1915,12 +1930,15 @@ where let sequence = self .reserve_sequence(&request.scope, &request.thread_id) .await?; + let now = Utc::now(); let message = ThreadMessageRecord { message_id, thread_id: request.thread_id.clone(), sequence, kind: MessageKind::CapabilityDisplayPreview, status: MessageStatus::Finalized, + created_at: Some(now), + updated_at: Some(now), actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -2011,6 +2029,7 @@ where let content = serde_json::to_string(&envelope) .map_err(|error| SessionThreadError::Serialization(error.to_string()))?; message.content = Some(content.clone()); + message.updated_at = Some(Utc::now()); Ok(()) }, ) @@ -2036,6 +2055,7 @@ where // Keep content and attachments in lockstep (as redaction does): // a content update must not leave stale attachment refs behind. message.attachments = Vec::new(); + message.updated_at = Some(Utc::now()); Ok(()) }, ) @@ -2060,6 +2080,7 @@ where message.status = MessageStatus::Finalized; message.content = Some(content.clone().into_text()); message.attachments = Vec::new(); + message.updated_at = Some(Utc::now()); Ok(()) }) .await?; @@ -2091,6 +2112,7 @@ where message.attachments = Vec::new(); message.tool_result_provider_call = None; message.redaction_ref = Some(request.redaction_ref.clone()); + message.updated_at = Some(Utc::now()); Ok(()) }, ) @@ -2927,6 +2949,8 @@ fn history_message(message: &ThreadMessageRecord) -> ThreadMessageRecord { reply_target_binding_id: message.reply_target_binding_id.clone(), turn_id: message.turn_id.clone(), turn_run_id: message.turn_run_id.clone(), + created_at: message.created_at, + updated_at: message.updated_at, tool_result_ref: message.tool_result_ref.clone(), tool_result_provider_call: None, content: message.content.clone(), diff --git a/crates/ironclaw_threads/src/in_memory.rs b/crates/ironclaw_threads/src/in_memory.rs index 161cb32f5ec..df8a547f4b7 100644 --- a/crates/ironclaw_threads/src/in_memory.rs +++ b/crates/ironclaw_threads/src/in_memory.rs @@ -159,13 +159,16 @@ impl SessionThreadService for InMemorySessionThreadService { thread.next_sequence += 1; // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. - thread.record.updated_at = Some(Utc::now()); + let now = Utc::now(); + thread.record.updated_at = Some(now); thread.messages.push(ThreadMessageRecord { message_id, thread_id: request.thread_id.clone(), sequence, kind: MessageKind::User, status: MessageStatus::Accepted, + created_at: Some(now), + updated_at: Some(now), actor_id: Some(request.actor_id), source_binding_id: request.source_binding_id.clone(), reply_target_binding_id: request.reply_target_binding_id, @@ -277,12 +280,15 @@ impl SessionThreadService for InMemorySessionThreadService { return Ok(existing.clone()); } let message_id = ThreadMessageId::new(); + let now = Utc::now(); let message = ThreadMessageRecord { message_id, thread_id: request.thread_id.clone(), sequence: thread.next_sequence, kind: MessageKind::Assistant, status: MessageStatus::Draft, + created_at: Some(now), + updated_at: Some(now), actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -297,7 +303,7 @@ impl SessionThreadService for InMemorySessionThreadService { thread.next_sequence += 1; // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. - thread.record.updated_at = Some(Utc::now()); + thread.record.updated_at = Some(now); thread.messages.push(message.clone()); Ok(message) } @@ -313,18 +319,24 @@ impl SessionThreadService for InMemorySessionThreadService { && message.turn_run_id.as_deref() == Some(request.turn_run_id.as_str()) }) { if existing.status == MessageStatus::Draft { + let now = Utc::now(); existing.status = MessageStatus::Finalized; existing.content = Some(request.content.into_text()); existing.attachments = Vec::new(); + existing.updated_at = Some(now); + thread.record.updated_at = Some(now); } return Ok(existing.clone()); } + let now = Utc::now(); let message = ThreadMessageRecord { message_id: ThreadMessageId::new(), thread_id: request.thread_id.clone(), sequence: thread.next_sequence, kind: MessageKind::Assistant, status: MessageStatus::Finalized, + created_at: Some(now), + updated_at: Some(now), actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -339,7 +351,7 @@ impl SessionThreadService for InMemorySessionThreadService { thread.next_sequence += 1; // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. - thread.record.updated_at = Some(Utc::now()); + thread.record.updated_at = Some(now); thread.messages.push(message.clone()); Ok(message) } @@ -392,6 +404,7 @@ impl SessionThreadService for InMemorySessionThreadService { .map_err(SessionThreadError::Serialization)? { existing.content = Some(content); + existing.updated_at = Some(Utc::now()); } } return Ok(existing.clone()); @@ -403,12 +416,15 @@ impl SessionThreadService for InMemorySessionThreadService { } let content = serde_json::to_string(&envelope) .map_err(|error| SessionThreadError::Serialization(error.to_string()))?; + let now = Utc::now(); let message = ThreadMessageRecord { message_id: ThreadMessageId::new(), thread_id: request.thread_id.clone(), sequence: thread.next_sequence, kind: MessageKind::ToolResultReference, status: MessageStatus::Finalized, + created_at: Some(now), + updated_at: Some(now), actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -423,7 +439,7 @@ impl SessionThreadService for InMemorySessionThreadService { thread.next_sequence += 1; // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. - thread.record.updated_at = Some(Utc::now()); + thread.record.updated_at = Some(now); thread.messages.push(message.clone()); Ok(message) } @@ -454,12 +470,15 @@ impl SessionThreadService for InMemorySessionThreadService { } let content = serde_json::to_string(&request.preview) .map_err(|error| SessionThreadError::Serialization(error.to_string()))?; + let now = Utc::now(); let message = ThreadMessageRecord { message_id: ThreadMessageId::new(), thread_id: request.thread_id.clone(), sequence: thread.next_sequence, kind: MessageKind::CapabilityDisplayPreview, status: MessageStatus::Finalized, + created_at: Some(now), + updated_at: Some(now), actor_id: None, source_binding_id: None, reply_target_binding_id: None, @@ -474,7 +493,7 @@ impl SessionThreadService for InMemorySessionThreadService { thread.next_sequence += 1; // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. - thread.record.updated_at = Some(Utc::now()); + thread.record.updated_at = Some(now); thread.messages.push(message.clone()); Ok(message) } @@ -512,6 +531,7 @@ impl SessionThreadService for InMemorySessionThreadService { serde_json::to_string(&envelope) .map_err(|error| SessionThreadError::Serialization(error.to_string()))?, ); + message.updated_at = Some(Utc::now()); Ok(message.clone()) } @@ -532,6 +552,7 @@ impl SessionThreadService for InMemorySessionThreadService { // assistant draft carries no attachments, so a content update must not // leave stale refs behind if a future draft path ever sets them. message.attachments = Vec::new(); + message.updated_at = Some(Utc::now()); Ok(message.clone()) } @@ -543,11 +564,19 @@ impl SessionThreadService for InMemorySessionThreadService { content: MessageContent, ) -> Result { let mut state = self.state.lock().await; - let message = get_message_mut(&mut state, scope, thread_id, message_id)?; + let thread = get_thread_mut(&mut state, scope, thread_id)?; + let message = thread + .messages + .iter_mut() + .find(|message| message.message_id == message_id) + .ok_or(SessionThreadError::UnknownMessage { message_id })?; ensure_draft(message)?; + let now = Utc::now(); message.status = MessageStatus::Finalized; message.content = Some(content.into_text()); message.attachments = Vec::new(); + message.updated_at = Some(now); + thread.record.updated_at = Some(now); Ok(message.clone()) } @@ -567,6 +596,7 @@ impl SessionThreadService for InMemorySessionThreadService { message.attachments = Vec::new(); message.tool_result_provider_call = None; message.redaction_ref = Some(request.redaction_ref); + message.updated_at = Some(Utc::now()); Ok(message.clone()) } @@ -1092,6 +1122,8 @@ fn history_message(message: &ThreadMessageRecord) -> ThreadMessageRecord { reply_target_binding_id: message.reply_target_binding_id.clone(), turn_id: message.turn_id.clone(), turn_run_id: message.turn_run_id.clone(), + created_at: message.created_at, + updated_at: message.updated_at, tool_result_ref: message.tool_result_ref.clone(), tool_result_provider_call: None, content: message.content.clone(), diff --git a/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs b/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs index 2ef6976d60a..4612c554eca 100644 --- a/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs +++ b/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs @@ -207,6 +207,91 @@ async fn filesystem_finalized_assistant_lookup_by_run_uses_persisted_message() { assert_eq!(finalized.content.as_deref(), Some("final")); } +#[tokio::test] +async fn filesystem_list_thread_history_returns_durable_message_timestamps() { + let backend = Arc::new(InMemoryBackend::new()); + let scoped = scoped_threads_fs_at(backend, "tenant-message-timestamps", "alice"); + let service = FilesystemSessionThreadService::new(scoped); + let scope = scope("message-timestamps"); + let thread = service + .ensure_thread(EnsureThreadRequest { + scope: scope.clone(), + thread_id: Some(ThreadId::new("thread-message-timestamps").unwrap()), + created_by_actor_id: "actor-a".into(), + title: None, + metadata_json: None, + }) + .await + .unwrap(); + + let before_user = Utc::now(); + let accepted = service + .accept_inbound_message(AcceptInboundMessageRequest { + scope: scope.clone(), + thread_id: thread.thread_id.clone(), + actor_id: "actor-a".into(), + source_binding_id: None, + reply_target_binding_id: None, + external_event_id: None, + content: MessageContent::text("hello"), + }) + .await + .unwrap(); + let after_user = Utc::now(); + + let draft = service + .append_assistant_draft(AppendAssistantDraftRequest { + scope: scope.clone(), + thread_id: thread.thread_id.clone(), + turn_run_id: "run-message-timestamps".into(), + content: MessageContent::text("working"), + }) + .await + .unwrap(); + let draft_created_at = draft.created_at.expect("draft has created_at"); + let draft_updated_at = draft.updated_at.expect("draft has updated_at"); + wait_until_after(draft_updated_at).await; + + let finalized = service + .finalize_assistant_message( + &scope, + &thread.thread_id, + draft.message_id, + MessageContent::text("done"), + ) + .await + .unwrap(); + let final_updated_at = finalized.updated_at.expect("finalized has updated_at"); + assert!( + final_updated_at > draft_updated_at, + "finalizing a draft must move updated_at forward" + ); + + let history = service + .list_thread_history(ThreadHistoryRequest { + scope, + thread_id: thread.thread_id, + }) + .await + .unwrap(); + let user = history + .messages + .iter() + .find(|message| message.message_id == accepted.message_id) + .expect("user message is in history"); + let user_created_at = user.created_at.expect("user message has created_at"); + assert!(user_created_at >= before_user && user_created_at <= after_user); + assert_eq!(user.updated_at, Some(user_created_at)); + + let assistant = history + .messages + .iter() + .find(|message| message.message_id == draft.message_id) + .expect("assistant message is in history"); + assert_eq!(assistant.created_at, Some(draft_created_at)); + assert_eq!(assistant.updated_at, Some(final_updated_at)); +} + #[tokio::test] async fn filesystem_append_finalized_assistant_message_is_finalized_and_idempotent_by_turn_run() { let backend = Arc::new(InMemoryBackend::new()); diff --git a/crates/ironclaw_threads/tests/session_thread_contract.rs b/crates/ironclaw_threads/tests/session_thread_contract.rs index cef370ad19f..b0022f7f131 100644 --- a/crates/ironclaw_threads/tests/session_thread_contract.rs +++ b/crates/ironclaw_threads/tests/session_thread_contract.rs @@ -954,6 +954,89 @@ async fn rejected_busy_rejects_non_user_message() { ); } +#[tokio::test] +async fn list_thread_history_returns_durable_message_timestamps() { + let service = InMemorySessionThreadService::default(); + let test_scope = scope("message-timestamps"); + let thread = service + .ensure_thread(EnsureThreadRequest { + scope: test_scope.clone(), + thread_id: Some(ThreadId::new("thread-message-timestamps").unwrap()), + created_by_actor_id: "actor-a".into(), + title: None, + metadata_json: None, + }) + .await + .unwrap(); + + let before_user = Utc::now(); + let accepted = service + .accept_inbound_message(AcceptInboundMessageRequest { + scope: test_scope.clone(), + thread_id: thread.thread_id.clone(), + actor_id: "actor-a".into(), + source_binding_id: None, + reply_target_binding_id: None, + external_event_id: None, + content: user_message("hello"), + }) + .await + .unwrap(); + let after_user = Utc::now(); + + let draft = service + .append_assistant_draft(AppendAssistantDraftRequest { + scope: test_scope.clone(), + thread_id: thread.thread_id.clone(), + turn_run_id: "run-message-timestamps".into(), + content: MessageContent::text("working"), + }) + .await + .unwrap(); + let draft_created_at = draft.created_at.expect("draft has created_at"); + let draft_updated_at = draft.updated_at.expect("draft has updated_at"); + wait_until_after(draft_updated_at).await; + + let finalized = service + .finalize_assistant_message( + &test_scope, + &thread.thread_id, + draft.message_id, + MessageContent::text("done"), + ) + .await + .unwrap(); + let final_updated_at = finalized.updated_at.expect("finalized has updated_at"); + assert!( + final_updated_at > draft_updated_at, + "finalizing a draft must move updated_at forward" + ); + + let history = service + .list_thread_history(ThreadHistoryRequest { + scope: test_scope, + thread_id: thread.thread_id, + }) + .await + .unwrap(); + let user = history + .messages + .iter() + .find(|message| message.message_id == accepted.message_id) + .expect("user message is in history"); + let user_created_at = user.created_at.expect("user message has created_at"); + assert!(user_created_at >= before_user && user_created_at <= after_user); + assert_eq!(user.updated_at, Some(user_created_at)); + + let assistant = history + .messages + .iter() + .find(|message| message.message_id == draft.message_id) + .expect("assistant message is in history"); + assert_eq!(assistant.created_at, Some(draft_created_at)); + assert_eq!(assistant.updated_at, Some(final_updated_at)); +} + #[tokio::test] async fn rejected_busy_rejects_already_finalized_user_message() { let service = InMemorySessionThreadService::default(); diff --git a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js index 9723e720977..9ec094b6bc6 100644 --- a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js +++ b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js @@ -158,10 +158,19 @@ function roleForRecord(record) { } function timestampForRecord(record) { - // ThreadMessageRecord has no top-level timestamp; surfaces use - // the sequence ordering for now. Browsers render the wall-clock - // when an event arrives (FinalReplyView.generated_at). - return record.received_at || record.created_at || null; + // Prefer durable server-side timestamps from the timeline. Finalized + // assistant rows can be born as drafts and finalized later, so their + // completion/update stamp is the user-visible receive time. + if (record.received_at) return record.received_at; + const kind = record.kind || ""; + const status = record.status || ""; + if ( + (kind === "assistant" || kind === "assistant_message") && + status === "finalized" + ) { + return record.updated_at || record.created_at || null; + } + return record.created_at || record.updated_at || null; } function toolCardFromPreviewRecord(record) { diff --git a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs index cdde7d2dc4a..790ab06111c 100644 --- a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs +++ b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs @@ -174,6 +174,38 @@ test("messagesFromTimeline: finalized assistant records are marked as final repl assert.equal(messages[1].isFinalReply, false); }); +test("messagesFromTimeline: user records keep durable created_at timestamps", () => { + const messages = messagesFromTimeline([ + { + message_id: "message-1", + kind: "user", + status: "submitted", + content: "check my calendar", + created_at: "2026-05-12T17:08:00Z", + updated_at: "2026-05-12T17:09:00Z", + }, + ]); + + assert.equal(messages.length, 1); + assert.equal(messages[0].timestamp, "2026-05-12T17:08:00Z"); +}); + +test("messagesFromTimeline: finalized assistant records prefer durable updated_at timestamps", () => { + const messages = messagesFromTimeline([ + { + message_id: "reply-1", + kind: "assistant", + status: "finalized", + content: "Done.", + created_at: "2026-05-12T17:06:00Z", + updated_at: "2026-05-12T17:08:00Z", + }, + ]); + + assert.equal(messages.length, 1); + assert.equal(messages[0].timestamp, "2026-05-12T17:08:00Z"); +}); + test("messagesFromTimeline: tool previews use timeline sequence as activity order", () => { const messages = messagesFromTimeline([ { From 59003fbec4e32a84e850546f16b1e162e46fc21a Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Tue, 7 Jul 2026 19:21:14 +0800 Subject: [PATCH 02/16] fix(reborn): harden chat timestamp persistence --- crates/ironclaw_threads/src/contract.rs | 83 +++++++ .../src/filesystem_service.rs | 140 ++++++++---- crates/ironclaw_threads/src/in_memory.rs | 205 +++++++++++++----- .../filesystem_session_thread_contract.rs | 7 +- .../tests/session_thread_contract.rs | 7 +- .../js/pages/chat/lib/history-messages.js | 12 +- .../pages/chat/lib/history-messages.test.mjs | 24 ++ 7 files changed, 353 insertions(+), 125 deletions(-) diff --git a/crates/ironclaw_threads/src/contract.rs b/crates/ironclaw_threads/src/contract.rs index b07eb7c5681..17a13d6993f 100644 --- a/crates/ironclaw_threads/src/contract.rs +++ b/crates/ironclaw_threads/src/contract.rs @@ -231,6 +231,43 @@ pub struct ThreadMessageRecord { pub redaction_ref: Option, } +/// New transcript rows must carry durable timestamps. Legacy rows persisted +/// before this field existed still deserialize with `None` and remain readable, +/// but all current writers should fail loudly if they try to append a fresh +/// row without both stamps. +pub(crate) fn validate_new_message_timestamps( + message: &ThreadMessageRecord, + context: &'static str, +) -> Result<(), crate::error::SessionThreadError> { + if message.created_at.is_none() || message.updated_at.is_none() { + return Err(crate::error::SessionThreadError::Backend(format!( + "new {context} message {} is missing durable timestamps", + message.message_id + ))); + } + Ok(()) +} + +/// Existing timestamp values are append-only metadata: once a row has either +/// stamp, mutation paths may update it but must not erase it. This keeps +/// legacy `None` values compatible while catching accidental nullification of +/// newly-written rows. +pub(crate) fn validate_message_timestamps_not_cleared( + before: &ThreadMessageRecord, + after: &ThreadMessageRecord, + context: &'static str, +) -> Result<(), crate::error::SessionThreadError> { + let created_at_cleared = before.created_at.is_some() && after.created_at.is_none(); + let updated_at_cleared = before.updated_at.is_some() && after.updated_at.is_none(); + if created_at_cleared || updated_at_cleared { + return Err(crate::error::SessionThreadError::Backend(format!( + "{context} cleared durable timestamps on message {}", + after.message_id + ))); + } + Ok(()) +} + /// Summary artifact over a stable transcript sequence range. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct SummaryArtifact { @@ -592,6 +629,29 @@ mod tests { } } + fn sample_message() -> ThreadMessageRecord { + let now = Utc::now(); + ThreadMessageRecord { + message_id: ThreadMessageId::new(), + thread_id: ThreadId::new("thread-contract").unwrap(), + sequence: 1, + kind: MessageKind::User, + status: MessageStatus::Accepted, + created_at: Some(now), + updated_at: Some(now), + actor_id: Some("actor-a".to_string()), + source_binding_id: None, + reply_target_binding_id: None, + turn_id: None, + turn_run_id: None, + tool_result_ref: None, + tool_result_provider_call: None, + content: Some("hello".to_string()), + attachments: Vec::new(), + redaction_ref: None, + } + } + #[test] fn text_constructor_carries_no_attachments() { let content = MessageContent::text("hello"); @@ -644,6 +704,29 @@ mod tests { assert!(validate_attachment_refs(&[at_cap]).is_ok()); } + #[test] + fn validate_new_message_timestamps_rejects_missing_stamps() { + let mut message = sample_message(); + message.created_at = None; + + let err = validate_new_message_timestamps(&message, "test message") + .expect_err("new messages without created_at must be rejected"); + + assert!(matches!(err, crate::error::SessionThreadError::Backend(_))); + } + + #[test] + fn validate_message_timestamps_not_cleared_rejects_nullification() { + let before = sample_message(); + let mut after = before.clone(); + after.updated_at = None; + + let err = validate_message_timestamps_not_cleared(&before, &after, "test update") + .expect_err("message updates must not clear existing timestamps"); + + assert!(matches!(err, crate::error::SessionThreadError::Backend(_))); + } + #[test] fn into_text_keeps_text_when_attachment_free() { // The non-lossy use: `into_text` is the text-only accessor for content diff --git a/crates/ironclaw_threads/src/filesystem_service.rs b/crates/ironclaw_threads/src/filesystem_service.rs index 225a84c0735..1e72f35fdaf 100644 --- a/crates/ironclaw_threads/src/filesystem_service.rs +++ b/crates/ironclaw_threads/src/filesystem_service.rs @@ -38,7 +38,7 @@ mod message_sequence_index; use std::sync::Arc; use async_trait::async_trait; -use chrono::Utc; +use chrono::{DateTime, Utc}; use futures::{StreamExt, future::join_all}; use ironclaw_filesystem::{ CasApply, CasExpectation, CasUpdateError, ContentType, Entry, FileType, FilesystemError, @@ -322,6 +322,7 @@ where message: &ThreadMessageRecord, description: &'static str, ) -> Result<(), SessionThreadError> { + crate::contract::validate_new_message_timestamps(message, description)?; let path = message_record_path(scope, thread_id, message.message_id)?; let entry = Self::message_entry(message)?; match put_with_cas( @@ -359,6 +360,7 @@ where thread_id: &ThreadId, message: &ThreadMessageRecord, ) -> Result { + crate::contract::validate_new_message_timestamps(message, "message append event")?; let path = message_append_log_path(scope, thread_id)?; let payload = serialize_pretty(&StoredThreadMessageRecord::from(message))?; match self @@ -455,6 +457,7 @@ where message: &mut ThreadMessageRecord, idempotency_record: Option<(&ScopedPath, &Entry)>, ) -> Result { + crate::contract::validate_new_message_timestamps(message, "transactional message")?; let resource_scope = scope.to_resource_scope(); let txn_prefix = scoped_path(THREADS_PREFIX)?; let thread_path = thread_record_path(scope, thread_id)?; @@ -1142,6 +1145,7 @@ where &self, scope: &ThreadScope, thread_id: &ThreadId, + updated_at: DateTime, ) -> Result<(), SessionThreadError> { let path = thread_record_path(scope, thread_id)?; for _ in 0..FILESYSTEM_CAS_RETRIES { @@ -1152,7 +1156,7 @@ where else { return Ok(()); }; - stored.record.updated_at = Some(Utc::now()); + stored.record.updated_at = Some(updated_at); let entry = Self::thread_entry(&stored)?; match put_with_cas( self.filesystem.as_ref(), @@ -1186,8 +1190,16 @@ where /// un-idempotent caller retry and duplicate the message. Logs and /// continues; the advisory `updated_at` stamp simply stays at its prior /// value until the next activity. - async fn touch_thread_updated_at_best_effort(&self, scope: &ThreadScope, thread_id: &ThreadId) { - if let Err(error) = self.touch_thread_updated_at(scope, thread_id).await { + async fn touch_thread_updated_at_best_effort_at( + &self, + scope: &ThreadScope, + thread_id: &ThreadId, + updated_at: DateTime, + ) { + if let Err(error) = self + .touch_thread_updated_at(scope, thread_id, updated_at) + .await + { // silent-ok: recency stamp is advisory after the message is durable. tracing::debug!( ?error, @@ -1248,7 +1260,13 @@ where (record, CasExpectation::Absent) } }; + let before = message.clone(); mutate(&mut message)?; + crate::contract::validate_message_timestamps_not_cleared( + &before, + &message, + "filesystem message update", + )?; let entry = Self::message_entry(&message)?; match put_with_cas( self.filesystem.as_ref(), @@ -1527,7 +1545,7 @@ where // sidebar surfaces this thread first. Best-effort: the message is // already durable, so a touch failure must not fail (and retry) the // accept. - self.touch_thread_updated_at_best_effort(&scope, &thread_id) + self.touch_thread_updated_at_best_effort_at(&scope, &thread_id, now) .await; Ok(AcceptedInboundMessage { @@ -1707,6 +1725,7 @@ where return Ok(existing); } let content = request.content.clone(); + let now = Utc::now(); let finalized = self .apply_message_update( &request.scope, @@ -1714,7 +1733,6 @@ where existing.message_id, |message| { ensure_draft(message)?; - let now = Utc::now(); message.status = MessageStatus::Finalized; message.content = Some(content.clone().into_text()); message.attachments = Vec::new(); @@ -1725,7 +1743,7 @@ where .await?; // Finalizing the in-flight draft is thread activity — stamp recency // (best-effort; the draft update above is already durable). - self.touch_thread_updated_at_best_effort(&request.scope, &request.thread_id) + self.touch_thread_updated_at_best_effort_at(&request.scope, &request.thread_id, now) .await; return Ok(finalized); } @@ -1781,7 +1799,7 @@ where } // Finalized assistant reply is thread activity — stamp recency // (best-effort; the append above is already durable). - self.touch_thread_updated_at_best_effort(&request.scope, &request.thread_id) + self.touch_thread_updated_at_best_effort_at(&request.scope, &request.thread_id, now) .await; Ok(message) } @@ -1828,14 +1846,17 @@ where }; let model_observation = envelope.model_observation.clone(); if provider_call_update.is_some() || model_observation.is_some() { - return self + let now = Utc::now(); + let updated = self .apply_message_update( &request.scope, &request.thread_id, existing.message_id, |message| { + let mut changed = false; if let Some(provider_call) = provider_call_update.as_ref() { message.tool_result_provider_call = Some(provider_call.clone()); + changed = true; } if let Some(model_observation) = model_observation.as_ref() { let content = message.content.as_deref().ok_or_else(|| { @@ -1850,13 +1871,25 @@ where .map_err(SessionThreadError::Serialization)? { message.content = Some(content); - message.updated_at = Some(Utc::now()); + changed = true; } } + if changed { + message.updated_at = Some(now); + } Ok(()) }, ) + .await?; + if updated != existing { + self.touch_thread_updated_at_best_effort_at( + &request.scope, + &request.thread_id, + now, + ) .await; + } + return Ok(updated); } return Ok(existing); } @@ -1951,6 +1984,7 @@ where redaction_ref: None, }; let path = message_record_path(&request.scope, &request.thread_id, message.message_id)?; + crate::contract::validate_new_message_timestamps(&message, "capability display preview")?; let entry = Self::message_entry(&message)?; match put_with_cas( self.filesystem.as_ref(), @@ -2008,7 +2042,9 @@ where let result_ref = request.result_ref.clone(); let thread_id_for_error = request.thread_id.clone(); let safe_summary = request.safe_summary; - self.apply_message_update( + let now = Utc::now(); + let updated = self + .apply_message_update( &request.scope, &request.thread_id, message.message_id, @@ -2029,11 +2065,14 @@ where let content = serde_json::to_string(&envelope) .map_err(|error| SessionThreadError::Serialization(error.to_string()))?; message.content = Some(content.clone()); - message.updated_at = Some(Utc::now()); + message.updated_at = Some(now); Ok(()) }, ) - .await + .await?; + self.touch_thread_updated_at_best_effort_at(&request.scope, &request.thread_id, now) + .await; + Ok(updated) } async fn update_assistant_draft( @@ -2045,21 +2084,26 @@ where .ok_or_else(|| SessionThreadError::UnknownThread { thread_id: request.thread_id.clone(), })?; - self.apply_message_update( - &request.scope, - &request.thread_id, - request.message_id, - |message| { - ensure_draft(message)?; - message.content = Some(request.content.clone().into_text()); - // Keep content and attachments in lockstep (as redaction does): - // a content update must not leave stale attachment refs behind. - message.attachments = Vec::new(); - message.updated_at = Some(Utc::now()); - Ok(()) - }, - ) - .await + let now = Utc::now(); + let updated = self + .apply_message_update( + &request.scope, + &request.thread_id, + request.message_id, + |message| { + ensure_draft(message)?; + message.content = Some(request.content.clone().into_text()); + // Keep content and attachments in lockstep (as redaction does): + // a content update must not leave stale attachment refs behind. + message.attachments = Vec::new(); + message.updated_at = Some(now); + Ok(()) + }, + ) + .await?; + self.touch_thread_updated_at_best_effort_at(&request.scope, &request.thread_id, now) + .await; + Ok(updated) } async fn finalize_assistant_message( @@ -2074,13 +2118,14 @@ where .ok_or_else(|| SessionThreadError::UnknownThread { thread_id: thread_id.clone(), })?; + let now = Utc::now(); let finalized = self .apply_message_update(scope, thread_id, message_id, |message| { ensure_draft(message)?; message.status = MessageStatus::Finalized; message.content = Some(content.clone().into_text()); message.attachments = Vec::new(); - message.updated_at = Some(Utc::now()); + message.updated_at = Some(now); Ok(()) }) .await?; @@ -2088,7 +2133,7 @@ where // (best-effort; the finalize above is already durable). Without this, // the draft/update/finalize path would leave active threads stale in // the `updated_at`-sorted sidebar. - self.touch_thread_updated_at_best_effort(scope, thread_id) + self.touch_thread_updated_at_best_effort_at(scope, thread_id, now) .await; Ok(finalized) } @@ -2102,21 +2147,26 @@ where .ok_or_else(|| SessionThreadError::UnknownThread { thread_id: request.thread_id.clone(), })?; - self.apply_message_update( - &request.scope, - &request.thread_id, - request.message_id, - |message| { - message.status = MessageStatus::Redacted; - message.content = None; - message.attachments = Vec::new(); - message.tool_result_provider_call = None; - message.redaction_ref = Some(request.redaction_ref.clone()); - message.updated_at = Some(Utc::now()); - Ok(()) - }, - ) - .await + let now = Utc::now(); + let updated = self + .apply_message_update( + &request.scope, + &request.thread_id, + request.message_id, + |message| { + message.status = MessageStatus::Redacted; + message.content = None; + message.attachments = Vec::new(); + message.tool_result_provider_call = None; + message.redaction_ref = Some(request.redaction_ref.clone()); + message.updated_at = Some(now); + Ok(()) + }, + ) + .await?; + self.touch_thread_updated_at_best_effort_at(&request.scope, &request.thread_id, now) + .await; + Ok(updated) } async fn load_context_window( diff --git a/crates/ironclaw_threads/src/in_memory.rs b/crates/ironclaw_threads/src/in_memory.rs index df8a547f4b7..01d403b7ae5 100644 --- a/crates/ironclaw_threads/src/in_memory.rs +++ b/crates/ironclaw_threads/src/in_memory.rs @@ -161,7 +161,7 @@ impl SessionThreadService for InMemorySessionThreadService { // stamp so activity-ordered listings surface this thread first. let now = Utc::now(); thread.record.updated_at = Some(now); - thread.messages.push(ThreadMessageRecord { + let message = ThreadMessageRecord { message_id, thread_id: request.thread_id.clone(), sequence, @@ -179,7 +179,9 @@ impl SessionThreadService for InMemorySessionThreadService { content: Some(content_text), attachments, redaction_ref: None, - }); + }; + crate::contract::validate_new_message_timestamps(&message, "inbound user")?; + thread.messages.push(message); if let Some(key) = key { state.inbound_idempotency.insert( @@ -246,9 +248,15 @@ impl SessionThreadService for InMemorySessionThreadService { let mut state = self.state.lock().await; let message = get_message_mut(&mut state, scope, thread_id, message_id)?; ensure_user_accepted(message, "mark_message_submitted")?; + let before = message.clone(); message.status = MessageStatus::Submitted; message.turn_id = Some(turn_id); message.turn_run_id = Some(turn_run_id); + crate::contract::validate_message_timestamps_not_cleared( + &before, + message, + "mark_message_submitted", + )?; Ok(message.clone()) } @@ -261,9 +269,15 @@ impl SessionThreadService for InMemorySessionThreadService { let mut state = self.state.lock().await; let message = get_message_mut(&mut state, scope, thread_id, message_id)?; ensure_user_accepted(message, "mark_message_rejected_busy")?; + let before = message.clone(); message.status = MessageStatus::RejectedBusy; message.turn_id = None; message.turn_run_id = None; + crate::contract::validate_message_timestamps_not_cleared( + &before, + message, + "mark_message_rejected_busy", + )?; Ok(message.clone()) } @@ -304,6 +318,7 @@ impl SessionThreadService for InMemorySessionThreadService { // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. thread.record.updated_at = Some(now); + crate::contract::validate_new_message_timestamps(&message, "assistant draft")?; thread.messages.push(message.clone()); Ok(message) } @@ -320,10 +335,16 @@ impl SessionThreadService for InMemorySessionThreadService { }) { if existing.status == MessageStatus::Draft { let now = Utc::now(); + let before = existing.clone(); existing.status = MessageStatus::Finalized; existing.content = Some(request.content.into_text()); existing.attachments = Vec::new(); existing.updated_at = Some(now); + crate::contract::validate_message_timestamps_not_cleared( + &before, + existing, + "append_finalized_assistant_message", + )?; thread.record.updated_at = Some(now); } return Ok(existing.clone()); @@ -352,6 +373,7 @@ impl SessionThreadService for InMemorySessionThreadService { // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. thread.record.updated_at = Some(now); + crate::contract::validate_new_message_timestamps(&message, "finalized assistant")?; thread.messages.push(message.clone()); Ok(message) } @@ -375,12 +397,18 @@ impl SessionThreadService for InMemorySessionThreadService { && message.turn_run_id.as_deref() == Some(request.turn_run_id.as_str()) && message.tool_result_ref.as_deref() == Some(envelope.result_ref.as_str()) }) { + let now = Utc::now(); + let before = existing.clone(); + let mut changed = false; if let Some(provider_call) = provider_call.as_ref() { provider_call .validate() .map_err(SessionThreadError::Serialization)?; match existing.tool_result_provider_call.as_ref() { - None => existing.tool_result_provider_call = Some(provider_call.clone()), + None => { + existing.tool_result_provider_call = Some(provider_call.clone()); + changed = true; + } Some(existing_provider_call) if existing_provider_call == provider_call => {} Some(_) => { return Err(SessionThreadError::Serialization( @@ -404,10 +432,22 @@ impl SessionThreadService for InMemorySessionThreadService { .map_err(SessionThreadError::Serialization)? { existing.content = Some(content); - existing.updated_at = Some(Utc::now()); + changed = true; } } - return Ok(existing.clone()); + if changed { + existing.updated_at = Some(now); + } + crate::contract::validate_message_timestamps_not_cleared( + &before, + existing, + "append_tool_result_reference", + )?; + let updated = existing.clone(); + if changed { + thread.record.updated_at = Some(now); + } + return Ok(updated); } if let Some(provider_call) = &provider_call { provider_call @@ -440,6 +480,7 @@ impl SessionThreadService for InMemorySessionThreadService { // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. thread.record.updated_at = Some(now); + crate::contract::validate_new_message_timestamps(&message, "tool result reference")?; thread.messages.push(message.clone()); Ok(message) } @@ -494,6 +535,7 @@ impl SessionThreadService for InMemorySessionThreadService { // Appending a message is thread activity; bump the last-activity // stamp so activity-ordered listings surface this thread first. thread.record.updated_at = Some(now); + crate::contract::validate_new_message_timestamps(&message, "capability display preview")?; thread.messages.push(message.clone()); Ok(message) } @@ -502,58 +544,82 @@ impl SessionThreadService for InMemorySessionThreadService { &self, request: UpdateToolResultReferenceRequest, ) -> Result { + let now = Utc::now(); let mut state = self.state.lock().await; let thread = get_thread_mut(&mut state, &request.scope, &request.thread_id)?; - let message = thread - .messages - .iter_mut() - .find(|message| { - message.kind == MessageKind::ToolResultReference - && message.status == MessageStatus::Finalized - && message.turn_run_id.as_deref() == Some(request.turn_run_id.as_str()) - && message.tool_result_ref.as_deref() == Some(request.result_ref.as_str()) - }) - .ok_or_else(|| { - SessionThreadError::Backend(format!( - "tool result reference {} was not found in thread {}", - request.result_ref, request.thread_id - )) + let updated = { + let message = thread + .messages + .iter_mut() + .find(|message| { + message.kind == MessageKind::ToolResultReference + && message.status == MessageStatus::Finalized + && message.turn_run_id.as_deref() == Some(request.turn_run_id.as_str()) + && message.tool_result_ref.as_deref() == Some(request.result_ref.as_str()) + }) + .ok_or_else(|| { + SessionThreadError::Backend(format!( + "tool result reference {} was not found in thread {}", + request.result_ref, request.thread_id + )) + })?; + let before = message.clone(); + let content = message.content.as_deref().ok_or_else(|| { + SessionThreadError::Serialization( + "tool result reference content is missing".to_string(), + ) })?; - let content = message.content.as_deref().ok_or_else(|| { - SessionThreadError::Serialization( - "tool result reference content is missing".to_string(), - ) - })?; - let envelope = ToolResultReferenceEnvelope::from_json_str(content) - .map_err(SessionThreadError::Serialization)? - .with_safe_summary(request.safe_summary); - message.content = Some( - serde_json::to_string(&envelope) - .map_err(|error| SessionThreadError::Serialization(error.to_string()))?, - ); - message.updated_at = Some(Utc::now()); - Ok(message.clone()) + let envelope = ToolResultReferenceEnvelope::from_json_str(content) + .map_err(SessionThreadError::Serialization)? + .with_safe_summary(request.safe_summary); + message.content = Some( + serde_json::to_string(&envelope) + .map_err(|error| SessionThreadError::Serialization(error.to_string()))?, + ); + message.updated_at = Some(now); + crate::contract::validate_message_timestamps_not_cleared( + &before, + message, + "update_tool_result_reference", + )?; + message.clone() + }; + thread.record.updated_at = Some(now); + Ok(updated) } async fn update_assistant_draft( &self, request: UpdateAssistantDraftRequest, ) -> Result { + let now = Utc::now(); let mut state = self.state.lock().await; - let message = get_message_mut( - &mut state, - &request.scope, - &request.thread_id, - request.message_id, - )?; - ensure_draft(message)?; - message.content = Some(request.content.into_text()); - // Keep content and attachments in lockstep (as redaction does): an - // assistant draft carries no attachments, so a content update must not - // leave stale refs behind if a future draft path ever sets them. - message.attachments = Vec::new(); - message.updated_at = Some(Utc::now()); - Ok(message.clone()) + let thread = get_thread_mut(&mut state, &request.scope, &request.thread_id)?; + let updated = { + let message = thread + .messages + .iter_mut() + .find(|message| message.message_id == request.message_id) + .ok_or(SessionThreadError::UnknownMessage { + message_id: request.message_id, + })?; + ensure_draft(message)?; + let before = message.clone(); + message.content = Some(request.content.into_text()); + // Keep content and attachments in lockstep (as redaction does): an + // assistant draft carries no attachments, so a content update must not + // leave stale refs behind if a future draft path ever sets them. + message.attachments = Vec::new(); + message.updated_at = Some(now); + crate::contract::validate_message_timestamps_not_cleared( + &before, + message, + "update_assistant_draft", + )?; + message.clone() + }; + thread.record.updated_at = Some(now); + Ok(updated) } async fn finalize_assistant_message( @@ -572,10 +638,16 @@ impl SessionThreadService for InMemorySessionThreadService { .ok_or(SessionThreadError::UnknownMessage { message_id })?; ensure_draft(message)?; let now = Utc::now(); + let before = message.clone(); message.status = MessageStatus::Finalized; message.content = Some(content.into_text()); message.attachments = Vec::new(); message.updated_at = Some(now); + crate::contract::validate_message_timestamps_not_cleared( + &before, + message, + "finalize_assistant_message", + )?; thread.record.updated_at = Some(now); Ok(message.clone()) } @@ -584,20 +656,33 @@ impl SessionThreadService for InMemorySessionThreadService { &self, request: RedactMessageRequest, ) -> Result { + let now = Utc::now(); let mut state = self.state.lock().await; - let message = get_message_mut( - &mut state, - &request.scope, - &request.thread_id, - request.message_id, - )?; - message.status = MessageStatus::Redacted; - message.content = None; - message.attachments = Vec::new(); - message.tool_result_provider_call = None; - message.redaction_ref = Some(request.redaction_ref); - message.updated_at = Some(Utc::now()); - Ok(message.clone()) + let thread = get_thread_mut(&mut state, &request.scope, &request.thread_id)?; + let updated = { + let message = thread + .messages + .iter_mut() + .find(|message| message.message_id == request.message_id) + .ok_or(SessionThreadError::UnknownMessage { + message_id: request.message_id, + })?; + let before = message.clone(); + message.status = MessageStatus::Redacted; + message.content = None; + message.attachments = Vec::new(); + message.tool_result_provider_call = None; + message.redaction_ref = Some(request.redaction_ref); + message.updated_at = Some(now); + crate::contract::validate_message_timestamps_not_cleared( + &before, + message, + "redact_message", + )?; + message.clone() + }; + thread.record.updated_at = Some(now); + Ok(updated) } async fn load_context_window( diff --git a/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs b/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs index 4612c554eca..1e8f3f3c77d 100644 --- a/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs +++ b/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs @@ -249,8 +249,7 @@ async fn filesystem_list_thread_history_returns_durable_message_timestamps() { .await .unwrap(); let draft_created_at = draft.created_at.expect("draft has created_at"); - let draft_updated_at = draft.updated_at.expect("draft has updated_at"); - wait_until_after(draft_updated_at).await; + assert_eq!(draft.updated_at, Some(draft_created_at)); let finalized = service .finalize_assistant_message( @@ -262,10 +261,6 @@ async fn filesystem_list_thread_history_returns_durable_message_timestamps() { .await .unwrap(); let final_updated_at = finalized.updated_at.expect("finalized has updated_at"); - assert!( - final_updated_at > draft_updated_at, - "finalizing a draft must move updated_at forward" - ); let history = service .list_thread_history(ThreadHistoryRequest { diff --git a/crates/ironclaw_threads/tests/session_thread_contract.rs b/crates/ironclaw_threads/tests/session_thread_contract.rs index b0022f7f131..04205f336a2 100644 --- a/crates/ironclaw_threads/tests/session_thread_contract.rs +++ b/crates/ironclaw_threads/tests/session_thread_contract.rs @@ -994,8 +994,7 @@ async fn list_thread_history_returns_durable_message_timestamps() { .await .unwrap(); let draft_created_at = draft.created_at.expect("draft has created_at"); - let draft_updated_at = draft.updated_at.expect("draft has updated_at"); - wait_until_after(draft_updated_at).await; + assert_eq!(draft.updated_at, Some(draft_created_at)); let finalized = service .finalize_assistant_message( @@ -1007,10 +1006,6 @@ async fn list_thread_history_returns_durable_message_timestamps() { .await .unwrap(); let final_updated_at = finalized.updated_at.expect("finalized has updated_at"); - assert!( - final_updated_at > draft_updated_at, - "finalizing a draft must move updated_at forward" - ); let history = service .list_thread_history(ThreadHistoryRequest { diff --git a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js index 9ec094b6bc6..2716f9ca714 100644 --- a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js +++ b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js @@ -158,16 +158,12 @@ function roleForRecord(record) { } function timestampForRecord(record) { - // Prefer durable server-side timestamps from the timeline. Finalized - // assistant rows can be born as drafts and finalized later, so their - // completion/update stamp is the user-visible receive time. + // Prefer durable server-side timestamps from the timeline. Finalized rows + // can be updated after creation, so their latest durable update is the + // user-visible receive/completion time. if (record.received_at) return record.received_at; - const kind = record.kind || ""; const status = record.status || ""; - if ( - (kind === "assistant" || kind === "assistant_message") && - status === "finalized" - ) { + if (status === "finalized") { return record.updated_at || record.created_at || null; } return record.created_at || record.updated_at || null; diff --git a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs index 790ab06111c..9cbfca3409e 100644 --- a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs +++ b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs @@ -206,6 +206,30 @@ test("messagesFromTimeline: finalized assistant records prefer durable updated_a assert.equal(messages[0].timestamp, "2026-05-12T17:08:00Z"); }); +test("messagesFromTimeline: finalized non-assistant records prefer durable updated_at timestamps", () => { + const messages = messagesFromTimeline([ + { + message_id: "tool-preview-1", + kind: "capability_display_preview", + status: "finalized", + turn_run_id: "run-tools", + content: JSON.stringify({ + version: 1, + invocation_id: "invocation-extension-a", + capability_id: "builtin.extension_search", + status: "completed", + title: "extension_search", + }), + created_at: "2026-05-12T17:06:00Z", + updated_at: "2026-05-12T17:10:00Z", + }, + ]); + + assert.equal(messages.length, 1); + assert.equal(messages[0].role, "tool_activity"); + assert.equal(messages[0].timestamp, "2026-05-12T17:10:00Z"); +}); + test("messagesFromTimeline: tool previews use timeline sequence as activity order", () => { const messages = messagesFromTimeline([ { From 4d516b755cf151c09427d8be8a5b69f9e9618c60 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Tue, 7 Jul 2026 20:05:39 +0800 Subject: [PATCH 03/16] fix(reborn): address timestamp review feedback --- crates/ironclaw_threads/src/contract.rs | 36 +++++++++++++------ crates/ironclaw_threads/src/error.rs | 24 +++++++++++++ .../pages/chat/lib/history-messages.test.mjs | 17 +++++++++ 3 files changed, 67 insertions(+), 10 deletions(-) diff --git a/crates/ironclaw_threads/src/contract.rs b/crates/ironclaw_threads/src/contract.rs index 17a13d6993f..14cbaffb7d8 100644 --- a/crates/ironclaw_threads/src/contract.rs +++ b/crates/ironclaw_threads/src/contract.rs @@ -240,10 +240,11 @@ pub(crate) fn validate_new_message_timestamps( context: &'static str, ) -> Result<(), crate::error::SessionThreadError> { if message.created_at.is_none() || message.updated_at.is_none() { - return Err(crate::error::SessionThreadError::Backend(format!( - "new {context} message {} is missing durable timestamps", - message.message_id - ))); + return Err(crate::error::SessionThreadError::InvalidMessageTimestamp { + message_id: message.message_id, + context, + violation: crate::error::TimestampViolation::MissingDurableTimestamps, + }); } Ok(()) } @@ -260,10 +261,11 @@ pub(crate) fn validate_message_timestamps_not_cleared( let created_at_cleared = before.created_at.is_some() && after.created_at.is_none(); let updated_at_cleared = before.updated_at.is_some() && after.updated_at.is_none(); if created_at_cleared || updated_at_cleared { - return Err(crate::error::SessionThreadError::Backend(format!( - "{context} cleared durable timestamps on message {}", - after.message_id - ))); + return Err(crate::error::SessionThreadError::InvalidMessageTimestamp { + message_id: after.message_id, + context, + violation: crate::error::TimestampViolation::ClearedDurableTimestamps, + }); } Ok(()) } @@ -712,7 +714,14 @@ mod tests { let err = validate_new_message_timestamps(&message, "test message") .expect_err("new messages without created_at must be rejected"); - assert!(matches!(err, crate::error::SessionThreadError::Backend(_))); + assert!(matches!( + err, + crate::error::SessionThreadError::InvalidMessageTimestamp { + context: "test message", + violation: crate::error::TimestampViolation::MissingDurableTimestamps, + .. + } + )); } #[test] @@ -724,7 +733,14 @@ mod tests { let err = validate_message_timestamps_not_cleared(&before, &after, "test update") .expect_err("message updates must not clear existing timestamps"); - assert!(matches!(err, crate::error::SessionThreadError::Backend(_))); + assert!(matches!( + err, + crate::error::SessionThreadError::InvalidMessageTimestamp { + context: "test update", + violation: crate::error::TimestampViolation::ClearedDurableTimestamps, + .. + } + )); } #[test] diff --git a/crates/ironclaw_threads/src/error.rs b/crates/ironclaw_threads/src/error.rs index 8a2cf584f70..428cb6b05ee 100644 --- a/crates/ironclaw_threads/src/error.rs +++ b/crates/ironclaw_threads/src/error.rs @@ -1,8 +1,26 @@ +use std::fmt; + use ironclaw_host_api::ThreadId; use thiserror::Error; use crate::{MessageStatus, ThreadMessageId}; +/// Specific timestamp contract violation on a transcript message. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TimestampViolation { + MissingDurableTimestamps, + ClearedDurableTimestamps, +} + +impl fmt::Display for TimestampViolation { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::MissingDurableTimestamps => write!(f, "missing durable timestamps"), + Self::ClearedDurableTimestamps => write!(f, "cleared durable timestamps"), + } + } +} + /// Canonical thread/transcript service errors. #[derive(Debug, Error)] pub enum SessionThreadError { @@ -48,6 +66,12 @@ pub enum SessionThreadError { }, #[error("invalid attachment on inbound message: {0}")] InvalidAttachment(String), + #[error("{context} timestamp violation on message {message_id}: {violation}")] + InvalidMessageTimestamp { + message_id: ThreadMessageId, + context: &'static str, + violation: TimestampViolation, + }, #[error("failed to create generated thread id: {0}")] GeneratedThreadId(String), #[error("serialization error: {0}")] diff --git a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs index 9cbfca3409e..d697cedd229 100644 --- a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs +++ b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.test.mjs @@ -206,6 +206,23 @@ test("messagesFromTimeline: finalized assistant records prefer durable updated_a assert.equal(messages[0].timestamp, "2026-05-12T17:08:00Z"); }); +test("messagesFromTimeline: received_at takes precedence over durable timestamps", () => { + const messages = messagesFromTimeline([ + { + message_id: "reply-received", + kind: "assistant", + status: "finalized", + content: "Done.", + received_at: "2026-05-12T17:11:00Z", + created_at: "2026-05-12T17:06:00Z", + updated_at: "2026-05-12T17:08:00Z", + }, + ]); + + assert.equal(messages.length, 1); + assert.equal(messages[0].timestamp, "2026-05-12T17:11:00Z"); +}); + test("messagesFromTimeline: finalized non-assistant records prefer durable updated_at timestamps", () => { const messages = messagesFromTimeline([ { From cade5470cc17b9da52687741ee061029dfe58627 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Tue, 7 Jul 2026 20:10:20 +0800 Subject: [PATCH 04/16] fix(reborn): map timestamp thread errors --- crates/ironclaw_product_workflow/src/reborn_services.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/crates/ironclaw_product_workflow/src/reborn_services.rs b/crates/ironclaw_product_workflow/src/reborn_services.rs index 97e7aca0f04..3c6916a5c0c 100644 --- a/crates/ironclaw_product_workflow/src/reborn_services.rs +++ b/crates/ironclaw_product_workflow/src/reborn_services.rs @@ -6078,6 +6078,7 @@ fn map_timeline_probe_error(error: SessionThreadError) -> RebornServicesError { match error { SessionThreadError::Serialization(_) | SessionThreadError::Deserialization(_) + | SessionThreadError::InvalidMessageTimestamp { .. } | SessionThreadError::Backend(_) => RebornServicesError::from_status_kind( RebornServicesErrorCode::Unavailable, RebornServicesErrorKind::TimelineUnavailable, @@ -6118,6 +6119,7 @@ fn map_thread_error(error: SessionThreadError) -> RebornServicesError { SessionThreadError::GeneratedThreadId(_) | SessionThreadError::Serialization(_) | SessionThreadError::Deserialization(_) + | SessionThreadError::InvalidMessageTimestamp { .. } | SessionThreadError::Backend(_) => RebornServicesError::service_unavailable(true), } } From 6f187d4f17d6a01322f7e20390be17521a2f87b4 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Tue, 7 Jul 2026 20:13:37 +0800 Subject: [PATCH 05/16] fix(reborn): derive timestamp violation errors --- crates/ironclaw_threads/src/error.rs | 15 +++------------ 1 file changed, 3 insertions(+), 12 deletions(-) diff --git a/crates/ironclaw_threads/src/error.rs b/crates/ironclaw_threads/src/error.rs index 428cb6b05ee..f975b9cd900 100644 --- a/crates/ironclaw_threads/src/error.rs +++ b/crates/ironclaw_threads/src/error.rs @@ -1,26 +1,17 @@ -use std::fmt; - use ironclaw_host_api::ThreadId; use thiserror::Error; use crate::{MessageStatus, ThreadMessageId}; /// Specific timestamp contract violation on a transcript message. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] pub enum TimestampViolation { + #[error("missing durable timestamps")] MissingDurableTimestamps, + #[error("cleared durable timestamps")] ClearedDurableTimestamps, } -impl fmt::Display for TimestampViolation { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match self { - Self::MissingDurableTimestamps => write!(f, "missing durable timestamps"), - Self::ClearedDurableTimestamps => write!(f, "cleared durable timestamps"), - } - } -} - /// Canonical thread/transcript service errors. #[derive(Debug, Error)] pub enum SessionThreadError { From c6baafd6265705151d7dabb1088c89e3b5bf9587 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Tue, 7 Jul 2026 20:26:43 +0800 Subject: [PATCH 06/16] fix(stress): classify timestamp thread errors --- tools/ironclaw_stress/src/user_turn.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/tools/ironclaw_stress/src/user_turn.rs b/tools/ironclaw_stress/src/user_turn.rs index 24a9f89e8ab..50380991670 100644 --- a/tools/ironclaw_stress/src/user_turn.rs +++ b/tools/ironclaw_stress/src/user_turn.rs @@ -1869,7 +1869,9 @@ fn thread_failure(stage: impl Into, error: SessionThreadError) -> Operat SessionThreadError::Serialization(_) | SessionThreadError::Deserialization(_) => { "thread_serialization" } - SessionThreadError::Backend(_) => "thread_backend", + SessionThreadError::InvalidMessageTimestamp { .. } | SessionThreadError::Backend(_) => { + "thread_backend" + } }; OperationFailure::new(bucket, stage, error) } From c51b3caa1e83e052c201960448ae43f4751db2f6 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Tue, 7 Jul 2026 20:35:53 +0800 Subject: [PATCH 07/16] fix(reborn): optimize timestamp update validation --- crates/ironclaw_threads/src/contract.rs | 42 +++++--- .../src/filesystem_service.rs | 12 ++- crates/ironclaw_threads/src/in_memory.rs | 96 ++++++++++++------- .../js/pages/chat/lib/history-messages.js | 6 +- 4 files changed, 104 insertions(+), 52 deletions(-) diff --git a/crates/ironclaw_threads/src/contract.rs b/crates/ironclaw_threads/src/contract.rs index 14cbaffb7d8..8cf74e5dcf0 100644 --- a/crates/ironclaw_threads/src/contract.rs +++ b/crates/ironclaw_threads/src/contract.rs @@ -188,9 +188,11 @@ pub struct SessionThreadRecord { /// before activity timestamps existed; such records sort oldest. #[serde(default, skip_serializing_if = "Option::is_none")] pub created_at: Option>, - /// Last time the thread saw activity (a message was appended). Drives - /// the sidebar "Recent" ordering — newest activity first. Bumped on - /// every append; `None` for legacy records (sorts oldest). + /// Last time the thread saw durable user-visible transcript activity. + /// Drives the sidebar "Recent" ordering — newest activity first. Writers + /// that touch both a message and the thread activity stamp reuse one + /// timestamp for that logical change; pure idempotent replays remain + /// side-effect-free. `None` for legacy records (sorts oldest). #[serde(default, skip_serializing_if = "Option::is_none")] pub updated_at: Option>, } @@ -209,7 +211,8 @@ pub struct ThreadMessageRecord { pub created_at: Option>, /// Last time the message row was materially changed. Draft finalization /// updates this so UI refreshes can show the reply completion time instead - /// of the browser refresh time. + /// of the browser refresh time. Writers capture one timestamp per logical + /// message mutation so related message/thread stamps stay atomic. #[serde(default, skip_serializing_if = "Option::is_none")] pub updated_at: Option>, pub actor_id: Option, @@ -253,16 +256,22 @@ pub(crate) fn validate_new_message_timestamps( /// stamp, mutation paths may update it but must not erase it. This keeps /// legacy `None` values compatible while catching accidental nullification of /// newly-written rows. -pub(crate) fn validate_message_timestamps_not_cleared( - before: &ThreadMessageRecord, - after: &ThreadMessageRecord, +/// Lightweight timestamp-only variant for hot mutation paths. Callers usually +/// already hold a mutable message and should not clone large payloads just to +/// prove timestamp metadata was not erased. +pub(crate) fn validate_message_timestamp_fields_not_cleared( + message_id: ThreadMessageId, + before_created_at: Option>, + before_updated_at: Option>, + after_created_at: Option>, + after_updated_at: Option>, context: &'static str, ) -> Result<(), crate::error::SessionThreadError> { - let created_at_cleared = before.created_at.is_some() && after.created_at.is_none(); - let updated_at_cleared = before.updated_at.is_some() && after.updated_at.is_none(); + let created_at_cleared = before_created_at.is_some() && after_created_at.is_none(); + let updated_at_cleared = before_updated_at.is_some() && after_updated_at.is_none(); if created_at_cleared || updated_at_cleared { return Err(crate::error::SessionThreadError::InvalidMessageTimestamp { - message_id: after.message_id, + message_id, context, violation: crate::error::TimestampViolation::ClearedDurableTimestamps, }); @@ -725,13 +734,20 @@ mod tests { } #[test] - fn validate_message_timestamps_not_cleared_rejects_nullification() { + fn validate_message_timestamp_fields_not_cleared_rejects_nullification() { let before = sample_message(); let mut after = before.clone(); after.updated_at = None; - let err = validate_message_timestamps_not_cleared(&before, &after, "test update") - .expect_err("message updates must not clear existing timestamps"); + let err = validate_message_timestamp_fields_not_cleared( + after.message_id, + before.created_at, + before.updated_at, + after.created_at, + after.updated_at, + "test update", + ) + .expect_err("message updates must not clear existing timestamps"); assert!(matches!( err, diff --git a/crates/ironclaw_threads/src/filesystem_service.rs b/crates/ironclaw_threads/src/filesystem_service.rs index 1e72f35fdaf..88a2ad9b155 100644 --- a/crates/ironclaw_threads/src/filesystem_service.rs +++ b/crates/ironclaw_threads/src/filesystem_service.rs @@ -1260,11 +1260,15 @@ where (record, CasExpectation::Absent) } }; - let before = message.clone(); + let before_created_at = message.created_at; + let before_updated_at = message.updated_at; mutate(&mut message)?; - crate::contract::validate_message_timestamps_not_cleared( - &before, - &message, + crate::contract::validate_message_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, "filesystem message update", )?; let entry = Self::message_entry(&message)?; diff --git a/crates/ironclaw_threads/src/in_memory.rs b/crates/ironclaw_threads/src/in_memory.rs index 01d403b7ae5..2e471eb9e75 100644 --- a/crates/ironclaw_threads/src/in_memory.rs +++ b/crates/ironclaw_threads/src/in_memory.rs @@ -248,13 +248,17 @@ impl SessionThreadService for InMemorySessionThreadService { let mut state = self.state.lock().await; let message = get_message_mut(&mut state, scope, thread_id, message_id)?; ensure_user_accepted(message, "mark_message_submitted")?; - let before = message.clone(); + let before_created_at = message.created_at; + let before_updated_at = message.updated_at; message.status = MessageStatus::Submitted; message.turn_id = Some(turn_id); message.turn_run_id = Some(turn_run_id); - crate::contract::validate_message_timestamps_not_cleared( - &before, - message, + crate::contract::validate_message_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, "mark_message_submitted", )?; Ok(message.clone()) @@ -269,13 +273,17 @@ impl SessionThreadService for InMemorySessionThreadService { let mut state = self.state.lock().await; let message = get_message_mut(&mut state, scope, thread_id, message_id)?; ensure_user_accepted(message, "mark_message_rejected_busy")?; - let before = message.clone(); + let before_created_at = message.created_at; + let before_updated_at = message.updated_at; message.status = MessageStatus::RejectedBusy; message.turn_id = None; message.turn_run_id = None; - crate::contract::validate_message_timestamps_not_cleared( - &before, - message, + crate::contract::validate_message_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, "mark_message_rejected_busy", )?; Ok(message.clone()) @@ -335,14 +343,18 @@ impl SessionThreadService for InMemorySessionThreadService { }) { if existing.status == MessageStatus::Draft { let now = Utc::now(); - let before = existing.clone(); + let before_created_at = existing.created_at; + let before_updated_at = existing.updated_at; existing.status = MessageStatus::Finalized; existing.content = Some(request.content.into_text()); existing.attachments = Vec::new(); existing.updated_at = Some(now); - crate::contract::validate_message_timestamps_not_cleared( - &before, - existing, + crate::contract::validate_message_timestamp_fields_not_cleared( + existing.message_id, + before_created_at, + before_updated_at, + existing.created_at, + existing.updated_at, "append_finalized_assistant_message", )?; thread.record.updated_at = Some(now); @@ -398,7 +410,8 @@ impl SessionThreadService for InMemorySessionThreadService { && message.tool_result_ref.as_deref() == Some(envelope.result_ref.as_str()) }) { let now = Utc::now(); - let before = existing.clone(); + let before_created_at = existing.created_at; + let before_updated_at = existing.updated_at; let mut changed = false; if let Some(provider_call) = provider_call.as_ref() { provider_call @@ -438,9 +451,12 @@ impl SessionThreadService for InMemorySessionThreadService { if changed { existing.updated_at = Some(now); } - crate::contract::validate_message_timestamps_not_cleared( - &before, - existing, + crate::contract::validate_message_timestamp_fields_not_cleared( + existing.message_id, + before_created_at, + before_updated_at, + existing.created_at, + existing.updated_at, "append_tool_result_reference", )?; let updated = existing.clone(); @@ -563,7 +579,8 @@ impl SessionThreadService for InMemorySessionThreadService { request.result_ref, request.thread_id )) })?; - let before = message.clone(); + let before_created_at = message.created_at; + let before_updated_at = message.updated_at; let content = message.content.as_deref().ok_or_else(|| { SessionThreadError::Serialization( "tool result reference content is missing".to_string(), @@ -577,9 +594,12 @@ impl SessionThreadService for InMemorySessionThreadService { .map_err(|error| SessionThreadError::Serialization(error.to_string()))?, ); message.updated_at = Some(now); - crate::contract::validate_message_timestamps_not_cleared( - &before, - message, + crate::contract::validate_message_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, "update_tool_result_reference", )?; message.clone() @@ -604,16 +624,20 @@ impl SessionThreadService for InMemorySessionThreadService { message_id: request.message_id, })?; ensure_draft(message)?; - let before = message.clone(); + let before_created_at = message.created_at; + let before_updated_at = message.updated_at; message.content = Some(request.content.into_text()); // Keep content and attachments in lockstep (as redaction does): an // assistant draft carries no attachments, so a content update must not // leave stale refs behind if a future draft path ever sets them. message.attachments = Vec::new(); message.updated_at = Some(now); - crate::contract::validate_message_timestamps_not_cleared( - &before, - message, + crate::contract::validate_message_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, "update_assistant_draft", )?; message.clone() @@ -638,14 +662,18 @@ impl SessionThreadService for InMemorySessionThreadService { .ok_or(SessionThreadError::UnknownMessage { message_id })?; ensure_draft(message)?; let now = Utc::now(); - let before = message.clone(); + let before_created_at = message.created_at; + let before_updated_at = message.updated_at; message.status = MessageStatus::Finalized; message.content = Some(content.into_text()); message.attachments = Vec::new(); message.updated_at = Some(now); - crate::contract::validate_message_timestamps_not_cleared( - &before, - message, + crate::contract::validate_message_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, "finalize_assistant_message", )?; thread.record.updated_at = Some(now); @@ -667,16 +695,20 @@ impl SessionThreadService for InMemorySessionThreadService { .ok_or(SessionThreadError::UnknownMessage { message_id: request.message_id, })?; - let before = message.clone(); + let before_created_at = message.created_at; + let before_updated_at = message.updated_at; message.status = MessageStatus::Redacted; message.content = None; message.attachments = Vec::new(); message.tool_result_provider_call = None; message.redaction_ref = Some(request.redaction_ref); message.updated_at = Some(now); - crate::contract::validate_message_timestamps_not_cleared( - &before, - message, + crate::contract::validate_message_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, "redact_message", )?; message.clone() diff --git a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js index 2716f9ca714..c4b2e2e8395 100644 --- a/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js +++ b/crates/ironclaw_webui_v2/static/js/pages/chat/lib/history-messages.js @@ -158,9 +158,9 @@ function roleForRecord(record) { } function timestampForRecord(record) { - // Prefer durable server-side timestamps from the timeline. Finalized rows - // can be updated after creation, so their latest durable update is the - // user-visible receive/completion time. + // Prefer durable server-side timestamps from the timeline. For finalized + // rows, `updated_at` captures the durable receive/finalization time shown in + // chat history. if (record.received_at) return record.received_at; const status = record.status || ""; if (status === "finalized") { From 3cd0db8bf44ead1b10cd5f88ac73d458e5294ca5 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Tue, 7 Jul 2026 21:02:05 +0800 Subject: [PATCH 08/16] fix(stress): separate timestamp failure bucket --- tools/ironclaw_stress/src/user_turn.rs | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/tools/ironclaw_stress/src/user_turn.rs b/tools/ironclaw_stress/src/user_turn.rs index 50380991670..ac7a4039864 100644 --- a/tools/ironclaw_stress/src/user_turn.rs +++ b/tools/ironclaw_stress/src/user_turn.rs @@ -1869,9 +1869,8 @@ fn thread_failure(stage: impl Into, error: SessionThreadError) -> Operat SessionThreadError::Serialization(_) | SessionThreadError::Deserialization(_) => { "thread_serialization" } - SessionThreadError::InvalidMessageTimestamp { .. } | SessionThreadError::Backend(_) => { - "thread_backend" - } + SessionThreadError::InvalidMessageTimestamp { .. } => "thread_timestamp_invalid", + SessionThreadError::Backend(_) => "thread_backend", }; OperationFailure::new(bucket, stage, error) } From d358f670864268a834cb21b61062b96d89f4fd20 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 8 Jul 2026 17:33:12 +0800 Subject: [PATCH 09/16] test(reborn): avoid outbound target coverage stack overflow --- tests/integration/outbound_target.rs | 127 ++++++++++++++++----------- 1 file changed, 76 insertions(+), 51 deletions(-) diff --git a/tests/integration/outbound_target.rs b/tests/integration/outbound_target.rs index b0823175a33..57354b44011 100644 --- a/tests/integration/outbound_target.rs +++ b/tests/integration/outbound_target.rs @@ -178,60 +178,65 @@ async fn target_set_unknown_target_routes_to_invalid_input() { .expect("an unknown target surfaces as Failed(InvalidInput)"); } -#[tokio::test] -async fn target_set_disabled_by_settings_routes_to_policy_denied() { - let group = RebornIntegrationGroup::outbound_target_tools() - .await - .expect("outbound-target-tools group builds"); - let harness = group - .thread("conv-outbound-set-denied") - .script([ - RebornScriptedReply::tool_call( - "builtin.outbound_delivery_target_set", - serde_json::json!({ "target_id": KNOWN_TARGET_ID }), - ), - RebornScriptedReply::text("that tool is disabled"), - ]) - .build() - .await - .expect("thread builds"); +#[test] +fn target_set_disabled_by_settings_routes_to_policy_denied() { + run_async_test_with_stack( + "target_set_disabled_by_settings_routes_to_policy_denied", + async { + let group = RebornIntegrationGroup::outbound_target_tools() + .await + .expect("outbound-target-tools group builds"); + let harness = group + .thread("conv-outbound-set-denied") + .script([ + RebornScriptedReply::tool_call( + "builtin.outbound_delivery_target_set", + serde_json::json!({ "target_id": KNOWN_TARGET_ID }), + ), + RebornScriptedReply::text("that tool is disabled"), + ]) + .build() + .await + .expect("thread builds"); - // Persist a `Disabled` per-tool override for the run's effective dispatch - // user (the thread binding actor), driving the settings decision to Deny. - group - .capability_harness() - .expect("outbound_target_tools always uses HostRuntime") - .disable_outbound_target_set_tool( - harness.binding.tenant_id.clone(), - harness.binding.actor_user_id.clone(), - ) - .await - .expect("tool override persists"); + // Persist a `Disabled` per-tool override for the run's effective dispatch + // user (the thread binding actor), driving the settings decision to Deny. + group + .capability_harness() + .expect("outbound_target_tools always uses HostRuntime") + .disable_outbound_target_set_tool( + harness.binding.tenant_id.clone(), + harness.binding.actor_user_id.clone(), + ) + .await + .expect("tool override persists"); - harness - .submit_turn("send my replies to slack dm alpha") - .await - .expect("turn completes despite the policy-denied target_set"); + harness + .submit_turn("send my replies to slack dm alpha") + .await + .expect("turn completes despite the policy-denied target_set"); - harness - .assert_tool_invoked("builtin.outbound_delivery_target_set") - .await - .expect("target_set dispatched through the synthetic-capability port"); - harness - .assert_tool_error(ToolErrorClass::Failed, "policy_denied") - .await - .expect("a disabled tool surfaces as Failed(PolicyDenied)"); - // A policy-denied dispatch must short-circuit before ever reaching the - // facade set seam — proves the deny happened at the settings-decision gate, - // not merely that the model observed a policy_denied error string. - let facade = group - .capability_harness() - .expect("outbound_target_tools always uses HostRuntime") - .outbound_preferences_facade_for_test() - .expect("outbound_target_tools always wires a facade double"); - assert!( - facade.recorded_set_target_ids().is_empty(), - "a policy-denied target_set must not reach the facade set seam" + harness + .assert_tool_invoked("builtin.outbound_delivery_target_set") + .await + .expect("target_set dispatched through the synthetic-capability port"); + harness + .assert_tool_error(ToolErrorClass::Failed, "policy_denied") + .await + .expect("a disabled tool surfaces as Failed(PolicyDenied)"); + // A policy-denied dispatch must short-circuit before ever reaching the + // facade set seam — proves the deny happened at the settings-decision gate, + // not merely that the model observed a policy_denied error string. + let facade = group + .capability_harness() + .expect("outbound_target_tools always uses HostRuntime") + .outbound_preferences_facade_for_test() + .expect("outbound_target_tools always wires a facade double"); + assert!( + facade.recorded_set_target_ids().is_empty(), + "a policy-denied target_set must not reach the facade set seam" + ); + }, ); } @@ -349,3 +354,23 @@ async fn target_set_approval_gate_deny_leaves_preference_unchanged() { "a denied target_set must not reach the facade set seam" ); } + +fn run_async_test_with_stack(name: &'static str, test: F) +where + F: std::future::Future + Send + 'static, +{ + let handle = std::thread::Builder::new() + .name(name.to_string()) + .stack_size(16 * 1024 * 1024) + .spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("tokio test runtime") + .block_on(test); + }) + .expect("spawn stack-sized test thread"); + if let Err(panic) = handle.join() { + std::panic::resume_unwind(panic); + } +} From bff89fa64178629f8ca3ac85dc24bf0f54cbc79b Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 8 Jul 2026 17:38:32 +0800 Subject: [PATCH 10/16] test(reborn): avoid skill activate coverage stack overflow --- tests/integration/skill_activate.rs | 159 ++++++++++++++++------------ 1 file changed, 93 insertions(+), 66 deletions(-) diff --git a/tests/integration/skill_activate.rs b/tests/integration/skill_activate.rs index 0cd02430ae4..4e54f9a3b85 100644 --- a/tests/integration/skill_activate.rs +++ b/tests/integration/skill_activate.rs @@ -185,74 +185,101 @@ async fn skill_activate_over_budget_surfaces_recoverable_failed() { /// (`AmbiguousSkill` → `Failed(InvalidInput)`), distinct from the /// `ContextBudgetExceeded` arm above. Proves neither candidate's instructions /// leak into a later model request despite the name matching both. -#[tokio::test] -async fn skill_activate_ambiguous_name_surfaces_recoverable_failed() { - let group = RebornIntegrationGroup::skill_activation_tools() - .await - .expect("skill-activation group builds"); - let capability_harness = group - .capability_harness() - .expect("skill-activation group has a host-runtime capability harness"); - capability_harness - .seed_system_skill_for_test( - "duplicate", - "a system-scoped skill", - "SYSTEM_DUPLICATE_SKILL_SENTINEL", - ) - .expect("system-scoped duplicate skill seeds"); +#[test] +fn skill_activate_ambiguous_name_surfaces_recoverable_failed() { + run_async_test_with_stack( + "skill_activate_ambiguous_name_surfaces_recoverable_failed", + async { + let group = RebornIntegrationGroup::skill_activation_tools() + .await + .expect("skill-activation group builds"); + let capability_harness = group + .capability_harness() + .expect("skill-activation group has a host-runtime capability harness"); + capability_harness + .seed_system_skill_for_test( + "duplicate", + "a system-scoped skill", + "SYSTEM_DUPLICATE_SKILL_SENTINEL", + ) + .expect("system-scoped duplicate skill seeds"); - let harness = group - .thread("conv-skill-activate-ambiguous") - .script([ - RebornScriptedReply::tool_call( - "builtin.skill_activate", - json!({"names": ["duplicate"]}), - ), - RebornScriptedReply::text("that name is ambiguous"), - ]) - .build() - .await - .expect("thread builds"); - // The user-scoped root is seeded under the SAME (tenant, actor) the built - // thread's run resolves under (`harness.binding`) — only then does the - // user `/skills` mount the run actually reads from contain this file. - capability_harness - .seed_user_skill_for_test( - &harness.binding.tenant_id, - &harness.binding.actor_user_id, - "duplicate", - "a user-scoped skill", - "USER_DUPLICATE_SKILL_SENTINEL", - ) - .expect("user-scoped duplicate skill seeds"); + let harness = group + .thread("conv-skill-activate-ambiguous") + .script([ + RebornScriptedReply::tool_call( + "builtin.skill_activate", + json!({"names": ["duplicate"]}), + ), + RebornScriptedReply::text("that name is ambiguous"), + ]) + .build() + .await + .expect("thread builds"); + // The user-scoped root is seeded under the SAME (tenant, actor) the built + // thread's run resolves under (`harness.binding`) — only then does the + // user `/skills` mount the run actually reads from contain this file. + capability_harness + .seed_user_skill_for_test( + &harness.binding.tenant_id, + &harness.binding.actor_user_id, + "duplicate", + "a user-scoped skill", + "USER_DUPLICATE_SKILL_SENTINEL", + ) + .expect("user-scoped duplicate skill seeds"); - harness - .submit_turn("activate the duplicate skill") - .await - .expect("turn completes despite the ambiguous activation"); + harness + .submit_turn("activate the duplicate skill") + .await + .expect("turn completes despite the ambiguous activation"); - // Model-visible Failed, not a terminal driver_unavailable. - harness - .assert_tool_error(ToolErrorClass::Failed, "ambiguous skill") - .await - .expect("ambiguous skill name surfaced as a recoverable Failed tool error"); - harness - .assert_reply_contains("that name is ambiguous") - .await - .expect("run recovered and finalized"); - // Neither candidate's instructions may leak into a later model request — - // an ambiguous selection activates nothing. - for sentinel in [ - "SYSTEM_DUPLICATE_SKILL_SENTINEL", - "USER_DUPLICATE_SKILL_SENTINEL", - ] { - let err = harness - .assert_model_request_contains(sentinel) - .await - .expect_err("an ambiguous activation must not inject either candidate's instructions"); - assert!( - err.to_string().starts_with("no model request contained"), - "expected the intended \"not found\" assertion failure, got a different harness error: {err}" - ); + // Model-visible Failed, not a terminal driver_unavailable. + harness + .assert_tool_error(ToolErrorClass::Failed, "ambiguous skill") + .await + .expect("ambiguous skill name surfaced as a recoverable Failed tool error"); + harness + .assert_reply_contains("that name is ambiguous") + .await + .expect("run recovered and finalized"); + // Neither candidate's instructions may leak into a later model request — + // an ambiguous selection activates nothing. + for sentinel in [ + "SYSTEM_DUPLICATE_SKILL_SENTINEL", + "USER_DUPLICATE_SKILL_SENTINEL", + ] { + let err = harness + .assert_model_request_contains(sentinel) + .await + .expect_err( + "an ambiguous activation must not inject either candidate's instructions", + ); + assert!( + err.to_string().starts_with("no model request contained"), + "expected the intended \"not found\" assertion failure, got a different harness error: {err}" + ); + } + }, + ); +} + +fn run_async_test_with_stack(name: &'static str, test: F) +where + F: std::future::Future + Send + 'static, +{ + let handle = std::thread::Builder::new() + .name(name.to_string()) + .stack_size(16 * 1024 * 1024) + .spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("tokio test runtime") + .block_on(test); + }) + .expect("spawn stack-sized test thread"); + if let Err(panic) = handle.join() { + std::panic::resume_unwind(panic); } } From 7b83ef06b10db6e111a6d4e1487f8945fb35943f Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 8 Jul 2026 17:55:33 +0800 Subject: [PATCH 11/16] test(reborn): wait for trigger mutator gateway attempts --- .../tests/trigger_poller_e2e.rs | 31 ++++++++++++++++--- 1 file changed, 26 insertions(+), 5 deletions(-) diff --git a/crates/ironclaw_reborn_composition/tests/trigger_poller_e2e.rs b/crates/ironclaw_reborn_composition/tests/trigger_poller_e2e.rs index c70fa01a8ba..260a73cb8c0 100644 --- a/crates/ironclaw_reborn_composition/tests/trigger_poller_e2e.rs +++ b/crates/ironclaw_reborn_composition/tests/trigger_poller_e2e.rs @@ -194,6 +194,21 @@ impl TriggerMutatorAttemptGateway { self.registration_outcomes.lock().await.clone() } + async fn wait_for_registration_outcomes( + &self, + expected_len: usize, + timeout: Duration, + ) -> BTreeMap> { + let stop = Instant::now() + timeout; + loop { + let outcomes = self.registration_outcomes().await; + if outcomes.len() >= expected_len || Instant::now() >= stop { + return outcomes; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + } + /// One provider tool call per scheduled-trigger mutator capability id, /// paired with a payload shaped for that capability's real input schema /// (see `ironclaw_host_runtime::first_party_tools::trigger_management`'s @@ -1447,10 +1462,9 @@ async fn scheduled_trigger_denies_mutators_with_tool_disclosure( // legitimate trigger's id before the fire runs. gateway.set_mutator_target_trigger_id(trigger_id).await; - // Wait for the fire to settle. This is the model's ONLY turn where a - // capability call can be attempted (`invoke_trigger_create` above never - // touched the model), so once `last_status` is set, the mutator - // self-attempts inside `gateway` have already run to completion. + // Wait for the fire submission to settle. `last_status` is written when + // the turn is accepted/replayed by the poller, before the asynchronous + // agent loop necessarily reaches the model gateway under coverage builds. let settled = wait_for_settled( &repo, &tenant_id, @@ -1460,10 +1474,17 @@ async fn scheduled_trigger_denies_mutators_with_tool_disclosure( ) .await; + // Now wait for the fired run's only model turn to reach the gateway and + // attempt all four mutator registrations. Without this second wait the + // test can race on slower instrumented builds: the trigger record is + // already marked accepted, but the model-facing turn has not run yet. + let registration_outcomes = gateway + .wait_for_registration_outcomes(4, Duration::from_secs(15)) + .await; + runtime.shutdown().await.expect("runtime shutdown"); let captured_contents = gateway.captured_message_contents().await; - let registration_outcomes = gateway.registration_outcomes().await; assert_eq!( registration_outcomes.len(), From b4c5a7a1fed806cb4f50aae5716c0313b74e66a7 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 8 Jul 2026 17:59:36 +0800 Subject: [PATCH 12/16] test(reborn): align stack test helper signatures --- tests/integration/outbound_target.rs | 9 +++++---- tests/integration/skill_activate.rs | 9 +++++---- 2 files changed, 10 insertions(+), 8 deletions(-) diff --git a/tests/integration/outbound_target.rs b/tests/integration/outbound_target.rs index 57354b44011..a010706e4d4 100644 --- a/tests/integration/outbound_target.rs +++ b/tests/integration/outbound_target.rs @@ -182,7 +182,7 @@ async fn target_set_unknown_target_routes_to_invalid_input() { fn target_set_disabled_by_settings_routes_to_policy_denied() { run_async_test_with_stack( "target_set_disabled_by_settings_routes_to_policy_denied", - async { + || async { let group = RebornIntegrationGroup::outbound_target_tools() .await .expect("outbound-target-tools group builds"); @@ -355,9 +355,10 @@ async fn target_set_approval_gate_deny_leaves_preference_unchanged() { ); } -fn run_async_test_with_stack(name: &'static str, test: F) +fn run_async_test_with_stack(name: &'static str, test: F) where - F: std::future::Future + Send + 'static, + F: FnOnce() -> Fut + Send + 'static, + Fut: std::future::Future + 'static, { let handle = std::thread::Builder::new() .name(name.to_string()) @@ -367,7 +368,7 @@ where .enable_all() .build() .expect("tokio test runtime") - .block_on(test); + .block_on(test()); }) .expect("spawn stack-sized test thread"); if let Err(panic) = handle.join() { diff --git a/tests/integration/skill_activate.rs b/tests/integration/skill_activate.rs index 4e54f9a3b85..af9ace7a92f 100644 --- a/tests/integration/skill_activate.rs +++ b/tests/integration/skill_activate.rs @@ -189,7 +189,7 @@ async fn skill_activate_over_budget_surfaces_recoverable_failed() { fn skill_activate_ambiguous_name_surfaces_recoverable_failed() { run_async_test_with_stack( "skill_activate_ambiguous_name_surfaces_recoverable_failed", - async { + || async { let group = RebornIntegrationGroup::skill_activation_tools() .await .expect("skill-activation group builds"); @@ -264,9 +264,10 @@ fn skill_activate_ambiguous_name_surfaces_recoverable_failed() { ); } -fn run_async_test_with_stack(name: &'static str, test: F) +fn run_async_test_with_stack(name: &'static str, test: F) where - F: std::future::Future + Send + 'static, + F: FnOnce() -> Fut + Send + 'static, + Fut: std::future::Future + 'static, { let handle = std::thread::Builder::new() .name(name.to_string()) @@ -276,7 +277,7 @@ where .enable_all() .build() .expect("tokio test runtime") - .block_on(test); + .block_on(test()); }) .expect("spawn stack-sized test thread"); if let Err(panic) = handle.join() { From 7b3ecf1c2f3edafad5fa578e775acc94e45cc94a Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 8 Jul 2026 18:06:43 +0800 Subject: [PATCH 13/16] test(webui): assert reopened timeline timestamps --- tests/integration/webui_v2_product_api.rs | 24 ++++++++++++++++++----- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/tests/integration/webui_v2_product_api.rs b/tests/integration/webui_v2_product_api.rs index 91c1d76d2c5..9d913811082 100644 --- a/tests/integration/webui_v2_product_api.rs +++ b/tests/integration/webui_v2_product_api.rs @@ -85,14 +85,28 @@ async fn thread_history_cold_get_and_libsql_reopen() { .await; assert_eq!(status, StatusCode::OK, "body: {body}"); let messages = body["messages"].as_array().expect("messages array"); - assert!( - messages.iter().any(|message| message["kind"] == "assistant" + let finalized_reply = messages + .iter() + .find(|message| { + message["kind"] == "assistant" && message["status"] == "finalized" && message["content"] .as_str() - .is_some_and(|content| content.contains("pong"))), - "expected a finalized assistant message containing 'pong' after a fresh libsql reopen: {body}" - ); + .is_some_and(|content| content.contains("pong")) + }) + .unwrap_or_else(|| { + panic!( + "expected a finalized assistant message containing 'pong' after a fresh libsql reopen: {body}" + ) + }); + for field in ["created_at", "updated_at"] { + assert!( + finalized_reply + .get(field) + .is_some_and(|value| !value.is_null()), + "expected finalized assistant message {field} to be present after a fresh libsql reopen: {body}" + ); + } } /// InMemory sibling of the LibSql reopen above: proves service From 05c628a48294fe21983aaaa30681645fa5871df7 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 8 Jul 2026 18:08:52 +0800 Subject: [PATCH 14/16] test(webui): validate reopened timeline timestamp format --- tests/integration/webui_v2_product_api.rs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/tests/integration/webui_v2_product_api.rs b/tests/integration/webui_v2_product_api.rs index 9d913811082..565807a9f29 100644 --- a/tests/integration/webui_v2_product_api.rs +++ b/tests/integration/webui_v2_product_api.rs @@ -14,6 +14,7 @@ mod support; use std::sync::Arc; use axum::http::StatusCode; +use chrono::{DateTime, Utc}; use ironclaw_events::InMemoryDurableEventLog; use ironclaw_filesystem::{CompositeRootFilesystem, LibSqlRootFilesystem}; use ironclaw_host_api::{CapabilityId, EffectKind, ExtensionId, PermissionMode}; @@ -100,12 +101,11 @@ async fn thread_history_cold_get_and_libsql_reopen() { ) }); for field in ["created_at", "updated_at"] { - assert!( - finalized_reply - .get(field) - .is_some_and(|value| !value.is_null()), - "expected finalized assistant message {field} to be present after a fresh libsql reopen: {body}" - ); + finalized_reply[field] + .as_str() + .unwrap_or_else(|| panic!("{field} missing after libsql reopen: {body}")) + .parse::>() + .unwrap_or_else(|error| panic!("{field} not RFC3339 after libsql reopen: {error}")); } } From acac69af5f5ba6ecc3808a9fd726f240b7202644 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 8 Jul 2026 18:44:39 +0800 Subject: [PATCH 15/16] test(reborn): run over-budget skill activation on larger stack --- tests/integration/skill_activate.rs | 100 +++++++++++++++------------- 1 file changed, 54 insertions(+), 46 deletions(-) diff --git a/tests/integration/skill_activate.rs b/tests/integration/skill_activate.rs index af9ace7a92f..efc6ec0130a 100644 --- a/tests/integration/skill_activate.rs +++ b/tests/integration/skill_activate.rs @@ -127,55 +127,63 @@ async fn skill_criteria_auto_activation_stays_off_on_coordinator_path() { /// error. An oversized system skill (~10k tokens, over /// `DEFAULT_MAX_SKILL_CONTEXT_TOKENS = 4000`) drives the real selection path /// (`reserve_skill_budget` → `ContextBudgetExceeded` → `Failed`). -#[tokio::test] -async fn skill_activate_over_budget_surfaces_recoverable_failed() { - let group = RebornIntegrationGroup::skill_activation_tools() - .await - .expect("skill-activation group builds"); - let capability_harness = group - .capability_harness() - .expect("skill-activation group has a host-runtime capability harness"); - let oversized_prompt = "BLOAT_SKILL_FILLER ".repeat(2200); // ~41.8k chars ≈ 10k tokens > 4000 budget - capability_harness - .seed_system_skill_for_test("bloat", "an oversized skill", &oversized_prompt) - .expect("oversized system skill seeds"); +#[test] +fn skill_activate_over_budget_surfaces_recoverable_failed() { + run_async_test_with_stack( + "skill_activate_over_budget_surfaces_recoverable_failed", + || async { + let group = RebornIntegrationGroup::skill_activation_tools() + .await + .expect("skill-activation group builds"); + let capability_harness = group + .capability_harness() + .expect("skill-activation group has a host-runtime capability harness"); + let oversized_prompt = "BLOAT_SKILL_FILLER ".repeat(2200); // ~41.8k chars ≈ 10k tokens > 4000 budget + capability_harness + .seed_system_skill_for_test("bloat", "an oversized skill", &oversized_prompt) + .expect("oversized system skill seeds"); - let harness = group - .thread("conv-skill-activate-over-budget") - .script([ - RebornScriptedReply::tool_call("builtin.skill_activate", json!({"names": ["bloat"]})), - RebornScriptedReply::text("could not activate"), - ]) - .build() - .await - .expect("thread builds"); + let harness = group + .thread("conv-skill-activate-over-budget") + .script([ + RebornScriptedReply::tool_call( + "builtin.skill_activate", + json!({"names": ["bloat"]}), + ), + RebornScriptedReply::text("could not activate"), + ]) + .build() + .await + .expect("thread builds"); - harness - .submit_turn("activate the big one") - .await - .expect("turn completes despite the failed activation"); + harness + .submit_turn("activate the big one") + .await + .expect("turn completes despite the failed activation"); - // Model-visible Failed, not a terminal driver_unavailable. - harness - .assert_tool_error(ToolErrorClass::Failed, "skill context budget") - .await - .expect("over-budget activation surfaced as a recoverable Failed tool error"); - harness - .assert_reply_contains("could not activate") - .await - .expect("run recovered and finalized"); - // Specific error check (not generic `is_err()`): `assert_model_request_contains` - // has a second, unrelated `Err` path (JSON serialization failure of the - // captured request) — asserting the exact "not found" message text rules - // that out, so an infra-level failure can't masquerade as proof the - // skill's instructions were absent. - let err = harness - .assert_model_request_contains("BLOAT_SKILL_FILLER") - .await - .expect_err("a failed activation must not inject the skill's instructions"); - assert!( - err.to_string().starts_with("no model request contained"), - "expected the intended \"not found\" assertion failure, got a different harness error: {err}" + // Model-visible Failed, not a terminal driver_unavailable. + harness + .assert_tool_error(ToolErrorClass::Failed, "skill context budget") + .await + .expect("over-budget activation surfaced as a recoverable Failed tool error"); + harness + .assert_reply_contains("could not activate") + .await + .expect("run recovered and finalized"); + // Specific error check (not generic `is_err()`): `assert_model_request_contains` + // has a second, unrelated `Err` path (JSON serialization failure of the + // captured request) — asserting the exact "not found" message text rules + // that out, so an infra-level failure can't masquerade as proof the + // skill's instructions were absent. + let err = harness + .assert_model_request_contains("BLOAT_SKILL_FILLER") + .await + .expect_err("a failed activation must not inject the skill's instructions"); + assert!( + err.to_string().starts_with("no model request contained"), + "expected the intended \"not found\" assertion failure, got a different harness error: {err}" + ); + }, ); } From 396095b4dad148c911c83ff9ce0bb626b83de4a2 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 8 Jul 2026 20:32:08 +0800 Subject: [PATCH 16/16] test(reborn): avoid project create coverage stack overflow --- tests/integration/project_create.rs | 54 +++++++++++++++++++++++++---- 1 file changed, 48 insertions(+), 6 deletions(-) diff --git a/tests/integration/project_create.rs b/tests/integration/project_create.rs index 439521441b4..79c1f5d6bd3 100644 --- a/tests/integration/project_create.rs +++ b/tests/integration/project_create.rs @@ -24,8 +24,15 @@ use reborn_support::group::RebornIntegrationGroup; use reborn_support::project_service_fault::FAULT_INJECT_DENIED_PROJECT_NAME; use reborn_support::reply::RebornScriptedReply; -#[tokio::test] -async fn project_create_capability_dispatches_and_persists_project() { +#[test] +fn project_create_capability_dispatches_and_persists_project() { + run_async_test_with_stack( + "project_create_capability_dispatches_and_persists_project", + project_create_capability_dispatches_and_persists_project_impl, + ); +} + +async fn project_create_capability_dispatches_and_persists_project_impl() { let group = RebornIntegrationGroup::project_lifecycle() .await .expect("project-lifecycle group builds"); @@ -92,8 +99,15 @@ async fn project_create_capability_dispatches_and_persists_project() { /// `InvalidInput` error to a recoverable, model-visible `Failed` tool error — /// proving the reject path routes through Completed-turn/Failed-outcome /// plumbing rather than aborting the run. -#[tokio::test] -async fn project_create_invalid_input_routes_to_recoverable_tool_error() { +#[test] +fn project_create_invalid_input_routes_to_recoverable_tool_error() { + run_async_test_with_stack( + "project_create_invalid_input_routes_to_recoverable_tool_error", + project_create_invalid_input_routes_to_recoverable_tool_error_impl, + ); +} + +async fn project_create_invalid_input_routes_to_recoverable_tool_error_impl() { let group = RebornIntegrationGroup::project_lifecycle() .await .expect("project-lifecycle group builds"); @@ -140,8 +154,15 @@ async fn project_create_invalid_input_routes_to_recoverable_tool_error() { /// independently confirmed bug for provider-tool-call-originated invocations /// under local-dev composition (issue #5608) that collapses the "retry twice, /// then Failed" contract into an immediate `driver_unavailable`. -#[tokio::test] -async fn project_create_denied_fault_routes_to_recoverable_tool_error() { +#[test] +fn project_create_denied_fault_routes_to_recoverable_tool_error() { + run_async_test_with_stack( + "project_create_denied_fault_routes_to_recoverable_tool_error", + project_create_denied_fault_routes_to_recoverable_tool_error_impl, + ); +} + +async fn project_create_denied_fault_routes_to_recoverable_tool_error_impl() { let group = RebornIntegrationGroup::project_lifecycle_fault_injected() .await .expect("project-lifecycle fault-injection group builds"); @@ -176,3 +197,24 @@ async fn project_create_denied_fault_routes_to_recoverable_tool_error() { .await .expect("run recovered and finalized instead of dying at the fault"); } + +fn run_async_test_with_stack(name: &'static str, test: F) +where + F: FnOnce() -> Fut + Send + 'static, + Fut: std::future::Future + 'static, +{ + let handle = std::thread::Builder::new() + .name(name.to_string()) + .stack_size(16 * 1024 * 1024) + .spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("tokio test runtime") + .block_on(test()); + }) + .expect("spawn stack-sized test thread"); + if let Err(panic) = handle.join() { + std::panic::resume_unwind(panic); + } +}