pomelo-s3 0.12.2

S3-compatible ObjectSource for pomelo-data, reading objects over SigV4-presigned HTTPS GETs.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
//! Generic S3 `ObjectSource` for `pomelo-data` — reads objects from any
//! S3-compatible store (Cloudflare R2, AWS S3, MinIO, GCS-S3) over HTTPS with
//! SigV4-presigned GETs. Endpoint is configurable, so nothing here is
//! Cloudflare-specific; the container points it at R2, OSS users at anything.
//!
//! `pomelo-data` stays pure (trait + `LocalSource`); cloud access lives here.

use std::time::Duration;

use pomelo_data::error::DataError;
use pomelo_data::{ObjectLister, ObjectSink, ObjectSource};
use rusty_s3::actions::ListObjectsV2;
use rusty_s3::{Bucket, Credentials, S3Action, UrlStyle};

/// Cap on the XML `ListObjectsV2` response body per page — generous for even a
/// full 1000-key page of long keys.
const LIST_RESPONSE_LIMIT: u64 = 16 * 1024 * 1024;

/// How long a presigned GET URL stays valid — generous, single-shot use.
const SIGN_TTL: Duration = Duration::from_secs(300);

/// An `ObjectSource` backed by an S3-compatible bucket.
pub struct S3Source {
    bucket: Bucket,
    creds: Credentials,
}

impl S3Source {
    /// `endpoint` is the host root (no bucket), e.g.
    /// `https://<acct>.r2.cloudflarestorage.com`. Path-style addressing, so the
    /// bucket is appended to the path (works with custom endpoints).
    /// `session_token` is `Some` for **temporary** credentials (AWS STS / IAM
    /// role); it's signed into requests as `X-Amz-Security-Token`. R2 and other
    /// static-key stores pass `None`.
    pub fn new(
        endpoint: &str,
        bucket: &str,
        access_key: &str,
        secret_key: &str,
        session_token: Option<&str>,
        region: &str,
    ) -> Result<Self, DataError> {
        let url = endpoint
            .parse()
            .map_err(|e| DataError::Io(format!("bad S3 endpoint {endpoint:?}: {e}")))?;
        let bucket = Bucket::new(url, UrlStyle::Path, bucket.to_string(), region.to_string())
            .map_err(|e| DataError::Io(format!("bad S3 bucket: {e}")))?;
        let creds = match session_token {
            Some(t) => Credentials::new_with_token(access_key, secret_key, t),
            None => Credentials::new(access_key, secret_key),
        };
        Ok(Self { bucket, creds })
    }

    /// Build from `S3_ENDPOINT` / `S3_BUCKET` / `S3_ACCESS_KEY_ID` /
    /// `S3_SECRET_ACCESS_KEY` (+ optional `S3_SESSION_TOKEN` for temporary
    /// credentials, and `S3_REGION`, default `auto`).
    pub fn from_env() -> Result<Self, DataError> {
        let var = |k: &str| std::env::var(k).map_err(|_| DataError::Io(format!("missing env {k}")));
        let region = std::env::var("S3_REGION").unwrap_or_else(|_| "auto".to_string());
        let token = std::env::var("S3_SESSION_TOKEN").ok();
        Self::new(
            &var("S3_ENDPOINT")?,
            &var("S3_BUCKET")?,
            &var("S3_ACCESS_KEY_ID")?,
            &var("S3_SECRET_ACCESS_KEY")?,
            token.as_deref(),
            &region,
        )
    }
}

impl ObjectSource for S3Source {
    fn get(&self, key: &str) -> Result<Option<Vec<u8>>, DataError> {
        let url = self
            .bucket
            .get_object(Some(&self.creds), key)
            .sign(SIGN_TTL);
        match ureq::get(url.as_str()).call() {
            Ok(resp) => {
                // Per-symbol gzip files are small, but lift the default body cap
                // so a large universe file can't get rejected mid-load.
                let bytes = resp
                    .into_body()
                    .with_config()
                    .limit(256 * 1024 * 1024)
                    .read_to_vec()
                    .map_err(|e| DataError::Io(format!("read {key}: {e}")))?;
                Ok(Some(bytes))
            }
            Err(ureq::Error::StatusCode(404)) => Ok(None),
            Err(e) => Err(DataError::Io(format!("GET {key}: {e}"))),
        }
    }
}

