Skip to main content

lumen_server/service/
artifact.rs

1use 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}