Skip to content

Commit 26aed20

Browse files
committed
Pool opt-in plugin instances for parallel enrichment
Compile components once and bound deferred calls across files, reserving an initial-presentation instance so model latency cannot block initial folds. Serial plugins retain one instance. AI-assisted implementation with Codex.
1 parent fdf7d0b commit 26aed20

6 files changed

Lines changed: 298 additions & 19 deletions

File tree

‎docs/plugins.md‎

Lines changed: 24 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -133,11 +133,28 @@ diffr builds these records once per call, from the file's manifest entry and
133133
its trees as the plugins before left them, and hands the same records to a
134134
native plugin and to a component. Nothing else reaches a plugin.
135135

136-
diffr makes one instance of each plugin per run and calls it for every file.
137-
A native plugin's instance is called from several files' workers in parallel,
138-
so it must be `Send + Sync`; interior mutability is the plugin's own
139-
business. A component's instance lives in one wasmtime store, which runs one
140-
call at a time: diffr calls a component plugin one file at a time.
136+
By default diffr makes one instance of each plugin and serializes calls to it.
137+
A plugin author can declare `parallel = true` at the top of `plugin.toml` to
138+
permit independent deferred-work instances. This promises that calls do not depend on shared
139+
mutable state, call order, or unique external side effects. Initial presentation and enrichment may use different instances. Plugin ordering within a file
140+
is still sequential.
141+
142+
Opted-in plugins expose a host-owned `instances` setting (default 4, range 1–64):
143+
144+
```toml
145+
[plugins.bundled.summarize]
146+
instances = 4
147+
```
148+
149+
The pool is shared across every file and phase of that plugin, so the bound is
150+
per comparison, not per file. A separate initial-presentation instance keeps
151+
queries, classification, and mutation responsive while the deferred pool is busy.
152+
With `instances = 1`, all phases share the same instance for serial debugging. The host compiles a WASM component once and creates
153+
independent stores. Each store still executes only one call at a time. Leases
154+
return on success and errors. Cancellation skips queued enrichment; active calls
155+
retain their configured request timeouts. Native plugins use the same pool.
156+
The summarizer batches selected folds into one request per file, so this
157+
parallelizes files; it does not split a file's batch into per-fold requests.
141158

142159
### Fresh ids
143160

@@ -319,8 +336,8 @@ next run.
319336

320337
`plugins/summarize/src/http.rs` is a complete outgoing-request example for
321338
plugin authors. The host provides TLS; components do not need to embed a
322-
TLS implementation. WASM instances process files serially; the legacy
323-
`max_concurrency` option is accepted but has no effect. Timeouts apply to
339+
TLS implementation. Each WASM instance processes calls serially; `instances` controls the host pool.
340+
The legacy `max_concurrency` option is still accepted but has no effect. Timeouts apply to
324341
connection, first-byte and between-byte waits.
325342
`cargo test --features wasm-plugin-tests --bin diffr external_component_summarizes_over_http`
326343
checks the external component against a local model endpoint, including retries.

‎plugins/summarize/plugin.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
name = "summarize"
2+
parallel = true
23
title = "Summaries"
34
description = "Pseudocode summaries for large new function bodies and tests."
45

‎src/plugin/config.rs‎

Lines changed: 42 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ pub(crate) const PATH: &str = "path";
2222

2323
/// The keys in a plugin entry that diffr owns: a plugin's options may not
2424
/// use them.
25-
pub(crate) const RESERVED: [&str; 2] = [ENABLED, PATH];
25+
pub(crate) const RESERVED: [&str; 3] = [ENABLED, PATH, "instances"];
2626