impl ObjectSink for S3Source {
    fn put(&self, key: &str, bytes: &[u8]) -> Result<(), DataError> {
        let url = self
            .bucket
            .put_object(Some(&self.creds), key)
            .sign(SIGN_TTL);
        match ureq::put(url.as_str()).send(bytes) {
            Ok(_) => Ok(()),
            Err(e) => Err(DataError::Io(format!("PUT {key}: {e}"))),
        }
    }
}

impl ObjectLister for S3Source {
    /// `ListObjectsV2`, paginating on `NextContinuationToken` (S3 returns at
    /// most 1000 keys per page). Returns full keys (the `prefix` as stored),
    /// matching what [`ObjectSource::get`] expects.
    fn list(&self, prefix: &str) -> Result<Vec<String>, DataError> {
        let mut keys = Vec::new();
        let mut continuation_token: Option<String> = None;
        loop {
            let mut action = ListObjectsV2::new(&self.bucket, Some(&self.creds));
            action.with_prefix(prefix);
            if let Some(tok) = &continuation_token {
                action.with_continuation_token(tok.as_str());
            }
            let url = action.sign(SIGN_TTL);
            let body = match ureq::get(url.as_str()).call() {
                Ok(resp) => resp
                    .into_body()
                    .with_config()
                    .limit(LIST_RESPONSE_LIMIT)
                    .read_to_string()
                    .map_err(|e| DataError::Io(format!("read LIST {prefix}: {e}")))?,
                Err(e) => return Err(DataError::Io(format!("LIST {prefix}: {e}"))),
            };
            let parsed = ListObjectsV2::parse_response(&body)
                .map_err(|e| DataError::Io(format!("parse LIST {prefix} response: {e}")))?;
            keys.extend(parsed.contents.into_iter().map(|c| c.key));
            continuation_token = parsed.next_continuation_token;
            if continuation_token.is_none() {
                break;
            }
        }
        Ok(keys)
    }
}

/// Resolved S3/R2 connection details from the environment.
pub struct S3Conn {
    pub endpoint: String,
    pub access_key: String,
    pub secret_key: String,
    /// `Some` for temporary (AWS STS / IAM role) credentials.
    pub session_token: Option<String>,
    pub region: String,
}

/// Resolve S3/R2 credentials + endpoint from the environment, trying the `S3_`
/// prefix first (Cloudflare R2 / static keys) then `AWS_` (AWS S3, incl. IAM-role
/// temporary credentials via `AWS_SESSION_TOKEN`). Per prefix `P` it reads
/// `{P}ACCESS_KEY_ID`, `{P}SECRET_ACCESS_KEY`, optional `{P}SESSION_TOKEN`,
/// `{P}ENDPOINT`/`{P}ENDPOINT_URL` (AWS: derived from `{P}REGION` if unset), and
/// `{P}REGION` (default `auto`). `get` is injected so the chain is unit-testable.
/// A different deployment prefix should be mapped onto `S3_*`/`AWS_*` in the env.
pub fn resolve_s3_conn(get: impl Fn(&str) -> Option<String>) -> Result<S3Conn, String> {
    for p in ["S3_", "AWS_"] {
        let Some(access_key) = get(&format!("{p}ACCESS_KEY_ID")) else {
            continue;
        };
        let secret_key = get(&format!("{p}SECRET_ACCESS_KEY")).ok_or_else(|| {
            format!("{p}ACCESS_KEY_ID is set but {p}SECRET_ACCESS_KEY is missing")
        })?;
        let session_token = get(&format!("{p}SESSION_TOKEN"));
        let region = get(&format!("{p}REGION")).unwrap_or_else(|| "auto".to_string());
        let endpoint = get(&format!("{p}ENDPOINT"))
            .or_else(|| get(&format!("{p}ENDPOINT_URL")))
            .or_else(|| {
                // AWS with a real region but no explicit endpoint → derive it.
                (p == "AWS_" && region != "auto")
                    .then(|| format!("https://s3.{region}.amazonaws.com"))
            })
            .ok_or_else(|| {
                format!("set {p}ENDPOINT (R2: https://<acct>.r2.cloudflarestorage.com) or {p}REGION (AWS)")
            })?;
        return Ok(S3Conn {
            endpoint,
            access_key,
            secret_key,
            session_token,
            region,
        });
    }
    Err("no S3 credentials in env: set S3_ACCESS_KEY_ID + S3_SECRET_ACCESS_KEY (+ S3_ENDPOINT) for R2, or AWS_ACCESS_KEY_ID + AWS_SECRET_ACCESS_KEY (+ AWS_SESSION_TOKEN) for an AWS IAM role".into())
}

