Skip to content

Commit fdf7d0b

Browse files
committed
Stream initial files before deferred annotation updates
Opt-in NDJSON v4 preserves existing v3 consumers. Separate bounded enrichment workers emit label-only updates and local errors without invalidating the initial diff. AI-assisted implementation with Codex.
1 parent e876467 commit fdf7d0b

6 files changed

Lines changed: 306 additions & 38 deletions

File tree

‎docs/streaming.md‎

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -415,3 +415,34 @@ The bundled plugins, in their default order, each a folder under `plugins/`:
415415
A binary diff has no regions, so only `hide-files` can change how it starts
416416
out.
417417

418+
419+
## Progressive annotations (opt-in v4)
420+
421+
`diffr main HEAD --format ndjson --stream-annotations` emits `start.version = 4`.
422+
Without the flag the v3 stream continues to emit fully enriched files, for
423+
existing consumers including the TUI.
424+
425+
In v4, each successful `file` contains the complete source, initial folds,
426+
alignment, and authoritative structural changed-line ranges. It is flushed
427+
before enrichment starts. Text files subsequently receive an `annotations`
428+
event, identified by the same file pair:
429+
430+
```json
431+
{"type":"annotations","file":{"rhs":{"path":"a.rs","oid":"...","mode":"100644"}},"annotations":[{"region_id":12,"label":"validate input\nwrite result"}]}
432+
```
433+
434+
Apply each label to the existing region. Do not rebuild the editor or replace
435+
collapsed state, selection, or scroll position. The update never changes region
436+
IDs, fold-state IDs, line ranges, alignment, or counts. An empty annotations list
437+
is valid. If enrichment fails, the event instead has an empty list and an
438+
`error` with code `enrichment_failed`; keep showing the initial file.
439+
440+
Events for different files may interleave, but a file always precedes its
441+
annotations. `complete` follows all annotation work. Its succeeded/failed totals
442+
count initial file results, not annotation updates. An annotation error still
443+
sets process exit status 2, independently of `--exit-code`.
444+
445+
Parsing and enrichment use separate workers with a bounded queue. This keeps
446+
slow model requests off parsing workers while bounding retained file trees;
447+
when enrichment falls behind, the queue applies backpressure. Disconnecting
448+
stops queued work; already-running plugin calls finish under their own timeouts.

