Skip to main content

flatland_client_lib/
asset_backend.rs

1//! Pluggable published-asset backends: local disk, GCS/Firebase, or S3-compatible.
2
3use std::path::{Path, PathBuf};
4
5use anyhow::Context;
6use async_trait::async_trait;
7use reqwest::header::{AUTHORIZATION, CONTENT_LENGTH, CONTENT_TYPE};
8
9use crate::assets::{firebase_download_url, urlencoding_encode_path, DEFAULT_FIREBASE_BUCKET};
10
11/// How published client/sim packs are stored and fetched.
12#[derive(Debug, Clone, Copy, PartialEq, Eq)]
13pub enum AssetBackendKind {
14    /// Files under `FLATLAND_ASSETS_LOCAL_ROOT` (single-machine / OSS laptop).
15    Local,
16    /// Google Cloud Storage / Firebase Storage JSON + media APIs.
17    Gcs,
18    /// S3-compatible object store (AWS, MinIO, R2) via path-style HTTPS.
19    S3,
20}
21
22impl AssetBackendKind {
23    pub fn parse(raw: &str) -> anyhow::Result<Self> {
24        match raw.trim().to_ascii_lowercase().as_str() {
25            "local" | "file" | "disk" => Ok(Self::Local),
26            "gcs" | "firebase" | "gs" => Ok(Self::Gcs),
27            "s3" | "minio" | "r2" => Ok(Self::S3),
28            other => {
29                anyhow::bail!("unknown FLATLAND_ASSETS_BACKEND={other:?} (expected local|gcs|s3)")
30            }
31        }
32    }
33}
34
35#[derive(Debug, Clone)]
36pub struct AssetBackendConfig {
37    pub kind: AssetBackendKind,
38    /// GCS/Firebase or S3 bucket name.
39    pub bucket: String,
40    /// Object key prefix for **client** gfx packs (default `flatland3/client-assets`).
41    pub client_prefix: String,
42    /// Object key prefix for **sim** content packs (default `flatland3/sim-content`).
43    pub sim_prefix: String,
44    /// Local published root when `kind == Local`.
45    pub local_root: PathBuf,
46    /// S3 region (default `us-east-1`).
47    pub s3_region: String,
48    /// Optional S3 endpoint (MinIO/R2), e.g. `https://minio.example:9000`.
49    pub s3_endpoint: Option<String>,
50    /// Override for client `latest.json` fetch URL.
51    pub index_url_override: Option<String>,
52}
53
54impl Default for AssetBackendConfig {
55    fn default() -> Self {
56        Self::from_env()
57    }
58}
59
60impl AssetBackendConfig {
61    pub fn from_env() -> Self {
62        let kind = std::env::var("FLATLAND_ASSETS_BACKEND")
63            .ok()
64            .and_then(|s| AssetBackendKind::parse(&s).ok())
65            .unwrap_or(AssetBackendKind::Gcs);
66        let bucket = std::env::var("FLATLAND_ASSETS_BUCKET")
67            .unwrap_or_else(|_| DEFAULT_FIREBASE_BUCKET.to_string());
68        let client_prefix = std::env::var("FLATLAND_ASSETS_PREFIX")
69            .unwrap_or_else(|_| "flatland3/client-assets".to_string());
70        let sim_prefix = std::env::var("FLATLAND_SIM_ASSETS_PREFIX")
71            .unwrap_or_else(|_| "flatland3/sim-content".to_string());
72        let local_root = std::env::var("FLATLAND_ASSETS_LOCAL_ROOT")
73            .map(PathBuf::from)
74            .unwrap_or_else(|_| {
75                dirs::home_dir()
76                    .unwrap_or_else(|| PathBuf::from("."))
77                    .join(".flatland3")
78                    .join("published")
79            });
80        let s3_region =
81            std::env::var("FLATLAND_ASSETS_S3_REGION").unwrap_or_else(|_| "us-east-1".to_string());
82        let s3_endpoint = std::env::var("FLATLAND_ASSETS_S3_ENDPOINT")
83            .ok()
84            .filter(|s| !s.is_empty());
85        let index_url_override = std::env::var("FLATLAND_ASSETS_INDEX_URL")
86            .ok()
87            .filter(|s| !s.is_empty());
88        Self {
89            kind,
90            bucket,
91            client_prefix,
92            sim_prefix,
93            local_root,
94            s3_region,
95            s3_endpoint,
96            index_url_override,
97        }
98    }
99
100    pub fn client_index_object(&self) -> String {
101        format!("{}/latest.json", self.client_prefix.trim_end_matches('/'))
102    }
103
104    pub fn sim_index_object(&self) -> String {
105        format!("{}/latest.json", self.sim_prefix.trim_end_matches('/'))
106    }
107
108    /// HTTPS (or `file://`) URL for the client `latest.json` index.
109    pub fn client_index_url(&self) -> String {
110        if let Some(url) = &self.index_url_override {
111            return url.clone();
112        }
113        self.object_url(&self.client_index_object())
114    }
115
116    pub fn sim_index_url(&self) -> String {
117        self.object_url(&self.sim_index_object())
118    }
119
120    pub fn client_object_key(&self, publish_rev: u64, relative: &str) -> String {
121        format!(
122            "{}/rev-{}/{}",
123            self.client_prefix.trim_end_matches('/'),
124            publish_rev,
125            relative
126        )
127    }
128
129    pub fn sim_object_key(&self, publish_rev: u64, relative: &str) -> String {
130        format!(
131            "{}/rev-{}/{}",
132            self.sim_prefix.trim_end_matches('/'),
133            publish_rev,
134            relative
135        )
136    }
137
138    /// Public fetch URL for an object key.
139    pub fn object_url(&self, object_key: &str) -> String {
140        match self.kind {
141            AssetBackendKind::Local => {
142                let path = self.local_root.join(object_key);
143                format!("file://{}", path.display())
144            }
145            AssetBackendKind::Gcs => firebase_download_url(&self.bucket, object_key),
146            AssetBackendKind::S3 => self.s3_object_url(object_key),
147        }
148    }
149
150    fn s3_object_url(&self, object_key: &str) -> String {
151        let key = object_key.trim_start_matches('/');
152        if let Some(endpoint) = &self.s3_endpoint {
153            let base = endpoint.trim_end_matches('/');
154            format!("{base}/{}/{key}", self.bucket)
155        } else {
156            format!(
157                "https://{}.s3.{}.amazonaws.com/{key}",
158                self.bucket, self.s3_region
159            )
160        }
161    }
162
163    pub fn local_object_path(&self, object_key: &str) -> PathBuf {
164        self.local_root.join(object_key)
165    }
166}
167
168/// Put/get published bytes (admin upload + worker/client download).
169#[async_trait]
170pub trait AssetStore: Send + Sync {
171    async fn put(&self, object_key: &str, content_type: &str, bytes: &[u8]) -> anyhow::Result<()>;
172
173    async fn get(&self, object_key: &str) -> anyhow::Result<Vec<u8>>;
174}
175
176pub fn store_from_config(cfg: &AssetBackendConfig) -> anyhow::Result<Box<dyn AssetStore>> {
177    match cfg.kind {
178        AssetBackendKind::Local => Ok(Box::new(LocalAssetStore {
179            root: cfg.local_root.clone(),
180        })),
181        AssetBackendKind::Gcs => Ok(Box::new(GcsAssetStore {
182            bucket: cfg.bucket.clone(),
183            client: reqwest::Client::new(),
184        })),
185        AssetBackendKind::S3 => Ok(Box::new(S3AssetStore {
186            bucket: cfg.bucket.clone(),
187            region: cfg.s3_region.clone(),
188            endpoint: cfg.s3_endpoint.clone(),
189            client: reqwest::Client::new(),
190        })),
191    }
192}
193
194pub struct LocalAssetStore {
195    pub root: PathBuf,
196}
197
198#[async_trait]
199impl AssetStore for LocalAssetStore {
200    async fn put(&self, object_key: &str, _content_type: &str, bytes: &[u8]) -> anyhow::Result<()> {
201        let path = self.root.join(object_key);
202        if let Some(parent) = path.parent() {
203            std::fs::create_dir_all(parent)
204                .with_context(|| format!("mkdir {}", parent.display()))?;
205        }
206        std::fs::write(&path, bytes).with_context(|| format!("write {}", path.display()))?;
207        Ok(())
208    }
209
210    async fn get(&self, object_key: &str) -> anyhow::Result<Vec<u8>> {
211        let path = self.root.join(object_key);
212        std::fs::read(&path).with_context(|| format!("read {}", path.display()))
213    }
214}
215
216pub struct GcsAssetStore {
217    pub bucket: String,
218    pub client: reqwest::Client,
219}
220
221#[async_trait]
222impl AssetStore for GcsAssetStore {
223    async fn put(&self, object_key: &str, content_type: &str, bytes: &[u8]) -> anyhow::Result<()> {
224        let token = gcs_upload_bearer_token().await?;
225        let url = format!(
226            "https://storage.googleapis.com/upload/storage/v1/b/{}/o?uploadType=media&name={}",
227            self.bucket,
228            urlencoding_encode_path(object_key)
229        );
230        let response = self
231            .client
232            .post(&url)
233            .header(AUTHORIZATION, format!("Bearer {token}"))
234            .header(CONTENT_TYPE, content_type)
235            .header(CONTENT_LENGTH, bytes.len())
236            .body(bytes.to_vec())
237            .send()
238            .await?;
239        if !response.status().is_success() {
240            let status = response.status();
241            let body = response.text().await.unwrap_or_default();
242            anyhow::bail!("GCS upload {object_key} failed: {status} {body}");
243        }
244        Ok(())
245    }
246
247    async fn get(&self, object_key: &str) -> anyhow::Result<Vec<u8>> {
248        let url = firebase_download_url(&self.bucket, object_key);
249        let response = self.client.get(&url).send().await?.error_for_status()?;
250        Ok(response.bytes().await?.to_vec())
251    }
252}
253
254pub struct S3AssetStore {
255    pub bucket: String,
256    pub region: String,
257    pub endpoint: Option<String>,
258    pub client: reqwest::Client,
259}
260
261#[async_trait]
262impl AssetStore for S3AssetStore {
263    async fn put(&self, object_key: &str, content_type: &str, bytes: &[u8]) -> anyhow::Result<()> {
264        // Prefer AWS CLI when present — avoids embedding SigV4 in the published client crate.
265        if which_aws_cli() {
266            return s3_cli_put(
267                &self.bucket,
268                object_key,
269                content_type,
270                bytes,
271                self.endpoint.as_deref(),
272                &self.region,
273            )
274            .await;
275        }
276        anyhow::bail!(
277            "S3 upload requires the `aws` CLI (aws s3 cp) or set FLATLAND_ASSETS_BACKEND=gcs|local"
278        )
279    }
280
281    async fn get(&self, object_key: &str) -> anyhow::Result<Vec<u8>> {
282        let url = if let Some(endpoint) = &self.endpoint {
283            format!(
284                "{}/{}/{}",
285                endpoint.trim_end_matches('/'),
286                self.bucket,
287                object_key
288            )
289        } else {
290            format!(
291                "https://{}.s3.{}.amazonaws.com/{}",
292                self.bucket, self.region, object_key
293            )
294        };
295        let response = self.client.get(&url).send().await?.error_for_status()?;
296        Ok(response.bytes().await?.to_vec())
297    }
298}
299
300fn which_aws_cli() -> bool {
301    std::process::Command::new("aws")
302        .arg("--version")
303        .output()
304        .map(|o| o.status.success())
305        .unwrap_or(false)
306}
307
308async fn s3_cli_put(
309    bucket: &str,
310    object_key: &str,
311    content_type: &str,
312    bytes: &[u8],
313    endpoint: Option<&str>,
314    region: &str,
315) -> anyhow::Result<()> {
316    let tmp = tempfile_path(object_key)?;
317    if let Some(parent) = tmp.parent() {
318        std::fs::create_dir_all(parent)?;
319    }
320    std::fs::write(&tmp, bytes)?;
321    let uri = format!("s3://{bucket}/{object_key}");
322    let mut cmd = tokio::process::Command::new("aws");
323    cmd.args(["s3", "cp", tmp.to_str().unwrap(), &uri]);
324    cmd.args(["--content-type", content_type]);
325    cmd.args(["--region", region]);
326    if let Some(ep) = endpoint {
327        cmd.args(["--endpoint-url", ep]);
328    }
329    let output = cmd.output().await?;
330    let _ = std::fs::remove_file(&tmp);
331    if !output.status.success() {
332        anyhow::bail!(
333            "aws s3 cp failed: {}",
334            String::from_utf8_lossy(&output.stderr)
335        );
336    }
337    Ok(())
338}
339
340fn tempfile_path(object_key: &str) -> anyhow::Result<PathBuf> {
341    let name = object_key.replace('/', "_");
342    let dir = std::env::temp_dir().join("flatland-asset-upload");
343    std::fs::create_dir_all(&dir)?;
344    Ok(dir.join(name))
345}
346
347async fn gcs_upload_bearer_token() -> anyhow::Result<String> {
348    if let Ok(token) = std::env::var("FLATLAND_ASSETS_UPLOAD_TOKEN") {
349        if !token.is_empty() {
350            return Ok(token);
351        }
352    }
353    let output = tokio::process::Command::new("gcloud")
354        .args(["auth", "application-default", "print-access-token"])
355        .output()
356        .await
357        .map_err(|err| {
358            anyhow::anyhow!(
359                "gcloud ADC token failed ({err}); set FLATLAND_ASSETS_UPLOAD_TOKEN or run gcloud auth application-default login"
360            )
361        })?;
362    if !output.status.success() {
363        anyhow::bail!(
364            "gcloud auth application-default print-access-token failed: {}",
365            String::from_utf8_lossy(&output.stderr)
366        );
367    }
368    Ok(String::from_utf8(output.stdout)?.trim().to_string())
369}
370
371pub fn mime_for_path(path: &Path) -> &'static str {
372    match path
373        .extension()
374        .and_then(|e| e.to_str())
375        .unwrap_or("")
376        .to_ascii_lowercase()
377        .as_str()
378    {
379        "png" => "image/png",
380        "jpg" | "jpeg" => "image/jpeg",
381        "webp" => "image/webp",
382        "mp3" => "audio/mpeg",
383        "wav" => "audio/wav",
384        "ogg" => "audio/ogg",
385        "json" => "application/json",
386        "yaml" | "yml" => "application/x-yaml",
387        _ => "application/octet-stream",
388    }
389}
390
391#[cfg(test)]
392mod tests {
393    use super::*;
394
395    #[test]
396    fn parse_backend_kinds() {
397        assert_eq!(
398            AssetBackendKind::parse("firebase").unwrap(),
399            AssetBackendKind::Gcs
400        );
401        assert_eq!(
402            AssetBackendKind::parse("local").unwrap(),
403            AssetBackendKind::Local
404        );
405        assert_eq!(
406            AssetBackendKind::parse("minio").unwrap(),
407            AssetBackendKind::S3
408        );
409    }
410
411    #[test]
412    fn gcs_object_urls() {
413        let cfg = AssetBackendConfig {
414            kind: AssetBackendKind::Gcs,
415            bucket: "flatland-8911e.appspot.com".into(),
416            client_prefix: "flatland3/client-assets".into(),
417            sim_prefix: "flatland3/sim-content".into(),
418            local_root: PathBuf::from("/tmp"),
419            s3_region: "us-east-1".into(),
420            s3_endpoint: None,
421            index_url_override: None,
422        };
423        let url = cfg.client_index_url();
424        assert!(url.contains("flatland-8911e.appspot.com"));
425        assert!(url.contains("latest.json"));
426        assert_eq!(
427            cfg.client_index_object(),
428            "flatland3/client-assets/latest.json"
429        );
430    }
431}