Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
125 changes: 58 additions & 67 deletions crates/ironclaw_loop_support/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -125,10 +125,11 @@ use tokio::sync::{Mutex, OnceCell};

use async_trait::async_trait;
use ironclaw_threads::{
AppendAssistantDraftRequest, AppendToolResultReferenceRequest, ContextMessage,
AppendAssistantDraftRequest, AppendFinalizedAssistantMessageRequest,
AppendToolResultReferenceRequest, ContextMessage, FinalizedAssistantMessageByRunRequest,
LoadContextMessagesRequest, LoadContextWindowRequest, MessageContent, MessageKind,
MessageStatus, ProviderToolCallReferenceEnvelope, SessionThreadError, SessionThreadService,
SummaryArtifact, ThreadHistoryRequest, ThreadMessageId, ThreadMessageRecord, ThreadScope,
ProviderToolCallReferenceEnvelope, SessionThreadError, SessionThreadService, SummaryArtifact,
ThreadHistoryRequest, ThreadMessageId, ThreadMessageRecord, ThreadScope,
ToolResultReferenceEnvelope, ToolResultSafeSummary, UpdateAssistantDraftRequest,
};
use ironclaw_turns::{
Expand Down Expand Up @@ -608,57 +609,54 @@ where
) -> Result<LoopMessageRef, AgentLoopHostError> {
validate_thread_scope_for_run(&self.thread_scope, &self.run_context)?;
let reply_content = request.reply.content;
let draft = self
// One-shot finalize keyed on `turn_run_id`. This collapses the former
// append-draft-then-CAS-finalize two-write sequence into a single
// call: on the filesystem backend with no prior draft it takes the
// append-only fast path (one log append + sequence index, no
// per-message file rewrite); if the loop streamed a draft first
// (`begin_assistant_draft`), it resolves that draft by run and
// finalizes it in place. It is idempotent for a resumed/retried turn —
// a second call returns the already-finalized message.
Comment on lines +612 to +619

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Soften the filesystem-backend guarantee in this adapter comment.

This module does not enforce “one log append + sequence index” or “no per-message file rewrite”; keep the comment to the one-shot thread-service contract, or move backend internals to the backend test/docs.

Suggested wording
-        // One-shot finalize keyed on `turn_run_id`. This collapses the former
-        // append-draft-then-CAS-finalize two-write sequence into a single
-        // call: on the filesystem backend with no prior draft it takes the
-        // append-only fast path (one log append + sequence index, no
-        // per-message file rewrite); if the loop streamed a draft first
-        // (`begin_assistant_draft`), it resolves that draft by run and
-        // finalizes it in place. It is idempotent for a resumed/retried turn —
-        // a second call returns the already-finalized message.
+        // One-shot finalize keyed on `turn_run_id`. This collapses the former
+        // append-draft-then-CAS-finalize sequence into a single thread-service
+        // call. If the loop streamed a draft first (`begin_assistant_draft`),
+        // the thread service resolves that draft by run and finalizes it in
+        // place. It is idempotent for a resumed/retried turn: a second call
+        // returns the already-finalized message.

As per coding guidelines, “Comments that promise guarantees across layers must either be enforced by code/tests or softened to describe intent.”

📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
// One-shot finalize keyed on `turn_run_id`. This collapses the former
// append-draft-then-CAS-finalize two-write sequence into a single
// call: on the filesystem backend with no prior draft it takes the
// append-only fast path (one log append + sequence index, no
// per-message file rewrite); if the loop streamed a draft first
// (`begin_assistant_draft`), it resolves that draft by run and
// finalizes it in place. It is idempotent for a resumed/retried turn —
// a second call returns the already-finalized message.
// One-shot finalize keyed on `turn_run_id`. This collapses the former
// append-draft-then-CAS-finalize sequence into a single thread-service
// call. If the loop streamed a draft first (`begin_assistant_draft`),
// the thread service resolves that draft by run and finalizes it in
// place. It is idempotent for a resumed/retried turn: a second call
// returns the already-finalized message.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@crates/ironclaw_loop_support/src/lib.rs` around lines 612 - 619, Soften the
adapter comment in the one-shot finalize flow so it only describes the
thread-service contract, not filesystem-backend implementation guarantees. In
the comment near the finalize logic in ironclaw_loop_support::lib, remove claims
like “one log append + sequence index” and “no per-message file rewrite,” and
keep the wording focused on `turn_run_id`, `begin_assistant_draft`, and the
idempotent finalize behavior. If backend details are important, move them to the
backend implementation tests/docs instead of this adapter comment.

Source: Coding guidelines

let finalized = match self
.thread_service
.append_assistant_draft(AppendAssistantDraftRequest {
.append_finalized_assistant_message(AppendFinalizedAssistantMessageRequest {
scope: self.thread_scope.clone(),
thread_id: self.run_context.thread_id.clone(),
turn_run_id: self.run_context.run_id.to_string(),
content: MessageContent::text(reply_content.clone()),
})
.await
.map_err(transcript_write_error)?;
if draft.status == MessageStatus::Finalized {
if draft.content.as_deref() == Some(reply_content.as_str()) {
let message_ref = message_ref(draft.message_id)?;
self.emit_assistant_reply_finalized(message_ref.clone())
.await?;
return Ok(message_ref);
{
Ok(message) => message,
Err(append_error) => {
// Concurrency convergence: a racing finalize for the same run
// can win the CAS first, leaving this call's finalize to fail.
// If a finalized reply for this run already exists with matching
// content, converge on it rather than failing the turn — this
// preserves the idempotent-under-concurrent-duplicate contract
// the former append-draft-then-finalize path provided.
match self
.finalized_reply_for_run_matching(&reply_content)
.await?
{
Some(existing) => existing,
None => return Err(transcript_write_error(append_error)),
}
Comment on lines +638 to +644

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

If append_finalized_assistant_message fails (e.g., due to a unique constraint violation on turn_run_id because another concurrent finalize won the race), but the already finalized message has different content (a divergent re-run), the current implementation of finalized_reply_for_run_matching will filter it out and return None.

As a result, the caller will fall back to returning transcript_write_error(append_error), which logs a warning with the raw database unique constraint violation details. This is noisy and misleading because a divergent re-run is a known business logic error that should be handled cleanly.

By fetching the existing finalized message regardless of content (removing the filter from the helper), and letting the subsequent content-comparison check (finalized.content.as_deref() != Some(reply_content.as_str())) handle it, we cleanly return the expected TranscriptWriteFailed error for divergent content without logging the database error as a warning/error.

Suggested change
match self
.finalized_reply_for_run_matching(&reply_content)
.await?
{
Some(existing) => existing,
None => return Err(transcript_write_error(append_error)),
}
match self
.finalized_reply_for_run()
.await?
{
Some(existing) => existing,
None => return Err(transcript_write_error(append_error)),
}

}
};
// Preserve the prior contract: if a finalized reply already exists for
// this run with different content (a divergent re-run), surface a
// transcript write failure rather than silently keeping stale content.
if finalized.content.as_deref() != Some(reply_content.as_str()) {
return Err(AgentLoopHostError::new(
AgentLoopHostErrorKind::TranscriptWriteFailed,
"assistant transcript write failed",
));
}
let finalized = self
.thread_service
.finalize_assistant_message(
&self.thread_scope,
&self.run_context.thread_id,
draft.message_id,
MessageContent::text(reply_content.clone()),
)
.await;
match finalized {
Ok(message) => {
let message_ref = message_ref(message.message_id)?;
self.emit_assistant_reply_finalized(message_ref.clone())
.await?;
Ok(message_ref)
}
Err(error) => {
if let Some(message_id) = self
.already_finalized_matching_reply(draft.message_id, &reply_content)
.await?
{
let message_ref = message_ref(message_id)?;
self.emit_assistant_reply_finalized(message_ref.clone())
.await?;
return Ok(message_ref);
}
Err(transcript_write_error(error))
}
}
let message_ref = message_ref(finalized.message_id)?;
self.emit_assistant_reply_finalized(message_ref.clone())
.await?;
Ok(message_ref)
}

async fn append_capability_result_ref(
Expand Down Expand Up @@ -723,6 +721,25 @@ impl<S> ThreadBackedLoopTranscriptPort<S>
where
S: SessionThreadService + ?Sized + Send + Sync,
{
/// Look up the finalized assistant reply for this run, returning it only
/// when its content matches `reply_content`. Used to converge a finalize
/// that lost a concurrent race onto the winner's message.
async fn finalized_reply_for_run_matching(
&self,
reply_content: &str,
) -> Result<Option<ThreadMessageRecord>, AgentLoopHostError> {
let existing = self
.thread_service
.finalized_assistant_message_by_run(FinalizedAssistantMessageByRunRequest {
scope: self.thread_scope.clone(),
thread_id: self.run_context.thread_id.clone(),
turn_run_id: self.run_context.run_id.to_string(),
})
.await
.map_err(transcript_write_error)?;
Ok(existing.filter(|message| message.content.as_deref() == Some(reply_content)))
}
Comment on lines +727 to +741

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

Simplify this helper to return the existing finalized message regardless of content. This allows the caller to cleanly handle divergent content checks and avoid logging misleading database unique constraint warnings under concurrent races.

    async fn finalized_reply_for_run(
        &self,
    ) -> Result<Option<ThreadMessageRecord>, AgentLoopHostError> {
        self.thread_service
            .finalized_assistant_message_by_run(FinalizedAssistantMessageByRunRequest {
                scope: self.thread_scope.clone(),
                thread_id: self.run_context.thread_id.clone(),
                turn_run_id: self.run_context.run_id.to_string(),
            })
            .await
            .map_err(transcript_write_error)
    }


async fn emit_assistant_reply_finalized(
&self,
message_ref: LoopMessageRef,
Expand Down Expand Up @@ -779,32 +796,6 @@ where
}
Ok(message)
}

async fn already_finalized_matching_reply(
&self,
message_id: ThreadMessageId,
reply_content: &str,
) -> Result<Option<ThreadMessageId>, AgentLoopHostError> {
let history = self
.thread_service
.list_thread_history(ThreadHistoryRequest {
scope: self.thread_scope.clone(),
thread_id: self.run_context.thread_id.clone(),
})
.await
.map_err(transcript_write_error)?;
let expected_run_id = self.run_context.run_id.to_string();
Ok(history.messages.into_iter().find_map(|message| {
let belongs_to_run = message.turn_run_id.as_deref() == Some(expected_run_id.as_str());
let matches_reply = message.status == MessageStatus::Finalized
&& message.content.as_deref() == Some(reply_content);
if message.message_id == message_id && belongs_to_run && matches_reply {
Some(message.message_id)
} else {
None
}
}))
}
}

/// Empty capability surface for the text-only loop-support MVP.
Expand Down
120 changes: 120 additions & 0 deletions crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2252,6 +2252,126 @@ async fn transcript_port_finalizes_assistant_reply_into_durable_thread_history()
);
}

/// Finalizing twice for the same run (e.g. a resumed/retried turn) must be
/// idempotent: one finalized assistant row, same message ref both times. The
/// one-shot `append_finalized_assistant_message` path the port now uses keys on
/// `turn_run_id`, so the second call resolves the existing finalized message
/// rather than appending a duplicate.
#[tokio::test]
async fn transcript_port_finalize_assistant_reply_is_idempotent_for_one_run() {
let fixture = ThreadFixture::new().await;
let adapter = ThreadBackedLoopTranscriptPort::new(
Arc::clone(&fixture.thread_service),
fixture.thread_scope.clone(),
fixture.run_context.clone(),
);

let first = adapter
.finalize_assistant_message(FinalizeAssistantMessage {
reply: AssistantReply {
content: "final answer".to_string(),
},
})
.await
.unwrap();
let second = adapter
.finalize_assistant_message(FinalizeAssistantMessage {
reply: AssistantReply {
content: "final answer".to_string(),
},
})
.await
.unwrap();

assert_eq!(first.as_str(), second.as_str());
let history = fixture
.thread_service
.list_thread_history(ThreadHistoryRequest {
scope: fixture.thread_scope.clone(),
thread_id: fixture.thread_id.clone(),
})
.await
.unwrap();
let assistant_rows: Vec<_> = history
.messages
.iter()
.filter(|message| message.kind == MessageKind::Assistant)
.collect();
assert_eq!(
assistant_rows.len(),
1,
"double finalize must not materialize a second assistant row"
);
assert_eq!(assistant_rows[0].status, MessageStatus::Finalized);
assert_eq!(assistant_rows[0].content.as_deref(), Some("final answer"));
}

/// When the loop streamed a draft first (`begin_assistant_draft` +
/// `update_assistant_draft`), finalizing must finalize that existing draft IN
/// PLACE — same message id, final content — not append a second message. The
/// one-shot path resolves the streamed draft by `turn_run_id` and finalizes it.
#[tokio::test]
async fn transcript_port_finalize_finalizes_existing_streamed_draft_in_place() {
let fixture = ThreadFixture::new().await;
let adapter = ThreadBackedLoopTranscriptPort::new(
Arc::clone(&fixture.thread_service),
fixture.thread_scope.clone(),
fixture.run_context.clone(),
);

let draft_ref = adapter
.begin_assistant_draft(BeginAssistantDraft {
reply: AssistantReply {
content: "partial".to_string(),
},
})
.await
.unwrap();
adapter
.update_assistant_draft(UpdateAssistantDraft {
message_ref: draft_ref.clone(),
reply: AssistantReply {
content: "partial answer".to_string(),
},
})
.await
.unwrap();

let finalized_ref = adapter
.finalize_assistant_message(FinalizeAssistantMessage {
reply: AssistantReply {
content: "partial answer complete".to_string(),
},
})
.await
.unwrap();

assert_eq!(
finalized_ref.as_str(),
draft_ref.as_str(),
"finalize must finalize the streamed draft in place, not create a new message"
);
let history = fixture
.thread_service
.list_thread_history(ThreadHistoryRequest {
scope: fixture.thread_scope.clone(),
thread_id: fixture.thread_id.clone(),
})
.await
.unwrap();
let assistant_rows: Vec<_> = history
.messages
.iter()
.filter(|message| message.kind == MessageKind::Assistant)
.collect();
assert_eq!(assistant_rows.len(), 1);
assert_eq!(assistant_rows[0].status, MessageStatus::Finalized);
assert_eq!(
assistant_rows[0].content.as_deref(),
Some("partial answer complete")
);
}
Comment on lines +2255 to +2373

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Verify whether caller-level coverage already exists for append-error convergence
# and divergent duplicate finalize behavior.
rg -n -C3 'finalized_reply_for_run_matching|TranscriptWriteFailed|divergent|append_finalized_assistant_message' crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

Repository: nearai/ironclaw

Length of output: 556


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Inspect the production finalize path and existing contract coverage around it.
rg -n -C4 'finalize_assistant_message|append_finalized_assistant_message|finalized_reply_for_run_matching|TranscriptWriteFailed|divergent' \
  crates/ironclaw_loop_support/src/lib.rs \
  crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

# Map the relevant test file structure so we can read only the pertinent sections.
ast-grep outline crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs --view expanded

# Read the production finalize path around the reported lines.
sed -n '600,690p' crates/ironclaw_loop_support/src/lib.rs

# Read the nearby test region that already covers finalize behavior.
sed -n '2220,2395p' crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

Repository: nearai/ironclaw

Length of output: 42985


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Search for any caller-level tests that exercise the append-failure fallback or divergence rejection.
rg -n -C3 'finalized_reply_for_run_matching|TranscriptWriteFailed|append_finalized_assistant_message|finalize_assistant_message\(' crates/ironclaw_loop_support/tests crates/ironclaw_loop_support/src

# Inspect the concrete error handling branch in the production code.
sed -n '620,675p' crates/ironclaw_loop_support/src/lib.rs

Repository: nearai/ironclaw

Length of output: 21931


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Inspect the production finalize path and nearby tests for the fallback branch.
sed -n '620,675p' crates/ironclaw_loop_support/src/lib.rs
sed -n '2250,2390p' crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

# Find any other tests that mention the fallback/error-convergence path.
rg -n -C3 'finalized_reply_for_run_matching|TranscriptWriteFailed|append_finalized_assistant_message|divergent' crates/ironclaw_loop_support/tests

Repository: nearai/ironclaw

Length of output: 8517


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Inspect the production branch that handles append failure and divergence.
sed -n '620,675p' crates/ironclaw_loop_support/src/lib.rs

Repository: nearai/ironclaw

Length of output: 2790


🏁 Script executed:

#!/bin/bash
set -euo pipefail

sed -n '2660,2805p' crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

Repository: nearai/ironclaw

Length of output: 5173


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Inspect the candidate finalize tests that may already cover the fallback and divergence paths.
sed -n '2660,2805p' crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

Repository: nearai/ironclaw

Length of output: 5173


🏁 Script executed:

#!/bin/bash
set -euo pipefail

sed -n '2658,2798p' crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

Repository: nearai/ironclaw

Length of output: 4969


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Look for a test that exercises divergent finalize content or the append-failure fallback.
rg -n -C4 'divergent|different content|TranscriptWriteFailed|append_finalized_assistant_message|finalized_reply_for_run_matching|mismatch' \
  crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

# Read the concurrent-duplicate test continuation and the gated service implementation.
sed -n '2790,2865p' crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs
sed -n '4268,4298p' crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs

Repository: nearai/ironclaw

Length of output: 5899


Add a divergent-content finalize test. The concurrent-duplicate case already covers the append-failure convergence path; add a caller-level test that finalizing the same run with different content returns TranscriptWriteFailed and leaves history unchanged.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@crates/ironclaw_loop_support/tests/thread_loop_support_contract.rs` around
lines 2255 - 2373, The finalize tests in
transcript_port_finalize_assistant_reply_is_idempotent_for_one_run and
transcript_port_finalize_finalizes_existing_streamed_draft_in_place only cover
duplicate-content convergence; add a new caller-level test that finalizing the
same run with different assistant content returns TranscriptWriteFailed. Reuse
ThreadFixture and ThreadBackedLoopTranscriptPort, invoke
finalize_assistant_message twice with mismatched content for the same
turn_run_id, and assert history stays unchanged with no extra assistant row.

Source: Path instructions


#[tokio::test]
async fn transcript_port_appends_tool_result_reference_envelope_idempotently() {
let fixture = ThreadFixture::new().await;
Expand Down
Loading