1use super::{
4 AuthMode, HookSet, Pipeline, meta_error, ms,
5 revocation::{RevokeBudget, RevokeProgress},
6};
7use crate::{
8 authority::FenceKind,
9 error::ServerError,
10 policy::NamespacePolicy,
11 repo::{Addressing, NamespaceKey},
12 store::{MultipartBlobStore, NamespaceStore, codec, keys},
13};
14use mkit_attest::grant::EpochTransition;
15use mkit_core::repo_identity::Namespace;
16
17impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
18 pub(super) async fn stored_authority_generation(
19 &self,
20 ns: &NamespaceKey,
21 ) -> Result<u64, ServerError> {
22 self.meta
23 .get(&self.shards.coordinator(ns), &keys::authority_generation())
24 .await
25 .map_err(meta_error)?
26 .as_ref()
27 .map(codec::decode_u64)
28 .transpose()
29 .map_err(meta_error)
30 .map(|generation| generation.unwrap_or(0))
31 }
32
33 pub(super) async fn ensure_authority_activation(
36 &self,
37 ns: &NamespaceKey,
38 ) -> Result<Option<u64>, ServerError> {
39 let p = self.shards.coordinator(ns);
40 let wanted = [keys::authority_generation(), keys::lease_recovery()];
41 for _ in 0..1 {
42 let rows = self.meta.get_many(&p, &wanted).await.map_err(meta_error)?;
43 let [generation, raw_mode] = rows.as_slice() else {
44 return Err(super::internal("authority activation row count"));
45 };
46 let current = generation
47 .as_ref()
48 .map(codec::decode_u64)
49 .transpose()
50 .map_err(meta_error)?
51 .unwrap_or(0);
52 let mut mode = raw_mode
53 .as_ref()
54 .map(codec::decode_lease_recovery)
55 .transpose()
56 .map_err(meta_error)?
57 .unwrap_or(codec::LeaseRecovery {
58 resumed_at_ms: 0,
59 authority_fence: None,
60 authority_ready: None,
61 activation_only: Some(true),
62 });
63 if self.cfg.authority_fence.is_none() {
64 if generation.is_some() || mode.authority_fence == Some(true) {
65 return Err(ServerError::unavailable(
66 "persisted authority fence requires enabled executor",
67 ));
68 }
69 return Ok(None);
70 }
71 if mode.authority_fence == Some(true) && generation.is_none() {
72 return Err(ServerError::unavailable(
73 "authority generation missing from fenced namespace",
74 ));
75 }
76 if mode.authority_fence != Some(true) {
77 mode.authority_fence = Some(true);
78 mode.authority_ready = Some(false);
79 let batch = crate::store::Batch::new()
80 .require(super::lease::observed_guard(
81 wanted[0].clone(),
82 generation.as_ref(),
83 ))
84 .require(super::lease::observed_guard(
85 wanted[1].clone(),
86 raw_mode.as_ref(),
87 ))
88 .put(wanted[0].clone(), codec::encode_u64(current))
89 .put(wanted[1].clone(), codec::encode_lease_recovery(&mode));
90 match self.apply_meta(&p, batch).await? {
91 crate::store::BatchOutcome::Committed => {}
92 crate::store::BatchOutcome::PreconditionFailed { .. } => continue,
93 crate::store::BatchOutcome::DeadlinePassed { .. } => {
94 return Err(super::internal("activation had no deadline"));
95 }
96 }
97 }
100 if mode.authority_ready == Some(true) {
101 return Ok(None);
102 }
103 if self
104 .revoke_fence_step(ns, &RevokeBudget::default(), FenceKind::Authority)
105 .await?
106 != RevokeProgress::Complete
107 {
108 return Err(
109 ServerError::unavailable("authority activation pending; retry")
110 .with_header("Retry-After", "1"),
111 );
112 }
113 let prior = codec::encode_lease_recovery(&mode);
114 mode.authority_ready = Some(true);
115 let batch = crate::store::Batch::new()
116 .require(crate::store::Precondition::Equals(
117 wanted[0].clone(),
118 codec::encode_u64(current),
119 ))
120 .require(crate::store::Precondition::Equals(wanted[1].clone(), prior))
121 .put(wanted[1].clone(), codec::encode_lease_recovery(&mode));
122 match self.apply_meta(&p, batch).await? {
123 crate::store::BatchOutcome::Committed => return Ok(Some(current)),
124 crate::store::BatchOutcome::PreconditionFailed { .. } => {}
125 crate::store::BatchOutcome::DeadlinePassed { .. } => {
126 return Err(super::internal("activation ready had no deadline"));
127 }
128 }
129 }
130 Err(
131 ServerError::unavailable("authority activation contention; retry")
132 .with_header("Retry-After", "1"),
133 )
134 }
135
136 pub(super) fn check_ticket_generation<'a>(
138 &'a self,
139 ns: &'a NamespaceKey,
140 generation: Option<u64>,
141 ) -> crate::BoxFuture<'a, Result<(), ServerError>> {
142 Box::pin(async move {
143 if self.cfg.authority_fence.is_none() {
144 if generation.is_some() {
145 return Err(ServerError::unavailable(
146 "fenced ticket requires enabled executor",
147 ));
148 }
149 Box::pin(self.ensure_authority_activation(ns)).await?;
150 } else {
151 let current = self
154 .meta
155 .get(&self.shards.coordinator(ns), &keys::authority_generation())
156 .await
157 .map_err(meta_error)?;
158 let current = current
159 .as_ref()
160 .map(codec::decode_u64)
161 .transpose()
162 .map_err(meta_error)?
163 .ok_or_else(|| {
164 ServerError::unavailable(
165 "authority generation missing from fenced namespace",
166 )
167 })?;
168 if generation != Some(current) {
169 return Err(crate::authority::moved());
170 }
171 }
172 Ok(())
173 })
174 }
175
176 pub async fn get_authority_generation(&self, namespace: &str) -> Result<u64, ServerError> {
180 if self.cfg.authority_fence.is_none() {
181 return Err(ServerError::unimplemented("authority fencing is disabled"));
182 }
183 let ns = Namespace::parse(namespace)
184 .map_err(|_| ServerError::invalid_argument("invalid namespace"))?;
185 self.authority_namespace(&ns).await?;
186 self.stored_authority_generation(&NamespaceKey::from_namespace(&ns))
187 .await
188 }
189
190 async fn authority_namespace(&self, ns: &Namespace) -> Result<(), ServerError> {
191 let Addressing::Multi(multi) = &self.cfg.addressing else {
192 return Err(ServerError::unimplemented(
193 "authority fencing requires multi addressing",
194 ));
195 };
196 match &multi.namespace_policy {
197 NamespacePolicy::Allowlist(allowed) if allowed.contains(ns) => Ok(()),
198 NamespacePolicy::Any { .. } => {
199 let key = NamespaceKey::from_namespace(ns);
200 if self
201 .meta
202 .get(&self.shards.coordinator(&key), &keys::namespace_record())
203 .await
204 .map_err(meta_error)?
205 .is_some()
206 {
207 Ok(())
208 } else {
209 Err(ServerError::permission_denied("namespace not served"))
210 }
211 }
212 _ => Err(ServerError::permission_denied("namespace not served")),
213 }
214 }
215
216 pub async fn set_authority_generation(&self, wire: &str) -> Result<u64, ServerError> {
221 let fence = self
222 .cfg
223 .authority_fence
224 .as_ref()
225 .ok_or_else(|| ServerError::unimplemented("authority fencing is disabled"))?;
226 let AuthMode::AuthV2(auth) = &self.cfg.auth else {
227 return Err(ServerError::unimplemented(
228 "authority fencing requires auth-v2",
229 ));
230 };
231 let statement = fence.verify(wire, auth.audience(), self.clock.now_ms())?;
232 self.authority_namespace(&statement.namespace).await?;
233 let ns = NamespaceKey::from_namespace(&statement.namespace);
234 if self
235 .transition_fence(&ns, statement.generation, FenceKind::Authority)
236 .await?
237 == EpochTransition::Reject
238 {
239 return Err(ServerError::permission_denied(
240 "authority generation step rejected",
241 ));
242 }
243 if let Some(completed) = Box::pin(self.ensure_authority_activation(&ns)).await? {
244 if self.stored_authority_generation(&ns).await? == completed {
245 return Ok(completed);
246 }
247 return Err(
248 ServerError::unavailable("authority revocation pending; retry")
249 .with_header("Retry-After", "1"),
250 );
251 }
252 let start = ms(self.clock.now_ms());
253 for _ in 0..1 {
254 let elapsed = ms(self.clock.now_ms()).saturating_sub(start);
255 if elapsed >= 5_000 {
256 break;
257 }
258 let before = self.stored_authority_generation(&ns).await?;
259 let budget = RevokeBudget::new((5_000 - elapsed).min(1_000));
260 if self
261 .revoke_fence_step(&ns, &budget, FenceKind::Authority)
262 .await?
263 == RevokeProgress::Complete
264 {
265 let after = self.stored_authority_generation(&ns).await?;
266 if before == after {
267 return Ok(after);
268 }
269 }
270 }
271 Err(
272 ServerError::unavailable("authority revocation pending; retry")
273 .with_header("Retry-After", "1"),
274 )
275 }
276}