2727
/// A plugin folder's description, and its component when it has one.
2828
pub(crate) const MANIFEST_FILE: &str = "plugin.toml";
@@ -35,6 +35,10 @@ pub(crate) struct Manifest {
3535
/// The plugin's entry name in `[plugins]`, and the prefix of every tag
3636
/// its queries set: `<name>:<tag>`.
3737
pub(crate) name: String,
38+
/// Opt in only when independent instances can process different files.
39+
/// No cross-call state, ordering, or unique external side effects may be required.
40+
#[serde(default)]
41+
pub(crate) parallel: bool,
3842
/// The human name settings screens group the plugin's settings under.
3943
pub(crate) title: String,
4044
#[serde(default)]
@@ -176,6 +180,13 @@ impl Manifest {
176180
"x-group": group,
177181
}),
178182
);
183+
if self.parallel {
184+
properties.insert("instances".into(), json!({
185+
"type": "integer", "minimum": 1, "maximum": 64, "default": 4,
186+
"title": "Parallel instances", "description": "Maximum simultaneous deferred plugin calls across all files.",
187+
"x-group": group,
188+
}));
189+
}
179190
for (key, option) in &self.options {
180191
let mut option = option.clone();
181192
if let Some(option) = option.as_object_mut() {
@@ -301,6 +312,9 @@ pub(crate) struct Entry {
301312
/// plugin's default.
302313
#[serde(skip_serializing_if = "Option::is_none")]
303314
pub(crate) enabled: Option<bool>,
315+
/// Host-owned deferred-work bound, shared by all files of this plugin.
316+
#[serde(skip_serializing_if = "Option::is_none")]
317+
pub(crate) instances: Option<usize>,
304318
/// The plugin's folder on disk, as written: relative to the
305319
/// configuration file's directory, or absolute.
306320
#[serde(skip_serializing_if = "Option::is_none")]
@@ -422,6 +436,12 @@ impl PluginsConfig {
422436
};
423437
let manifest = &folder.manifest;
424438
entry.enabled.get_or_insert(manifest.enabled_by_default());
439+
let instances = entry
440+
.instances
441+
.unwrap_or(if manifest.parallel { 4 } else { 1 });
442+
if !(1..=64).contains(&instances) || (!manifest.parallel && instances != 1) {
443+
return Err(ConfigError(format!("plugins.{reference}.instances: expected 1..=64 for a parallel plugin, or 1 for a serial plugin")));
444+
}
425445
manifest
426446
.validate(&entry.options)
427447
.map_err(|error| ConfigError(format!("plugins.{reference}: {error}")))?;
@@ -742,3 +762,24 @@ mod tests {
742762
.starts_with("plugins.bundled.summarize: provider: "));
743763
}
744764
}
765+
766+
#[cfg(test)]
767+
mod concurrency_tests {
768+
use crate::config::Config;
769+
770+
#[test]
771+
fn concurrency_requires_plugin_opt_in_and_a_bounded_positive_count() {
772+
for count in [0, 65] {
773+
assert!(Config::from_toml(&format!(
774+
"[plugins.bundled.summarize]\ninstances = {count}\n"
775+
))
776+
.is_err());
777+
}
778+
assert!(Config::from_toml("[plugins.bundled.context]\ninstances = 2\n").is_err());
779+
let config = Config::from_toml("[plugins.bundled.summarize]\ninstances = 3\n").unwrap();
780+
assert_eq!(
781+
config.plugins.entries["bundled.summarize"].instances,
782+
Some(3)
783+
);
784+
}
785+
}

‎src/plugin/mod.rs‎

Lines changed: 46 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@
88
//! `plugin.toml` (name, title, options schema) and, when it has
99
//! one, its component. [`Pipeline::from_config`] then takes the component
1010
//! ([`wasm`]) or the bundled native implementation registered under
11-
//! the plugin's name ([`native`]), and makes the plugin's one instance for
11+
//! the plugin's name ([`native`]), and makes the plugin's instance pool for
1212
//! the run from its options. From there a [`Runner`] is a [`Runner`]: each
1313
//! call gets the contract's records, built once per call from the file's
1414
//! manifest entry and its current trees, and the host functions of [`host`].
@@ -32,6 +32,7 @@ pub(crate) mod builtin;
3232
pub(crate) mod config;
3333
pub(crate) mod host;
3434
pub(crate) mod native;
35+
mod pool;
3536
pub(crate) mod queries;
3637
pub(crate) mod wasm;
3738

@@ -137,19 +138,31 @@ impl Pipeline {
137138
Some(engine) => engine,
138139
None => engine.insert(wasm::engine()?),
139140
};
140-
pipeline.push(name, options, &|host, options| {
141-
wasm::WasmPlugin::load(engine, &path)?.create(host, options)
142-
})?
141+
let plugin = wasm::WasmPlugin::load(engine, &path)
142+
.with_context(|| format!("plugins.{name}"))?;
143+
pipeline.push_instances(
144+
name,
145+
options,
146+
entry
147+
.instances
148+
.unwrap_or(if folder.manifest.parallel { 4 } else { 1 }),
149+
&|host, options| plugin.create(host, options),
150+
)?
143151
}
144152
None => {
145153
let create = native::lookup(&folder.manifest.name)?.ok_or_else(|| {
146154
anyhow!(
147155
"plugins.{name}: no native implementation is registered for this bundled plugin"
148156
)
149157
})?;
150-
pipeline.push(name, options, &|host, options| {
151-
native::registered(create, host, options)
152-
})?
158+
pipeline.push_instances(
159+
name,
160+
options,
161+
entry
162+
.instances
163+
.unwrap_or(if folder.manifest.parallel { 4 } else { 1 }),
164+
&|host, options| native::registered(create, host, options),
165+
)?
153166
}
154167
}
155168
}
@@ -158,11 +171,35 @@ impl Pipeline {
158171

159172
/// Make the plugin `name` with `create` from `options` and add it to the
160173
/// end of the pipeline.
174+
#[cfg(test)]
161175
fn push(&mut self, name: &str, options: Value, create: Create<'_>) -> anyhow::Result<()> {
176+
self.push_instances(name, options, 1, create)
177+
}
178+
179+
fn push_instances(
180+
&mut self,
181+
name: &str,
182+
options: Value,
183+
instances: usize,
184+
create: Create<'_>,
185+
) -> anyhow::Result<()> {
162186
let reference = name;
163187
let name: Arc<str> = name.split_once('.').map_or(name, |(_, name)| name).into();
164-
let runner = create(self.host(&name), &options.to_string())
165-
.with_context(|| format!("plugins.{reference}"))?;
188+
let runners = (0..instances)
189+
.map(|_| {
190+
create(self.host(&name), &options.to_string())
191+
.with_context(|| format!("plugins.{reference}"))
192+
})
193+
.collect::<anyhow::Result<Vec<_>>>()?;
194+
let initial = if instances > 1 {
195+
Some(
196+
create(self.host(&name), &options.to_string())
197+
.with_context(|| format!("plugins.{reference}"))?,
198+
)
199+
} else {
200+
None
201+
};
202+
let runner = Box::new(pool::Pool::with_initial(runners, initial));
166203
self.plugins.push(Loaded { name, runner });
167204
Ok(())
168205
}

0 commit comments

Comments
 (0)