/// The destination for an `fmp-sync` run: a local path or an S3/R2 bucket. Both
/// implement [`ObjectSink`]/[`ObjectSource`], so `sync_into` writes a
/// byte-identical `data-layout.md` tree either way (the S3 variant prepends an
/// optional key prefix from the URL path).
pub enum OutStore {
    Local(pomelo_data::LocalSource),
    // Box the S3 client: it's much larger than the Local variant (clippy::large_enum_variant).
    S3 { src: Box<S3Source>, prefix: String },
}

impl OutStore {
    /// Parse `--out`: an `s3://bucket[/prefix]` URL builds an [`S3Source`]
    /// from the standard `$S3_*` env credentials; anything else is a local path.
    pub fn parse(out: &str) -> Result<Self, String> {
        let Some(rest) = out.strip_prefix("s3://") else {
            return Ok(OutStore::Local(pomelo_data::LocalSource::new(out)));
        };
        let (bucket, prefix) = match rest.split_once('/') {
            Some((b, p)) => (b, p.trim_matches('/')),
            None => (rest, ""),
        };
        if bucket.is_empty() {
            return Err("s3:// URL needs a bucket: s3://bucket[/prefix]".into());
        }
        let conn = resolve_s3_conn(|k| std::env::var(k).ok())?;
        let src = S3Source::new(
            &conn.endpoint,
            bucket,
            &conn.access_key,
            &conn.secret_key,
            conn.session_token.as_deref(),
            &conn.region,
        )
        .map_err(|e| e.to_string())?;
        Ok(OutStore::S3 {
            src: Box::new(src),
            prefix: prefix.to_string(),
        })
    }

    pub fn is_s3(&self) -> bool {
        matches!(self, OutStore::S3 { .. })
    }

    /// Prepend the S3 key prefix (if any) to a data-layout key.
    fn prefixed(prefix: &str, key: &str) -> String {
        if prefix.is_empty() {
            key.to_string()
        } else {
            format!("{prefix}/{key}")
        }
    }
}

impl ObjectSource for OutStore {
    fn get(&self, key: &str) -> Result<Option<Vec<u8>>, DataError> {
        match self {
            OutStore::Local(l) => l.get(key),
            OutStore::S3 { src, prefix } => src.get(&OutStore::prefixed(prefix, key)),
        }
    }
}

impl ObjectSink for OutStore {
    fn put(&self, key: &str, bytes: &[u8]) -> Result<(), DataError> {
        match self {
            OutStore::Local(l) => l.put(key, bytes),
            OutStore::S3 { src, prefix } => src.put(&OutStore::prefixed(prefix, key), bytes),
        }
    }
}

