Repository navigation
fix(claude_code): do not re-send already delivered user turns to a shared session #364
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
fbfbfb5
10cc302
52a8d4d
b182b21
2a12217
878530e
bb9a441
8377ab6
122928e
b05f0f1
964f3d8
2b0e8df
2ba27e9
d083449
9eebb4b
508d472
cff8643
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
|
|
@@ -144,7 +144,7 @@ fn nonzero_exit_message( | |||||
|
|
||||||
| use super::bridge::{ChatMessage, ChatResponse, ProviderDelta}; | ||||||
| use super::event_mapper::EventMapper; | ||||||
| use super::input_builder::build_stdin; | ||||||
| use super::input_builder::{build_stdin_with_delivered, pending_fingerprints}; | ||||||
| use super::session_store::{SessionStore, generate_uuid_v4, is_uuid_v4}; | ||||||
| use super::stream_parser::{ClaudeCodeEvent, StreamJsonParser}; | ||||||
|
|
||||||
|
|
@@ -532,7 +532,17 @@ pub(crate) async fn run_turn(ctx: TurnContext<'_>) -> anyhow::Result<ChatRespons | |||||
|
|
||||||
| // Validate input *before* spawning so we don't launch a process we | ||||||
| // can't feed (CodeRabbit: validate before spawn). | ||||||
| let stdin_bytes = build_stdin(ctx.messages, is_new); | ||||||
| // Reserve this call's pending turns atomically (per store, re-read from | ||||||
| // disk) before spawning, so two calls on one resumed thread cannot both | ||||||
| // send the same turn. The claim is released if the turn fails. | ||||||
| let pending = pending_fingerprints(ctx.messages); | ||||||
| let (delivered, claim) = if is_new || !ctx.persist_session { | ||||||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Release claims before retrying a missing session The missing-session recovery path recursively calls [RULE] stale-delivery-claim · |
||||||
| (std::collections::HashSet::new(), None) | ||||||
| } else { | ||||||
| let (already, claim) = ctx.session_store.claim_delivered(&ctx.thread_id, &pending); | ||||||
|
senamakel marked this conversation as resolved.
|
||||||
| (already, Some(claim)) | ||||||
| }; | ||||||
| let stdin_bytes = build_stdin_with_delivered(ctx.messages, is_new, &delivered); | ||||||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Do not replace an already delivered turn with a user notice When all pending fingerprints are marked delivered, this builder emits a synthetic notice as stdin rather than sending no user turn. Claude receives that notice as a new user message, which can create an extra assistant response or cause the notice to be treated as user content. The resumed-session path should continue the existing session without injecting a replacement user turn, or use a protocol mechanism that does not add user input. [RULE] synthetic-user-input · |
||||||
| if stdin_bytes.is_empty() { | ||||||
| anyhow::bail!("[claude-code][driver] no input messages to deliver"); | ||||||
| } | ||||||
|
|
@@ -736,6 +746,25 @@ pub(crate) async fn run_turn(ctx: TurnContext<'_>) -> anyhow::Result<ChatRespons | |||||
| } | ||||||
| } | ||||||
|
|
||||||
| // The session now holds every pending user turn of this call. Remember | ||||||
| // them so another service or loop iteration resuming the same thread does | ||||||
| // not deliver them again (concurrent callers are covered by the claim above). | ||||||
| // One-shot (non-durable) calls never resume, so they record nothing. | ||||||
| if ctx.persist_session { | ||||||
| let accepted_id = mapper.session_id.as_deref().unwrap_or(&cc_session_id); | ||||||
| if let Err(error) = ctx.session_store.record_delivered(&ctx.thread_id, &pending) { | ||||||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
When two provider instances resume the same thread concurrently, caller A can reserve a pending turn while caller B sees it as in-flight and sends only the already-delivered notice. If B succeeds first, this call records the entire Useful? React with 👍 / 👎. There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Propagate delivery recording failures A successful Claude invocation can reach this branch while Additional
|
||||||
| if let Err(error) = ctx.session_store.record_delivered(&ctx.thread_id, &pending) { | |
| ctx.session_store.record_delivered(&ctx.thread_id, &pending)?; |
[RULE] ignored-errors ·
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -15,6 +15,13 @@ | |
| use super::bridge::ChatMessage; | ||
| use base64::Engine as _; | ||
| use serde_json::{Value, json}; | ||
| use sha2::{Digest, Sha256}; | ||
| use std::collections::HashSet; | ||
|
|
||
| /// Sent on resume when every pending user turn is already in the Claude | ||
| /// session (another service or loop iteration on the same thread delivered it). | ||
| /// Re-sending the text would make the model see the same message twice. | ||
| const ALREADY_DELIVERED_NOTICE: &str = "[The user's latest message was already delivered earlier in this session. Respond to it now; do not treat it as a new, repeated message.]"; | ||
|
|
||
| /// Upper bound on a single decoded inline image's byte size; larger images | ||
| /// are dropped rather than inlined (see [`image_block`]). | ||
|
|
@@ -23,6 +30,17 @@ const MAX_IMAGE_BYTES: usize = 5 * 1024 * 1024; | |
| /// Build the bytes to write to claude's stdin. Returns an empty `Vec` | ||
| /// when there is nothing to send (caller should abort). | ||
| pub fn build_stdin(messages: &[ChatMessage], is_new_session: bool) -> Vec<u8> { | ||
| build_stdin_with_delivered(messages, is_new_session, &HashSet::new()) | ||
| } | ||
|
|
||
| /// Like [`build_stdin`], but on a resumed session skips pending user turns whose | ||
| /// fingerprint (see [`pending_fingerprints`]) is in `delivered`: the session | ||
| /// already received them. Ignored for a new session, which replays context. | ||
| pub fn build_stdin_with_delivered( | ||
| messages: &[ChatMessage], | ||
| is_new_session: bool, | ||
| delivered: &HashSet<String>, | ||
| ) -> Vec<u8> { | ||
| // Resolve any `[Image: … #att:<id>]` placeholders to on-disk `[IMAGE:<path>]` | ||
| // markers so pasted images can be inlined below. No-op for messages that | ||
| // carry no image placeholder, so plain-text turns are unaffected. | ||
|
|
@@ -38,7 +56,30 @@ pub fn build_stdin(messages: &[ChatMessage], is_new_session: bool) -> Vec<u8> { | |
| // context on a new session; never resubmit an earlier answered prompt. | ||
| let last_user_pos = non_system.iter().rposition(|m| m.role == "user"); | ||
| let active_user_pos = last_user_pos.filter(|&pos| pos == non_system.len() - 1); | ||
| let active_user_content = active_user_pos.and_then(|_| pending_user_content(&non_system)); | ||
| let active_user_content = active_user_pos.and_then(|_| { | ||
| let skip = if is_new_session { | ||
| &HashSet::new() | ||
| } else { | ||
| delivered | ||
| }; | ||
| let turns = pending_user_turns(&non_system); | ||
| let fresh: Vec<&str> = turns | ||
|
senamakel marked this conversation as resolved.
|
||
| .iter() | ||
| .filter(|(fingerprint, _)| !skip.contains(fingerprint)) | ||
|
senamakel marked this conversation as resolved.
|
||
| .map(|(_, content)| *content) | ||
| .collect(); | ||
| if turns.is_empty() { | ||
| None | ||
| } else if fresh.is_empty() { | ||
|
senamakel marked this conversation as resolved.
|
||
| tracing::debug!( | ||
| "[claude-code][input] all {} pending user turn(s) already delivered to session", | ||
| turns.len() | ||
| ); | ||
| Some(ALREADY_DELIVERED_NOTICE.to_string()) | ||
| } else { | ||
| Some(fresh.join("\n\n")) | ||
| } | ||
| }); | ||
|
|
||
| let mut content: Vec<Value> = Vec::new(); | ||
| if is_new_session { | ||
|
|
@@ -86,21 +127,66 @@ pub fn build_stdin(messages: &[ChatMessage], is_new_session: bool) -> Vec<u8> { | |
| out.into_bytes() | ||
| } | ||
|
|
||
| /// Join user turns that arrived after the most recent assistant response. A | ||
| /// single Claude input message is required, but queued steering/user messages | ||
| /// must retain their order instead of silently dropping every turn except the | ||
| /// last one. | ||
| fn pending_user_content(non_system: &[&ChatMessage]) -> Option<String> { | ||
| let after_assistant = non_system | ||
| /// User turns that arrived after the most recent assistant response, each with | ||
| /// its delivery fingerprint. A single Claude input message is required, but | ||
| /// queued steering/user messages must retain their order instead of silently | ||
| /// dropping every turn except the last one. | ||
| fn pending_user_turns<'a>(non_system: &[&'a ChatMessage]) -> Vec<(String, &'a str)> { | ||
| let last_assistant = non_system | ||
| .iter() | ||
| .rposition(|message| message.role == "assistant") | ||
| .map_or(0, |position| position + 1); | ||
| let pending: Vec<&str> = non_system[after_assistant..] | ||
| .rposition(|message| message.role == "assistant"); | ||
| let anchor = last_assistant.map_or("", |position| non_system[position].content.as_str()); | ||
|
senamakel marked this conversation as resolved.
senamakel marked this conversation as resolved.
|
||
| // Replies are not unique ("Done" twice), so the reply's ordinal among the | ||
| // assistant turns marks which exchange the pending turns belong to. | ||
| let reply_ordinal = non_system | ||
|
senamakel marked this conversation as resolved.
|
||
| .iter() | ||
| .filter(|message| message.role == "assistant") | ||
| .count(); | ||
|
Comment on lines
+141
to
+144
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
For provider instances holding different slices of the same thread, counting every assistant message makes the same pending turn hash differently. For example, Useful? React with 👍 / 👎. |
||
| let after_assistant = last_assistant.map_or(0, |position| position + 1); | ||
| let mut seen: std::collections::HashMap<&str, usize> = std::collections::HashMap::new(); | ||
| non_system[after_assistant..] | ||
| .iter() | ||
| .filter(|message| message.role == "user" && !message.content.is_empty()) | ||
| .map(|message| message.content.as_str()) | ||
| .collect(); | ||
| (!pending.is_empty()).then(|| pending.join("\n\n")) | ||
| .map(|message| { | ||
| // Occurrence of this exact text among the pending turns: distinct | ||
| // for a repeated message, yet unchanged when an earlier, different | ||
| // pending turn is absent from another service's slice. | ||
| let occurrence = seen.entry(message.content.as_str()).or_insert(0); | ||
|
senamakel marked this conversation as resolved.
|
||
| let fp = fingerprint(anchor, reply_ordinal, *occurrence, &message.content); | ||
| *occurrence += 1; | ||
| (fp, message.content.as_str()) | ||
| }) | ||
| .collect() | ||
| } | ||
|
|
||
| /// Identity of one pending user turn: the assistant reply it follows (text and | ||
| /// ordinal, so two identical replies are still different boundaries), how many | ||
| /// times the same text already appeared among the pending turns, and the text. | ||
| /// User-turn history position is deliberately not part of it, so services that | ||
| /// hold different slices of the same thread still agree on it, while a user | ||
| /// repeating the same words after a new reply gets a new identity. | ||
| fn fingerprint(anchor: &str, reply_ordinal: usize, occurrence: usize, content: &str) -> String { | ||
| let mut hasher = Sha256::new(); | ||
| hasher.update(anchor.as_bytes()); | ||
| hasher.update([0u8]); | ||
| hasher.update(reply_ordinal.to_le_bytes()); | ||
| hasher.update(occurrence.to_le_bytes()); | ||
| hasher.update(content.as_bytes()); | ||
| let digest = hasher.finalize(); | ||
| digest[..12].iter().map(|b| format!("{b:02x}")).collect() | ||
| } | ||
|
|
||
| /// Fingerprints of the user turns a resumed call would deliver for `messages` | ||
| /// (empty unless the final non-system turn is from the user). | ||
| pub fn pending_fingerprints(messages: &[ChatMessage]) -> Vec<String> { | ||
| let non_system: Vec<&ChatMessage> = messages.iter().filter(|m| m.role != "system").collect(); | ||
| if non_system.last().is_none_or(|m| m.role != "user") { | ||
| return Vec::new(); | ||
| } | ||
| pending_user_turns(&non_system) | ||
| .into_iter() | ||
| .map(|(fingerprint, _)| fingerprint) | ||
| .collect() | ||
| } | ||
|
|
||
| /// Render the turns before `end` (the latest user turn) as a plain-text preamble | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Qualify the no-duplicate-delivery claim as process-local
The comment asserts "two calls on one resumed thread cannot both send the same turn", but the reservation is explicitly process-local (
IN_FLIGHTis a static, andsession_store.rsdocuments "Separate OS processes are not locked against each other"). Two OS processes sharing the same workspace and thread — the exact scenario this feature exists for — can still both spawn the CLI and deliver the same turn, and no test in the diff pins the multi-process case. The store-level serialization was added in this revision, which resolved the earlier store findings; what remains is this overclaiming comment at the driver. Soften it to match what the code guarantees, or add an inter-process lock (e.g. an flock on<store>.lockheld across claim→record) if cross-process callers are real.[RULE] unpinned-invariant ·