Skip to content
Open
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
52 changes: 40 additions & 12 deletions crates/jp_conversation/src/stream/projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -244,10 +244,13 @@ pub(super) fn apply(events: &mut Vec<InternalEvent>, event_ids: &mut EventIds) -
&mut projected,
&mut event_origins,
event_ids,
summary_marker(turn > 0, run_to + 1 < policies.len()),
&summary.text,
conv_event.timestamp,
turn,
run_to,
TurnOrigin::Summary {
from: turn,
to: run_to,
},
);
}
// Drop the original event — it's covered by the summary.
Expand Down Expand Up @@ -555,6 +558,35 @@ fn resolve_policies(max_turn: usize, compactions: &[crate::Compaction]) -> Vec<T
policies
}

/// The synthetic request text that introduces an injected summary.
///
/// The marker says which part of the conversation the summary stands for, and
/// that the messages around it are still in the model's context.
/// A single "summary of previous conversation" marker is only true for a
/// summary that opens the conversation: placed after kept turns, it tells the
/// model everything above it was replaced, and the model then disowns turns it
/// can plainly read.
///
/// `has_before` and `has_after` say whether any turn precedes or follows the
/// summarized run in the projected view.
const fn summary_marker(has_before: bool, has_after: bool) -> &'static str {
match (has_before, has_after) {
(false, true) => {
"[Summary of the earlier part of this conversation. The messages after it are still in \
your context.]"
}
(true, true) => {
"[Summary of a middle part of this conversation. The messages before and after it are \
still in your context.]"
}
(true, false) => {
"[Summary of the latest part of this conversation. The messages before it are still in \
your context.]"
}
(false, false) => "[Summary of this conversation so far.]",
}
}

