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