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/src/reborn_services.rs b/crates/ironclaw_product_workflow/src/reborn_services.rs index a0c5f4565d8..bc5269e43d5 100644 --- a/crates/ironclaw_product_workflow/src/reborn_services.rs +++ b/crates/ironclaw_product_workflow/src/reborn_services.rs @@ -6607,6 +6607,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, @@ -6647,6 +6648,7 @@ fn map_thread_error(error: SessionThreadError) -> RebornServicesError { SessionThreadError::GeneratedThreadId(_) | SessionThreadError::Serialization(_) | SessionThreadError::Deserialization(_) + | SessionThreadError::InvalidMessageTimestamp { .. } | SessionThreadError::Backend(_) => RebornServicesError::service_unavailable(true), } } diff --git a/crates/ironclaw_product_workflow/tests/reborn_services_contract.rs b/crates/ironclaw_product_workflow/tests/reborn_services_contract.rs index 281a600f768..3f94ccbd810 100644 --- a/crates/ironclaw_product_workflow/tests/reborn_services_contract.rs +++ b/crates/ironclaw_product_workflow/tests/reborn_services_contract.rs @@ -192,6 +192,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 ad25bfcac8f..99b2fcb4bbb 100644 --- a/crates/ironclaw_reborn/src/loop_exit_applier/tests/mod.rs +++ b/crates/ironclaw_reborn/src/loop_exit_applier/tests/mod.rs @@ -839,6 +839,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, @@ -856,6 +858,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 7b7c81a8308..d83132ef1cc 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, @@ -549,6 +551,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, @@ -697,6 +701,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..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>, } @@ -203,6 +205,16 @@ 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. 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, pub source_binding_id: Option, pub reply_target_binding_id: Option, @@ -222,6 +234,51 @@ 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::InvalidMessageTimestamp { + message_id: message.message_id, + context, + violation: crate::error::TimestampViolation::MissingDurableTimestamps, + }); + } + 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. +/// 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(); + if created_at_cleared || updated_at_cleared { + return Err(crate::error::SessionThreadError::InvalidMessageTimestamp { + message_id, + context, + violation: crate::error::TimestampViolation::ClearedDurableTimestamps, + }); + } + Ok(()) +} + /// Summary artifact over a stable transcript sequence range. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct SummaryArtifact { @@ -583,6 +640,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"); @@ -635,6 +715,50 @@ 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::InvalidMessageTimestamp { + context: "test message", + violation: crate::error::TimestampViolation::MissingDurableTimestamps, + .. + } + )); + } + + #[test] + 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_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, + crate::error::SessionThreadError::InvalidMessageTimestamp { + context: "test update", + violation: crate::error::TimestampViolation::ClearedDurableTimestamps, + .. + } + )); + } + #[test] fn into_text_keeps_text_when_attachment_free() { // The non-lossy use: `into_text` is the text-only accessor for content @@ -717,6 +841,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/error.rs b/crates/ironclaw_threads/src/error.rs index 8a2cf584f70..f975b9cd900 100644 --- a/crates/ironclaw_threads/src/error.rs +++ b/crates/ironclaw_threads/src/error.rs @@ -3,6 +3,15 @@ use thiserror::Error; use crate::{MessageStatus, ThreadMessageId}; +/// Specific timestamp contract violation on a transcript message. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] +pub enum TimestampViolation { + #[error("missing durable timestamps")] + MissingDurableTimestamps, + #[error("cleared durable timestamps")] + ClearedDurableTimestamps, +} + /// Canonical thread/transcript service errors. #[derive(Debug, Error)] pub enum SessionThreadError { @@ -48,6 +57,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_threads/src/filesystem_service.rs b/crates/ironclaw_threads/src/filesystem_service.rs index e553a88ed8b..1d041f23327 100644 --- a/crates/ironclaw_threads/src/filesystem_service.rs +++ b/crates/ironclaw_threads/src/filesystem_service.rs @@ -43,7 +43,7 @@ use std::{ }; 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, FilesystemError, @@ -355,6 +355,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( @@ -392,6 +393,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 @@ -488,6 +490,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)?; @@ -1201,15 +1204,16 @@ where /// recency authority for listing; avoiding a full `thread.json` CAS here /// keeps activity writes row-shaped. /// - /// Best-effort under contention: a `VersionMismatch` means a concurrent /// The message itself is already durably written before most callers reach /// this path, so best-effort wrappers log and continue on failure. async fn touch_thread_updated_at( &self, scope: &ThreadScope, thread_id: &ThreadId, + updated_at: DateTime, ) -> Result<(), SessionThreadError> { - self.touch_thread_index_updated_at(scope, thread_id).await + self.touch_thread_index_updated_at(scope, thread_id, updated_at) + .await } /// Best-effort recency stamp for after-commit call sites. The message is @@ -1219,8 +1223,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, @@ -1281,7 +1293,17 @@ where (record, CasExpectation::Absent) } }; + let before_created_at = message.created_at; + let before_updated_at = message.updated_at; mutate(&mut 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)?; match put_with_cas( self.filesystem.as_ref(), @@ -1473,12 +1495,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(), @@ -1562,7 +1587,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 { @@ -1683,12 +1708,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, @@ -1739,6 +1767,7 @@ where return Ok(existing); } let content = request.content.clone(); + let now = Utc::now(); let finalized = self .apply_message_update( &request.scope, @@ -1749,25 +1778,29 @@ where message.status = MessageStatus::Finalized; message.content = Some(content.clone().into_text()); message.attachments = Vec::new(); + message.updated_at = Some(now); Ok(()) }, ) .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); } 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, @@ -1808,7 +1841,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) } @@ -1855,14 +1888,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(|| { @@ -1877,12 +1913,25 @@ where .map_err(SessionThreadError::Serialization)? { message.content = Some(content); + 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); } @@ -1896,12 +1945,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, @@ -1953,12 +2005,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, @@ -1971,6 +2026,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(), @@ -2028,7 +2084,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, @@ -2049,10 +2107,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(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( @@ -2064,20 +2126,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(); - 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( @@ -2092,12 +2160,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(now); Ok(()) }) .await?; @@ -2105,7 +2175,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) } @@ -2119,20 +2189,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()); - 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( @@ -2885,6 +2961,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/filesystem_service/thread_index.rs b/crates/ironclaw_threads/src/filesystem_service/thread_index.rs index 6e7197c0491..cfb2f4824f9 100644 --- a/crates/ironclaw_threads/src/filesystem_service/thread_index.rs +++ b/crates/ironclaw_threads/src/filesystem_service/thread_index.rs @@ -3,7 +3,7 @@ use std::{ sync::Arc, }; -use chrono::Utc; +use chrono::{DateTime, Utc}; use futures::StreamExt; use ironclaw_filesystem::{ CasApply, ContentType, Entry, FileType, Filter, Page, RecordKind, RootFilesystem, cas_update, @@ -275,8 +275,8 @@ where &self, scope: &ThreadScope, thread_id: &ThreadId, + updated_at: DateTime, ) -> Result<(), SessionThreadError> { - let now = Utc::now(); let path = thread_index_record_path(scope, thread_id)?; let resource_scope = scope.to_resource_scope(); let scope_for_retry = scope.clone(); @@ -307,11 +307,11 @@ where false, )); }; - stored.record.updated_at = Some(now); + stored.record.updated_at = Some(updated_at); Self::thread_index_record(&stored) } }; - index.record.updated_at = Some(now); + index.record.updated_at = Some(updated_at); Ok(CasApply::new(index, true)) } }, @@ -880,7 +880,7 @@ mod tests { let thread_id = ThreadId::new("thread-missing-touch").unwrap(); service - .touch_thread_index_updated_at(&request_scope, &thread_id) + .touch_thread_index_updated_at(&request_scope, &thread_id, chrono::Utc::now()) .await .expect("missing touch is a no-op"); service.mark_thread_index_scope_complete(&request_scope); diff --git a/crates/ironclaw_threads/src/in_memory.rs b/crates/ironclaw_threads/src/in_memory.rs index 161cb32f5ec..2e471eb9e75 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()); - thread.messages.push(ThreadMessageRecord { + let now = Utc::now(); + thread.record.updated_at = Some(now); + let message = 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, @@ -176,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( @@ -243,9 +248,19 @@ 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_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_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()) } @@ -258,9 +273,19 @@ 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_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_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()) } @@ -277,12 +302,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 +325,8 @@ 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); + crate::contract::validate_new_message_timestamps(&message, "assistant draft")?; thread.messages.push(message.clone()); Ok(message) } @@ -313,18 +342,34 @@ 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(); + 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_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); } 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 +384,8 @@ 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); + crate::contract::validate_new_message_timestamps(&message, "finalized assistant")?; thread.messages.push(message.clone()); Ok(message) } @@ -363,12 +409,19 @@ 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_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 .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( @@ -392,9 +445,25 @@ impl SessionThreadService for InMemorySessionThreadService { .map_err(SessionThreadError::Serialization)? { existing.content = Some(content); + changed = true; } } - return Ok(existing.clone()); + if changed { + existing.updated_at = Some(now); + } + 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(); + if changed { + thread.record.updated_at = Some(now); + } + return Ok(updated); } if let Some(provider_call) = &provider_call { provider_call @@ -403,12 +472,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 +495,8 @@ 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); + crate::contract::validate_new_message_timestamps(&message, "tool result reference")?; thread.messages.push(message.clone()); Ok(message) } @@ -454,12 +527,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 +550,8 @@ 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); + crate::contract::validate_new_message_timestamps(&message, "capability display preview")?; thread.messages.push(message.clone()); Ok(message) } @@ -483,56 +560,90 @@ 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_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(), + ) })?; - 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()))?, - ); - 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_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() + }; + 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(); - 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_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_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, + "update_assistant_draft", + )?; + message.clone() + }; + thread.record.updated_at = Some(now); + Ok(updated) } async fn finalize_assistant_message( @@ -543,11 +654,29 @@ 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(); + 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_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); Ok(message.clone()) } @@ -555,19 +684,37 @@ 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); - 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_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_timestamp_fields_not_cleared( + message.message_id, + before_created_at, + before_updated_at, + message.created_at, + message.updated_at, + "redact_message", + )?; + message.clone() + }; + thread.record.updated_at = Some(now); + Ok(updated) } async fn load_context_window( @@ -1092,6 +1239,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 f675d089970..b85dedd0a55 100644 --- a/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs +++ b/crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs @@ -454,6 +454,86 @@ 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"); + assert_eq!(draft.updated_at, Some(draft_created_at)); + + 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"); + + 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..04205f336a2 100644 --- a/crates/ironclaw_threads/tests/session_thread_contract.rs +++ b/crates/ironclaw_threads/tests/session_thread_contract.rs @@ -954,6 +954,84 @@ 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"); + assert_eq!(draft.updated_at, Some(draft_created_at)); + + 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"); + + 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/frontend/src/pages/chat/lib/history-messages.test.ts b/crates/ironclaw_webui_v2/frontend/src/pages/chat/lib/history-messages.test.ts index f08c566d250..08c1fb3ba3e 100644 --- a/crates/ironclaw_webui_v2/frontend/src/pages/chat/lib/history-messages.test.ts +++ b/crates/ironclaw_webui_v2/frontend/src/pages/chat/lib/history-messages.test.ts @@ -174,6 +174,79 @@ 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: 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([ + { + 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([ { diff --git a/crates/ironclaw_webui_v2/frontend/src/pages/chat/lib/history-messages.ts b/crates/ironclaw_webui_v2/frontend/src/pages/chat/lib/history-messages.ts index 8e723aaf60a..19d8e9539f4 100644 --- a/crates/ironclaw_webui_v2/frontend/src/pages/chat/lib/history-messages.ts +++ b/crates/ironclaw_webui_v2/frontend/src/pages/chat/lib/history-messages.ts @@ -158,10 +158,15 @@ 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. 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") { + return record.updated_at || record.created_at || null; + } + return record.created_at || record.updated_at || null; } function toolCardFromPreviewRecord(record) { diff --git a/tests/integration/outbound_target.rs b/tests/integration/outbound_target.rs index b0823175a33..a010706e4d4 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,24 @@ 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: 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); + } +} 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); + } +} diff --git a/tests/integration/skill_activate.rs b/tests/integration/skill_activate.rs index 0cd02430ae4..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}" + ); + }, ); } @@ -185,74 +193,102 @@ 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: 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); } } diff --git a/tests/integration/webui_v2_product_api.rs b/tests/integration/webui_v2_product_api.rs index 91c1d76d2c5..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}; @@ -85,14 +86,27 @@ 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"] { + 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}")); + } } /// InMemory sibling of the LibSql reopen above: proves service diff --git a/tools/ironclaw_stress/src/user_turn.rs b/tools/ironclaw_stress/src/user_turn.rs index 12b63baf1ee..7ea17c4717d 100644 --- a/tools/ironclaw_stress/src/user_turn.rs +++ b/tools/ironclaw_stress/src/user_turn.rs @@ -2352,6 +2352,7 @@ fn thread_failure(stage: impl Into, error: SessionThreadError) -> Operat SessionThreadError::Serialization(_) | SessionThreadError::Deserialization(_) => { "thread_serialization" } + SessionThreadError::InvalidMessageTimestamp { .. } => "thread_timestamp_invalid", SessionThreadError::Backend(_) => "thread_backend", }; OperationFailure::new(bucket, stage, error)