1use mkit_core::hash::{Hash, from_hex};
10use mkit_core::protocol::{PackKey, RefWriteCondition};
11use mkit_core::write_auth::{Authorized, is_hex};
12
13use crate::error::ServerError;
14use crate::principal::Principal;
15use crate::repo::RepoId;
16
17const SERVICE_PREFIX: &str = "/mkit.transport.v1.TransportService/";
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
22#[non_exhaustive]
23pub enum Procedure {
24 ListRepos,
26 ListRefs,
28 ReadRef,
30 UpdateRef,
32 AdvanceRefs,
34 BeginUpload,
36 UploadPart,
38 CompleteUpload,
40 PackExists,
42 UploadPack,
44 DownloadPack,
46 GetReceipt,
48 SetRepoVisibility,
50 IssueObjectUrl,
52 HttpGetObject,
55 HttpGetRefPath,
57}
58
59impl Procedure {
60 #[must_use]
63 pub const fn connect_path(self) -> &'static str {
64 match self {
65 Self::ListRepos => "/mkit.transport.v1.TransportService/ListRepos",
66 Self::ListRefs => "/mkit.transport.v1.TransportService/ListRefs",
67 Self::ReadRef => "/mkit.transport.v1.TransportService/ReadRef",
68 Self::UpdateRef => "/mkit.transport.v1.TransportService/UpdateRef",
69 Self::AdvanceRefs => "/mkit.transport.v1.TransportService/AdvanceRefs",
70 Self::BeginUpload => "/mkit.transport.v1.TransportService/BeginUpload",
71 Self::UploadPart => "/mkit.transport.v1.TransportService/UploadPart",
72 Self::CompleteUpload => "/mkit.transport.v1.TransportService/CompleteUpload",
73 Self::PackExists => "/mkit.transport.v1.TransportService/PackExists",
74 Self::UploadPack => "/mkit.transport.v1.TransportService/UploadPack",
75 Self::DownloadPack => "/mkit.transport.v1.TransportService/DownloadPack",
76 Self::GetReceipt => "/mkit.transport.v1.TransportService/GetReceipt",
77 Self::SetRepoVisibility => "/mkit.transport.v1.TransportService/SetRepoVisibility",
78 Self::IssueObjectUrl => "/mkit.transport.v1.TransportService/IssueObjectUrl",
79 Self::HttpGetObject => "/mkit.http.v1/GetObject",
80 Self::HttpGetRefPath => "/mkit.http.v1/GetRefPath",
81 }
82 }
83
84 #[must_use]
88 pub fn from_connect_path(path: &str) -> Option<Self> {
89 Some(match path.strip_prefix(SERVICE_PREFIX)? {
90 "ListRepos" => Self::ListRepos,
91 "ListRefs" => Self::ListRefs,
92 "ReadRef" => Self::ReadRef,
93 "UpdateRef" => Self::UpdateRef,
94 "AdvanceRefs" => Self::AdvanceRefs,
95 "BeginUpload" => Self::BeginUpload,
96 "UploadPart" => Self::UploadPart,
97 "CompleteUpload" => Self::CompleteUpload,
98 "PackExists" => Self::PackExists,
99 "UploadPack" => Self::UploadPack,
100 "DownloadPack" => Self::DownloadPack,
101 "GetReceipt" => Self::GetReceipt,
102 "SetRepoVisibility" => Self::SetRepoVisibility,
103 "IssueObjectUrl" => Self::IssueObjectUrl,
104 _ => return None,
105 })
106 }
107
108 #[must_use]
112 pub const fn is_write(self) -> bool {
113 match self {
114 Self::UpdateRef
115 | Self::AdvanceRefs
116 | Self::BeginUpload
117 | Self::UploadPack
118 | Self::UploadPart
119 | Self::CompleteUpload
120 | Self::SetRepoVisibility => true,
121 Self::ListRepos
122 | Self::ListRefs
123 | Self::ReadRef
124 | Self::PackExists
125 | Self::DownloadPack
126 | Self::GetReceipt
127 | Self::IssueObjectUrl
128 | Self::HttpGetObject
129 | Self::HttpGetRefPath => false,
130 }
131 }
132
133 #[must_use]
135 pub const fn is_streaming(self) -> bool {
136 matches!(
137 self,
138 Self::UploadPack | Self::UploadPart | Self::DownloadPack
139 )
140 }
141}
142
143#[derive(Debug, Clone, Copy, PartialEq, Eq)]
145#[non_exhaustive]
146pub enum Commitment {
147 Body(Hash),
149 Pack {
151 id: Hash,
153 len: u64,
155 },
156 Part {
158 ticket: Hash,
160 index: u32,
162 subtree: Hash,
164 len: u64,
166 },
167}
168
169impl Commitment {
170 pub fn parse(text: &str) -> Result<Self, ServerError> {
176 let invalid = || ServerError::invalid_argument("invalid content commitment");
177 if let Some(digest) = text.strip_prefix("body:") {
178 return Ok(Self::Body(canonical_hash(digest).ok_or_else(invalid)?));
179 }
180 if let Some(part) = text.strip_prefix("part:") {
181 let mut fields = part.split(':');
182 let ticket = canonical_hash(fields.next().ok_or_else(invalid)?).ok_or_else(invalid)?;
183 let index =
184 canonical_decimal::<u32>(fields.next().ok_or_else(invalid)?).ok_or_else(invalid)?;
185 let subtree = canonical_hash(fields.next().ok_or_else(invalid)?).ok_or_else(invalid)?;
186 let len =
187 canonical_decimal::<u64>(fields.next().ok_or_else(invalid)?).ok_or_else(invalid)?;
188 if fields.next().is_some() {
189 return Err(invalid());
190 }
191 return Ok(Self::Part {
192 ticket,
193 index,
194 subtree,
195 len,
196 });
197 }
198 let (digest, len) = text
199 .strip_prefix("pack:")
200 .and_then(|pack| pack.split_once(':'))
201 .ok_or_else(invalid)?;
202 let len = canonical_decimal::<u64>(len).ok_or_else(invalid)?;
203 Ok(Self::Pack {
204 id: canonical_hash(digest).ok_or_else(invalid)?,
205 len,
206 })
207 }
208}
209
210fn canonical_decimal<T: core::str::FromStr + ToString>(text: &str) -> Option<T> {
211 text.parse::<T>().ok().filter(|n| n.to_string() == text)
212}
213
214fn canonical_hash(text: &str) -> Option<Hash> {
216 if is_hex(text, 32) {
217 from_hex(text).ok()
218 } else {
219 None
220 }
221}
222
223#[derive(Debug, Clone, PartialEq, Eq)]
229#[non_exhaustive]
230pub struct VerifiedAuth {
231 pub signer: [u8; 32],
233 pub replay_scope: Hash,
235 pub fingerprint: Hash,
238 pub nonce: String,
240 pub commitment: Commitment,
242 pub expires_at_ms: i64,
244 pub created_at_ms: i64,
249}
250
251impl TryFrom<&Authorized> for VerifiedAuth {
252 type Error = ServerError;
253
254 fn try_from(auth: &Authorized) -> Result<Self, Self::Error> {
257 let malformed = || ServerError::unauthenticated("malformed auth v2 authorization");
258 if !is_hex(&auth.nonce, 32) {
259 return Err(malformed());
260 }
261 Ok(Self {
262 signer: canonical_hash(&auth.public_key).ok_or_else(malformed)?,
263 replay_scope: canonical_hash(&auth.scope).ok_or_else(malformed)?,
264 fingerprint: canonical_hash(&auth.fingerprint).ok_or_else(malformed)?,
265 nonce: auth.nonce.clone(),
266 commitment: Commitment::parse(&auth.commitment).map_err(|_| malformed())?,
267 expires_at_ms: auth.expires_at,
268 created_at_ms: 0,
269 })
270 }
271}
272
273#[derive(Debug, Clone, PartialEq, Eq)]
275pub struct RefUpdate {
276 pub name: String,
278 pub condition: RefWriteCondition,
280 pub new: Option<Hash>,
283}
284
285#[derive(Debug, Clone, PartialEq, Eq)]
287#[non_exhaustive]
288pub enum OpKind {
289 ListRepos {
291 name_prefix: String,
293 },
294 ListRefs {
296 prefix: String,
298 },
299 ReadRef {
301 name: String,
303 },
304 UpdateRef(RefUpdate),
306 AdvanceRefs {
308 head: RefUpdate,
310 packmap: RefUpdate,
312 tickets: Vec<Hash>,
314 },
315 BeginUpload {
317 ref_name: String,
319 key: PackKey,
321 bytes: u64,
323 },
324 PackExists {
326 key: PackKey,
328 },
329 UploadPack {
331 key: PackKey,
333 declared_len: u64,
335 },
336 DownloadPack {
338 key: PackKey,
340 },
341 SetRepoVisibility {
344 visibility: mkit_attest::grant::Visibility,
346 },
347 IssueObjectUrl {
350 target: crate::url_token::UrlTarget,
352 ttl_seconds: u32,
354 },
355 HttpGet {
359 ref_name: Option<String>,
361 },
362}
363
364impl OpKind {
365 #[must_use]
367 pub const fn procedure(&self) -> Procedure {
368 match self {
369 Self::ListRepos { .. } => Procedure::ListRepos,
370 Self::ListRefs { .. } => Procedure::ListRefs,
371 Self::ReadRef { .. } => Procedure::ReadRef,
372 Self::UpdateRef(_) => Procedure::UpdateRef,
373 Self::AdvanceRefs { .. } => Procedure::AdvanceRefs,
374 Self::BeginUpload { .. } => Procedure::BeginUpload,
375 Self::PackExists { .. } => Procedure::PackExists,
376 Self::UploadPack { .. } => Procedure::UploadPack,
377 Self::DownloadPack { .. } => Procedure::DownloadPack,
378 Self::SetRepoVisibility { .. } => Procedure::SetRepoVisibility,
379 Self::IssueObjectUrl { .. } => Procedure::IssueObjectUrl,
380 Self::HttpGet { ref_name: None } => Procedure::HttpGetObject,
381 Self::HttpGet { ref_name: Some(_) } => Procedure::HttpGetRefPath,
382 }
383 }
384}
385
386#[derive(Debug, Clone, PartialEq, Eq)]
389pub struct GrantRef {
390 pub id: Hash,
392 pub epoch: u64,
394 pub presence_requirement: Option<PresenceRequirement>,
396}
397
398#[derive(Debug, Clone, PartialEq, Eq)]
401pub enum PresenceRequirement {
402 Absent(String),
404 Present(String),
406}
407
408#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
411#[non_exhaustive]
412pub enum CallerView {
413 #[default]
415 Anonymous,
416 Reader,
418 Writer,
421}
422
423#[derive(Debug, Clone, Default, PartialEq, Eq)]
429#[non_exhaustive]
430pub struct AuthzFacts {
431 pub grant: Option<GrantRef>,
433 pub owner: bool,
435 pub caller_view: CallerView,
438 pub authority_generation: Option<u64>,
440}
441
442#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
444#[non_exhaustive]
445pub struct Creation {
446 pub namespace: bool,
448 pub repo: bool,
450}
451
452#[derive(Debug, Clone, PartialEq, Eq)]
456#[non_exhaustive]
457pub struct Operation {
458 pub repo: RepoId,
460 pub principal: Principal,
462 pub auth: Option<VerifiedAuth>,
464 pub write_grant: Option<crate::Redacted>,
466 pub kind: OpKind,
468 pub leased_epoch: Option<u64>,
470 pub observed_epoch: Option<u64>,
472 pub business_now_ms: Option<i64>,
474 pub authz: AuthzFacts,
476 pub creation: Creation,
478 pub created: Creation,
480}
481
482impl Operation {
483 #[must_use]
485 pub fn new(
486 repo: RepoId,
487 principal: Principal,
488 auth: Option<VerifiedAuth>,
489 kind: OpKind,
490 ) -> Self {
491 Self {
492 repo,
493 principal,
494 auth,
495 write_grant: None,
496 kind,
497 leased_epoch: None,
498 observed_epoch: None,
499 business_now_ms: None,
500 authz: AuthzFacts::default(),
501 creation: Creation::default(),
502 created: Creation::default(),
503 }
504 }
505
506 #[must_use]
508 pub const fn procedure(&self) -> Procedure {
509 self.kind.procedure()
510 }
511}
512
513#[cfg(test)]
514mod tests {
515 use mkit_core::hash::to_hex;
516 use mkit_core::write_auth::{Context, Headers, verify_headers};
517
518 use super::*;
519 use crate::error::Code;
520 use crate::repo::{NamespaceKey, RepoName};
521
522 const ALL: [(Procedure, &str); 14] = [
523 (Procedure::ListRefs, "ListRefs"),
524 (Procedure::ListRepos, "ListRepos"),
525 (Procedure::ReadRef, "ReadRef"),
526 (Procedure::UpdateRef, "UpdateRef"),
527 (Procedure::AdvanceRefs, "AdvanceRefs"),
528 (Procedure::BeginUpload, "BeginUpload"),
529 (Procedure::UploadPart, "UploadPart"),
530 (Procedure::CompleteUpload, "CompleteUpload"),
531 (Procedure::PackExists, "PackExists"),
532 (Procedure::UploadPack, "UploadPack"),
533 (Procedure::DownloadPack, "DownloadPack"),
534 (Procedure::GetReceipt, "GetReceipt"),
535 (Procedure::SetRepoVisibility, "SetRepoVisibility"),
536 (Procedure::IssueObjectUrl, "IssueObjectUrl"),
537 ];
538
539 const EXEMPT: [&str; 5] = [
542 "GetServerInfo",
543 "GetGrantEpoch",
544 "SetGrantEpoch",
545 "GetAuthorityGeneration",
546 "SetAuthorityGeneration",
547 ];
548
549 #[test]
550 fn every_transport_rpc_is_classified_or_exempt() {
551 let proto = include_str!("../../../../proto/mkit/transport/v1/transport.proto");
552 let service = proto
553 .split("service TransportService")
554 .nth(1)
555 .expect("TransportService");
556 let mut names = Vec::new();
557 for line in service.lines() {
558 let line = line.trim_start();
559 if let Some(rest) = line.strip_prefix("rpc ")
560 && let Some(name) = rest.split('(').next()
561 {
562 names.push(name.trim());
563 }
564 }
565 assert_eq!(names.len(), ALL.len() + EXEMPT.len());
566 for name in &names {
567 let path = format!("/mkit.transport.v1.TransportService/{name}");
568 if EXEMPT.contains(name) {
569 assert_eq!(Procedure::from_connect_path(&path), None, "{name}");
570 } else {
571 assert!(
572 Procedure::from_connect_path(&path).is_some(),
573 "{name} is not classified"
574 );
575 }
576 }
577 for (_, name) in ALL {
578 assert!(names.contains(&name), "{name} missing from the proto");
579 }
580 }
581
582 #[test]
583 fn http_procedures_are_hook_names_only() {
584 for (procedure, path) in [
585 (Procedure::HttpGetObject, "/mkit.http.v1/GetObject"),
586 (Procedure::HttpGetRefPath, "/mkit.http.v1/GetRefPath"),
587 ] {
588 assert_eq!(procedure.connect_path(), path);
589 assert_eq!(Procedure::from_connect_path(path), None);
590 assert!(!procedure.is_write());
591 assert!(!procedure.is_streaming());
592 }
593 assert_eq!(
594 OpKind::HttpGet { ref_name: None }.procedure(),
595 Procedure::HttpGetObject
596 );
597 assert_eq!(
598 OpKind::HttpGet {
599 ref_name: Some("refs/heads/main".into())
600 }
601 .procedure(),
602 Procedure::HttpGetRefPath
603 );
604 }
605
606 #[test]
607 fn procedure_paths_roundtrip() {
608 for (procedure, method) in ALL {
610 let path = procedure.connect_path();
611 assert_eq!(
612 path,
613 format!("/mkit.transport.v1.TransportService/{method}")
614 );
615 assert_eq!(Procedure::from_connect_path(path), Some(procedure));
616 }
617 for bad in [
618 "",
619 "/mkit.transport.v1.TransportService/",
620 "/mkit.transport.v1.TransportService/updateref",
621 "/mkit.transport.v1.TransportService/UpdateRef/",
622 "/mkit.repo.v1.RepoService/UpdateRef",
623 "mkit.transport.v1.TransportService/UpdateRef",
624 ] {
625 assert_eq!(Procedure::from_connect_path(bad), None, "{bad}");
626 }
627 }
628
629 #[test]
630 fn procedure_write_and_streaming_classes() {
631 let writes: Vec<_> = ALL
632 .iter()
633 .filter(|(p, _)| p.is_write())
634 .map(|(p, _)| *p)
635 .collect();
636 assert_eq!(
637 writes,
638 [
639 Procedure::UpdateRef,
640 Procedure::AdvanceRefs,
641 Procedure::BeginUpload,
642 Procedure::UploadPart,
643 Procedure::CompleteUpload,
644 Procedure::UploadPack,
645 Procedure::SetRepoVisibility
646 ]
647 );
648 let streams: Vec<_> = ALL
649 .iter()
650 .filter(|(p, _)| p.is_streaming())
651 .map(|(p, _)| *p)
652 .collect();
653 assert_eq!(
654 streams,
655 [
656 Procedure::UploadPart,
657 Procedure::UploadPack,
658 Procedure::DownloadPack
659 ]
660 );
661 }
662
663 fn golden() -> serde_json::Value {
664 serde_json::from_str(include_str!("../../../tests/golden/auth-v2/unary.json")).unwrap()
665 }
666
667 fn authorized(commitment: &str) -> Authorized {
668 Authorized {
669 scope: "11".repeat(32),
670 public_key: "22".repeat(32),
671 nonce: "ab".repeat(32),
672 fingerprint: "33".repeat(32),
673 commitment: commitment.to_owned(),
674 expires_at: 1_700_000_300_000,
675 }
676 }
677
678 #[test]
679 fn verified_auth_from_authorized_parses_body_and_pack_commitments() {
680 let fixture = golden();
681 let field = |name: &str| fixture[name].as_str().unwrap().to_owned();
682 let created_at = fixture["created_at"].as_i64().unwrap();
683 let expires_at = fixture["expires_at"].as_i64().unwrap();
684 let headers = Headers {
685 version: Some("2".into()),
686 audience: Some(field("audience")),
687 repository: Some(field("repository")),
688 public_key: Some(field("public_key")),
689 signature: Some(field("signature")),
690 commitment: Some(field("commitment")),
691 digest: Some(field("body_digest")),
692 created_at: Some(created_at.to_string()),
693 expires_at: Some(expires_at.to_string()),
694 idempotency_key: Some(field("nonce")),
695 };
696 let audience = field("audience");
697 let repository = field("repository");
698 let commitment = field("commitment");
699 let auth = verify_headers(
700 Context {
701 audience: &audience,
702 repository: &repository,
703 },
704 &field("procedure"),
705 Some(&commitment),
706 created_at + 1,
707 &headers,
708 )
709 .unwrap();
710
711 let verified = VerifiedAuth::try_from(&auth).unwrap();
712 assert_eq!(
713 verified.commitment,
714 Commitment::Body(from_hex(&field("body_digest")).unwrap())
715 );
716 assert_eq!(to_hex(&verified.signer), field("public_key"));
717 assert_eq!(to_hex(&verified.fingerprint), field("signing_digest"));
718 assert_eq!(to_hex(&verified.replay_scope), auth.scope);
719 assert_eq!(verified.nonce, field("nonce"));
720 assert_eq!(verified.expires_at_ms, expires_at);
721
722 let pack = authorized(&format!("pack:{}:12", "cd".repeat(32)));
723 let verified = VerifiedAuth::try_from(&pack).unwrap();
724 assert_eq!(
725 verified.commitment,
726 Commitment::Pack {
727 id: [0xcd; 32],
728 len: 12
729 }
730 );
731 assert_eq!(verified.signer, [0x22; 32]);
732 assert_eq!(verified.replay_scope, [0x11; 32]);
733 }
734
735 #[test]
736 fn verified_auth_rejects_noncanonical_fields() {
737 let digest = "cd".repeat(32);
738 for commitment in [
739 format!("body:{}", digest.to_uppercase()),
740 format!("body:{}", &digest[..62]),
741 format!("pack:{digest}:012"),
742 format!("pack:{digest}:+12"),
743 format!("pack:{digest}:"),
744 format!("pack:{digest}"),
745 format!("pack:{digest}:18446744073709551616"),
746 format!("part:{digest}:1"),
747 String::new(),
748 ] {
749 assert_eq!(
750 Commitment::parse(&commitment).unwrap_err().code(),
751 Code::InvalidArgument,
752 "{commitment}"
753 );
754 let err = VerifiedAuth::try_from(&authorized(&commitment)).unwrap_err();
755 assert_eq!(err.code(), Code::Unauthenticated, "{commitment}");
756 }
757 let body = format!("body:{digest}");
758 for broken in [
759 Authorized {
760 public_key: "AA".repeat(32),
761 ..authorized(&body)
762 },
763 Authorized {
764 scope: "11".repeat(31),
765 ..authorized(&body)
766 },
767 Authorized {
768 fingerprint: "zz".repeat(32),
769 ..authorized(&body)
770 },
771 Authorized {
772 nonce: "AB".repeat(32),
773 ..authorized(&body)
774 },
775 ] {
776 let err = VerifiedAuth::try_from(&broken).unwrap_err();
777 assert_eq!(err.code(), Code::Unauthenticated);
778 }
779 }
780
781 #[test]
782 fn part_commitment_parses_canonically() {
783 let text = format!("part:{}:2:{}:8388608", "ab".repeat(32), "cd".repeat(32));
784 assert_eq!(
785 Commitment::parse(&text).unwrap(),
786 Commitment::Part {
787 ticket: [0xab; 32],
788 index: 2,
789 subtree: [0xcd; 32],
790 len: 8_388_608,
791 }
792 );
793 for wrong in [
794 text.replace(":2:", ":02:"),
795 text.replace(":8388608", ":08388608"),
796 text.to_uppercase(),
797 format!("{text}:extra"),
798 ] {
799 assert_eq!(
800 Commitment::parse(&wrong).unwrap_err().code(),
801 Code::InvalidArgument
802 );
803 }
804 }
805
806 fn update(name: &str) -> RefUpdate {
807 RefUpdate {
808 name: name.to_owned(),
809 condition: RefWriteCondition::Missing,
810 new: Some([1; 32]),
811 }
812 }
813
814 #[test]
815 fn op_kind_maps_to_its_procedure() {
816 let key = PackKey::new([9; 32]);
817 let cases = [
818 (
819 OpKind::ListRefs {
820 prefix: String::new(),
821 },
822 Procedure::ListRefs,
823 ),
824 (
825 OpKind::ReadRef {
826 name: "refs/heads/main".into(),
827 },
828 Procedure::ReadRef,
829 ),
830 (
831 OpKind::UpdateRef(update("refs/heads/main")),
832 Procedure::UpdateRef,
833 ),
834 (
835 OpKind::AdvanceRefs {
836 head: update("refs/heads/main"),
837 packmap: update("refs/packmap/main"),
838 tickets: Vec::new(),
839 },
840 Procedure::AdvanceRefs,
841 ),
842 (OpKind::PackExists { key }, Procedure::PackExists),
843 (
844 OpKind::UploadPack {
845 key,
846 declared_len: 12,
847 },
848 Procedure::UploadPack,
849 ),
850 (OpKind::DownloadPack { key }, Procedure::DownloadPack),
851 (
852 OpKind::SetRepoVisibility {
853 visibility: mkit_attest::grant::Visibility::Private,
854 },
855 Procedure::SetRepoVisibility,
856 ),
857 (
858 OpKind::IssueObjectUrl {
859 target: crate::url_token::UrlTarget::Object([0xaa; 32]),
860 ttl_seconds: 60,
861 },
862 Procedure::IssueObjectUrl,
863 ),
864 ];
865 for (kind, procedure) in cases {
866 assert_eq!(kind.procedure(), procedure);
867 }
868 }
869
870 #[test]
871 fn operation_authz_defaults_empty() {
872 let op = Operation::new(
873 RepoId {
874 namespace: NamespaceKey::deployment_default(),
875 name: RepoName::new("room-a").unwrap(),
876 },
877 Principal::Anonymous,
878 None,
879 OpKind::ReadRef {
880 name: "refs/heads/main".into(),
881 },
882 );
883 assert_eq!(op.authz, AuthzFacts::default());
884 assert_eq!(op.authz.grant, None);
885 assert!(!op.authz.owner);
886 assert_eq!(op.procedure(), Procedure::ReadRef);
887 assert!(!op.procedure().is_write());
888 }
889}