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: 24 additions & 7 deletions docs/plugins.md
Original file line number Diff line number Diff line change
Expand Up @@ -133,11 +133,28 @@ diffr builds these records once per call, from the file's manifest entry and
its trees as the plugins before left them, and hands the same records to a
native plugin and to a component. Nothing else reaches a plugin.

diffr makes one instance of each plugin per run and calls it for every file.
A native plugin's instance is called from several files' workers in parallel,
so it must be `Send + Sync`; interior mutability is the plugin's own
business. A component's instance lives in one wasmtime store, which runs one
call at a time: diffr calls a component plugin one file at a time.
By default diffr makes one instance of each plugin and serializes calls to it.
A plugin author can declare `parallel = true` at the top of `plugin.toml` to
permit independent deferred-work instances. This promises that calls do not depend on shared
mutable state, call order, or unique external side effects. Initial presentation and enrichment may use different instances. Plugin ordering within a file
is still sequential.

Opted-in plugins expose a host-owned `instances` setting (default 4, range 1–64):

```toml
[plugins.bundled.summarize]
instances = 4
```

The pool is shared across every file and phase of that plugin, so the bound is
per comparison, not per file. A separate initial-presentation instance keeps
queries, classification, and mutation responsive while the deferred pool is busy.
With `instances = 1`, all phases share the same instance for serial debugging. The host compiles a WASM component once and creates
independent stores. Each store still executes only one call at a time. Leases
return on success and errors. Cancellation skips queued enrichment; active calls
retain their configured request timeouts. Native plugins use the same pool.
The summarizer batches selected folds into one request per file, so this
parallelizes files; it does not split a file's batch into per-fold requests.

### Fresh ids

Expand Down Expand Up @@ -319,8 +336,8 @@ next run.

`plugins/summarize/src/http.rs` is a complete outgoing-request example for
plugin authors. The host provides TLS; components do not need to embed a
TLS implementation. WASM instances process files serially; the legacy
`max_concurrency` option is accepted but has no effect. Timeouts apply to
TLS implementation. Each WASM instance processes calls serially; `instances` controls the host pool.
The legacy `max_concurrency` option is still accepted but has no effect. Timeouts apply to
connection, first-byte and between-byte waits.
`cargo test --features wasm-plugin-tests --bin diffr external_component_summarizes_over_http`
checks the external component against a local model endpoint, including retries.
Expand Down
1 change: 1 addition & 0 deletions plugins/summarize/plugin.toml
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
name = "summarize"
parallel = true
title = "Summaries"
description = "Pseudocode summaries for large new function bodies and tests."

Expand Down
43 changes: 42 additions & 1 deletion src/plugin/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ pub(crate) const PATH: &str = "path";

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

/// A plugin folder's description, and its component when it has one.
pub(crate) const MANIFEST_FILE: &str = "plugin.toml";
Expand All @@ -35,6 +35,10 @@ pub(crate) struct Manifest {
/// The plugin's entry name in `[plugins]`, and the prefix of every tag
/// its queries set: `<name>:<tag>`.
pub(crate) name: String,
/// Opt in only when independent instances can process different files.
/// No cross-call state, ordering, or unique external side effects may be required.
#[serde(default)]
pub(crate) parallel: bool,
/// The human name settings screens group the plugin's settings under.
pub(crate) title: String,
#[serde(default)]
Expand Down Expand Up @@ -176,6 +180,13 @@ impl Manifest {
"x-group": group,
}),
);
if self.parallel {
properties.insert("instances".into(), json!({
"type": "integer", "minimum": 1, "maximum": 64, "default": 4,
"title": "Parallel instances", "description": "Maximum simultaneous deferred plugin calls across all files.",
"x-group": group,
}));
}
for (key, option) in &self.options {
let mut option = option.clone();
if let Some(option) = option.as_object_mut() {
Expand Down Expand Up @@ -301,6 +312,9 @@ pub(crate) struct Entry {
/// plugin's default.
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) enabled: Option<bool>,
/// Host-owned deferred-work bound, shared by all files of this plugin.
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) instances: Option<usize>,
/// The plugin's folder on disk, as written: relative to the
/// configuration file's directory, or absolute.
#[serde(skip_serializing_if = "Option::is_none")]
Expand Down Expand Up @@ -422,6 +436,12 @@ impl PluginsConfig {
};
let manifest = &folder.manifest;
entry.enabled.get_or_insert(manifest.enabled_by_default());
let instances = entry
.instances
.unwrap_or(if manifest.parallel { 4 } else { 1 });
if !(1..=64).contains(&instances) || (!manifest.parallel && instances != 1) {
return Err(ConfigError(format!("plugins.{reference}.instances: expected 1..=64 for a parallel plugin, or 1 for a serial plugin")));
}
manifest
.validate(&entry.options)
.map_err(|error| ConfigError(format!("plugins.{reference}: {error}")))?;
Expand Down Expand Up @@ -742,3 +762,24 @@ mod tests {
.starts_with("plugins.bundled.summarize: provider: "));
}
}