‎src/cli.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,7 @@ pub(crate) fn run() -> Result<i32> {
6161
.arg(flag("find-renames").short('M').conflicts_with("no-renames"))
6262
.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"))
6363
.arg(Arg::new("format").long("format").value_parser(["ndjson"]).help("Write the event stream to stdout instead of opening the terminal UI"))
64+
.arg(flag("stream-annotations").requires("format").help("Emit initial files followed by deferred annotations (NDJSON v4)"))
6465
.arg(flag("syntax").help("Include every token's tree-sitter capture name in --format ndjson output"))
6566
.arg(Arg::new("width").long("width").value_parser(clap::value_parser!(usize)).help("Columns for --stat; defaults to the terminal's width"))
6667
.arg(flag("ignore-comments"))
@@ -104,6 +105,7 @@ pub(crate) fn run() -> Result<i32> {
104105
}
105106
let stream_options = crate::protocol::stream::Options {
106107
syntax: args.get_flag("syntax"),
108+
updates: args.get_flag("stream-annotations"),
107109
};
108110
let items: Vec<OsString> = args
109111
.get_many::<OsString>("items")

‎src/plugin/tests/deferred.rs‎

Lines changed: 132 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,132 @@
1+
use super::*;
2+
use std::io::Write;
3+
use std::sync::atomic::{AtomicBool, Ordering};
4+
use std::sync::Arc;
5+
6+
struct WatchingWriter {
7+
bytes: Vec<u8>,
8+
file_flushed: Arc<AtomicBool>,
9+
}
10+
impl Write for WatchingWriter {
11+
fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
12+
self.bytes.extend_from_slice(bytes);
13+
Ok(bytes.len())
14+
}
15+
fn flush(&mut self) -> std::io::Result<()> {
16+
if String::from_utf8_lossy(&self.bytes)
17+
.lines()
18+
.any(|line| serde_json::from_str::<serde_json::Value>(line).unwrap()["type"] == "file")
19+
{
20+
self.file_flushed.store(true, Ordering::SeqCst);
21+
}
22+
Ok(())
23+
}
24+
}
25+
26+
struct Deferred {
27+
file_flushed: Arc<AtomicBool>,
28+
fail: bool,
29+
}
30+
impl Runner for Deferred {
31+
fn queries(&self, _: Host) -> anyhow::Result<Vec<types::QuerySource>> {
32+
Ok(vec![])
33+
}
34+
fn classify(&self, _: Host, _: &types::FileEntry) -> anyhow::Result<Vec<String>> {
35+
Ok(vec![])
36+
}
37+
fn mutate(
38+
&self,
39+
_: Host,
40+
_: &types::FileEntry,
41+
_: &types::SourceSides,
42+
) -> anyhow::Result<Vec<types::Move>> {
43+
Ok(vec![])
44+
}
45+
fn enrich(
46+
&self,
47+
_: Host,
48+
_: &types::FileEntry,
49+
sides: &types::SourceSides,
50+
) -> anyhow::Result<Vec<types::Annotation>> {
51+
assert!(
52+
self.file_flushed.load(Ordering::SeqCst),
53+
"initial file must be flushed before slow enrichment starts"
54+
);
55+
if self.fail {
56+
anyhow::bail!("model unavailable");
57+
}
58+
let sides = tree::sides(sides)?;
59+
let id = sides.rhs().unwrap().regions[0].id;
60+
Ok(vec![types::Annotation {
61+
region_id: id,
62+
label: "do work".into(),
63+
}])
64+
}
65+
}
66+
67+
#[test]
68+
fn initial_file_is_flushed_before_enrichment_and_survives_failure() {
69+
for fail in [false, true] {
70+
let file_flushed = Arc::new(AtomicBool::new(false));
71+
let mut pipeline = Pipeline::default();
72+
pipeline
73+
.push("deferred", json!({}), &|_, _| {
74+
Ok(Box::new(Deferred {
75+
file_flushed: file_flushed.clone(),
76+
fail,
77+
}))
78+
})
79+
.unwrap();
80+
let params = Config::from_toml("").unwrap().compile().unwrap();
81+
let mut output = WatchingWriter {
82+
bytes: Vec::new(),
83+
file_flushed,
84+
};
85+
let ended = crate::protocol::stream::write_file(
86+
"a.rs",
87+
"b.rs",
88+
(0, 20),
89+
|| {
90+
crate::summary::DiffResult::try_from_sources_with_params(
91+
"a.rs",
92+
"",
93+
"fn f() { work(); }\n",
94+
&params,
95+
)
96+
},
97+
&params,
98+
&pipeline,
99+
crate::protocol::stream::Options {
100+
syntax: false,
101+
updates: true,
102+
},
103+
&mut output,
104+
)
105+
.unwrap();
106+
assert_eq!(ended.failed, fail);
107+
assert!(!ended.aborted);
108+
let events: Vec<serde_json::Value> = String::from_utf8(output.bytes)
109+
.unwrap()
110+
.lines()
111+
.map(|line| serde_json::from_str(line).unwrap())
112+
.collect();
113+
assert_eq!(
114+
events
115+
.iter()
116+
.map(|event| event["type"].as_str().unwrap())
117+
.collect::<Vec<_>>(),
118+
["start", "file", "annotations", "complete"]
119+
);
120+
assert_eq!(events[0]["version"], 4);
121+
assert!(
122+
events[1]["diff"]["structural_changes"]["head"]
123+
.as_array()
124+
.unwrap()
125+
.len()
126+
> 0
127+
);
128+
assert_eq!(events[2].get("error").is_some(), fail);
129+
assert_eq!(events[3]["succeeded"], 1);
130+
assert_eq!(events[3]["failed"], 0);
131+
}
132+
}

‎src/plugin/tests/mod.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -484,3 +484,5 @@ fn a_subset_of_bundled_plugins_can_use_shared_query_tags() {
484484
.compile()
485485
.unwrap();
486486
}
487+
488+
mod deferred;

‎src/protocol/mod.rs‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,14 @@ pub enum Event {
5656
#[serde(flatten)]
5757
outcome: Outcome,
5858
},
59+
/// Deferred labels for a previously emitted successful file. A failure
60+
/// affects enrichment only; the initial file and its counts stay valid.
61+
Annotations {
62+
file: Pairing<FileRef>,
63+
annotations: Vec<Annotation>,
64+
#[serde(default, skip_serializing_if = "Option::is_none")]
65+
error: Option<Problem>,
66+
},
5967
/// The footer. `aborted` is present when a run-level failure stopped
6068
/// the comparison early; every file already emitted stays valid.
6169
Complete {

0 commit comments

Comments
 (0)