diff --git a/docs/plugins.md b/docs/plugins.md index 468166f95..76c8651f9 100644 --- a/docs/plugins.md +++ b/docs/plugins.md @@ -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 @@ -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. diff --git a/plugins/summarize/plugin.toml b/plugins/summarize/plugin.toml index 50a073898..b493f37f0 100644 --- a/plugins/summarize/plugin.toml +++ b/plugins/summarize/plugin.toml @@ -1,4 +1,5 @@ name = "summarize" +parallel = true title = "Summaries" description = "Pseudocode summaries for large new function bodies and tests." diff --git a/src/plugin/config.rs b/src/plugin/config.rs index ea2917486..8293e3e05 100644 --- a/src/plugin/config.rs +++ b/src/plugin/config.rs @@ -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"; @@ -35,6 +35,10 @@ pub(crate) struct Manifest { /// The plugin's entry name in `[plugins]`, and the prefix of every tag /// its queries set: `:`. 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)] @@ -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() { @@ -301,6 +312,9 @@ pub(crate) struct Entry { /// plugin's default. #[serde(skip_serializing_if = "Option::is_none")] pub(crate) enabled: Option, + /// Host-owned deferred-work bound, shared by all files of this plugin. + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) instances: Option, /// The plugin's folder on disk, as written: relative to the /// configuration file's directory, or absolute. #[serde(skip_serializing_if = "Option::is_none")] @@ -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}")))?; @@ -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) + ); + } +} diff --git a/src/plugin/mod.rs b/src/plugin/mod.rs index 6a48f908a..b68cde4fd 100644 --- a/src/plugin/mod.rs +++ b/src/plugin/mod.rs @@ -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`]. @@ -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; @@ -137,9 +138,16 @@ 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(|| { @@ -147,9 +155,14 @@ impl Pipeline { "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), + )? } } } @@ -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 = 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::>>()?; + 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(()) } diff --git a/src/plugin/pool.rs b/src/plugin/pool.rs new file mode 100644 index 000000000..619e3cf27 --- /dev/null +++ b/src/plugin/pool.rs @@ -0,0 +1,183 @@ +//! A per-plugin deferred-work bound shared by every file. Independent +//! instances are opt-in; a serial plugin keeps one instance for its whole run. +use super::{host::Host, Runner}; +use diffr_plugin_sdk::types; +use std::sync::{Condvar, Mutex}; + +pub(super) struct Pool { + initial: Option>>, + idle: Mutex>>, + available: Condvar, +} + +impl Pool { + #[cfg(test)] + fn new(instances: Vec>) -> Self { + Self::with_initial(instances, None) + } + + pub(super) fn with_initial( + instances: Vec>, + initial: Option>, + ) -> Self { + assert!( + !instances.is_empty(), + "a plugin needs at least one instance" + ); + Self { + initial: initial.map(Mutex::new), + idle: Mutex::new(instances), + available: Condvar::new(), + } + } + + // Fast presentation never waits for a pool full of network requests. + // Serial plugins share their sole instance across both phases. + fn initial(&self, call: impl FnOnce(&dyn Runner) -> anyhow::Result) -> anyhow::Result { + match &self.initial { + Some(instance) => call(instance.lock().expect("initial plugin instance").as_ref()), + None => self.with(call), + } + } + + fn with(&self, call: impl FnOnce(&dyn Runner) -> anyhow::Result) -> anyhow::Result { + let mut idle = self.idle.lock().expect("instance pool"); + while idle.is_empty() { + idle = self.available.wait(idle).expect("instance pool"); + } + let lease = Lease { + pool: self, + instance: idle.pop(), + }; + drop(idle); + call(lease.instance.as_deref().expect("leased instance")) + } +} + +/// Return the slot on both errors and unwinding; never hold the queue lock +/// while executing a guest or waiting on its HTTP response. +struct Lease<'a> { + pool: &'a Pool, + instance: Option>, +} +impl Drop for Lease<'_> { + fn drop(&mut self) { + self.pool + .idle + .lock() + .expect("instance pool") + .push(self.instance.take().expect("leased instance")); + self.pool.available.notify_one(); + } +} + +impl Runner for Pool { + fn queries(&self, host: Host) -> anyhow::Result> { + self.initial(|runner| runner.queries(host)) + } + fn classify(&self, host: Host, file: &types::FileEntry) -> anyhow::Result> { + self.initial(|runner| runner.classify(host, file)) + } + fn mutate( + &self, + host: Host, + file: &types::FileEntry, + sides: &types::SourceSides, + ) -> anyhow::Result> { + self.initial(|runner| runner.mutate(host, file, sides)) + } + fn enrich( + &self, + host: Host, + file: &types::FileEntry, + sides: &types::SourceSides, + ) -> anyhow::Result> { + self.with(|runner| runner.enrich(host, file, sides)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::{mpsc, Arc}; + use std::time::Duration; + + struct Empty; + impl Runner for Empty { + 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![]) + } + } + + #[test] + fn independent_calls_overlap_but_never_exceed_the_bound() { + let pool = Pool::new(vec![Box::new(Empty), Box::new(Empty)]); + let active = AtomicUsize::new(0); + let peak = AtomicUsize::new(0); + let gate = Arc::new((Mutex::new(false), Condvar::new())); + let (entered, arrivals) = mpsc::channel(); + std::thread::scope(|scope| { + for _ in 0..6 { + let (pool, active, peak, entered, gate) = + (&pool, &active, &peak, entered.clone(), gate.clone()); + scope.spawn(move || { + pool.with(|_| { + let count = active.fetch_add(1, Ordering::SeqCst) + 1; + peak.fetch_max(count, Ordering::SeqCst); + entered.send(()).unwrap(); + let (lock, ready) = &*gate; + let mut open = lock.lock().unwrap(); + while !*open { + open = ready.wait(open).unwrap(); + } + active.fetch_sub(1, Ordering::SeqCst); + Ok(()) + }) + .unwrap() + }); + } + let first = arrivals.recv_timeout(Duration::from_secs(5)); + let second = arrivals.recv_timeout(Duration::from_secs(5)); + // Always release the workers, even if an assertion will fail. + *gate.0.lock().unwrap() = true; + gate.1.notify_all(); + assert!( + first.is_ok() && second.is_ok(), + "two requests should enter before either returns" + ); + }); + assert_eq!(peak.load(Ordering::SeqCst), 2); + } + + #[test] + fn initial_presentation_does_not_wait_for_busy_enrichment() { + let pool = Pool::with_initial(vec![Box::new(Empty)], Some(Box::new(Empty))); + pool.with(|_| { + // Holding the only enrichment lease must not block initial calls. + assert_eq!(pool.initial(|_| Ok(7))?, 7); + Ok(()) + }) + .unwrap(); + } + + #[test] + fn errors_return_the_slot_to_a_serial_pool() { + let pool = Pool::new(vec![Box::new(Empty)]); + assert!(pool + .with::<()>(|_| anyhow::bail!("request failed")) + .is_err()); + assert_eq!(pool.with(|_| Ok(42)).unwrap(), 42); + } +} diff --git a/src/plugin/wasm.rs b/src/plugin/wasm.rs index b2d3cc545..33266d50f 100644 --- a/src/plugin/wasm.rs +++ b/src/plugin/wasm.rs @@ -4,10 +4,10 @@ //! Each component is compiled once, when the pipeline is built (with //! wasmtime's disk cache, so an unchanged component is not recompiled on the //! next run), linked against WASI and the `host` imports, and instantiated -//! once in its own [`Store`]. The plugin's `new` makes the one `plugin` +//! in independent [`Store`]s when the plugin opts into pooling. The plugin's `new` makes the one `plugin` //! resource of the run in that instance, and every `classify` and `mutate` //! calls that resource. A store runs one call at a time, so it sits behind a -//! mutex: the rayon workers call one component plugin one file at a time. +//! mutex: each pooled instance accepts one call at a time. //! The records a call gets are the ones a native plugin gets, lowered field //! for field into the generated bindings. //!