impl ObjectLister for OutStore {
    fn list(&self, prefix: &str) -> Result<Vec<String>, DataError> {
        match self {
            OutStore::Local(l) => l.list(prefix),
            OutStore::S3 {
                src,
                prefix: store_prefix,
            } => {
                let full_prefix = OutStore::prefixed(store_prefix, prefix);
                let keys = src.list(&full_prefix)?;
                // Strip the store's own key prefix back off, so callers see the
                // same data-layout-relative keys as with a local/unprefixed store.
                let strip_from = if store_prefix.is_empty() {
                    0
                } else {
                    store_prefix.len() + 1
                };
                Ok(keys
                    .into_iter()
                    .map(|k| k.get(strip_from..).unwrap_or(&k).to_string())
                    .collect())
            }
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use std::fs;
    use std::io::{Read, Write};
    use std::net::TcpListener;
    use std::thread;

    /// Minimal stub S3: serves "hi" for any key, 404 for keys containing
    /// "missing". Handles exactly `n` connections then exits.
    fn spawn_stub(n: usize) -> String {
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        thread::spawn(move || {
            for _ in 0..n {
                let (mut sock, _) = listener.accept().unwrap();
                let mut buf = [0u8; 2048];
                let read = sock.read(&mut buf).unwrap();
                let line = String::from_utf8_lossy(&buf[..read]);
                let first = line.lines().next().unwrap_or("");
                let resp = if first.contains("missing") {
                    "HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
                        .to_string()
                } else {
                    "HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nhi"
                        .to_string()
                };
                sock.write_all(resp.as_bytes()).unwrap();
            }
        });
        format!("http://{addr}")
    }

    /// Serves `body` once with a `Content-Encoding: gzip` header — to prove the
    /// client does NOT transparently decompress (objects are `.csv.gz` and the
    /// loaders gunzip them; auto-decompress would corrupt that).
    fn spawn_encoded_stub(body: Vec<u8>) -> String {
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        thread::spawn(move || {
            let (mut sock, _) = listener.accept().unwrap();
            let mut buf = [0u8; 2048];
            let _ = sock.read(&mut buf).unwrap();
            let head = format!(
                "HTTP/1.1 200 OK\r\nContent-Encoding: gzip\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
                body.len()
            );
            sock.write_all(head.as_bytes()).unwrap();
            sock.write_all(&body).unwrap();
        });
        format!("http://{addr}")
    }

    /// A `ListObjectsV2` XML response body for `keys`, with `next_token` set
    /// when the caller should page again.
    fn list_objects_v2_xml(keys: &[String], next_token: Option<&str>) -> String {
        let contents: String = keys
            .iter()
            .map(|k| {
                format!(
                    "<Contents><Key>{k}</Key><LastModified>2020-01-01T00:00:00.000Z</LastModified>\
                     <ETag>\"e\"</ETag><Size>1</Size><StorageClass>STANDARD</StorageClass></Contents>"
                )
            })
            .collect();
        let token = next_token
            .map(|t| format!("<NextContinuationToken>{t}</NextContinuationToken>"))
            .unwrap_or_default();
        format!(
            "<?xml version=\"1.0\" encoding=\"UTF-8\"?>\
             <ListBucketResult xmlns=\"http://s3.amazonaws.com/doc/2006-03-01/\">\
             <Name>bucket</Name><Prefix>prices/</Prefix><KeyCount>{}</KeyCount>\
             <MaxKeys>1000</MaxKeys><IsTruncated>{}</IsTruncated>{contents}{token}\
             <EncodingType>url</EncodingType></ListBucketResult>",
            keys.len(),
            next_token.is_some(),
        )
    }

    /// Serves each of `bodies` in order, one per accepted connection, as a
    /// `200 OK` response. Used to stub successive `ListObjectsV2` pages.
    fn spawn_body_pages_stub(bodies: Vec<String>) -> String {
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        thread::spawn(move || {
            for body in bodies {
                let (mut sock, _) = listener.accept().unwrap();
                let mut buf = [0u8; 4096];
                let _ = sock.read(&mut buf).unwrap();
                let resp = format!(
                    "HTTP/1.1 200 OK\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
                    body.len()
                );
                sock.write_all(resp.as_bytes()).unwrap();
            }
        });
        format!("http://{addr}")
    }

    #[test]
    fn get_returns_bytes_and_none_on_404() {
        let endpoint = spawn_stub(2);
        let src = S3Source::new(&endpoint, "bucket", "ak", "sk", None, "auto").unwrap();
        assert_eq!(src.get("prices/AAPL.csv.gz").unwrap(), Some(b"hi".to_vec()));
        assert_eq!(src.get("prices/missing.csv.gz").unwrap(), None);
    }

    #[test]
    fn get_does_not_decompress_content_encoding_gzip() {
        use flate2::write::GzEncoder;
        use flate2::Compression;
        let mut enc = GzEncoder::new(Vec::new(), Compression::default());
        enc.write_all(b"hello,world\n").unwrap();
        let gz = enc.finish().unwrap();

        let endpoint = spawn_encoded_stub(gz.clone());
        let src = S3Source::new(&endpoint, "bucket", "ak", "sk", None, "auto").unwrap();
        let got = src.get("fundamentals/AAA.csv.gz").unwrap().unwrap();
        // Raw gzip bytes must come back verbatim (magic 0x1f 0x8b intact), NOT decompressed.
        assert_eq!(got, gz, "S3Source must return raw stored bytes");
        assert_eq!(&got[..2], &[0x1f, 0x8b]);
    }

    #[test]
    fn put_returns_ok_on_2xx() {
        let endpoint = spawn_stub(1); // existing stub: 200 for non-"missing" keys
        let src = S3Source::new(&endpoint, "bucket", "ak", "sk", None, "auto").unwrap();
        src.put("panels/close.csv.gz", b"gzip-bytes").unwrap();
    }

    #[test]
    fn new_rejects_a_malformed_endpoint() {
        let err = match S3Source::new("not a url", "bucket", "ak", "sk", None, "auto") {
            Err(e) => e,
            Ok(_) => panic!("expected a malformed-endpoint error"),
        };
        assert!(matches!(err, DataError::Io(_)));
    }

    /// A bound-then-closed port: connections are refused, so both GET and PUT hit
    /// the non-404 error arms (transport error, not a status code).
    fn dead_endpoint() -> String {
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        drop(listener); // free the port so connects are refused
        format!("http://{addr}")
    }

    #[test]
    fn get_and_put_surface_transport_errors() {
        let src = S3Source::new(&dead_endpoint(), "bucket", "ak", "sk", None, "auto").unwrap();
        assert!(matches!(
            src.get("prices/AAA.csv.gz"),
            Err(DataError::Io(_))
        ));
        assert!(matches!(
            src.put("panels/close.csv.gz", b"x"),
            Err(DataError::Io(_))
        ));
    }

    #[test]
    fn from_env_reads_vars_and_reports_missing() {
        // These S3_* vars are used by no other test in this crate, so mutating the
        // process environment here is safe within this test binary.
        for k in [
            "S3_ENDPOINT",
            "S3_BUCKET",
            "S3_ACCESS_KEY_ID",
            "S3_SECRET_ACCESS_KEY",
            "S3_REGION",
        ] {
            std::env::remove_var(k);
        }
        // Missing required var → Err naming the var.
        assert!(matches!(
            S3Source::from_env(),
            Err(DataError::Io(ref m)) if m.contains("S3_ENDPOINT")
        ));

        std::env::set_var("S3_ENDPOINT", "https://example.r2.cloudflarestorage.com");
        std::env::set_var("S3_BUCKET", "bucket");
        std::env::set_var("S3_ACCESS_KEY_ID", "ak");
        std::env::set_var("S3_SECRET_ACCESS_KEY", "sk");
        // S3_REGION left unset → defaults to "auto".
        S3Source::from_env().expect("from_env builds when all required vars are present");

        for k in [
            "S3_ENDPOINT",
            "S3_BUCKET",
            "S3_ACCESS_KEY_ID",
            "S3_SECRET_ACCESS_KEY",
        ] {
            std::env::remove_var(k);
        }
    }

    #[test]
    fn parse_local_vs_s3_and_key_prefixing() {
        // Non-s3 → local path (no env needed).
        assert!(!OutStore::parse("./mydata").unwrap().is_s3());
        assert!(!OutStore::parse("/tmp/x").unwrap().is_s3());
        // Key prefixing: empty prefix is a passthrough; a prefix joins with '/'.
        assert_eq!(
            OutStore::prefixed("", "prices/AAPL.csv.gz"),
            "prices/AAPL.csv.gz"
        );
        assert_eq!(
            OutStore::prefixed("mirror/v1", "panels/piotroski_score.csv.gz"),
            "mirror/v1/panels/piotroski_score.csv.gz"
        );
    }

    /// Build an env lookup over a fixed table (no real env — deterministic under
    /// parallel tests).
    fn env_of(pairs: &[(&str, &str)]) -> impl Fn(&str) -> Option<String> {
        let owned: Vec<(String, String)> = pairs
            .iter()
            .map(|(k, v)| (k.to_string(), v.to_string()))
            .collect();
        move |k: &str| owned.iter().find(|(kk, _)| kk == k).map(|(_, v)| v.clone())
    }

    #[test]
    fn resolve_s3_conn_prefers_s3_then_falls_back_to_aws() {
        // S3_ wins when present (R2 static keys, explicit endpoint, no token).
        let c = resolve_s3_conn(env_of(&[
            ("S3_ENDPOINT", "https://acct.r2.cloudflarestorage.com"),
            ("S3_ACCESS_KEY_ID", "r2ak"),
            ("S3_SECRET_ACCESS_KEY", "r2sk"),
            ("AWS_ACCESS_KEY_ID", "awsak"),
            ("AWS_SECRET_ACCESS_KEY", "awssk"),
        ]))
        .unwrap();
        assert_eq!(c.access_key, "r2ak");
        assert_eq!(c.endpoint, "https://acct.r2.cloudflarestorage.com");
        assert_eq!(c.region, "auto");
        assert!(c.session_token.is_none());

        // AWS_ fallback: IAM temp creds (session token) + endpoint derived from region.
        let c = resolve_s3_conn(env_of(&[
            ("AWS_ACCESS_KEY_ID", "ASIAEXAMPLE"),
            ("AWS_SECRET_ACCESS_KEY", "sk"),
            ("AWS_SESSION_TOKEN", "tok"),
            ("AWS_REGION", "us-east-1"),
        ]))
        .unwrap();
        assert_eq!(c.access_key, "ASIAEXAMPLE");
        assert_eq!(c.session_token.as_deref(), Some("tok"));
        assert_eq!(c.endpoint, "https://s3.us-east-1.amazonaws.com");

        // Access key without its secret → a clear error.
        assert!(resolve_s3_conn(env_of(&[("S3_ACCESS_KEY_ID", "x")])).is_err());
        // Nothing set → error.
        assert!(resolve_s3_conn(env_of(&[])).is_err());
    }

    #[test]
    fn list_paginates_past_a_thousand_keys_via_continuation_token() {
        // First page: a full 1000-key page (S3's per-page cap) + a continuation
        // token; second page: the remainder, no token. `list` must concatenate
        // both pages into a single result.
        let page1: Vec<String> = (0..1000)
            .map(|i| format!("prices/SYM{i:04}.csv.gz"))
            .collect();
        let page2: Vec<String> = (1000..1200)
            .map(|i| format!("prices/SYM{i:04}.csv.gz"))
            .collect();
        let endpoint = spawn_body_pages_stub(vec![
            list_objects_v2_xml(&page1, Some("tok1")),
            list_objects_v2_xml(&page2, None),
        ]);
        let src = S3Source::new(&endpoint, "bucket", "ak", "sk", None, "auto").unwrap();
        let keys = src.list("prices").unwrap();
        assert_eq!(keys.len(), 1200);
        assert_eq!(keys[0], "prices/SYM0000.csv.gz");
        assert_eq!(keys[999], "prices/SYM0999.csv.gz");
        assert_eq!(keys[1199], "prices/SYM1199.csv.gz");
    }

    #[test]
    fn list_returns_empty_on_no_contents() {
        let endpoint = spawn_body_pages_stub(vec![list_objects_v2_xml(&[], None)]);
        let src = S3Source::new(&endpoint, "bucket", "ak", "sk", None, "auto").unwrap();
        assert_eq!(src.list("panels").unwrap(), Vec::<String>::new());
    }

    #[test]
    fn list_surfaces_transport_errors() {
        let src = S3Source::new(&dead_endpoint(), "bucket", "ak", "sk", None, "auto").unwrap();
        assert!(matches!(src.list("prices"), Err(DataError::Io(_))));
    }

    #[test]
    fn outstore_list_strips_the_store_prefix_from_s3_keys() {
        let body = list_objects_v2_xml(&["mirror/v1/prices/AAPL.csv.gz".to_string()], None);
        let endpoint = spawn_body_pages_stub(vec![body]);
        let src = S3Source::new(&endpoint, "bucket", "ak", "sk", None, "auto").unwrap();
        let out = OutStore::S3 {
            src: Box::new(src),
            prefix: "mirror/v1".to_string(),
        };
        assert_eq!(
            out.list("prices").unwrap(),
            vec!["prices/AAPL.csv.gz".to_string()]
        );
    }

    #[test]
    fn outstore_list_delegates_to_local_source() {
        let dir = std::env::temp_dir().join("pomelo_s3_outstore_list_test");
        let _ = fs::remove_dir_all(&dir);
        fs::create_dir_all(dir.join("prices")).unwrap();
        fs::write(dir.join("prices/AAPL.csv.gz"), b"x").unwrap();
        let out = OutStore::parse(dir.to_str().unwrap()).unwrap();
        assert_eq!(
            out.list("prices").unwrap(),
            vec!["prices/AAPL.csv.gz".to_string()]
        );
    }
}