#[cfg(test)]
mod concurrency_tests {
use crate::config::Config;

#[test]
fn concurrency_requires_plugin_opt_in_and_a_bounded_positive_count() {
for count in [0, 65] {
assert!(Config::from_toml(&format!(
"[plugins.bundled.summarize]\ninstances = {count}\n"
))
.is_err());
}
assert!(Config::from_toml("[plugins.bundled.context]\ninstances = 2\n").is_err());
let config = Config::from_toml("[plugins.bundled.summarize]\ninstances = 3\n").unwrap();
assert_eq!(
config.plugins.entries["bundled.summarize"].instances,
Some(3)
);
}
}
55 changes: 46 additions & 9 deletions src/plugin/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
//! `plugin.toml` (name, title, options schema) and, when it has
//! one, its component. [`Pipeline::from_config`] then takes the component
//! ([`wasm`]) or the bundled native implementation registered under
//! the plugin's name ([`native`]), and makes the plugin's one instance for
//! the plugin's name ([`native`]), and makes the plugin's instance pool for
//! the run from its options. From there a [`Runner`] is a [`Runner`]: each
//! call gets the contract's records, built once per call from the file's
//! manifest entry and its current trees, and the host functions of [`host`].
Expand All @@ -32,6 +32,7 @@ pub(crate) mod builtin;
pub(crate) mod config;
pub(crate) mod host;
pub(crate) mod native;
mod pool;
pub(crate) mod queries;
pub(crate) mod wasm;

Expand Down Expand Up @@ -137,19 +138,31 @@ impl Pipeline {
Some(engine) => engine,
None => engine.insert(wasm::engine()?),
};
pipeline.push(name, options, &|host, options| {
wasm::WasmPlugin::load(engine, &path)?.create(host, options)
})?
let plugin = wasm::WasmPlugin::load(engine, &path)
.with_context(|| format!("plugins.{name}"))?;
pipeline.push_instances(
name,
options,
entry
.instances
.unwrap_or(if folder.manifest.parallel { 4 } else { 1 }),
&|host, options| plugin.create(host, options),
)?
}
None => {
let create = native::lookup(&folder.manifest.name)?.ok_or_else(|| {
anyhow!(
"plugins.{name}: no native implementation is registered for this bundled plugin"
)
})?;
pipeline.push(name, options, &|host, options| {
native::registered(create, host, options)
})?
pipeline.push_instances(
name,
options,
entry
.instances
.unwrap_or(if folder.manifest.parallel { 4 } else { 1 }),
&|host, options| native::registered(create, host, options),
)?
}
}
}
Expand All @@ -158,11 +171,35 @@ impl Pipeline {

/// Make the plugin `name` with `create` from `options` and add it to the
/// end of the pipeline.
#[cfg(test)]
fn push(&mut self, name: &str, options: Value, create: Create<'_>) -> anyhow::Result<()> {
self.push_instances(name, options, 1, create)
}

fn push_instances(
&mut self,
name: &str,
options: Value,
instances: usize,
create: Create<'_>,
) -> anyhow::Result<()> {
let reference = name;
let name: Arc<str> = name.split_once('.').map_or(name, |(_, name)| name).into();
let runner = create(self.host(&name), &options.to_string())
.with_context(|| format!("plugins.{reference}"))?;
let runners = (0..instances)
.map(|_| {
create(self.host(&name), &options.to_string())
.with_context(|| format!("plugins.{reference}"))
})
.collect::<anyhow::Result<Vec<_>>>()?;
let initial = if instances > 1 {
Some(
create(self.host(&name), &options.to_string())
.with_context(|| format!("plugins.{reference}"))?,
)
} else {
None
};
let runner = Box::new(pool::Pool::with_initial(runners, initial));
self.plugins.push(Loaded { name, runner });
Ok(())
}
Expand Down
Loading
Loading