1use 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 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 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 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 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}