diff --git a/docs/streaming.md b/docs/streaming.md index 8b6781587..1ee98af24 100644 --- a/docs/streaming.md +++ b/docs/streaming.md @@ -415,3 +415,34 @@ The bundled plugins, in their default order, each a folder under `plugins/`: A binary diff has no regions, so only `hide-files` can change how it starts out. + +## Progressive annotations (opt-in v4) + +`diffr main HEAD --format ndjson --stream-annotations` emits `start.version = 4`. +Without the flag the v3 stream continues to emit fully enriched files, for +existing consumers including the TUI. + +In v4, each successful `file` contains the complete source, initial folds, +alignment, and authoritative structural changed-line ranges. It is flushed +before enrichment starts. Text files subsequently receive an `annotations` +event, identified by the same file pair: + +```json +{"type":"annotations","file":{"rhs":{"path":"a.rs","oid":"...","mode":"100644"}},"annotations":[{"region_id":12,"label":"validate input\nwrite result"}]} +``` + +Apply each label to the existing region. Do not rebuild the editor or replace +collapsed state, selection, or scroll position. The update never changes region +IDs, fold-state IDs, line ranges, alignment, or counts. An empty annotations list +is valid. If enrichment fails, the event instead has an empty list and an +`error` with code `enrichment_failed`; keep showing the initial file. + +Events for different files may interleave, but a file always precedes its +annotations. `complete` follows all annotation work. Its succeeded/failed totals +count initial file results, not annotation updates. An annotation error still +sets process exit status 2, independently of `--exit-code`. + +Parsing and enrichment use separate workers with a bounded queue. This keeps +slow model requests off parsing workers while bounding retained file trees; +when enrichment falls behind, the queue applies backpressure. Disconnecting +stops queued work; already-running plugin calls finish under their own timeouts. diff --git a/src/cli.rs b/src/cli.rs index 9ff3a7200..1190cf292 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -61,6 +61,7 @@ pub(crate) fn run() -> Result { .arg(flag("find-renames").short('M').conflicts_with("no-renames")) .arg(Arg::new("unified").short('U').long("unified").value_parser(clap::value_parser!(u32)).help("Unchanged lines kept around each change; defaults to plugins.bundled.context.lines")) .arg(Arg::new("format").long("format").value_parser(["ndjson"]).help("Write the event stream to stdout instead of opening the terminal UI")) + .arg(flag("stream-annotations").requires("format").help("Emit initial files followed by deferred annotations (NDJSON v4)")) .arg(flag("syntax").help("Include every token's tree-sitter capture name in --format ndjson output")) .arg(Arg::new("width").long("width").value_parser(clap::value_parser!(usize)).help("Columns for --stat; defaults to the terminal's width")) .arg(flag("ignore-comments")) @@ -104,6 +105,7 @@ pub(crate) fn run() -> Result { } let stream_options = crate::protocol::stream::Options { syntax: args.get_flag("syntax"), + updates: args.get_flag("stream-annotations"), }; let items: Vec = args .get_many::("items") diff --git a/src/plugin/tests/deferred.rs b/src/plugin/tests/deferred.rs new file mode 100644 index 000000000..9b4f25b09 --- /dev/null +++ b/src/plugin/tests/deferred.rs @@ -0,0 +1,132 @@ +use super::*; +use std::io::Write; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; + +struct WatchingWriter { + bytes: Vec, + file_flushed: Arc, +} +impl Write for WatchingWriter { + fn write(&mut self, bytes: &[u8]) -> std::io::Result { + self.bytes.extend_from_slice(bytes); + Ok(bytes.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + if String::from_utf8_lossy(&self.bytes) + .lines() + .any(|line| serde_json::from_str::(line).unwrap()["type"] == "file") + { + self.file_flushed.store(true, Ordering::SeqCst); + } + Ok(()) + } +} + +struct Deferred { + file_flushed: Arc, + fail: bool, +} +impl Runner for Deferred { + fn queries(&self, _: Host) -> anyhow::Result> { + Ok(vec![]) + } + fn classify(&self, _: Host, _: &types::FileEntry) -> anyhow::Result> { + Ok(vec![]) + } + fn mutate( + &self, + _: Host, + _: &types::FileEntry, + _: &types::SourceSides, + ) -> anyhow::Result> { + Ok(vec![]) + } + fn enrich( + &self, + _: Host, + _: &types::FileEntry, + sides: &types::SourceSides, + ) -> anyhow::Result> { + assert!( + self.file_flushed.load(Ordering::SeqCst), + "initial file must be flushed before slow enrichment starts" + ); + if self.fail { + anyhow::bail!("model unavailable"); + } + let sides = tree::sides(sides)?; + let id = sides.rhs().unwrap().regions[0].id; + Ok(vec![types::Annotation { + region_id: id, + label: "do work".into(), + }]) + } +} + +#[test] +fn initial_file_is_flushed_before_enrichment_and_survives_failure() { + for fail in [false, true] { + let file_flushed = Arc::new(AtomicBool::new(false)); + let mut pipeline = Pipeline::default(); + pipeline + .push("deferred", json!({}), &|_, _| { + Ok(Box::new(Deferred { + file_flushed: file_flushed.clone(), + fail, + })) + }) + .unwrap(); + let params = Config::from_toml("").unwrap().compile().unwrap(); + let mut output = WatchingWriter { + bytes: Vec::new(), + file_flushed, + }; + let ended = crate::protocol::stream::write_file( + "a.rs", + "b.rs", + (0, 20), + || { + crate::summary::DiffResult::try_from_sources_with_params( + "a.rs", + "", + "fn f() { work(); }\n", + ¶ms, + ) + }, + ¶ms, + &pipeline, + crate::protocol::stream::Options { + syntax: false, + updates: true, + }, + &mut output, + ) + .unwrap(); + assert_eq!(ended.failed, fail); + assert!(!ended.aborted); + let events: Vec = String::from_utf8(output.bytes) + .unwrap() + .lines() + .map(|line| serde_json::from_str(line).unwrap()) + .collect(); + assert_eq!( + events + .iter() + .map(|event| event["type"].as_str().unwrap()) + .collect::>(), + ["start", "file", "annotations", "complete"] + ); + assert_eq!(events[0]["version"], 4); + assert!( + events[1]["diff"]["structural_changes"]["head"] + .as_array() + .unwrap() + .len() + > 0 + ); + assert_eq!(events[2].get("error").is_some(), fail); + assert_eq!(events[3]["succeeded"], 1); + assert_eq!(events[3]["failed"], 0); + } +} diff --git a/src/plugin/tests/mod.rs b/src/plugin/tests/mod.rs index a29e288b8..2e6030813 100644 --- a/src/plugin/tests/mod.rs +++ b/src/plugin/tests/mod.rs @@ -484,3 +484,5 @@ fn a_subset_of_bundled_plugins_can_use_shared_query_tags() { .compile() .unwrap(); } + +mod deferred; diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index d0dcf57a9..6b3707260 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -56,6 +56,14 @@ pub enum Event { #[serde(flatten)] outcome: Outcome, }, + /// Deferred labels for a previously emitted successful file. A failure + /// affects enrichment only; the initial file and its counts stay valid. + Annotations { + file: Pairing, + annotations: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + error: Option, + }, /// The footer. `aborted` is present when a run-level failure stopped /// the comparison early; every file already emitted stays valid. Complete { diff --git a/src/protocol/stream.rs b/src/protocol/stream.rs index ca5f50531..6dcdb8d86 100644 --- a/src/protocol/stream.rs +++ b/src/protocol/stream.rs @@ -26,6 +26,8 @@ use std::thread; pub(crate) struct Options { /// Emit every token's capture name (`--syntax`). pub(crate) syntax: bool, + /// Opt-in v4 stream; v3 consumers continue to receive finished files. + pub(crate) updates: bool, } /// What the stream ended with: whether any file failed, and whether a @@ -66,9 +68,12 @@ pub(crate) fn write( failed, aborted, .. } = &event { - ended.failed = *failed > 0; + ended.failed |= *failed > 0; ended.aborted = aborted.is_some(); } + if matches!(&event, Event::Annotations { error: Some(_), .. }) { + ended.failed = true; + } serde_json::to_writer(&mut output, &event)?; output.write_all(b"\n")?; output.flush()?; @@ -104,7 +109,7 @@ fn produce( ) -> Result<(), Disconnected> { let send = |event: Event| sender.send(event).map_err(|_| Disconnected); let start = Event::Start { - version: VERSION, + version: if options.updates { 4 } else { VERSION }, lhs: Snapshot::from(&session.comparison.before), rhs: Snapshot::from(&session.comparison.after), files: manifest, @@ -118,40 +123,81 @@ fn produce( session, cancelled: Arc::clone(&cancelled), }; - pool.install(|| { - loader.par_bridge().for_each(|(file, loaded)| { - let (visibility, outcome) = match loaded.and_then(|loaded| diffed(&loaded, options)) { - Ok((entry, diff)) => match shape(pipeline, &entry, diff) { - Ok((visibility, diff)) => (visibility, Outcome::Diff { diff }), - Err(error) => { - // A run-level failure: stop pulling files, let the ones - // in flight finish, and report why in the footer. - failed.fetch_add(1, Ordering::Relaxed); + let (enrich_sender, enrich_receiver) = + sync_channel::<(FileChange, Pairing)>(pool.current_num_threads()); + let enrich_receiver = Mutex::new(enrich_receiver); + thread::scope(|scope| { + let enrich_receiver = &enrich_receiver; + if options.updates { + for _ in 0..pool.current_num_threads() { + let sender = &sender; + let cancelled = &cancelled; + scope.spawn(move || loop { + let job = enrich_receiver.lock().expect("enrichment queue").recv(); + let Ok((entry, sides)) = job else { + break; + }; + if cancelled.load(Ordering::Relaxed) { + continue; + } + let event = enrich_event(pipeline, &entry, &sides); + if sender.send(event).is_err() { cancelled.store(true, Ordering::Relaxed); - aborted.lock().expect("abort reason").get_or_insert(error); - return; } - }, - Err(error) => ( - Visibility::default(), - Outcome::Error { - error: wire_error(&error), - }, - ), - }; - match &outcome { - Outcome::Diff { .. } => succeeded.fetch_add(1, Ordering::Relaxed), - Outcome::Error { .. } => failed.fetch_add(1, Ordering::Relaxed), - }; - let event = Event::File { - file: file.sides, - visibility, - outcome, - }; - if sender.send(event).is_err() { - cancelled.store(true, Ordering::Relaxed); + }); } + } + pool.install(|| { + loader.par_bridge().for_each(|(file, loaded)| { + let (visibility, outcome) = match loaded.and_then(|loaded| diffed(&loaded, options)) + { + Ok((entry, diff)) => match shape(pipeline, &entry, diff, options.updates) { + Ok((visibility, diff)) => (visibility, Outcome::Diff { diff }), + Err(error) => { + // A run-level failure: stop pulling files, let the ones + // in flight finish, and report why in the footer. + failed.fetch_add(1, Ordering::Relaxed); + cancelled.store(true, Ordering::Relaxed); + aborted.lock().expect("abort reason").get_or_insert(error); + return; + } + }, + Err(error) => ( + Visibility::default(), + Outcome::Error { + error: wire_error(&error), + }, + ), + }; + match &outcome { + Outcome::Diff { .. } => succeeded.fetch_add(1, Ordering::Relaxed), + Outcome::Error { .. } => failed.fetch_add(1, Ordering::Relaxed), + }; + let pending = if options.updates { + match &outcome { + Outcome::Diff { + diff: Diff::Text { sides, .. }, + } => Some((file.manifest_entry(), sides.clone())), + _ => None, + } + } else { + None + }; + let event = Event::File { + file: file.sides, + visibility, + outcome, + }; + if sender.send(event).is_err() { + cancelled.store(true, Ordering::Relaxed); + } else if let Some(pending) = pending { + if !cancelled.load(Ordering::Relaxed) && enrich_sender.send(pending).is_err() { + cancelled.store(true, Ordering::Relaxed); + } + } + }); }); + drop(enrich_sender); }); let aborted = aborted.into_inner().expect("abort reason"); send(Event::Complete { @@ -161,6 +207,25 @@ fn produce( }) } +/// An enrichment failure is local to this file's annotations. +fn enrich_event(pipeline: &Pipeline, entry: &FileChange, sides: &Pairing) -> Event { + let (annotations, error) = match pipeline.enrich(entry, sides) { + Ok(annotations) => (annotations, None), + Err(error) => ( + Vec::new(), + Some(Problem { + code: "enrichment_failed".into(), + message: format!("{error:#}"), + }), + ), + }; + Event::Annotations { + file: entry.file.clone(), + annotations, + error, + } +} + /// The wire record for an error, built as it is written. The code comes from /// the typed cause attached where the error arose; an error nothing /// classified is `internal`. @@ -220,6 +285,7 @@ fn shape( pipeline: &Pipeline, entry: &FileChange, diff: Diff, + updates: bool, ) -> anyhow::Result<(Visibility, Diff)> { match diff { Diff::Text { @@ -227,7 +293,11 @@ fn shape( mut stats, .. } => { - let visibility = pipeline.run(entry, &mut sides)?; + let visibility = if updates { + pipeline.prepare(entry, &mut sides)? + } else { + pipeline.run(entry, &mut sides)? + }; let coverage = change_coverage(&sides); stats.visible = coverage.initially_visible.counts(); Ok(( @@ -245,7 +315,11 @@ fn shape( syntax: Vec::new(), regions: Vec::new(), }); - let visibility = pipeline.run(entry, &mut empty)?; + let visibility = if updates { + pipeline.prepare(entry, &mut empty)? + } else { + pipeline.run(entry, &mut empty)? + }; Ok((visibility, Diff::Binary { sides })) } } @@ -283,7 +357,7 @@ pub(crate) fn write_file( let entry = file.manifest_entry(); let mut output = BufWriter::new(output); let start = Event::Start { - version: VERSION, + version: if options.updates { 4 } else { VERSION }, lhs: Snapshot::Path { path: before.to_owned(), }, @@ -319,7 +393,7 @@ pub(crate) fn write_file( }, }, ); - match shape(pipeline, &entry, projected) { + match shape(pipeline, &entry, projected, options.updates) { Ok((visibility, diff)) => ( Some(Event::File { file: file.sides, @@ -337,8 +411,24 @@ pub(crate) fn write_file( serde_json::to_writer(&mut output, record)?; output.write_all(b"\n")?; } + output.flush()?; + let mut enrichment_failed = false; + if options.updates { + if let Some(Event::File { + outcome: Outcome::Diff { + diff: Diff::Text { sides, .. }, + }, + .. + }) = &record + { + let event = enrich_event(pipeline, &entry, sides); + enrichment_failed = matches!(&event, Event::Annotations { error: Some(_), .. }); + serde_json::to_writer(&mut output, &event)?; + output.write_all(b"\n")?; + } + } let ended = Ended { - failed, + failed: failed || enrichment_failed, aborted: aborted.is_some(), }; serde_json::to_writer( @@ -634,7 +724,10 @@ mod conflict_tests { }, ¶ms, &Pipeline::default(), - Options { syntax: false }, + Options { + syntax: false, + updates: false, + }, &mut output, ) .unwrap();