Skip to main content

blut_graph_core/
plugin.rs

1// SPDX-License-Identifier: AGPL-3.0-or-later
2
3use alloc::string::String;
4use alloc::vec::Vec;
5use core::fmt;
6
7use serde::{Deserialize, Serialize};
8
9use crate::model::Capability;
10
11pub const PLUGIN_PROTOCOL_VERSION: u32 = 2;
12const CONTROL_MAGIC: &[u8; 4] = b"BPC2";
13
14#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
15#[serde(rename_all = "kebab-case")]
16pub enum ExecutableDigestAlgorithm {
17    /// BLAKE3 derive-key mode with context `blut.plugin-executable.v1`.
18    Blake3DomainV1,
19}
20
21#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
22#[serde(rename_all = "kebab-case")]
23pub enum SignatureAlgorithm {
24    Ed25519,
25}
26
27pub fn executable_digest(
28    algorithm: ExecutableDigestAlgorithm,
29    executable_bytes: &[u8],
30) -> [u8; 32] {
31    match algorithm {
32        ExecutableDigestAlgorithm::Blake3DomainV1 => {
33            let mut hasher = blake3::Hasher::new_derive_key("blut.plugin-executable.v1");
34            hasher.update(executable_bytes);
35            *hasher.finalize().as_bytes()
36        }
37    }
38}
39
40#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
41#[serde(rename_all = "kebab-case")]
42pub enum TeardownPolicy {
43    GracefulThenKill,
44    ImmediateKill,
45}
46
47#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
48pub struct ProcessContract {
49    pub startup_deadline_millis: u64,
50    pub request_deadline_millis: u64,
51    pub heartbeat_interval_millis: u64,
52    pub heartbeat_grace_millis: u64,
53    pub teardown_deadline_millis: u64,
54    pub kill_grace_millis: u64,
55    pub max_inflight: u32,
56    pub max_frame_bytes: u32,
57    pub teardown: TeardownPolicy,
58}
59
60impl ProcessContract {
61    pub fn is_valid(&self) -> bool {
62        self.startup_deadline_millis > 0
63            && self.request_deadline_millis > 0
64            && self.heartbeat_interval_millis > 0
65            && self.heartbeat_grace_millis >= self.heartbeat_interval_millis
66            && self.teardown_deadline_millis > 0
67            && self.kill_grace_millis <= self.teardown_deadline_millis
68            && self.max_inflight > 0
69            && self.max_frame_bytes >= 64
70    }
71}
72
73#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
74pub struct PluginManifest {
75    pub protocol_version: u32,
76    pub plugin_id: String,
77    pub executable_digest_algorithm: ExecutableDigestAlgorithm,
78    pub executable_digest: [u8; 32],
79    pub capabilities: Vec<Capability>,
80    pub process: ProcessContract,
81    pub signature_algorithm: SignatureAlgorithm,
82    /// Stable verifier key identifier; key resolution is a supervisor policy.
83    pub signing_key_id: String,
84    pub signature: Vec<u8>,
85}
86
87#[derive(Serialize)]
88struct UnsignedPluginManifest<'a> {
89    protocol_version: u32,
90    plugin_id: &'a str,
91    executable_digest_algorithm: ExecutableDigestAlgorithm,
92    executable_digest: [u8; 32],
93    capabilities: &'a [Capability],
94    process: &'a ProcessContract,
95    signature_algorithm: SignatureAlgorithm,
96    signing_key_id: &'a str,
97}
98
99impl PluginManifest {
100    pub fn normalize(&mut self) -> Result<(), PluginError> {
101        self.capabilities.sort_unstable();
102        self.capabilities.dedup();
103        if self.protocol_version != PLUGIN_PROTOCOL_VERSION {
104            return Err(PluginError::ProtocolVersion(self.protocol_version));
105        }
106        if self.plugin_id.is_empty()
107            || self.signing_key_id.is_empty()
108            || self.signature.len() != 64
109            || self.capabilities.is_empty()
110            || self.capabilities.iter().any(|item| item.0.is_empty())
111            || !self.process.is_valid()
112        {
113            return Err(PluginError::InvalidManifest);
114        }
115        Ok(())
116    }
117
118    /// Canonical signature preimage. The signature itself is excluded to avoid
119    /// circular signing; capabilities are normalized before serialization.
120    pub fn unsigned_signing_bytes(&self) -> Result<Vec<u8>, PluginError> {
121        let mut normalized = self.clone();
122        normalized.capabilities.sort_unstable();
123        normalized.capabilities.dedup();
124        if normalized.protocol_version != PLUGIN_PROTOCOL_VERSION
125            || normalized.plugin_id.is_empty()
126            || normalized.signing_key_id.is_empty()
127            || normalized.capabilities.is_empty()
128            || normalized.capabilities.iter().any(|item| item.0.is_empty())
129            || !normalized.process.is_valid()
130        {
131            return Err(PluginError::InvalidManifest);
132        }
133        postcard::to_allocvec(&UnsignedPluginManifest {
134            protocol_version: normalized.protocol_version,
135            plugin_id: &normalized.plugin_id,
136            executable_digest_algorithm: normalized.executable_digest_algorithm,
137            executable_digest: normalized.executable_digest,
138            capabilities: &normalized.capabilities,
139            process: &normalized.process,
140            signature_algorithm: normalized.signature_algorithm,
141            signing_key_id: &normalized.signing_key_id,
142        })
143        .map_err(|_| PluginError::MalformedFrame)
144    }
145
146    pub fn signing_digest(&self) -> Result<[u8; 32], PluginError> {
147        let bytes = self.unsigned_signing_bytes()?;
148        let mut hasher = blake3::Hasher::new_derive_key("blut.plugin-manifest.v1");
149        hasher.update(&bytes);
150        Ok(*hasher.finalize().as_bytes())
151    }
152}
153
154#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
155pub struct PluginRequest {
156    pub request_id: u64,
157    pub invocation_id: [u8; 32],
158    pub capability: Capability,
159    pub payload_content_id: [u8; 32],
160    pub shared_memory_lease: Option<String>,
161    /// Absolute supervisor-clock deadline. Zero is never interpreted as infinite.
162    pub deadline_millis: u64,
163}
164
165#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
166pub struct PluginFailure {
167    pub domain: String,
168    pub code: String,
169    pub message: String,
170    pub retryable: bool,
171}
172
173#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
174pub struct PluginResponse {
175    pub request_id: u64,
176    pub output_content_id: Option<[u8; 32]>,
177    pub receipt: Vec<u8>,
178    pub failure: Option<PluginFailure>,
179}
180
181#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
182#[serde(rename_all = "kebab-case")]
183pub enum PluginLifecycle {
184    Spawned,
185    Handshaking,
186    Ready,
187    Draining,
188    Terminated,
189}
190
191impl PluginLifecycle {
192    pub const fn permits(self, next: Self) -> bool {
193        matches!(
194            (self, next),
195            (Self::Spawned, Self::Handshaking | Self::Terminated)
196                | (
197                    Self::Handshaking,
198                    Self::Ready | Self::Draining | Self::Terminated
199                )
200                | (Self::Ready, Self::Draining | Self::Terminated)
201                | (Self::Draining, Self::Terminated)
202        )
203    }
204}
205
206#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
207#[serde(rename_all = "kebab-case")]
208pub enum PluginControlFrame {
209    Hello {
210        manifest: PluginManifest,
211        supervisor_nonce: [u8; 32],
212        deadline_millis: u64,
213    },
214    Ready {
215        supervisor_nonce: [u8; 32],
216        process_id: u64,
217    },
218    Invoke(PluginRequest),
219    Complete(PluginResponse),
220    Heartbeat {
221        sequence: u64,
222        monotonic_millis: u64,
223    },
224    Cancel {
225        request_id: u64,
226        deadline_millis: u64,
227    },
228    Shutdown {
229        reason: String,
230        deadline_millis: u64,
231    },
232    Ack {
233        request_id: Option<u64>,
234        lifecycle: PluginLifecycle,
235    },
236}
237
238#[derive(Clone, Copy, Debug, PartialEq, Eq)]
239pub struct PluginControlLimits {
240    pub max_frame_bytes: usize,
241    pub max_receipt_bytes: usize,
242    pub max_signature_bytes: usize,
243}
244
245impl Default for PluginControlLimits {
246    fn default() -> Self {
247        Self {
248            max_frame_bytes: 1024 * 1024,
249            max_receipt_bytes: 256 * 1024,
250            max_signature_bytes: 16 * 1024,
251        }
252    }
253}
254
255impl PluginControlLimits {
256    pub fn for_process(process: &ProcessContract) -> Self {
257        let defaults = Self::default();
258        let max_frame_bytes = process.max_frame_bytes as usize;
259        Self {
260            max_frame_bytes,
261            max_receipt_bytes: defaults.max_receipt_bytes.min(max_frame_bytes),
262            max_signature_bytes: defaults.max_signature_bytes.min(max_frame_bytes),
263        }
264    }
265}
266
267impl PluginControlFrame {
268    pub fn to_control_bytes(&self) -> Result<Vec<u8>, PluginError> {
269        self.to_control_bytes_with_limits(PluginControlLimits::default())
270    }
271
272    pub fn to_control_bytes_with_limits(
273        &self,
274        limits: PluginControlLimits,
275    ) -> Result<Vec<u8>, PluginError> {
276        validate_frame(self, limits)?;
277        let mut bytes = Vec::from(CONTROL_MAGIC.as_slice());
278        bytes.extend(postcard::to_allocvec(self).map_err(|_| PluginError::MalformedFrame)?);
279        if bytes.len() > limits.max_frame_bytes {
280            return Err(PluginError::FrameTooLarge);
281        }
282        Ok(bytes)
283    }
284
285    pub fn from_control_bytes(
286        bytes: &[u8],
287        limits: PluginControlLimits,
288    ) -> Result<Self, PluginError> {
289        if bytes.len() > limits.max_frame_bytes {
290            return Err(PluginError::FrameTooLarge);
291        }
292        let body = bytes
293            .strip_prefix(CONTROL_MAGIC)
294            .ok_or(PluginError::BadControlMagic)?;
295        let (frame, remainder): (Self, &[u8]) =
296            postcard::take_from_bytes(body).map_err(|_| PluginError::MalformedFrame)?;
297        if !remainder.is_empty() {
298            return Err(PluginError::MalformedFrame);
299        }
300        validate_frame(&frame, limits)?;
301        Ok(frame)
302    }
303}
304
305fn validate_frame(
306    frame: &PluginControlFrame,
307    limits: PluginControlLimits,
308) -> Result<(), PluginError> {
309    match frame {
310        PluginControlFrame::Hello {
311            manifest,
312            deadline_millis,
313            ..
314        } => {
315            let mut normalized = manifest.clone();
316            normalized.normalize()?;
317            if normalized != *manifest || manifest.signature.len() > limits.max_signature_bytes {
318                return Err(PluginError::InvalidManifest);
319            }
320            if manifest.process.max_frame_bytes as usize > limits.max_frame_bytes {
321                return Err(PluginError::FrameTooLarge);
322            }
323            if *deadline_millis == 0 {
324                return Err(PluginError::InvalidDeadline);
325            }
326        }
327        PluginControlFrame::Ready { process_id, .. } if *process_id == 0 => {
328            return Err(PluginError::InvalidLifecycle);
329        }
330        PluginControlFrame::Invoke(request) => validate_request(request)?,
331        PluginControlFrame::Complete(response) => {
332            validate_response(response, limits.max_receipt_bytes)?;
333        }
334        PluginControlFrame::Heartbeat {
335            monotonic_millis, ..
336        } if *monotonic_millis == 0 => return Err(PluginError::InvalidDeadline),
337        PluginControlFrame::Cancel {
338            deadline_millis, ..
339        }
340        | PluginControlFrame::Shutdown {
341            deadline_millis, ..
342        } if *deadline_millis == 0 => return Err(PluginError::InvalidDeadline),
343        PluginControlFrame::Shutdown { reason, .. } if reason.is_empty() => {
344            return Err(PluginError::MalformedFrame);
345        }
346        PluginControlFrame::Ack {
347            request_id: Some(_),
348            lifecycle: PluginLifecycle::Terminated,
349        } => return Err(PluginError::InvalidLifecycle),
350        _ => {}
351    }
352    Ok(())
353}
354
355fn validate_request(request: &PluginRequest) -> Result<(), PluginError> {
356    if request.capability.0.is_empty() || request.deadline_millis == 0 {
357        return Err(PluginError::InvalidDeadline);
358    }
359    Ok(())
360}
361
362fn validate_response(
363    response: &PluginResponse,
364    max_receipt_bytes: usize,
365) -> Result<(), PluginError> {
366    if response.receipt.len() > max_receipt_bytes
367        || response.failure.as_ref().is_some_and(|failure| {
368            failure.domain.is_empty()
369                || failure.code.is_empty()
370                || response.output_content_id.is_some()
371        })
372    {
373        return Err(PluginError::MalformedFrame);
374    }
375    Ok(())
376}
377
378#[derive(Clone, Debug, PartialEq, Eq)]
379pub enum PluginError {
380    ProtocolVersion(u32),
381    InvalidManifest,
382    ManifestUntrusted,
383    CapabilityUndeclared(String),
384    InvalidDeadline,
385    InvalidLifecycle,
386    BadControlMagic,
387    FrameTooLarge,
388    MalformedFrame,
389    Transport(String),
390    MismatchedResponse,
391}
392
393impl fmt::Display for PluginError {
394    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
395        write!(f, "{self:?}")
396    }
397}
398
399#[cfg(feature = "std")]
400impl std::error::Error for PluginError {}
401
402pub trait PluginHost {
403    fn verify_manifest(&self, manifest: &PluginManifest) -> Result<(), PluginError>;
404
405    fn exchange(
406        &mut self,
407        manifest: &PluginManifest,
408        request: &PluginRequest,
409    ) -> Result<PluginResponse, PluginError>;
410
411    fn call(
412        &mut self,
413        manifest: &PluginManifest,
414        request: &PluginRequest,
415    ) -> Result<PluginResponse, PluginError> {
416        let mut normalized = manifest.clone();
417        normalized.normalize()?;
418        if normalized != *manifest {
419            return Err(PluginError::InvalidManifest);
420        }
421        validate_request(request)?;
422        let process_limits = PluginControlLimits::for_process(&manifest.process);
423        PluginControlFrame::Invoke(request.clone()).to_control_bytes_with_limits(process_limits)?;
424        self.verify_manifest(manifest)?;
425        if !manifest.capabilities.contains(&request.capability) {
426            return Err(PluginError::CapabilityUndeclared(
427                request.capability.0.clone(),
428            ));
429        }
430        let response = self.exchange(manifest, request)?;
431        PluginControlFrame::Complete(response.clone())
432            .to_control_bytes_with_limits(process_limits)?;
433        if response.request_id != request.request_id {
434            return Err(PluginError::MismatchedResponse);
435        }
436        Ok(response)
437    }
438}
439
440#[cfg(test)]
441mod tests {
442    use alloc::string::ToString;
443    use alloc::vec;
444
445    use super::*;
446
447    struct Host;
448
449    impl PluginHost for Host {
450        fn verify_manifest(&self, _manifest: &PluginManifest) -> Result<(), PluginError> {
451            Ok(())
452        }
453
454        fn exchange(
455            &mut self,
456            _manifest: &PluginManifest,
457            request: &PluginRequest,
458        ) -> Result<PluginResponse, PluginError> {
459            Ok(PluginResponse {
460                request_id: request.request_id,
461                output_content_id: None,
462                receipt: vec![],
463                failure: None,
464            })
465        }
466    }
467
468    struct OversizedHost;
469
470    impl PluginHost for OversizedHost {
471        fn verify_manifest(&self, _manifest: &PluginManifest) -> Result<(), PluginError> {
472            Ok(())
473        }
474
475        fn exchange(
476            &mut self,
477            manifest: &PluginManifest,
478            request: &PluginRequest,
479        ) -> Result<PluginResponse, PluginError> {
480            Ok(PluginResponse {
481                request_id: request.request_id,
482                output_content_id: None,
483                receipt: vec![0; manifest.process.max_frame_bytes as usize],
484                failure: None,
485            })
486        }
487    }
488
489    fn manifest() -> PluginManifest {
490        PluginManifest {
491            protocol_version: PLUGIN_PROTOCOL_VERSION,
492            plugin_id: "adapter".to_string(),
493            executable_digest_algorithm: ExecutableDigestAlgorithm::Blake3DomainV1,
494            executable_digest: executable_digest(
495                ExecutableDigestAlgorithm::Blake3DomainV1,
496                b"executable",
497            ),
498            capabilities: vec![Capability("export".to_string())],
499            process: ProcessContract {
500                startup_deadline_millis: 1_000,
501                request_deadline_millis: 5_000,
502                heartbeat_interval_millis: 500,
503                heartbeat_grace_millis: 1_500,
504                teardown_deadline_millis: 1_000,
505                kill_grace_millis: 100,
506                max_inflight: 4,
507                max_frame_bytes: 65_536,
508                teardown: TeardownPolicy::GracefulThenKill,
509            },
510            signature_algorithm: SignatureAlgorithm::Ed25519,
511            signing_key_id: "test-key".to_string(),
512            signature: vec![1; 64],
513        }
514    }
515
516    fn request(capability: &str) -> PluginRequest {
517        PluginRequest {
518            request_id: 1,
519            invocation_id: [7; 32],
520            capability: Capability(capability.to_string()),
521            payload_content_id: [0; 32],
522            shared_memory_lease: None,
523            deadline_millis: 123,
524        }
525    }
526
527    #[test]
528    fn undeclared_capability_never_reaches_transport() {
529        assert!(matches!(
530            Host.call(&manifest(), &request("import")),
531            Err(PluginError::CapabilityUndeclared(_))
532        ));
533    }
534
535    #[test]
536    fn executable_digest_is_domain_separated_and_stable() {
537        assert_eq!(
538            executable_digest(ExecutableDigestAlgorithm::Blake3DomainV1, b"abc"),
539            [
540                0xed, 0x25, 0x1c, 0x0f, 0x1e, 0xb9, 0xc7, 0x7c, 0x6b, 0x0e, 0x39, 0x1d, 0x4c, 0x62,
541                0xba, 0x66, 0x29, 0xc1, 0x91, 0x56, 0x9c, 0x3e, 0x5f, 0xa6, 0xce, 0xee, 0x22, 0x59,
542                0x6a, 0x82, 0x18, 0xa4,
543            ]
544        );
545    }
546
547    #[test]
548    fn manifest_signature_preimage_and_bpc2_wire_are_literal_goldens() {
549        assert_eq!(
550            manifest().signing_digest().unwrap(),
551            [
552                35, 152, 182, 151, 168, 246, 64, 146, 153, 175, 245, 69, 198, 185, 247, 4, 17, 128,
553                1, 174, 112, 76, 208, 232, 80, 232, 42, 245, 207, 34, 39, 128,
554            ]
555        );
556        let mut invoke_golden = vec![66, 80, 67, 50, 2, 1];
557        invoke_golden.extend([7; 32]);
558        invoke_golden.extend([6, b'e', b'x', b'p', b'o', b'r', b't']);
559        invoke_golden.extend([0; 32]);
560        invoke_golden.extend([0, 123]);
561        assert_eq!(
562            PluginControlFrame::Invoke(request("export"))
563                .to_control_bytes()
564                .unwrap(),
565            invoke_golden
566        );
567    }
568
569    #[test]
570    fn ed25519_manifest_rejects_noncanonical_signature_lengths() {
571        for length in [0, 63, 65] {
572            let mut invalid = manifest();
573            invalid.signature = vec![0; length];
574            assert_eq!(invalid.normalize(), Err(PluginError::InvalidManifest));
575        }
576    }
577
578    #[test]
579    fn control_wire_rejects_trailing_bytes_and_zero_deadlines() {
580        let frame = PluginControlFrame::Invoke(request("export"));
581        let mut bytes = frame.to_control_bytes().unwrap();
582        assert_eq!(
583            PluginControlFrame::from_control_bytes(&bytes, PluginControlLimits::default()),
584            Ok(frame)
585        );
586        bytes.push(0);
587        assert_eq!(
588            PluginControlFrame::from_control_bytes(&bytes, PluginControlLimits::default()),
589            Err(PluginError::MalformedFrame)
590        );
591
592        let mut invalid = request("export");
593        invalid.deadline_millis = 0;
594        assert_eq!(
595            PluginControlFrame::Invoke(invalid).to_control_bytes(),
596            Err(PluginError::InvalidDeadline)
597        );
598    }
599
600    #[test]
601    fn lifecycle_never_reenters_ready_after_draining_or_termination() {
602        assert!(PluginLifecycle::Spawned.permits(PluginLifecycle::Handshaking));
603        assert!(PluginLifecycle::Handshaking.permits(PluginLifecycle::Ready));
604        assert!(PluginLifecycle::Ready.permits(PluginLifecycle::Draining));
605        assert!(PluginLifecycle::Draining.permits(PluginLifecycle::Terminated));
606        assert!(!PluginLifecycle::Draining.permits(PluginLifecycle::Ready));
607        assert!(!PluginLifecycle::Terminated.permits(PluginLifecycle::Spawned));
608    }
609
610    #[test]
611    fn response_cannot_claim_an_output_and_failure_together() {
612        let response = PluginResponse {
613            request_id: 1,
614            output_content_id: Some([1; 32]),
615            receipt: vec![],
616            failure: Some(PluginFailure {
617                domain: "plugin.test".into(),
618                code: "failed".into(),
619                message: "failed".into(),
620                retryable: false,
621            }),
622        };
623        assert_eq!(
624            PluginControlFrame::Complete(response).to_control_bytes(),
625            Err(PluginError::MalformedFrame)
626        );
627    }
628
629    #[test]
630    fn control_encoder_enforces_the_same_frame_bound_as_the_decoder() {
631        let frame = PluginControlFrame::Shutdown {
632            reason: "x".repeat(PluginControlLimits::default().max_frame_bytes),
633            deadline_millis: 1,
634        };
635        assert_eq!(frame.to_control_bytes(), Err(PluginError::FrameTooLarge));
636        assert_eq!(
637            OversizedHost.call(&manifest(), &request("export")),
638            Err(PluginError::FrameTooLarge)
639        );
640    }
641}