1use std::{
2 error::Error,
3 fmt,
4 fmt::Write as _,
5 fs,
6 path::{Path, PathBuf},
7 str,
8 str::FromStr,
9 sync::{Arc, Mutex},
10};
11
12use crate::host_authoring::HostInput;
13use anyhow::{Context, bail};
14use lenso_app_plan::authoring::{
15 PluginInstanceId, PluginRootInstance, PluginRootResolutionError, PluginRootSnapshot,
16 ResolvedApp,
17};
18use sha2::{Digest, Sha256};
19
20use super::{
21 MAX_CONFIGURATION_BYTES, PLUGIN_ROOT, PluginRootAuthoringState, atomic_write,
22 inspect_plugin_root, load_host_catalog, lock_plugin_root, snapshot_plugin_root,
23 validate_existing_plugin_id, validate_instance_filename,
24};
25
26const PROPOSAL_SCHEMA: &str = "lenso.plugin-configuration-proposal.v1";
27const PUBLICATION_SCHEMA: &str = "lenso.plugin-configuration-publication.v1";
28const SOURCE_DIGEST_SCHEMA: &str = "lenso.plugin-configuration-source.v1";
29
30#[derive(Clone, Debug, Eq, PartialEq)]
32pub struct PluginConfigurationAuthoritySource {
33 kind: String,
34 reference: String,
35}
36
37impl PluginConfigurationAuthoritySource {
38 pub fn new(kind: impl Into<String>, reference: impl Into<String>) -> anyhow::Result<Self> {
40 let kind = kind.into();
41 let reference = reference.into();
42 if kind.is_empty()
43 || kind.len() > 64
44 || !kind.bytes().all(|byte| {
45 byte.is_ascii_lowercase()
46 || byte.is_ascii_digit()
47 || matches!(byte, b'_' | b'-' | b'.')
48 })
49 {
50 bail!("Plugin configuration authority kind is invalid");
51 }
52 if reference.is_empty() || reference.len() > 256 || reference.chars().any(char::is_control)
53 {
54 bail!("Plugin configuration authority reference is invalid");
55 }
56 Ok(Self { kind, reference })
57 }
58
59 pub fn kind(&self) -> &str {
60 &self.kind
61 }
62
63 pub fn reference(&self) -> &str {
64 &self.reference
65 }
66}
67
68pub trait PluginConfigurationAuthority: fmt::Debug + Send + Sync {
73 fn source(&self) -> PluginConfigurationAuthoritySource;
74
75 fn inspect(&self) -> anyhow::Result<PluginRootAuthoringState>;
76
77 fn propose(
78 &self,
79 expected_revision: &PluginRootRevision,
80 plugin_id: &str,
81 instance: &str,
82 bytes: &[u8],
83 ) -> anyhow::Result<PluginConfigurationProposal>;
84
85 fn publish(
86 &self,
87 proposal: &PluginConfigurationProposal,
88 ) -> anyhow::Result<PluginConfigurationPublication>;
89}
90
91#[derive(Clone, Debug)]
93pub struct LocalPluginRootAuthority {
94 root: PathBuf,
95 access: Arc<Mutex<()>>,
96}
97
98impl LocalPluginRootAuthority {
99 pub fn new(root: impl Into<PathBuf>) -> Self {
100 Self {
101 root: root.into(),
102 access: Arc::new(Mutex::new(())),
103 }
104 }
105
106 pub fn root(&self) -> &Path {
107 &self.root
108 }
109
110 pub(crate) fn lock(&self) -> anyhow::Result<std::sync::MutexGuard<'_, ()>> {
111 self.access
112 .lock()
113 .map_err(|_| anyhow::anyhow!("Plugin configuration authority lock is poisoned"))
114 }
115}
116
117impl PluginConfigurationAuthority for LocalPluginRootAuthority {
118 fn source(&self) -> PluginConfigurationAuthoritySource {
119 PluginConfigurationAuthoritySource {
120 kind: "local_plugin_root".to_owned(),
121 reference: "app".to_owned(),
122 }
123 }
124
125 fn inspect(&self) -> anyhow::Result<PluginRootAuthoringState> {
126 let _guard = self.lock()?;
127 inspect_plugin_root(&self.root)
128 }
129
130 fn propose(
131 &self,
132 expected_revision: &PluginRootRevision,
133 plugin_id: &str,
134 instance: &str,
135 bytes: &[u8],
136 ) -> anyhow::Result<PluginConfigurationProposal> {
137 let _guard = self.lock()?;
138 propose_instance_configuration(&self.root, expected_revision, plugin_id, instance, bytes)
139 }
140
141 fn publish(
142 &self,
143 proposal: &PluginConfigurationProposal,
144 ) -> anyhow::Result<PluginConfigurationPublication> {
145 let _guard = self.lock()?;
146 publish_instance_configuration(&self.root, proposal)
147 }
148}
149
150#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
155pub struct PluginRootRevision(String);
156
157impl PluginRootRevision {
158 pub fn as_str(&self) -> &str {
159 &self.0
160 }
161}
162
163impl fmt::Display for PluginRootRevision {
164 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
165 formatter.write_str(&self.0)
166 }
167}
168
169impl FromStr for PluginRootRevision {
170 type Err = PluginRootRevisionParseError;
171
172 fn from_str(value: &str) -> Result<Self, Self::Err> {
173 let Some(digest) = value.strip_prefix("sha256:") else {
174 return Err(PluginRootRevisionParseError);
175 };
176 if digest.len() != 64 || !digest.bytes().all(|byte| byte.is_ascii_hexdigit()) {
177 return Err(PluginRootRevisionParseError);
178 }
179 Ok(Self(format!("sha256:{}", digest.to_ascii_lowercase())))
180 }
181}
182
183#[derive(Clone, Copy, Debug, Eq, PartialEq)]
185pub struct PluginRootRevisionParseError;
186
187impl fmt::Display for PluginRootRevisionParseError {
188 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
189 formatter.write_str("Plugin Root revision must be `sha256:` followed by 64 hex digits")
190 }
191}
192
193impl Error for PluginRootRevisionParseError {}
194
195#[derive(Clone, Debug, Eq, PartialEq)]
197pub struct PluginRootRevisionConflict {
198 expected: PluginRootRevision,
199 current: PluginRootRevision,
200}
201
202impl PluginRootRevisionConflict {
203 pub const fn expected(&self) -> &PluginRootRevision {
204 &self.expected
205 }
206
207 pub const fn current(&self) -> &PluginRootRevision {
208 &self.current
209 }
210}
211
212impl fmt::Display for PluginRootRevisionConflict {
213 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
214 write!(
215 formatter,
216 "Plugin Root revision conflict: expected {}, current {}",
217 self.expected, self.current
218 )
219 }
220}
221
222impl Error for PluginRootRevisionConflict {}
223
224#[derive(Clone, Debug)]
226pub struct PluginConfigurationProposal {
227 schema: &'static str,
228 base_revision: PluginRootRevision,
229 base_source_digest: PluginConfigurationSourceDigest,
230 candidate_revision: PluginRootRevision,
231 digest: String,
232 status: PluginConfigurationProposalStatus,
233 application: PluginConfigurationApplication,
234 diagnostics: Vec<PluginConfigurationDiagnostic>,
235 plugin_id: String,
236 instance_key: String,
237 toml: Vec<u8>,
238}
239
240impl PluginConfigurationProposal {
241 pub const fn schema(&self) -> &str {
242 self.schema
243 }
244
245 pub const fn base_revision(&self) -> &PluginRootRevision {
246 &self.base_revision
247 }
248
249 pub const fn base_source_digest(&self) -> &PluginConfigurationSourceDigest {
250 &self.base_source_digest
251 }
252
253 pub const fn candidate_revision(&self) -> &PluginRootRevision {
254 &self.candidate_revision
255 }
256
257 pub fn digest(&self) -> &str {
258 &self.digest
259 }
260
261 pub const fn status(&self) -> PluginConfigurationProposalStatus {
262 self.status
263 }
264
265 pub const fn application(&self) -> PluginConfigurationApplication {
266 self.application
267 }
268
269 pub fn diagnostics(&self) -> &[PluginConfigurationDiagnostic] {
270 &self.diagnostics
271 }
272
273 pub fn plugin_id(&self) -> &str {
274 &self.plugin_id
275 }
276
277 pub fn instance_key(&self) -> &str {
278 &self.instance_key
279 }
280}
281
282#[derive(Clone, Copy, Debug, Eq, PartialEq)]
284pub enum PluginConfigurationProposalStatus {
285 Ready,
286 NeedsDecision,
287 Rejected,
288}
289
290#[derive(Clone, Copy, Debug, Eq, PartialEq)]
292pub enum PluginConfigurationApplication {
293 Noop,
294 AppGeneration,
295 Blocked,
296}
297
298#[derive(Clone, Debug, Eq, PartialEq)]
300pub struct PluginConfigurationDiagnostic {
301 code: &'static str,
302 detail: String,
303}
304
305impl PluginConfigurationDiagnostic {
306 pub const fn code(&self) -> &str {
307 self.code
308 }
309
310 pub fn detail(&self) -> &str {
311 &self.detail
312 }
313}
314
315#[derive(Clone, Debug)]
317pub struct PluginConfigurationPublication {
318 schema: &'static str,
319 base_revision: PluginRootRevision,
320 base_source_digest: PluginConfigurationSourceDigest,
321 revision: PluginRootRevision,
322 proposal_digest: String,
323 resolved: ResolvedApp,
324}
325
326impl PluginConfigurationPublication {
327 pub const fn schema(&self) -> &str {
328 self.schema
329 }
330
331 pub const fn base_revision(&self) -> &PluginRootRevision {
332 &self.base_revision
333 }
334
335 pub const fn base_source_digest(&self) -> &PluginConfigurationSourceDigest {
336 &self.base_source_digest
337 }
338
339 pub const fn revision(&self) -> &PluginRootRevision {
340 &self.revision
341 }
342
343 pub fn proposal_digest(&self) -> &str {
344 &self.proposal_digest
345 }
346
347 pub const fn resolved(&self) -> &ResolvedApp {
348 &self.resolved
349 }
350
351 pub fn into_resolved(self) -> ResolvedApp {
352 self.resolved
353 }
354}
355
356pub fn propose_instance_configuration(
358 root: &Path,
359 expected_revision: &PluginRootRevision,
360 plugin_id: &str,
361 instance: &str,
362 bytes: &[u8],
363) -> anyhow::Result<PluginConfigurationProposal> {
364 validate_existing_plugin_id(plugin_id)?;
365 validate_instance_filename(instance)?;
366 let _lock = lock_plugin_root(root)?;
367 let host = load_host_catalog(root)?;
368 let current = snapshot_plugin_root(root, &host)?;
369 let current_revision = revision_for_snapshot(¤t)?;
370 ensure_revision(expected_revision, ¤t_revision)?;
371 let base_source_digest = source_digest_for_instance(root, plugin_id, instance)?;
372 build_proposal(
373 &host,
374 ¤t,
375 current_revision,
376 base_source_digest,
377 plugin_id,
378 instance,
379 bytes,
380 )
381}
382
383pub fn publish_instance_configuration(
385 root: &Path,
386 proposal: &PluginConfigurationProposal,
387) -> anyhow::Result<PluginConfigurationPublication> {
388 let _lock = lock_plugin_root(root)?;
389 let host = load_host_catalog(root)?;
390 let current = snapshot_plugin_root(root, &host)?;
391 let current_revision = revision_for_snapshot(¤t)?;
392 ensure_revision(&proposal.base_revision, ¤t_revision)?;
393 let current_source_digest =
394 source_digest_for_instance(root, &proposal.plugin_id, &proposal.instance_key)?;
395 ensure_source_digest(&proposal.base_source_digest, ¤t_source_digest)?;
396
397 let verified = build_proposal(
398 &host,
399 ¤t,
400 current_revision,
401 current_source_digest,
402 &proposal.plugin_id,
403 &proposal.instance_key,
404 &proposal.toml,
405 )?;
406 if proposal.candidate_revision != verified.candidate_revision
407 || proposal.digest != verified.digest
408 {
409 bail!("Plugin configuration proposal no longer matches its reviewed candidate");
410 }
411 if verified.status != PluginConfigurationProposalStatus::Ready
412 || verified.application == PluginConfigurationApplication::Blocked
413 {
414 let detail = verified
415 .diagnostics
416 .as_slice()
417 .first()
418 .map_or("candidate did not pass the Ready Gate", |diagnostic| {
419 diagnostic.detail()
420 });
421 bail!("Plugin configuration proposal cannot be published: {detail}");
422 }
423
424 let path = root
425 .join(PLUGIN_ROOT)
426 .join(&proposal.plugin_id)
427 .join(format!("{}.toml", proposal.instance_key));
428 atomic_write(&path, &proposal.toml)?;
429 let published = snapshot_plugin_root(root, &host)?;
430 let revision = revision_for_snapshot(&published)?;
431 if revision != proposal.candidate_revision {
432 bail!("published Plugin Root does not match the reviewed candidate revision");
433 }
434 let resolved = super::inspect_plugin_root(root)?.resolved().clone();
435 Ok(PluginConfigurationPublication {
436 schema: PUBLICATION_SCHEMA,
437 base_revision: proposal.base_revision.clone(),
438 base_source_digest: proposal.base_source_digest.clone(),
439 revision,
440 proposal_digest: proposal.digest.clone(),
441 resolved,
442 })
443}
444
445fn build_proposal(
446 host: &HostInput,
447 current: &PluginRootSnapshot,
448 base_revision: PluginRootRevision,
449 base_source_digest: PluginConfigurationSourceDigest,
450 plugin_id: &str,
451 instance: &str,
452 bytes: &[u8],
453) -> anyhow::Result<PluginConfigurationProposal> {
454 let configuration = parse_configuration(bytes)?;
455 let id = PluginInstanceId::new(plugin_id, instance);
456 let mut instances = current
457 .instances()
458 .iter()
459 .filter(|item| item.id() != &id)
460 .cloned()
461 .collect::<Vec<_>>();
462 instances.push(PluginRootInstance::new(plugin_id, instance).with_configuration(configuration));
463 let candidate = crate::preserve_dependency_selections(
464 PluginRootSnapshot::new(
465 current.releases().iter().cloned(),
466 instances,
467 current.disabled().iter().cloned(),
468 ),
469 current,
470 );
471 let candidate_revision = revision_for_snapshot(&candidate)?;
472 let authority = serde_json::to_vec(&(
473 PROPOSAL_SCHEMA,
474 host,
475 current,
476 base_source_digest.as_str(),
477 &candidate,
478 bytes,
479 ))
480 .context("encode Plugin configuration proposal authority")?;
481 let digest = sha256_digest(&authority);
482 let (status, application, diagnostics) = match host.resolve(&candidate) {
483 Ok(_) => (
484 PluginConfigurationProposalStatus::Ready,
485 if candidate_revision == base_revision {
486 PluginConfigurationApplication::Noop
487 } else {
488 PluginConfigurationApplication::AppGeneration
489 },
490 Vec::new(),
491 ),
492 Err(error) => {
493 let status = if matches!(
494 error,
495 PluginRootResolutionError::AmbiguousSlot { .. }
496 | PluginRootResolutionError::AmbiguousCapability { .. }
497 ) {
498 PluginConfigurationProposalStatus::NeedsDecision
499 } else {
500 PluginConfigurationProposalStatus::Rejected
501 };
502 (
503 status,
504 PluginConfigurationApplication::Blocked,
505 vec![PluginConfigurationDiagnostic {
506 code: resolution_error_code(&error),
507 detail: error.to_string(),
508 }],
509 )
510 }
511 };
512 Ok(PluginConfigurationProposal {
513 schema: PROPOSAL_SCHEMA,
514 base_revision,
515 base_source_digest,
516 candidate_revision,
517 digest,
518 status,
519 application,
520 diagnostics,
521 plugin_id: plugin_id.to_owned(),
522 instance_key: instance.to_owned(),
523 toml: bytes.to_vec(),
524 })
525}
526
527fn resolution_error_code(error: &PluginRootResolutionError) -> &'static str {
528 match error {
529 PluginRootResolutionError::InvalidHostConfiguration(_) => "host_admission_denied",
530 PluginRootResolutionError::AmbiguousSlot { .. } => "ambiguous_slot",
531 PluginRootResolutionError::AmbiguousCapability { .. } => "ambiguous_capability",
532 PluginRootResolutionError::InvalidConfiguration { .. } => "invalid_configuration",
533 PluginRootResolutionError::MissingRequiredSlot(_) => "missing_required_slot",
534 PluginRootResolutionError::MissingCapability { .. } => "missing_capability",
535 PluginRootResolutionError::RequiredInstanceDisabled(_) => "required_instance_disabled",
536 PluginRootResolutionError::UnknownPlugin(_) => "unknown_plugin",
537 PluginRootResolutionError::UnknownDisabledInstance(_) => "unknown_disabled_instance",
538 _ => "invalid_plugin_root",
539 }
540}
541
542fn parse_configuration(bytes: &[u8]) -> anyhow::Result<serde_json::Value> {
543 let byte_count = u64::try_from(bytes.len()).context("Plugin configuration is too large")?;
544 if byte_count > MAX_CONFIGURATION_BYTES {
545 bail!("Plugin configuration exceeds 256 KiB");
546 }
547 let text = str::from_utf8(bytes).context("Plugin configuration must be UTF-8 TOML")?;
548 let table: toml::Table = toml::from_str(text).context("parse Plugin configuration TOML")?;
549 serde_json::to_value(table).context("convert Plugin configuration to portable values")
550}
551
552pub(crate) fn ensure_revision(
553 expected: &PluginRootRevision,
554 current: &PluginRootRevision,
555) -> anyhow::Result<()> {
556 if expected == current {
557 return Ok(());
558 }
559 Err(PluginRootRevisionConflict {
560 expected: expected.clone(),
561 current: current.clone(),
562 }
563 .into())
564}
565
566#[derive(Clone, Debug, Eq, PartialEq)]
568pub struct PluginConfigurationSourceDigest(String);
569
570impl PluginConfigurationSourceDigest {
571 pub fn as_str(&self) -> &str {
572 &self.0
573 }
574
575 pub fn for_source(
577 plugin_id: &str,
578 instance: &str,
579 bytes: Option<&[u8]>,
580 ) -> anyhow::Result<Self> {
581 validate_existing_plugin_id(plugin_id)?;
582 validate_instance_filename(instance)?;
583 Ok(source_digest_for_bytes(plugin_id, instance, bytes))
584 }
585}
586
587impl fmt::Display for PluginConfigurationSourceDigest {
588 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
589 formatter.write_str(&self.0)
590 }
591}
592
593#[derive(Clone, Debug, Eq, PartialEq)]
595pub struct PluginConfigurationSourceConflict {
596 expected: PluginConfigurationSourceDigest,
597 current: PluginConfigurationSourceDigest,
598}
599
600impl PluginConfigurationSourceConflict {
601 pub const fn expected(&self) -> &PluginConfigurationSourceDigest {
602 &self.expected
603 }
604
605 pub const fn current(&self) -> &PluginConfigurationSourceDigest {
606 &self.current
607 }
608}
609
610impl fmt::Display for PluginConfigurationSourceConflict {
611 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
612 write!(
613 formatter,
614 "Plugin configuration source conflict: expected {}, current {}",
615 self.expected, self.current
616 )
617 }
618}
619
620impl Error for PluginConfigurationSourceConflict {}
621
622pub(crate) fn source_digest_for_bytes(
623 plugin_id: &str,
624 instance: &str,
625 bytes: Option<&[u8]>,
626) -> PluginConfigurationSourceDigest {
627 let mut authority = Sha256::new();
628 update_digest_component(&mut authority, SOURCE_DIGEST_SCHEMA.as_bytes());
629 update_digest_component(&mut authority, plugin_id.as_bytes());
630 update_digest_component(&mut authority, instance.as_bytes());
631 match bytes {
632 Some(bytes) => {
633 authority.update([1]);
634 update_digest_component(&mut authority, bytes);
635 }
636 None => authority.update([0]),
637 }
638 PluginConfigurationSourceDigest(encode_sha256(authority.finalize()))
639}
640
641fn source_digest_for_instance(
642 root: &Path,
643 plugin_id: &str,
644 instance: &str,
645) -> anyhow::Result<PluginConfigurationSourceDigest> {
646 let path = root
647 .join(PLUGIN_ROOT)
648 .join(plugin_id)
649 .join(format!("{instance}.toml"));
650 let bytes = match fs::symlink_metadata(&path) {
651 Ok(metadata) if metadata.file_type().is_file() => Some(
652 fs::read(&path)
653 .with_context(|| format!("read Plugin configuration source {}", path.display()))?,
654 ),
655 Ok(_) => bail!(
656 "Plugin configuration source must be a regular file: {}",
657 path.display()
658 ),
659 Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
660 Err(error) => {
661 return Err(error).with_context(|| {
662 format!("inspect Plugin configuration source {}", path.display())
663 });
664 }
665 };
666 Ok(source_digest_for_bytes(
667 plugin_id,
668 instance,
669 bytes.as_deref(),
670 ))
671}
672
673fn ensure_source_digest(
674 expected: &PluginConfigurationSourceDigest,
675 current: &PluginConfigurationSourceDigest,
676) -> anyhow::Result<()> {
677 if expected == current {
678 return Ok(());
679 }
680 Err(PluginConfigurationSourceConflict {
681 expected: expected.clone(),
682 current: current.clone(),
683 }
684 .into())
685}
686
687fn update_digest_component(authority: &mut Sha256, bytes: &[u8]) {
688 authority.update(u64::try_from(bytes.len()).unwrap_or(u64::MAX).to_be_bytes());
689 authority.update(bytes);
690}
691
692pub(crate) fn revision_for_snapshot(
693 snapshot: &PluginRootSnapshot,
694) -> anyhow::Result<PluginRootRevision> {
695 let canonical =
696 serde_json::to_vec(snapshot).context("encode Plugin Root revision authority")?;
697 Ok(PluginRootRevision(sha256_digest(&canonical)))
698}
699
700fn sha256_digest(bytes: &[u8]) -> String {
701 encode_sha256(Sha256::digest(bytes))
702}
703
704fn encode_sha256(digest: impl AsRef<[u8]>) -> String {
705 let digest = digest.as_ref();
706 let mut encoded = String::with_capacity(7 + digest.len() * 2);
707 encoded.push_str("sha256:");
708 for byte in digest {
709 write!(&mut encoded, "{byte:02x}").expect("writing to a String cannot fail");
710 }
711 encoded
712}
713
714#[cfg(test)]
715mod tests {
716 use std::{fs, sync::Arc};
717
718 use lenso_app_plan::authoring::{
719 HostDefaultPlugin, HostPluginRelease, HostSlot, PluginDescriptor,
720 };
721
722 use super::*;
723 use crate::{HOST_CATALOG, inspect_plugin_root};
724
725 fn fixture_root() -> tempfile::TempDir {
726 let root = tempfile::tempdir().unwrap();
727 fs::create_dir_all(root.path().join(".lenso")).unwrap();
728 let descriptor = PluginDescriptor::new("example.agent", "1.0.0", "agent")
729 .with_configuration_schema(serde_json::json!({
730 "type": "object",
731 "properties": {
732 "greeting": { "type": "string" }
733 },
734 "additionalProperties": false
735 }));
736 let host = lenso_app_plan::authoring::HostCatalog::new(
737 [HostSlot::one("agent")],
738 [HostPluginRelease::new(descriptor)],
739 [HostDefaultPlugin::new("example.agent", "default")],
740 );
741 fs::write(
742 root.path().join(HOST_CATALOG),
743 serde_json::to_vec(&host).unwrap(),
744 )
745 .unwrap();
746 root
747 }
748
749 #[test]
750 fn proposal_is_read_only_and_publication_advances_the_revision() {
751 let root = fixture_root();
752 let base = inspect_plugin_root(root.path()).unwrap().revision().clone();
753 let proposal = propose_instance_configuration(
754 root.path(),
755 &base,
756 "example.agent",
757 "default",
758 b"greeting = \"hello\"\n",
759 )
760 .unwrap();
761
762 assert_eq!(proposal.status(), PluginConfigurationProposalStatus::Ready);
763 assert_eq!(
764 proposal.application(),
765 PluginConfigurationApplication::AppGeneration
766 );
767 assert_eq!(proposal.base_revision(), &base);
768 assert!(
769 proposal
770 .base_source_digest()
771 .as_str()
772 .starts_with("sha256:")
773 );
774 assert_ne!(proposal.candidate_revision(), &base);
775 assert!(proposal.digest().starts_with("sha256:"));
776 assert!(!configuration_path(root.path()).exists());
777
778 let publication = publish_instance_configuration(root.path(), &proposal).unwrap();
779 assert_eq!(publication.base_revision(), &base);
780 assert_eq!(
781 publication.base_source_digest(),
782 proposal.base_source_digest()
783 );
784 assert_eq!(publication.revision(), proposal.candidate_revision());
785 assert_eq!(publication.proposal_digest(), proposal.digest());
786 assert_eq!(
787 fs::read_to_string(configuration_path(root.path())).unwrap(),
788 "greeting = \"hello\"\n"
789 );
790 }
791
792 #[test]
793 fn local_authority_dispatches_through_the_host_port() {
794 let root = fixture_root();
795 let authority: Arc<dyn PluginConfigurationAuthority> =
796 Arc::new(LocalPluginRootAuthority::new(root.path()));
797 let source = authority.source();
798 let base = authority.inspect().unwrap().revision().clone();
799
800 let proposal = authority
801 .propose(
802 &base,
803 "example.agent",
804 "default",
805 b"greeting = \"authority\"\n",
806 )
807 .unwrap();
808 assert!(!configuration_path(root.path()).exists());
809
810 let publication = authority.publish(&proposal).unwrap();
811
812 assert_eq!(source.kind(), "local_plugin_root");
813 assert_eq!(source.reference(), "app");
814 assert_eq!(publication.revision(), proposal.candidate_revision());
815 assert_eq!(
816 authority.inspect().unwrap().revision(),
817 publication.revision()
818 );
819 }
820
821 #[test]
822 fn stale_publication_fails_with_a_typed_conflict_and_preserves_the_winner() {
823 let root = fixture_root();
824 let base = inspect_plugin_root(root.path()).unwrap().revision().clone();
825 let first = proposal(root.path(), &base, b"greeting = \"first\"\n");
826 let stale = proposal(root.path(), &base, b"greeting = \"stale\"\n");
827 let first_publication = publish_instance_configuration(root.path(), &first).unwrap();
828
829 let error = publish_instance_configuration(root.path(), &stale).unwrap_err();
830 let conflict = error.downcast_ref::<PluginRootRevisionConflict>().unwrap();
831
832 assert_eq!(conflict.expected(), &base);
833 assert_eq!(conflict.current(), first_publication.revision());
834 assert_eq!(
835 fs::read_to_string(configuration_path(root.path())).unwrap(),
836 "greeting = \"first\"\n"
837 );
838 }
839
840 #[test]
841 fn concurrent_publications_allow_exactly_one_winner() {
842 let root = fixture_root();
843 let base = inspect_plugin_root(root.path()).unwrap().revision().clone();
844 let first = proposal(root.path(), &base, b"greeting = \"first\"\n");
845 let second = proposal(root.path(), &base, b"greeting = \"second\"\n");
846 let path = root.path().to_path_buf();
847 let barrier = Arc::new(std::sync::Barrier::new(3));
848 let handles = [first, second].map(|proposal| {
849 let path = path.clone();
850 let barrier = Arc::clone(&barrier);
851 std::thread::spawn(move || {
852 barrier.wait();
853 publish_instance_configuration(&path, &proposal)
854 })
855 });
856 barrier.wait();
857 let outcomes = handles.map(|handle| handle.join().unwrap());
858
859 assert_eq!(outcomes.iter().filter(|outcome| outcome.is_ok()).count(), 1);
860 let error = outcomes.into_iter().find_map(Result::err).unwrap();
861 assert!(error.downcast_ref::<PluginRootRevisionConflict>().is_some());
862 let contents = fs::read_to_string(configuration_path(&path)).unwrap();
863 assert!(contents == "greeting = \"first\"\n" || contents == "greeting = \"second\"\n");
864 }
865
866 #[test]
867 fn rejected_proposal_keeps_structured_diagnostics_without_writing() {
868 let root = fixture_root();
869 let base = inspect_plugin_root(root.path()).unwrap().revision().clone();
870 let proposal = proposal(root.path(), &base, b"unexpected = true\n");
871
872 assert_eq!(
873 proposal.status(),
874 PluginConfigurationProposalStatus::Rejected
875 );
876 assert_eq!(
877 proposal.application(),
878 PluginConfigurationApplication::Blocked
879 );
880 assert_eq!(proposal.diagnostics()[0].code(), "invalid_configuration");
881 assert!(!configuration_path(root.path()).exists());
882 assert!(publish_instance_configuration(root.path(), &proposal).is_err());
883 }
884
885 #[test]
886 fn plugin_root_revision_is_semantic_not_toml_formatting() {
887 let root = fixture_root();
888 let base = inspect_plugin_root(root.path()).unwrap().revision().clone();
889 let proposal = proposal(root.path(), &base, b"greeting = \"hello\"\n");
890 let publication = publish_instance_configuration(root.path(), &proposal).unwrap();
891 fs::write(
892 configuration_path(root.path()),
893 b"# human note\n\ngreeting=\"hello\"\n",
894 )
895 .unwrap();
896
897 let reformatted = inspect_plugin_root(root.path()).unwrap();
898 assert_eq!(reformatted.revision(), publication.revision());
899 }
900
901 #[test]
902 fn formatting_only_source_change_rejects_stale_publication_without_overwrite() {
903 let root = fixture_root();
904 let base = inspect_plugin_root(root.path()).unwrap().revision().clone();
905 let initial = proposal(root.path(), &base, b"greeting = \"hello\"\n");
906 let initial = publish_instance_configuration(root.path(), &initial).unwrap();
907 let stale = proposal(root.path(), initial.revision(), b"greeting = \"goodbye\"\n");
908 let external = b"# keep this human note\n\ngreeting=\"hello\"\n";
909 fs::write(configuration_path(root.path()), external).unwrap();
910
911 assert_eq!(
912 inspect_plugin_root(root.path()).unwrap().revision(),
913 initial.revision()
914 );
915 let error = publish_instance_configuration(root.path(), &stale).unwrap_err();
916 let conflict = error
917 .downcast_ref::<PluginConfigurationSourceConflict>()
918 .unwrap();
919
920 assert_eq!(conflict.expected(), stale.base_source_digest());
921 assert_ne!(conflict.current(), stale.base_source_digest());
922 assert_eq!(fs::read(configuration_path(root.path())).unwrap(), external);
923 }
924
925 #[test]
926 fn proposal_digest_closes_the_exact_reviewed_toml() {
927 let root = fixture_root();
928 let base = inspect_plugin_root(root.path()).unwrap().revision().clone();
929 let compact = propose_instance_configuration(
930 root.path(),
931 &base,
932 "example.agent",
933 "default",
934 b"greeting=\"hello\"\n",
935 )
936 .unwrap();
937 let formatted = propose_instance_configuration(
938 root.path(),
939 &base,
940 "example.agent",
941 "default",
942 b"greeting = \"hello\"\n",
943 )
944 .unwrap();
945
946 assert_eq!(compact.candidate_revision(), formatted.candidate_revision());
947 assert_ne!(compact.digest(), formatted.digest());
948 }
949
950 #[test]
951 fn source_digest_domains_raw_bytes_absence_and_instance_identity() {
952 let absent =
953 PluginConfigurationSourceDigest::for_source("example.agent", "default", None).unwrap();
954 let empty =
955 PluginConfigurationSourceDigest::for_source("example.agent", "default", Some(b""))
956 .unwrap();
957 let other_instance =
958 PluginConfigurationSourceDigest::for_source("example.agent", "secondary", None)
959 .unwrap();
960
961 assert_ne!(absent, empty);
962 assert_ne!(absent, other_instance);
963 }
964
965 #[test]
966 fn plugin_root_revision_round_trips_for_http_preconditions() {
967 let root = fixture_root();
968 let revision = inspect_plugin_root(root.path()).unwrap().revision().clone();
969
970 assert_eq!(
971 revision.as_str().parse::<PluginRootRevision>().unwrap(),
972 revision
973 );
974 assert!("sha256:not-a-digest".parse::<PluginRootRevision>().is_err());
975 }
976
977 fn proposal(
978 root: &Path,
979 base: &PluginRootRevision,
980 toml: &[u8],
981 ) -> PluginConfigurationProposal {
982 propose_instance_configuration(root, base, "example.agent", "default", toml).unwrap()
983 }
984
985 fn configuration_path(root: &Path) -> std::path::PathBuf {
986 root.join("plugins/example.agent/default.toml")
987 }
988}