/// Inject a synthetic `ChatRequest`/`ChatResponse` pair for a summary.
///
/// A leading `TurnStart` keeps the synthetic pair as its own turn so that
Expand All @@ -563,25 +595,21 @@ fn resolve_policies(max_turn: usize, compactions: &[crate::Compaction]) -> Vec<T
/// The `TurnStart` is not provider-visible, so it is filtered out before the
/// LLM request is built.
///
/// `from`/`to` are the raw turn range this summary replaces; every injected
/// event records it as its [`TurnOrigin`] so the run stays in lockstep with
/// `events`.
/// `marker` is the synthetic request text (see [`summary_marker`]).
/// `origin` is the raw turn range this summary replaces; every injected event
/// records it so the run stays in lockstep with `events`.
fn inject_summary(
events: &mut Vec<InternalEvent>,
origins: &mut Vec<TurnOrigin>,
event_ids: &mut EventIds,
marker: &str,
summary: &str,
timestamp: DateTime<Utc>,
from: usize,
to: usize,
origin: TurnOrigin,
) {
let origin = TurnOrigin::Summary { from, to };
for event in [
ConversationEvent::new(TurnStart, timestamp),
ConversationEvent::new(
ChatRequest::from("[Summary of previous conversation]"),
timestamp,
),
ConversationEvent::new(ChatRequest::from(marker), timestamp),
ConversationEvent::new(ChatResponse::message(summary), timestamp),
] {
events.push(InternalEvent {
Expand Down
70 changes: 70 additions & 0 deletions crates/jp_conversation/src/stream/projection_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1133,6 +1133,76 @@ fn summary_is_injected_as_its_own_turn() {
)));
}

/// The text of the synthetic request introducing the (single) injected summary.
fn summary_marker_text(stream: &ConversationStream) -> String {
stream
.iter()
.find_map(|e| match &e.event.kind {
EventKind::ChatRequest(r) if r.content.starts_with("[Summary") => {
Some(r.content.clone())
}
_ => None,
})
.expect("projection injects a summary request")
}

#[test]
fn summary_marker_for_a_leading_summary() {
let mut stream = message_turns(3);
stream.add_compaction(Compaction::new(0, 1).with_summary(SummaryPolicy::generated("s")));

stream.apply_projection();

assert_eq!(
summary_marker_text(&stream),
"[Summary of the earlier part of this conversation. The messages after it are still in \
your context.]"
);
}

#[test]
fn summary_marker_for_a_middle_summary() {
// The case the marker exists for: a model told the summary replaces
// "previous conversation" treats the kept turns above it as trimmed.
let mut stream = message_turns(4);
stream.add_compaction(Compaction::new(1, 2).with_summary(SummaryPolicy::generated("s")));

stream.apply_projection();

assert_eq!(
summary_marker_text(&stream),
"[Summary of a middle part of this conversation. The messages before and after it are \
still in your context.]"
);
}

#[test]
fn summary_marker_for_a_trailing_summary() {
let mut stream = message_turns(3);
stream.add_compaction(Compaction::new(1, 2).with_summary(SummaryPolicy::generated("s")));

stream.apply_projection();

assert_eq!(
summary_marker_text(&stream),
"[Summary of the latest part of this conversation. The messages before it are still in \
your context.]"
);
}

#[test]
fn summary_marker_for_a_summary_of_every_turn() {
let mut stream = message_turns(2);
stream.add_compaction(Compaction::new(0, 1).with_summary(SummaryPolicy::generated("s")));

stream.apply_projection();

assert_eq!(
summary_marker_text(&stream),
"[Summary of this conversation so far.]"
);
}

#[test]
fn distinct_adjacent_summaries_with_identical_text_stay_separate() {
// Two distinct single-turn summary compactions over adjacent turns that
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ expression: prepared.history
"content": [
{
"type": "text",
"text": "[Summary of previous conversation]"
"text": "[Summary of the earlier part of this conversation. The messages after it are still in your context.]"
}
]
},
Expand Down
42 changes: 36 additions & 6 deletions crates/jp_llm/src/provider/anthropic/acp/transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ use serde::de::DeserializeOwned;
use serde_json::{Value, json};
use sha2::{Digest as _, Sha256};
use tokio::sync::{Notify, mpsc};
use tracing::{debug, instrument::WithSubscriber as _, warn};
use tracing::{debug, info, instrument::WithSubscriber as _, warn};
use uuid::Uuid;

use super::{
Expand Down Expand Up @@ -198,7 +198,8 @@ fn record_failure(state: &Mutex<State>, error: StreamError) -> RpcError {
/// What a connection needs on disk and in the child's environment before it can
/// reach the adapter.
pub(super) struct Launch {
/// The derived transcript the adapter resumes from, removed when dropped.
/// The derived transcript the adapter resumes from, removed when dropped
/// unless `JP_DEBUG` keeps it.
pub(super) artifact: NativeArtifact,

/// The environment the adapter is given, and that the session options are
Expand All @@ -217,6 +218,7 @@ pub(super) struct Launch {
/// Otherwise reads `HOME` and `CLAUDE_CONFIG_DIR`, and writes into the
/// directory they name, so a caller that has no adapter to run should build its
/// own pieces rather than call this.
/// Also reads `JP_DEBUG`, which keeps the transcript after the request.
pub(super) fn launch(
prepared: &PreparedRequest,
context: &QueryContext,
Expand All @@ -225,7 +227,13 @@ pub(super) fn launch(
) -> Result<Launch, Error> {
let directory = native_directory(login_directory)?;
let project = project_name(context);
let artifact = NativeArtifact::write(prepared, &context.root, &directory, &project)?;
let artifact = NativeArtifact::write(
prepared,
&context.root,
&directory,
&project,
keep_transcript(),
)?;
let mut environment = options::environment(prepared, cache);
configure_storage_environment(
&mut environment,
Expand Down Expand Up @@ -570,23 +578,40 @@ fn native_directory(configured: Option<&Utf8Path>) -> Result<Utf8PathBuf, Error>
Ok(directory)
}

/// Whether `JP_DEBUG` asks for the derived transcript to outlive the request.
///
/// The transcript is the only record of what JP handed Claude Code, so keeping
/// it is what lets a debugging session compare it with the conversation.
/// Read the same way `jp` reads the variable: only `1` and `true` turn it on.
fn keep_transcript() -> bool {
env::var("JP_DEBUG")
.as_deref()
.is_ok_and(|value| value == "1" || value == "true")
}

/// The transcript a connection resumes from.
///
/// `session` decides which request JP opens with: `session/load` when there is
/// one to resume, `session/new` otherwise.
/// `path` is the file backing it, which only a run with an adapter to read it
/// needs.
/// `path` is the file removed when this is dropped: `None` when nothing was
/// written, or when the file is kept for debugging.
pub(super) struct NativeArtifact {
pub(super) session: Option<Uuid>,
pub(super) path: Option<Utf8PathBuf>,
}

impl NativeArtifact {
/// Write the transcript under `directory`.
///
/// With `keep`, the file stays on disk after the request and its path is
/// logged at `info`; otherwise the path is logged at `debug` and the file
/// is removed on drop.
fn write(
prepared: &PreparedRequest,
root: &Utf8Path,
directory: &Utf8Path,
project: &str,
keep: bool,
) -> Result<Self, Error> {
if prepared.history.is_empty() {
return Ok(Self {
Expand All @@ -609,9 +634,14 @@ impl NativeArtifact {
options.mode(0o600);
}
let mut file = options.open(&path).map_err(Error::NativeIo)?;
if keep {
info!(%path, "Keeping the derived Claude transcript because JP_DEBUG is set.");
} else {
debug!(%path, "Wrote the derived Claude transcript.");
}
let artifact = Self {
session: Some(session),
path: Some(path),
path: (!keep).then_some(path),
};
for record in prepared.records(session, root, Utc::now()) {
serde_json::to_writer(&mut file, &record).map_err(Error::NativeJson)?;
Expand Down
47 changes: 47 additions & 0 deletions crates/jp_llm/src/provider/anthropic/acp/transport_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,53 @@ fn project_directory_is_scoped_by_host_identity_not_worktree_path() {
assert_eq!(project_name(&context), "jp-c987654321-otvo8");
}

/// Where [`NativeArtifact::write`] puts the transcript for `session`.
fn transcript_path(directory: &Utf8Path, session: Uuid) -> Utf8PathBuf {
directory
.join("projects")
.join("jp-c123-otvo8")
.join(format!("{session}.jsonl"))
}

#[test]
fn the_transcript_is_removed_after_the_request() {
let directory = camino_tempfile::tempdir().unwrap();
let artifact = NativeArtifact::write(
&prepared(),
directory.path(),
directory.path(),
"jp-c123-otvo8",
false,
)
.unwrap();
let path = transcript_path(directory.path(), artifact.session.unwrap());
assert!(path.exists());

drop(artifact);

assert!(!path.exists());
}

#[test]
fn a_kept_transcript_outlives_the_request() {
// `JP_DEBUG` keeps the file: it is the only record of what JP handed
// Claude Code, which the trace log cannot show.
let directory = camino_tempfile::tempdir().unwrap();
let artifact = NativeArtifact::write(
&prepared(),
directory.path(),
directory.path(),
"jp-c123-otvo8",
true,
)
.unwrap();
let path = transcript_path(directory.path(), artifact.session.unwrap());

drop(artifact);

assert!(path.exists());
}

#[test]
fn storage_options_do_not_select_a_different_login_directory() {
let mut environment = BTreeMap::new();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ expression: request
"content": [
{
"type": "text",
"text": "[Summary of previous conversation]"
"text": "[Summary of the earlier part of this conversation. The messages after it are still in your context.]"
}
]
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ expression: request
"content": [
{
"type": "text",
"text": "[Summary of previous conversation]"
"text": "[Summary of the earlier part of this conversation. The messages after it are still in your context.]"
}
]
},
Expand All @@ -28,7 +28,7 @@ expression: request
"content": [
{
"type": "text",
"text": "[Summary of previous conversation]"
"text": "[Summary of a middle part of this conversation. The messages before and after it are still in your context.]"
}
]
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ expression: request
"messages": [
{
"role": "user",
"content": "[Summary of previous conversation]"
"content": "[Summary of the earlier part of this conversation. The messages after it are still in your context.]"
},
{
"role": "assistant",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,15 +7,15 @@ expression: request
"messages": [
{
"role": "user",
"content": "[Summary of previous conversation]"
"content": "[Summary of the earlier part of this conversation. The messages after it are still in your context.]"
},
{
"role": "assistant",
"content": "Summary A: France's capital and the start of the notes lookup."
},
{
"role": "user",
"content": "[Summary of previous conversation]"
"content": "[Summary of a middle part of this conversation. The messages before and after it are still in your context.]"
},
{
"role": "assistant",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ expression: request
{
"parts": [
{
"text": "[Summary of previous conversation]"
"text": "[Summary of the earlier part of this conversation. The messages after it are still in your context.]"
}
],
"role": "user"
Expand Down
Loading
Loading