lumen_server/service/
artifact.rs1use std::{sync::Arc, time::Duration};
2
3use super::{ArtifactRef, ArtifactStore, ArtifactWrite, BoxFuture, ServiceError, ServiceResult};
4
5pub trait PresignedUrlResolver: Send + Sync {
6 fn resolve<'a>(&'a self, artifact: &'a ArtifactWrite) -> BoxFuture<'a, ServiceResult<String>>;
7}
8
9#[derive(Debug, Clone)]
10pub struct StaticPresignedUrlResolver {
11 upload_url: String,
12}
13
14impl StaticPresignedUrlResolver {
15 pub fn new(upload_url: impl Into<String>) -> Self {
16 Self {
17 upload_url: upload_url.into(),
18 }
19 }
20}
21
22impl PresignedUrlResolver for StaticPresignedUrlResolver {
23 fn resolve<'a>(&'a self, _artifact: &'a ArtifactWrite) -> BoxFuture<'a, ServiceResult<String>> {
24 Box::pin(async move { Ok(self.upload_url.clone()) })
25 }
26}
27
28#[derive(Clone)]
29pub struct PresignedUrlArtifactStore<R> {
30 resolver: Arc<R>,
31 artifact_uri_prefix: String,
32 bearer_token: Option<String>,
33 timeout: Duration,
34}
35
36impl<R> PresignedUrlArtifactStore<R>
37where
38 R: PresignedUrlResolver,
39{
40 pub fn new(resolver: R, artifact_uri_prefix: impl Into<String>) -> Self {
41 Self {
42 resolver: Arc::new(resolver),
43 artifact_uri_prefix: artifact_uri_prefix.into(),
44 bearer_token: None,
45 timeout: Duration::from_secs(60),
46 }
47 }
48
49 pub fn with_bearer_token(mut self, bearer_token: impl Into<String>) -> Self {
50 self.bearer_token = Some(bearer_token.into());
51 self
52 }
53
54 pub fn with_timeout(mut self, timeout: Duration) -> Self {
55 self.timeout = timeout;
56 self
57 }
58}
59
60impl<R> ArtifactStore for PresignedUrlArtifactStore<R>
61where
62 R: PresignedUrlResolver + 'static,
63{
64 fn put<'a>(&'a self, artifact: ArtifactWrite) -> BoxFuture<'a, ServiceResult<ArtifactRef>> {
65 Box::pin(async move {
66 let upload_url = self.resolver.resolve(&artifact).await?;
67 let mut request = reqwest::Client::new()
68 .put(&upload_url)
69 .header(reqwest::header::CONTENT_TYPE, artifact.content_type)
70 .body(artifact.bytes.clone())
71 .timeout(self.timeout);
72 if let Some(token) = self.bearer_token.as_deref() {
73 request = request.bearer_auth(token);
74 }
75 let response = request.send().await.map_err(|err| ServiceError {
76 code: "artifact_upload_failed",
77 message: err.to_string(),
78 retryable: true,
79 })?;
80 if !response.status().is_success() {
81 return Err(ServiceError {
82 code: "artifact_upload_failed",
83 message: format!("artifact upload returned status {}", response.status()),
84 retryable: response.status().is_server_error()
85 || response.status().as_u16() == 429,
86 });
87 }
88
89 let id = artifact.job_id.0;
90 Ok(ArtifactRef {
91 id: id.clone(),
92 content_type: artifact.content_type,
93 bytes: artifact.bytes.len(),
94 uri: format!("{}/{}", self.artifact_uri_prefix.trim_end_matches('/'), id),
95 })
96 })
97 }
98}