1use std::num::NonZeroUsize;
15use std::sync::Arc;
16
17use async_trait::async_trait;
18use jiff::{SignedDuration, Timestamp};
19use tokio::sync::broadcast;
20
21use tollgate_core::{
22 AccountId, AccountSnapshot, CostUnits, FencingToken, LeaseId, Principal, PublishableSnapshot,
23 UsageEvent,
24};
25use tollgate_store::wire::{
26 API_PREFIX, AcquireRequest, ConsolidateRequest, IngestRequestRef, LeaseTtl, PrincipalsResponse,
27 Problem, ReleaseRequest,
28};
29use tollgate_store::{
30 AllocateError, Allocation, IngestError, IngestReport, LeaseAllocator, ReclaimBatch,
31 SnapshotPush, SnapshotResolution, SnapshotSource, StoreError, UsageSink,
32};
33
34pub struct HttpStore {
45 base: String,
46 transport: arc_swap::ArcSwap<HttpTransport>,
47}
48
49struct HttpTransport {
50 client: reqwest::Client,
51 bearer: Option<Arc<dyn crate::http_security::BearerProvider>>,
52 request_timeout: std::time::Duration,
53}
54
55impl HttpStore {
56 pub fn new(base_url: impl Into<String>) -> Result<Arc<Self>, StoreError> {
59 Self::with_config(base_url, crate::http_security::HttpStoreConfig::default())
60 }
61
62 pub fn with_timeouts(
71 base_url: impl Into<String>,
72 connect_timeout: std::time::Duration,
73 request_timeout: std::time::Duration,
74 ) -> Result<Arc<Self>, StoreError> {
75 Self::with_config(
76 base_url,
77 crate::http_security::HttpStoreConfig {
78 connect_timeout,
79 request_timeout,
80 ..Default::default()
81 },
82 )
83 }
84
85 pub fn with_config(
95 base_url: impl Into<String>,
96 config: crate::http_security::HttpStoreConfig,
97 ) -> Result<Arc<Self>, StoreError> {
98 let (base, client) = crate::http_security::client(&base_url.into(), &config)?;
99 Ok(Arc::new(Self {
100 base,
101 transport: arc_swap::ArcSwap::from_pointee(HttpTransport {
102 client,
103 bearer: config.bearer,
104 request_timeout: config.request_timeout,
105 }),
106 }))
107 }
108
109 pub fn reconfigure(
113 &self,
114 config: crate::http_security::HttpStoreConfig,
115 ) -> Result<(), StoreError> {
116 let (_, client) = crate::http_security::client(&self.base, &config)?;
117 self.transport.store(Arc::new(HttpTransport {
118 client,
119 bearer: config.bearer,
120 request_timeout: config.request_timeout,
121 }));
122 Ok(())
123 }
124
125 fn request(&self, method: reqwest::Method, path: &str) -> HttpRequest {
126 let transport = self.transport.load_full();
127 let request = transport
128 .client
129 .request(method, format!("{}{API_PREFIX}{path}", self.base));
130 HttpRequest { transport, request }
131 }
132}
133
134struct HttpRequest {
137 transport: Arc<HttpTransport>,
138 request: reqwest::RequestBuilder,
139}
140
141impl HttpRequest {
142 fn json<T: serde::Serialize + ?Sized>(mut self, body: &T) -> Self {
143 self.request = self.request.json(body);
144 self
145 }
146
147 async fn send(self) -> Result<reqwest::Response, StoreError> {
148 let mut request = self.request;
149 let deadline = tokio::time::Instant::now() + self.transport.request_timeout;
150 if let Some(provider) = &self.transport.bearer {
151 let token = tokio::time::timeout_at(deadline, provider.token())
152 .await
153 .map_err(|_| StoreError("control-plane credential deadline expired".into()))??;
154 request = request.header(reqwest::header::AUTHORIZATION, token.header());
155 }
156 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
157 if remaining.is_zero() {
158 return Err(StoreError("control-plane request deadline expired".into()));
159 }
160 request
163 .timeout(remaining)
164 .send()
165 .await
166 .map_err(|e| StoreError(format!("http: {}", e.without_url())))
167 }
168}
169
170fn transport_error(e: reqwest::Error) -> AllocateError {
171 AllocateError::Storage(StoreError(format!("http: {}", e.without_url())))
172}
173
174fn problem_to_allocate(problem: Problem) -> AllocateError {
178 match problem.code.as_str() {
179 "unknown-account" => AllocateError::UnknownAccount,
180 "account-inactive" => AllocateError::AccountInactive,
181 "insufficient-balance" => match problem.balance_shortfall {
185 Some(evidence) if problem.status == 409 && !evidence.remaining.is_zero() => {
186 AllocateError::BalanceInsufficient(evidence)
187 }
188 _ => AllocateError::InsufficientBalance,
189 },
190 "balance-exhausted" => match problem.balance_exhaustion {
191 Some(evidence) if problem.status == 409 => AllocateError::BalanceExhausted(evidence),
192 _ => AllocateError::Storage(StoreError(
193 "exhaustion response lacks valid evidence".into(),
194 )),
195 },
196 "balance-overflow" => AllocateError::BalanceOverflow,
197 "invalid-ttl" => AllocateError::InvalidTtl,
198 "unknown-lease" => AllocateError::UnknownLease,
199 "fenced" => AllocateError::Fenced,
200 "lease-not-active" => AllocateError::LeaseNotActive,
201 "invalid-release" => AllocateError::InvalidRelease,
202 _ => AllocateError::Storage(StoreError(format!(
203 "server {}: {}",
204 problem.status, problem.title
205 ))),
206 }
207}
208
209async fn problem_detail(response: reqwest::Response) -> String {
215 let status = response.status().as_u16();
216 match response.json::<Problem>().await {
217 Ok(problem) => format!("{status} {}: {}", problem.code, problem.title),
218 Err(_) => format!("server returned {status}"),
219 }
220}
221
222async fn read_problem(response: reqwest::Response) -> AllocateError {
223 let status = response.status().as_u16();
224 match response.json::<Problem>().await {
225 Ok(problem) => problem_to_allocate(problem),
226 Err(_) => AllocateError::Storage(StoreError(format!("server returned {status}"))),
227 }
228}
229
230async fn read_allocation(response: reqwest::Response) -> Result<Allocation, AllocateError> {
234 require_complete(response)?
235 .json()
236 .await
237 .map(credible_evidence)
238 .map_err(transport_error)
239}
240
241fn credible_evidence(mut allocation: Allocation) -> Allocation {
242 if allocation
243 .funding
244 .is_some_and(|funding| funding.remaining < allocation.grant.units)
245 {
246 allocation.funding = None;
247 }
248 allocation
249}
250
251#[async_trait]
252impl LeaseAllocator for HttpStore {
253 async fn acquire(
254 &self,
255 account: AccountId,
256 requested: CostUnits,
257 ttl: SignedDuration,
258 _now: Timestamp,
259 ) -> Result<Allocation, AllocateError> {
260 let ttl = LeaseTtl::try_from(ttl)?;
262 let response = self
263 .request(reqwest::Method::POST, "/leases/acquire")
264 .json(&AcquireRequest {
265 account_id: account,
266 requested,
267 ttl,
268 })
269 .send()
270 .await?;
271 if !response.status().is_success() {
272 return Err(read_problem(response).await);
273 }
274 read_allocation(response).await
275 }
276
277 async fn release(
278 &self,
279 lease_id: LeaseId,
280 fencing_token: FencingToken,
281 unspent: CostUnits,
282 _now: Timestamp,
283 ) -> Result<(), AllocateError> {
284 let response = self
285 .request(reqwest::Method::POST, "/leases/release")
286 .json(&ReleaseRequest {
287 lease_id,
288 fencing_token,
289 unspent,
290 })
291 .send()
292 .await?;
293 if !response.status().is_success() {
294 return Err(read_problem(response).await);
295 }
296 Ok(())
297 }
298
299 async fn consolidate(
300 &self,
301 lease_id: LeaseId,
302 fencing_token: FencingToken,
303 unspent: CostUnits,
304 requested: CostUnits,
305 needed: CostUnits,
306 ttl: SignedDuration,
307 _now: Timestamp,
308 ) -> Result<Allocation, AllocateError> {
309 let ttl = LeaseTtl::try_from(ttl)?;
314 let response = self
315 .request(reqwest::Method::POST, "/leases/consolidate")
316 .json(&ConsolidateRequest {
317 lease_id,
318 fencing_token,
319 unspent,
320 requested,
321 needed,
322 ttl,
323 })
324 .send()
325 .await?;
326 if !response.status().is_success() {
327 return Err(read_problem(response).await);
328 }
329 read_allocation(response).await
330 }
331
332 async fn reclaim_expired_batch(
333 &self,
334 _now: Timestamp,
335 _limit: NonZeroUsize,
336 ) -> Result<ReclaimBatch, StoreError> {
337 Err(StoreError(
340 "reclaim is server-side; not exposed through the instance transport".into(),
341 ))
342 }
343}
344
345#[async_trait]
346impl SnapshotSource for HttpStore {
347 async fn snapshot(&self, principal: Principal) -> Result<SnapshotResolution, StoreError> {
348 let response = self
349 .request(reqwest::Method::GET, &format!("/snapshots/{principal}"))
350 .send()
351 .await?;
352 if response.status() == reqwest::StatusCode::NOT_FOUND {
353 let problem: Problem = response.json().await.map_err(|e| {
354 StoreError(format!(
355 "snapshot endpoint returned an unstructured 404, not a confirmed unknown principal: {}", e.without_url()
356 ))
357 })?;
358 if problem.code == "unknown-principal" {
359 return Ok(SnapshotResolution::Unknown);
360 }
361 return Err(StoreError(format!(
362 "snapshot endpoint returned 404 with code {} instead of unknown-principal",
363 problem.code
364 )));
365 }
366 if response.status() == reqwest::StatusCode::GONE {
367 let problem: Problem = response
368 .json()
369 .await
370 .map_err(|e| StoreError(format!("http: {}", e.without_url())))?;
371 if problem.code != "revoked-principal" {
372 return Err(StoreError(format!(
373 "server {}: {}",
374 problem.status, problem.title
375 )));
376 }
377 let generation = problem.generation.ok_or_else(|| {
378 StoreError("revoked-principal response omitted generation".into())
379 })?;
380 return Ok(SnapshotResolution::Revoked { generation });
381 }
382 if !response.status().is_success() {
383 return Err(StoreError(problem_detail(response).await));
384 }
385 let snapshot: AccountSnapshot = require_complete(response)?
386 .json()
387 .await
388 .map_err(|e| StoreError(format!("http: {}", e.without_url())))?;
389 let snapshot = PublishableSnapshot::try_new(Arc::new(snapshot))
390 .map_err(|error| StoreError(format!("invalid snapshot from server: {error}")))?;
391 Ok(SnapshotResolution::Present(snapshot))
392 }
393
394 fn subscribe(&self) -> broadcast::Receiver<SnapshotPush> {
395 let (sender, receiver) = broadcast::channel(1);
396 drop(sender);
397 receiver
398 }
399
400 async fn principals(&self) -> Result<Option<Vec<Principal>>, StoreError> {
408 let response = self
409 .request(reqwest::Method::GET, "/snapshots")
410 .send()
411 .await?;
412 if response.status() == reqwest::StatusCode::NOT_IMPLEMENTED {
413 return Ok(None);
414 }
415 if !response.status().is_success() {
416 return Err(StoreError(problem_detail(response).await));
417 }
418 let body: PrincipalsResponse = require_complete(response)?
419 .json()
420 .await
421 .map_err(|e| StoreError(format!("http: {}", e.without_url())))?;
422 Ok(Some(body.principals))
423 }
424}
425
426#[async_trait]
427impl UsageSink for HttpStore {
428 async fn ingest(
429 &self,
430 events: &[UsageEvent],
431 _now: Timestamp,
432 ) -> Result<IngestReport, IngestError> {
433 let response = self
434 .request(reqwest::Method::POST, "/usage/ingest")
435 .json(&IngestRequestRef { events })
436 .send()
437 .await?;
438 let status = response.status();
439 if !status.is_success() {
440 let detail = problem_detail(response).await;
441 let terminal = status.is_client_error()
450 && status != reqwest::StatusCode::REQUEST_TIMEOUT
451 && status != reqwest::StatusCode::TOO_MANY_REQUESTS
452 && status != reqwest::StatusCode::UNAUTHORIZED
453 && status != reqwest::StatusCode::FORBIDDEN;
454 let error = StoreError(detail);
455 return Err(if terminal {
456 IngestError::Refused(error)
457 } else {
458 IngestError::Unavailable(error)
459 });
460 }
461 let bytes = bounded_body(
462 require_complete(response)?,
463 tollgate_store::wire::MAX_INGEST_REPORT_BYTES,
464 )
465 .await?;
466 let report: IngestReport = serde_json::from_slice(&bytes)
467 .map_err(|_| StoreError("invalid usage acknowledgement JSON".into()))?;
468 report.validate(events.len())?;
469 Ok(report)
470 }
471}
472
473fn require_complete(response: reqwest::Response) -> Result<reqwest::Response, StoreError> {
475 if response.status() != reqwest::StatusCode::OK
476 || response
477 .headers()
478 .contains_key(reqwest::header::CONTENT_RANGE)
479 {
480 return Err(StoreError(
481 "control-plane response requires a complete HTTP 200 response".into(),
482 ));
483 }
484 Ok(response)
485}
486
487async fn bounded_body(
488 mut response: reqwest::Response,
489 limit: usize,
490) -> Result<Vec<u8>, StoreError> {
491 if response
492 .content_length()
493 .is_some_and(|length| length > limit as u64)
494 {
495 return Err(StoreError(
496 "control-plane body exceeds its wire limit".into(),
497 ));
498 }
499 let mut bytes = Vec::new();
500 while let Some(chunk) = response
501 .chunk()
502 .await
503 .map_err(|e| StoreError(format!("control-plane transport: {}", e.without_url())))?
504 {
505 if chunk.len() > limit.saturating_sub(bytes.len()) {
506 return Err(StoreError(
507 "control-plane body exceeds its wire limit".into(),
508 ));
509 }
510 bytes.extend_from_slice(&chunk);
511 }
512 Ok(bytes)
513}
514
515#[async_trait]
516impl tollgate_store::KeySource for HttpStore {
517 async fn active_keys_page(
518 &self,
519 _now: Timestamp,
520 after: Option<tollgate_core::KeyId>,
521 limit: NonZeroUsize,
522 ) -> Result<tollgate_store::KeyPage, StoreError> {
523 tollgate_store::validate_key_page_limit(limit)?;
524 let mut request = self.request(reqwest::Method::GET, "/keys");
525 request.request = request.request.query(&[("limit", limit.get().to_string())]);
526 if let Some(after) = after {
527 request.request = request.request.query(&[("after", after.to_string())]);
528 }
529 let response = request.send().await?;
530 let status = response.status();
531 if !status.is_success() {
532 let bytes = bounded_body(response, tollgate_store::wire::MAX_KEYS_BODY_BYTES).await?;
533 let code = serde_json::from_slice::<Problem>(&bytes)
534 .ok()
535 .map(|p| p.code);
536 let safe_code = match code.as_deref() {
537 Some("authentication-required") => "authentication-required",
538 Some("scope-forbidden") => "scope-forbidden",
539 Some("invalid-query") => "invalid-query",
540 Some("invalid-limit") => "invalid-limit",
541 _ => "credential-source-unavailable",
542 };
543 return Err(StoreError(format!(
544 "credential read refused: {status} {safe_code}"
545 )));
546 }
547 let bytes = bounded_body(
548 require_complete(response)?,
549 tollgate_store::wire::MAX_KEYS_BODY_BYTES,
550 )
551 .await?;
552 let body: tollgate_store::wire::KeysResponse = serde_json::from_slice(&bytes)
553 .map_err(|_| StoreError("invalid credential page body".into()))?;
554 tollgate_store::KeyPage::try_new(
555 body.revision,
556 body.as_of,
557 after,
558 limit,
559 body.keys,
560 body.next_after,
561 )
562 }
563}
564
565#[cfg(test)]
566mod exhaustion_tests {
567 use super::*;
568 use tollgate_core::LeaseGrant;
569
570 #[test]
571 fn incomplete_exhaustion_responses_never_become_authoritative() {
572 for body in [
573 r#"{"status":409,"code":"balance-exhausted","title":"empty"}"#,
574 r#"{"status":409,"code":"balance-exhausted","title":"empty","balance_exhaustion":{}}"#,
575 r#"{"status":503,"code":"balance-exhausted","title":"empty","balance_exhaustion":{"period_end":null}}"#,
576 ] {
577 if let Ok(problem) = serde_json::from_str::<Problem>(body) {
578 assert!(matches!(
579 problem_to_allocate(problem),
580 AllocateError::Storage(_)
581 ));
582 }
583 }
584 let old = serde_json::from_str::<Problem>(
585 r#"{"status":409,"code":"insufficient-balance","title":"empty"}"#,
586 )
587 .unwrap();
588 assert_eq!(problem_to_allocate(old), AllocateError::InsufficientBalance);
589 }
590
591 #[test]
594 fn plain_refusal_codes_round_trip() {
595 for (code, expected) in [
596 ("unknown-account", AllocateError::UnknownAccount),
597 ("account-inactive", AllocateError::AccountInactive),
598 ("insufficient-balance", AllocateError::InsufficientBalance),
599 ("balance-overflow", AllocateError::BalanceOverflow),
600 ("invalid-ttl", AllocateError::InvalidTtl),
601 ("unknown-lease", AllocateError::UnknownLease),
602 ("fenced", AllocateError::Fenced),
603 ("lease-not-active", AllocateError::LeaseNotActive),
604 ("invalid-release", AllocateError::InvalidRelease),
605 ] {
606 let problem = Problem {
607 status: 409,
608 code: code.into(),
609 title: code.into(),
610 generation: None,
611 balance_exhaustion: None,
612 balance_shortfall: None,
613 };
614 assert_eq!(problem_to_allocate(problem), expected, "{code}");
615 }
616 }
617
618 #[test]
619 fn incomplete_shortfall_responses_never_become_authoritative() {
620 for body in [
621 r#"{"status":409,"code":"insufficient-balance","title":"short","balance_shortfall":{"remaining":0,"period_end":null}}"#,
623 r#"{"status":503,"code":"insufficient-balance","title":"short","balance_shortfall":{"remaining":5,"period_end":null}}"#,
624 ] {
625 let problem = serde_json::from_str::<Problem>(body).unwrap();
626 assert_eq!(
627 problem_to_allocate(problem),
628 AllocateError::InsufficientBalance,
629 "{body}"
630 );
631 }
632 assert!(
634 serde_json::from_str::<Problem>(
635 r#"{"status":409,"code":"insufficient-balance","title":"short","balance_shortfall":{"remaining":5}}"#
636 )
637 .is_err()
638 );
639 let attested = serde_json::from_str::<Problem>(
640 r#"{"status":409,"code":"insufficient-balance","title":"short","balance_shortfall":{"remaining":5,"period_end":null}}"#,
641 )
642 .unwrap();
643 assert_eq!(
644 problem_to_allocate(attested),
645 AllocateError::BalanceInsufficient(tollgate_core::BalanceShortfall {
646 remaining: CostUnits(5),
647 period_end: None,
648 })
649 );
650 }
651
652 #[test]
655 fn a_grant_without_evidence_parses_as_unattested() {
656 let grant = LeaseGrant {
657 lease_id: LeaseId(1),
658 account_id: AccountId(1),
659 fencing_token: FencingToken(1),
660 units: CostUnits(10),
661 expires_at: jiff::Timestamp::UNIX_EPOCH,
662 };
663 let old = serde_json::to_string(&grant).unwrap();
664 let parsed: Allocation = serde_json::from_str(&old).unwrap();
665 assert_eq!(
666 parsed,
667 Allocation {
668 grant,
669 funding: None
670 }
671 );
672 }
673
674 #[test]
675 fn grant_evidence_below_the_grant_itself_is_discarded() {
676 let grant = LeaseGrant {
677 lease_id: LeaseId(1),
678 account_id: AccountId(1),
679 fencing_token: FencingToken(1),
680 units: CostUnits(10),
681 expires_at: jiff::Timestamp::UNIX_EPOCH,
682 };
683 let with = |remaining| Allocation {
684 grant,
685 funding: Some(tollgate_core::BalanceShortfall {
686 remaining: CostUnits(remaining),
687 period_end: None,
688 }),
689 };
690 assert_eq!(credible_evidence(with(9)).funding, None);
691 assert_eq!(credible_evidence(with(10)), with(10));
692 }
693}