Skip to content
Merged
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
31 changes: 31 additions & 0 deletions docs/streaming.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
2 changes: 2 additions & 0 deletions src/cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ pub(crate) fn run() -> Result<i32> {
.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"))
Expand Down Expand Up @@ -104,6 +105,7 @@ pub(crate) fn run() -> Result<i32> {
}
let stream_options = crate::protocol::stream::Options {
syntax: args.get_flag("syntax"),
updates: args.get_flag("stream-annotations"),
};
let items: Vec<OsString> = args
.get_many::<OsString>("items")
Expand Down
132 changes: 132 additions & 0 deletions src/plugin/tests/deferred.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,132 @@
use super::*;
use std::io::Write;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;

struct WatchingWriter {
bytes: Vec<u8>,
file_flushed: Arc<AtomicBool>,
}
impl Write for WatchingWriter {
fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
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::<serde_json::Value>(line).unwrap()["type"] == "file")
{
self.file_flushed.store(true, Ordering::SeqCst);
}
Ok(())
}
}

struct Deferred {
file_flushed: Arc<AtomicBool>,
fail: bool,
}
impl Runner for Deferred {
fn queries(&self, _: Host) -> anyhow::Result<Vec<types::QuerySource>> {
Ok(vec![])
}
fn classify(&self, _: Host, _: &types::FileEntry) -> anyhow::Result<Vec<String>> {
Ok(vec![])
}
fn mutate(
&self,
_: Host,
_: &types::FileEntry,
_: &types::SourceSides,
) -> anyhow::Result<Vec<types::Move>> {
Ok(vec![])
}
fn enrich(
&self,
_: Host,
_: &types::FileEntry,
sides: &types::SourceSides,
) -> anyhow::Result<Vec<types::Annotation>> {
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",
&params,
)
},
&params,
&pipeline,
crate::protocol::stream::Options {
syntax: false,
updates: true,
},
&mut output,
)
.unwrap();
assert_eq!(ended.failed, fail);
assert!(!ended.aborted);
let events: Vec<serde_json::Value> = 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::<Vec<_>>(),
["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);
}
}
2 changes: 2 additions & 0 deletions src/plugin/tests/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -484,3 +484,5 @@ fn a_subset_of_bundled_plugins_can_use_shared_query_tags() {
.compile()
.unwrap();
}

mod deferred;
8 changes: 8 additions & 0 deletions src/protocol/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<FileRef>,
annotations: Vec<Annotation>,
#[serde(default, skip_serializing_if = "Option::is_none")]
error: Option<Problem>,
},
/// The footer. `aborted` is present when a run-level failure stopped
/// the comparison early; every file already emitted stays valid.
Complete {
Expand Down
Loading
Loading