Skip to main content

rc_core/admin/
on_demand_migration.rs

1//! On-Demand Migration administration, independent of HTTP transport.
2//!
3//! A bucket names an external S3-compatible source bucket. A GET that misses
4//! locally is served from the source and stored locally, and a background
5//! backfill job pulls the remainder. The wire contract is pinned by the
6//! vendored fixtures in `tests/fixtures/on_demand_migration/`.
7//!
8//! Two shapes live here on purpose. [`OnDemandMigrationConfigRequest`] is what
9//! `rc` sends: it carries the plaintext secret in zeroizing storage and is
10//! serialized once, into a zeroizing buffer. [`OnDemandMigrationConfigView`] is
11//! what the server returns: every field is optional with the server's own
12//! default, so an older server that omits a field still parses, and the
13//! credential values are dropped during deserialization so nothing downstream
14//! can print them.
15
16use crate::{Error, Result};
17use async_trait::async_trait;
18use serde::de::Deserializer;
19use serde::ser::{SerializeStruct, Serializer};
20use serde::{Deserialize, Serialize};
21use std::collections::BTreeMap;
22use zeroize::{Zeroize, Zeroizing};
23
24/// Capability label used in unsupported-feature diagnostics.
25pub const ON_DEMAND_MIGRATION_CAPABILITY: &str = "admin.on-demand-migration";
26
27/// Upper bound for one admin response body. A status document with latency
28/// histograms is a few kilobytes; a megabyte is generous without being unbounded.
29pub const MAX_ON_DEMAND_MIGRATION_RESPONSE_BYTES: usize = 1024 * 1024;
30
31/// Upper bound for a CA bundle passed with `--ca-cert`.
32pub const MAX_ON_DEMAND_MIGRATION_CA_CERT_BYTES: usize = 64 * 1024;
33
34/// The placeholder the server substitutes for a credential in every response.
35pub const REDACTED_SECRET: &str = "REDACTED";
36
37/// Largest value the server accepts for `policy.inline_max_bytes` (256 MiB).
38pub const MAX_INLINE_MAX_BYTES: u64 = 256 * 1024 * 1024;
39
40/// Largest value the server accepts for `policy.max_concurrent_pulls`.
41pub const MAX_CONCURRENT_PULLS: u32 = 256;
42
43// ---------------------------------------------------------------------------
44// Enumerations shared by the request and the view
45// ---------------------------------------------------------------------------
46
47/// Source vendor family for the S3-speaking providers `rc` can configure.
48#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
49#[serde(rename_all = "lowercase")]
50pub enum SourceProvider {
51    /// Generic S3-compatible endpoint.
52    S3,
53    Aws,
54    Minio,
55    Rustfs,
56    R2,
57    /// GCS XML interoperability API with HMAC keys.
58    Gcs,
59}
60
61impl SourceProvider {
62    pub const fn as_str(self) -> &'static str {
63        match self {
64            Self::S3 => "s3",
65            Self::Aws => "aws",
66            Self::Minio => "minio",
67            Self::Rustfs => "rustfs",
68            Self::R2 => "r2",
69            Self::Gcs => "gcs",
70        }
71    }
72
73    /// Only AWS derives its endpoint from the region.
74    pub const fn requires_endpoint(self) -> bool {
75        !matches!(self, Self::Aws)
76    }
77}
78
79/// Bucket addressing style; `auto` is resolved by the server.
80#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
81#[serde(rename_all = "lowercase")]
82pub enum PathStyle {
83    #[default]
84    Auto,
85    Path,
86    Virtual,
87}
88
89impl PathStyle {
90    pub const fn as_str(self) -> &'static str {
91        match self {
92            Self::Auto => "auto",
93            Self::Path => "path",
94            Self::Virtual => "virtual",
95        }
96    }
97}
98
99/// What a HEAD that misses locally does.
100#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
101#[serde(rename_all = "snake_case")]
102pub enum HeadPolicy {
103    #[default]
104    Proxy,
105    LocalOnly,
106}
107
108impl HeadPolicy {
109    pub const fn as_str(self) -> &'static str {
110        match self {
111            Self::Proxy => "proxy",
112            Self::LocalOnly => "local_only",
113        }
114    }
115}
116
117/// Whether a Range GET also queues a whole-object background pull.
118#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
119#[serde(rename_all = "snake_case")]
120pub enum RangeGetPolicy {
121    #[default]
122    ServeAndBackfill,
123    ServeOnly,
124}
125
126impl RangeGetPolicy {
127    pub const fn as_str(self) -> &'static str {
128        match self {
129            Self::ServeAndBackfill => "serve_and_backfill",
130            Self::ServeOnly => "serve_only",
131        }
132    }
133}
134
135/// How a source failure is answered to the client.
136#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
137#[serde(rename_all = "snake_case")]
138pub enum SourceErrorPolicy {
139    #[default]
140    Propagate,
141    NotFound,
142}
143
144impl SourceErrorPolicy {
145    pub const fn as_str(self) -> &'static str {
146        match self {
147            Self::Propagate => "propagate",
148            Self::NotFound => "not_found",
149        }
150    }
151}
152
153/// When the backfill job treats a local object as already migrated.
154#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
155#[serde(rename_all = "snake_case")]
156pub enum SkipExisting {
157    #[default]
158    Always,
159    EtagOrSize,
160}
161
162impl SkipExisting {
163    pub const fn as_str(self) -> &'static str {
164        match self {
165            Self::Always => "always",
166            Self::EtagOrSize => "etag_or_size",
167        }
168    }
169}
170
171// ---------------------------------------------------------------------------
172// Request shape
173// ---------------------------------------------------------------------------
174
175/// Static credentials for the source, with the secret in zeroizing storage.
176///
177/// `Debug` never prints the secret or the session token.
178#[derive(Clone)]
179pub struct SourceCredentialsRequest {
180    pub access_key: String,
181    pub secret_key: Zeroizing<String>,
182    pub session_token: Option<Zeroizing<String>>,
183}
184
185impl std::fmt::Debug for SourceCredentialsRequest {
186    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
187        formatter
188            .debug_struct("SourceCredentialsRequest")
189            .field("access_key", &self.access_key)
190            .field("secret_key", &REDACTED_SECRET)
191            .field(
192                "session_token",
193                &self.session_token.as_ref().map(|_| REDACTED_SECRET),
194            )
195            .finish()
196    }
197}
198
199impl Serialize for SourceCredentialsRequest {
200    fn serialize<S: Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
201        let mut state = serializer.serialize_struct("SourceCredentials", 3)?;
202        state.serialize_field("access_key", &self.access_key)?;
203        state.serialize_field("secret_key", self.secret_key.as_str())?;
204        state.serialize_field(
205            "session_token",
206            &self.session_token.as_deref().map(String::as_str),
207        )?;
208        state.end()
209    }
210}
211
212/// TLS settings for the source connection.
213#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize)]
214pub struct TlsRequest {
215    pub skip_verify: bool,
216    pub ca_cert_pem: Option<String>,
217}
218
219/// The external source bucket.
220#[derive(Clone, Debug, Serialize)]
221pub struct SourceRequest {
222    pub provider: SourceProvider,
223    pub endpoint: Option<String>,
224    pub region: String,
225    pub bucket: String,
226    pub path_style: PathStyle,
227    /// `None` means anonymous access to a public source bucket.
228    pub credentials: Option<SourceCredentialsRequest>,
229    pub tls: TlsRequest,
230}
231
232/// Key filters.
233#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize)]
234pub struct FilterRequest {
235    /// Only local keys with this prefix consult the source.
236    pub prefix: Option<String>,
237    /// Prepended to the local key to form the source key.
238    pub source_prefix: Option<String>,
239}
240
241/// Per-request source timeouts, in milliseconds.
242#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
243pub struct SourceTimeout {
244    #[serde(default = "default_connect_ms")]
245    pub connect_ms: u64,
246    #[serde(default = "default_first_byte_ms")]
247    pub first_byte_ms: u64,
248    #[serde(default = "default_idle_ms")]
249    pub idle_ms: u64,
250}
251
252impl Default for SourceTimeout {
253    fn default() -> Self {
254        Self {
255            connect_ms: default_connect_ms(),
256            first_byte_ms: default_first_byte_ms(),
257            idle_ms: default_idle_ms(),
258        }
259    }
260}
261
262/// Read-path policy.
263///
264/// The request sends every field explicitly, filled with the server defaults
265/// the fixtures pin, so the body `rc` produces is the documented wire shape
266/// rather than a partial document whose meaning depends on server defaults.
267#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
268pub struct PolicyConfig {
269    #[serde(default)]
270    pub head: HeadPolicy,
271    #[serde(default)]
272    pub range_get: RangeGetPolicy,
273    #[serde(default)]
274    pub source_error: SourceErrorPolicy,
275    #[serde(default)]
276    pub list_through: bool,
277    #[serde(default = "default_true")]
278    pub respect_local_delete_marker: bool,
279    #[serde(default = "default_true")]
280    pub preserve_etag: bool,
281    #[serde(default)]
282    pub copy_tags: bool,
283    #[serde(default = "default_true")]
284    pub emit_events: bool,
285    #[serde(default = "default_negative_cache_ttl_secs")]
286    pub negative_cache_ttl_secs: u64,
287    #[serde(default = "default_inline_max_bytes")]
288    pub inline_max_bytes: u64,
289    #[serde(default = "default_multipart_part_size_bytes")]
290    pub multipart_part_size_bytes: u64,
291    #[serde(default = "default_max_concurrent_pulls")]
292    pub max_concurrent_pulls: u32,
293    #[serde(default = "default_pull_queue_capacity")]
294    pub pull_queue_capacity: u32,
295    #[serde(default)]
296    pub source_timeout: SourceTimeout,
297    #[serde(default)]
298    pub bandwidth_limit_bytes_per_sec: Option<u64>,
299}
300
301impl Default for PolicyConfig {
302    fn default() -> Self {
303        Self {
304            head: HeadPolicy::default(),
305            range_get: RangeGetPolicy::default(),
306            source_error: SourceErrorPolicy::default(),
307            list_through: false,
308            respect_local_delete_marker: true,
309            preserve_etag: true,
310            copy_tags: false,
311            emit_events: true,
312            negative_cache_ttl_secs: default_negative_cache_ttl_secs(),
313            inline_max_bytes: default_inline_max_bytes(),
314            multipart_part_size_bytes: default_multipart_part_size_bytes(),
315            max_concurrent_pulls: default_max_concurrent_pulls(),
316            pull_queue_capacity: default_pull_queue_capacity(),
317            source_timeout: SourceTimeout::default(),
318            bandwidth_limit_bytes_per_sec: None,
319        }
320    }
321}
322
323const fn default_true() -> bool {
324    true
325}
326const fn default_version() -> u32 {
327    1
328}
329const fn default_negative_cache_ttl_secs() -> u64 {
330    30
331}
332const fn default_inline_max_bytes() -> u64 {
333    16 * 1024 * 1024
334}
335const fn default_multipart_part_size_bytes() -> u64 {
336    64 * 1024 * 1024
337}
338const fn default_max_concurrent_pulls() -> u32 {
339    8
340}
341const fn default_pull_queue_capacity() -> u32 {
342    1024
343}
344const fn default_connect_ms() -> u64 {
345    5000
346}
347const fn default_first_byte_ms() -> u64 {
348    15_000
349}
350const fn default_idle_ms() -> u64 {
351    30_000
352}
353
354/// The document `PUT .../on-demand-migration/{bucket}` accepts.
355#[derive(Clone, Debug, Serialize)]
356pub struct OnDemandMigrationConfigRequest {
357    pub version: u32,
358    pub enabled: bool,
359    pub source: SourceRequest,
360    pub filter: FilterRequest,
361    pub policy: PolicyConfig,
362}
363
364impl OnDemandMigrationConfigRequest {
365    /// A version-1, enabled configuration with default filter and policy.
366    pub fn new(source: SourceRequest) -> Self {
367        Self {
368            version: default_version(),
369            enabled: true,
370            source,
371            filter: FilterRequest::default(),
372            policy: PolicyConfig::default(),
373        }
374    }
375
376    /// Reject locally what the server would reject, before any network access
377    /// and before a secret is read. Every failure is a usage error.
378    pub fn validate(&self) -> Result<()> {
379        let source = &self.source;
380        match source.endpoint.as_deref() {
381            Some(endpoint) => validate_endpoint(endpoint)?,
382            None if source.provider.requires_endpoint() => {
383                return Err(Error::Config(format!(
384                    "--endpoint is required for provider {}",
385                    source.provider.as_str()
386                )));
387            }
388            None => {}
389        }
390        if source.region.trim().is_empty() {
391            return Err(Error::Config("--region must not be empty".into()));
392        }
393        validate_source_bucket(&source.bucket)?;
394        if let Some(credentials) = &source.credentials {
395            if credentials.access_key.is_empty() {
396                return Err(Error::Config("--access-key must not be empty".into()));
397            }
398            if credentials.secret_key.is_empty() {
399                return Err(Error::Config(
400                    "The source secret key must not be empty".into(),
401                ));
402            }
403            if credentials
404                .session_token
405                .as_ref()
406                .is_some_and(|token| token.is_empty())
407            {
408                return Err(Error::Config(
409                    "The source session token must not be empty".into(),
410                ));
411            }
412        }
413        if let Some(pem) = source.tls.ca_cert_pem.as_deref() {
414            validate_ca_cert_pem(pem)?;
415        }
416        for (flag, value) in [
417            ("--prefix", &self.filter.prefix),
418            ("--source-prefix", &self.filter.source_prefix),
419        ] {
420            if value.as_deref().is_some_and(str::is_empty) {
421                return Err(Error::Config(format!("{flag} must not be empty")));
422            }
423        }
424        let policy = &self.policy;
425        if policy.inline_max_bytes > MAX_INLINE_MAX_BYTES {
426            return Err(Error::Config(format!(
427                "--inline-max-bytes must be at most {MAX_INLINE_MAX_BYTES}"
428            )));
429        }
430        if !(1..=MAX_CONCURRENT_PULLS).contains(&policy.max_concurrent_pulls) {
431            return Err(Error::Config(format!(
432                "--max-concurrent-pulls must be between 1 and {MAX_CONCURRENT_PULLS}"
433            )));
434        }
435        Ok(())
436    }
437
438    /// Serialize into a zeroizing buffer. The result is the only copy of the
439    /// plaintext body; callers hand it to the transport without cloning.
440    pub fn to_wire_json(&self) -> Result<Zeroizing<Vec<u8>>> {
441        self.validate()?;
442        Ok(Zeroizing::new(serde_json::to_vec(self)?))
443    }
444}
445
446/// `http(s)://host[:port]` with nothing else: no path, query, fragment or
447/// userinfo. Userinfo would smuggle a credential into a log line.
448fn validate_endpoint(endpoint: &str) -> Result<()> {
449    let parsed = url::Url::parse(endpoint)
450        .map_err(|_| Error::Config("--endpoint must be an http(s) URL".into()))?;
451    if !matches!(parsed.scheme(), "http" | "https") {
452        return Err(Error::Config("--endpoint must use http or https".into()));
453    }
454    if parsed.host_str().is_none_or(str::is_empty) {
455        return Err(Error::Config("--endpoint must name a host".into()));
456    }
457    if !parsed.username().is_empty() || parsed.password().is_some() {
458        return Err(Error::Config(
459            "--endpoint must not embed credentials; pass them with --access-key".into(),
460        ));
461    }
462    if !matches!(parsed.path(), "" | "/") || parsed.query().is_some() || parsed.fragment().is_some()
463    {
464        return Err(Error::Config(
465            "--endpoint must be scheme://host[:port] with no path, query or fragment".into(),
466        ));
467    }
468    Ok(())
469}
470
471fn validate_source_bucket(bucket: &str) -> Result<()> {
472    if bucket.is_empty() || bucket.contains('/') || bucket.chars().any(char::is_whitespace) {
473        return Err(Error::Config(
474            "--source-bucket must be a non-empty bucket name without '/' or whitespace".into(),
475        ));
476    }
477    Ok(())
478}
479
480/// The server requires a PEM certificate block; checking here turns a wrong
481/// file into a usage error before the secret is prompted for.
482pub fn validate_ca_cert_pem(pem: &str) -> Result<()> {
483    if pem.len() > MAX_ON_DEMAND_MIGRATION_CA_CERT_BYTES {
484        return Err(Error::Config(format!(
485            "--ca-cert exceeds {MAX_ON_DEMAND_MIGRATION_CA_CERT_BYTES} bytes"
486        )));
487    }
488    if !pem.contains("-----BEGIN CERTIFICATE-----") {
489        return Err(Error::Config(
490            "--ca-cert must contain a PEM certificate (-----BEGIN CERTIFICATE-----)".into(),
491        ));
492    }
493    Ok(())
494}
495
496/// A local bucket name as it appears in the admin route. Anything that could
497/// change the route (a slash, a dot segment, whitespace) is refused here so the
498/// transport never has to reason about it.
499pub fn validate_local_bucket(bucket: &str) -> Result<()> {
500    if bucket.is_empty()
501        || bucket.len() > 255
502        || matches!(bucket, "." | "..")
503        || bucket
504            .chars()
505            .any(|c| c == '/' || c == '\\' || c == '%' || c == '?' || c == '#' || c.is_whitespace())
506    {
507        return Err(Error::InvalidPath(
508            "Expected alias/bucket with a plain bucket name".into(),
509        ));
510    }
511    Ok(())
512}
513
514/// Body of `POST .../backfill?op=start`; every field is optional.
515#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize)]
516pub struct BackfillStartRequest {
517    #[serde(skip_serializing_if = "Option::is_none")]
518    pub prefix: Option<String>,
519    #[serde(skip_serializing_if = "Option::is_none")]
520    pub skip_existing: Option<SkipExisting>,
521    /// List and count only; nothing is queued.
522    pub dry_run: bool,
523}
524
525impl BackfillStartRequest {
526    pub fn validate(&self) -> Result<()> {
527        if self.prefix.as_deref().is_some_and(str::is_empty) {
528            return Err(Error::Config("--prefix must not be empty".into()));
529        }
530        Ok(())
531    }
532}
533
534// ---------------------------------------------------------------------------
535// Response shapes
536// ---------------------------------------------------------------------------
537
538/// Redacted credential summary. Only presence survives deserialization: the
539/// server substitutes `REDACTED`, and a server that did not must still never
540/// reach stdout, so the values are discarded at the parsing boundary.
541#[derive(Clone, Debug, Default, PartialEq, Eq)]
542pub struct SourceCredentialsView {
543    pub access_key: String,
544    pub has_secret_key: bool,
545    pub has_session_token: bool,
546}
547
548impl<'de> Deserialize<'de> for SourceCredentialsView {
549    fn deserialize<D: Deserializer<'de>>(deserializer: D) -> std::result::Result<Self, D::Error> {
550        #[derive(Deserialize)]
551        struct Wire {
552            #[serde(default)]
553            access_key: String,
554            #[serde(default)]
555            secret_key: Option<String>,
556            #[serde(default)]
557            session_token: Option<String>,
558        }
559        let wire = Wire::deserialize(deserializer)?;
560        // Wipe whatever the server sent before it goes out of scope.
561        let mut secret_key = Zeroizing::new(wire.secret_key.unwrap_or_default());
562        let mut session_token = Zeroizing::new(wire.session_token.unwrap_or_default());
563        let view = Self {
564            access_key: wire.access_key,
565            has_secret_key: !secret_key.is_empty(),
566            has_session_token: !session_token.is_empty(),
567        };
568        secret_key.zeroize();
569        session_token.zeroize();
570        Ok(view)
571    }
572}
573
574impl Serialize for SourceCredentialsView {
575    fn serialize<S: Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
576        let mut state = serializer.serialize_struct("SourceCredentials", 3)?;
577        state.serialize_field("access_key", &self.access_key)?;
578        state.serialize_field(
579            "secret_key",
580            &self.has_secret_key.then_some(REDACTED_SECRET),
581        )?;
582        state.serialize_field(
583            "session_token",
584            &self.has_session_token.then_some(REDACTED_SECRET),
585        )?;
586        state.end()
587    }
588}
589
590/// TLS settings as returned by the server. The CA bundle is public material,
591/// but it is long; only its presence is kept for display.
592#[derive(Clone, Debug, Default, PartialEq, Eq)]
593pub struct TlsView {
594    pub skip_verify: bool,
595    pub has_ca_cert: bool,
596}
597
598impl<'de> Deserialize<'de> for TlsView {
599    fn deserialize<D: Deserializer<'de>>(deserializer: D) -> std::result::Result<Self, D::Error> {
600        #[derive(Deserialize)]
601        struct Wire {
602            #[serde(default)]
603            skip_verify: bool,
604            #[serde(default)]
605            ca_cert_pem: Option<String>,
606        }
607        let wire = Wire::deserialize(deserializer)?;
608        Ok(Self {
609            skip_verify: wire.skip_verify,
610            has_ca_cert: wire.ca_cert_pem.is_some_and(|pem| !pem.is_empty()),
611        })
612    }
613}
614
615impl Serialize for TlsView {
616    fn serialize<S: Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
617        let mut state = serializer.serialize_struct("Tls", 2)?;
618        state.serialize_field("skip_verify", &self.skip_verify)?;
619        state.serialize_field("has_ca_cert", &self.has_ca_cert)?;
620        state.end()
621    }
622}
623
624/// The source as returned by the server. `provider` stays a string so a
625/// provider this build does not know (for example `azure`) still displays.
626#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
627pub struct SourceView {
628    #[serde(default)]
629    pub provider: String,
630    #[serde(default)]
631    pub endpoint: Option<String>,
632    #[serde(default)]
633    pub region: String,
634    #[serde(default)]
635    pub bucket: String,
636    #[serde(default)]
637    pub path_style: PathStyle,
638    #[serde(default)]
639    pub credentials: Option<SourceCredentialsView>,
640    #[serde(default)]
641    pub tls: TlsView,
642}
643
644#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
645pub struct FilterView {
646    #[serde(default)]
647    pub prefix: Option<String>,
648    #[serde(default)]
649    pub source_prefix: Option<String>,
650}
651
652/// The redacted configuration document.
653#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
654pub struct OnDemandMigrationConfigView {
655    #[serde(default = "default_version")]
656    pub version: u32,
657    #[serde(default = "default_true")]
658    pub enabled: bool,
659    #[serde(default)]
660    pub source: SourceView,
661    #[serde(default)]
662    pub filter: FilterView,
663    #[serde(default)]
664    pub policy: PolicyConfig,
665}
666
667impl Default for OnDemandMigrationConfigView {
668    fn default() -> Self {
669        Self {
670            version: default_version(),
671            enabled: true,
672            source: SourceView::default(),
673            filter: FilterView::default(),
674            policy: PolicyConfig::default(),
675        }
676    }
677}
678
679/// What the `PUT` probe learned about the source.
680#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
681pub struct ProbeSummary {
682    #[serde(default)]
683    pub reachable: bool,
684    #[serde(default)]
685    pub listable: bool,
686    #[serde(default)]
687    pub sample_key: Option<String>,
688}
689
690/// Response of `PUT .../on-demand-migration/{bucket}`.
691#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
692pub struct OnDemandMigrationSetResult {
693    #[serde(default)]
694    pub bucket: String,
695    #[serde(default)]
696    pub dry_run: bool,
697    #[serde(default)]
698    pub config: Option<OnDemandMigrationConfigView>,
699    /// `None` for a dry run, which saves nothing.
700    #[serde(default)]
701    pub updated_at: Option<String>,
702    #[serde(default)]
703    pub probe: Option<ProbeSummary>,
704}
705
706/// Response of `GET .../on-demand-migration/{bucket}`.
707#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
708pub struct OnDemandMigrationConfigResult {
709    #[serde(default)]
710    pub bucket: String,
711    #[serde(default)]
712    pub config: Option<OnDemandMigrationConfigView>,
713    #[serde(default)]
714    pub updated_at: Option<String>,
715}
716
717#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
718pub struct BreakerStatus {
719    #[serde(default)]
720    pub state: String,
721    #[serde(default)]
722    pub opened_at: Option<String>,
723}
724
725#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
726pub struct LatencyBucket {
727    #[serde(default)]
728    pub le_ms: u64,
729    #[serde(default)]
730    pub count: u64,
731}
732
733#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
734pub struct SourceLatency {
735    #[serde(default)]
736    pub buckets: Vec<LatencyBucket>,
737    #[serde(default)]
738    pub count: u64,
739    #[serde(default)]
740    pub sum_ms: u64,
741}
742
743/// Per-node runtime counters. The nested maps keep the outcome and path
744/// labels open-ended so a new label on the server displays instead of failing.
745#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
746pub struct RuntimeCounters {
747    /// Operation (`get`, `head`) to outcome (`source_hit`, `source_miss`, ...) to count.
748    #[serde(default)]
749    pub requests_total: BTreeMap<String, BTreeMap<String, u64>>,
750    #[serde(default)]
751    pub pulled_bytes_total: u64,
752    /// Pull path (`inline`, `background`, `backfill`) to count.
753    #[serde(default)]
754    pub pulled_objects_total: BTreeMap<String, u64>,
755    /// Failure class to count.
756    #[serde(default)]
757    pub pull_failures_total: BTreeMap<String, u64>,
758    #[serde(default)]
759    pub source_latency: Option<SourceLatency>,
760}
761
762#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
763pub struct LastSourceError {
764    #[serde(default)]
765    pub class: String,
766    #[serde(default)]
767    pub at: Option<String>,
768}
769
770/// Counters of the bucket's backfill job as embedded in the status document.
771#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
772pub struct BackfillSummary {
773    #[serde(default)]
774    pub job_id: String,
775    #[serde(default)]
776    pub state: String,
777    #[serde(default)]
778    pub listed: u64,
779    #[serde(default)]
780    pub enqueued: u64,
781    #[serde(default)]
782    pub pulled: u64,
783    #[serde(default)]
784    pub skipped_existing: u64,
785    #[serde(default)]
786    pub failed: u64,
787    #[serde(default)]
788    pub bytes: u64,
789    #[serde(default)]
790    pub updated_at: Option<String>,
791}
792
793/// Response of `GET .../on-demand-migration/{bucket}/status`.
794///
795/// This is the answering node's view: counters, queue depth and breaker state
796/// are per node, while the configuration is cluster-wide.
797#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
798pub struct OnDemandMigrationStatus {
799    #[serde(default)]
800    pub configured: bool,
801    #[serde(default)]
802    pub enabled: bool,
803    #[serde(default)]
804    pub module_enabled: bool,
805    #[serde(default)]
806    pub provider: Option<String>,
807    #[serde(default)]
808    pub endpoint_host: Option<String>,
809    #[serde(default)]
810    pub breaker: Option<BreakerStatus>,
811    #[serde(default)]
812    pub counters: Option<RuntimeCounters>,
813    #[serde(default)]
814    pub last_source_error: Option<LastSourceError>,
815    #[serde(default)]
816    pub inflight_pulls: u64,
817    #[serde(default)]
818    pub queue_depth: u64,
819    /// Deliberately `null` on the server today. Rendered as an em dash, never
820    /// as zero: a missing ratio and a zero ratio mean different things.
821    #[serde(default)]
822    pub served_by_source_ratio: Option<f64>,
823    #[serde(default)]
824    pub updated_at: Option<String>,
825    #[serde(default)]
826    pub backfill: Option<BackfillSummary>,
827}
828
829#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
830pub struct BackfillLastError {
831    #[serde(default)]
832    pub class: String,
833    /// Hash of the failing key; the key itself never leaves the server.
834    #[serde(default)]
835    pub key_hash: Option<String>,
836    #[serde(default)]
837    pub at: Option<String>,
838}
839
840#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
841pub struct BackfillOwner {
842    #[serde(default)]
843    pub node: String,
844    #[serde(default)]
845    pub lease_until: Option<String>,
846}
847
848/// The backfill checkpoint document as stored on the server.
849#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
850pub struct BackfillJob {
851    #[serde(default = "default_version")]
852    pub format_version: u32,
853    #[serde(default)]
854    pub job_id: String,
855    #[serde(default)]
856    pub state: String,
857    #[serde(default)]
858    pub config_updated_at: Option<String>,
859    #[serde(default)]
860    pub prefix: Option<String>,
861    #[serde(default)]
862    pub skip_existing: SkipExisting,
863    #[serde(default)]
864    pub dry_run: bool,
865    #[serde(default)]
866    pub listed: u64,
867    #[serde(default)]
868    pub enqueued: u64,
869    #[serde(default)]
870    pub pulled: u64,
871    #[serde(default)]
872    pub skipped_existing: u64,
873    #[serde(default)]
874    pub failed: u64,
875    #[serde(default)]
876    pub bytes: u64,
877    #[serde(default)]
878    pub last_key: Option<String>,
879    #[serde(default)]
880    pub last_error: Option<BackfillLastError>,
881    #[serde(default)]
882    pub failed_keys: Vec<String>,
883    #[serde(default)]
884    pub started_at: Option<String>,
885    #[serde(default)]
886    pub updated_at: Option<String>,
887    #[serde(default)]
888    pub owner: Option<BackfillOwner>,
889}
890
891impl BackfillJob {
892    /// Whether the job can still change. `--watch` stops on a terminal state;
893    /// an unknown state from a newer server is treated as still running so the
894    /// watcher keeps refreshing rather than declaring victory early.
895    pub fn is_terminal(&self) -> bool {
896        is_terminal_backfill_state(&self.state)
897    }
898}
899
900pub fn is_terminal_backfill_state(state: &str) -> bool {
901    matches!(
902        state,
903        "cancelled" | "completed" | "completed_with_failures" | "failed"
904    )
905}
906
907/// Response of `POST`/`GET .../on-demand-migration/{bucket}/backfill`.
908#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
909pub struct BackfillJobResult {
910    #[serde(default)]
911    pub bucket: String,
912    #[serde(default)]
913    pub job: Option<BackfillJob>,
914}
915
916// ---------------------------------------------------------------------------
917// API
918// ---------------------------------------------------------------------------
919
920/// Administrative operations for on-demand migration.
921///
922/// Writes are never automatically retried: a `PUT` probes the source and a
923/// backfill start takes a lease, so a repeated request is a second decision.
924#[async_trait]
925pub trait OnDemandMigrationApi: Send + Sync {
926    /// Validate, probe and (unless `dry_run`) save the configuration.
927    async fn set_on_demand_migration(
928        &self,
929        bucket: &str,
930        config: &OnDemandMigrationConfigRequest,
931        dry_run: bool,
932    ) -> Result<OnDemandMigrationSetResult>;
933
934    async fn get_on_demand_migration(&self, bucket: &str) -> Result<OnDemandMigrationConfigResult>;
935
936    /// Idempotent; already-pulled objects stay in place.
937    async fn delete_on_demand_migration(&self, bucket: &str) -> Result<()>;
938
939    async fn on_demand_migration_status(&self, bucket: &str) -> Result<OnDemandMigrationStatus>;
940
941    async fn start_on_demand_migration_backfill(
942        &self,
943        bucket: &str,
944        request: &BackfillStartRequest,
945    ) -> Result<BackfillJobResult>;
946
947    async fn cancel_on_demand_migration_backfill(&self, bucket: &str) -> Result<BackfillJobResult>;
948
949    async fn on_demand_migration_backfill_status(&self, bucket: &str) -> Result<BackfillJobResult>;
950}
951
952#[cfg(test)]
953mod tests {
954    use super::*;
955    use serde_json::{Value, json};
956
957    const SET_REQUEST: &str =
958        include_str!("../../tests/fixtures/on_demand_migration/set_request.json");
959    const SET_RESPONSE: &str =
960        include_str!("../../tests/fixtures/on_demand_migration/set_response.json");
961    const GET_RESPONSE: &str =
962        include_str!("../../tests/fixtures/on_demand_migration/get_response.json");
963    const STATUS: &str = include_str!("../../tests/fixtures/on_demand_migration/status.json");
964    const STATUS_WITH_BACKFILL: &str =
965        include_str!("../../tests/fixtures/on_demand_migration/status_with_backfill.json");
966    const BACKFILL_JOB: &str =
967        include_str!("../../tests/fixtures/on_demand_migration/backfill_job.json");
968
969    fn fixture_request() -> OnDemandMigrationConfigRequest {
970        let mut request = OnDemandMigrationConfigRequest::new(SourceRequest {
971            provider: SourceProvider::Minio,
972            endpoint: Some("https://source.example.com:9000".into()),
973            region: "us-east-1".into(),
974            bucket: "legacy-photos".into(),
975            path_style: PathStyle::Auto,
976            credentials: Some(SourceCredentialsRequest {
977                access_key: "AKIASOURCE".into(),
978                secret_key: Zeroizing::new("sourceSecretKey123".into()),
979                session_token: None,
980            }),
981            tls: TlsRequest::default(),
982        });
983        request.filter.source_prefix = Some("photos/".into());
984        request
985    }
986
987    #[test]
988    fn set_request_matches_the_plaintext_wire_fixture() {
989        let body = fixture_request().to_wire_json().unwrap();
990        let actual: Value = serde_json::from_slice(&body).unwrap();
991        let expected: Value = serde_json::from_str(SET_REQUEST).unwrap();
992        assert_eq!(actual, expected);
993    }
994
995    #[test]
996    fn request_debug_never_prints_the_secret() {
997        let request = fixture_request();
998        let debug = format!("{request:?}");
999        assert!(debug.contains("AKIASOURCE"));
1000        assert!(!debug.contains("sourceSecretKey123"));
1001        assert!(debug.contains(REDACTED_SECRET));
1002    }
1003
1004    #[test]
1005    fn set_response_parses_and_drops_the_redacted_secret() {
1006        let result: OnDemandMigrationSetResult = serde_json::from_str(SET_RESPONSE).unwrap();
1007        assert_eq!(result.bucket, "photos");
1008        assert!(!result.dry_run);
1009        assert_eq!(result.updated_at.as_deref(), Some("2026-09-02T10:00:00Z"));
1010        let probe = result.probe.unwrap();
1011        assert!(probe.reachable && probe.listable);
1012        assert_eq!(probe.sample_key.as_deref(), Some("photos/2024/01.jpg"));
1013        let config = result.config.unwrap();
1014        assert_eq!(config.source.provider, "minio");
1015        let credentials = config.source.credentials.unwrap();
1016        assert_eq!(credentials.access_key, "AKIASOURCE");
1017        assert!(credentials.has_secret_key);
1018        assert!(!credentials.has_session_token);
1019        let serialized = serde_json::to_value(&credentials).unwrap();
1020        assert_eq!(serialized["secret_key"], REDACTED_SECRET);
1021        assert_eq!(serialized["session_token"], Value::Null);
1022    }
1023
1024    #[test]
1025    fn get_response_matches_the_fixture_and_redacts_a_leaked_secret() {
1026        let result: OnDemandMigrationConfigResult = serde_json::from_str(GET_RESPONSE).unwrap();
1027        assert_eq!(result.bucket, "photos");
1028        let config = result.config.unwrap();
1029        assert!(config.enabled);
1030        assert_eq!(config.version, 1);
1031        assert_eq!(config.filter.source_prefix.as_deref(), Some("photos/"));
1032        assert_eq!(config.policy, PolicyConfig::default());
1033        assert_eq!(config.source.path_style, PathStyle::Auto);
1034
1035        // A server that failed to redact must not make it to output either.
1036        let leaked = GET_RESPONSE.replace("\"REDACTED\"", "\"plaintext-secret\"");
1037        let result: OnDemandMigrationConfigResult = serde_json::from_str(&leaked).unwrap();
1038        let text = serde_json::to_string(&result).unwrap();
1039        assert!(!text.contains("plaintext-secret"));
1040        assert!(text.contains(REDACTED_SECRET));
1041    }
1042
1043    #[test]
1044    fn status_fixture_keeps_the_null_ratio_and_counters() {
1045        let status: OnDemandMigrationStatus = serde_json::from_str(STATUS).unwrap();
1046        assert!(status.configured && status.enabled && status.module_enabled);
1047        assert_eq!(status.provider.as_deref(), Some("minio"));
1048        assert_eq!(status.endpoint_host.as_deref(), Some("source.example.com"));
1049        assert_eq!(status.breaker.as_ref().unwrap().state, "half_open");
1050        assert_eq!(status.served_by_source_ratio, None);
1051        assert_eq!(status.inflight_pulls, 1);
1052        assert_eq!(status.queue_depth, 1);
1053        assert!(status.backfill.is_none());
1054        let counters = status.counters.unwrap();
1055        assert_eq!(counters.pulled_bytes_total, 4096);
1056        assert_eq!(counters.requests_total["get"]["source_hit"], 2);
1057        assert_eq!(counters.pull_failures_total["source_timeout"], 1);
1058        assert_eq!(counters.source_latency.unwrap().count, 3);
1059        assert_eq!(status.last_source_error.unwrap().class, "server_error");
1060    }
1061
1062    #[test]
1063    fn status_with_backfill_fixture_carries_the_summary() {
1064        let status: OnDemandMigrationStatus = serde_json::from_str(STATUS_WITH_BACKFILL).unwrap();
1065        let backfill = status.backfill.unwrap();
1066        assert_eq!(backfill.job_id, "11111111-1111-4111-8111-111111111111");
1067        assert_eq!(backfill.state, "running");
1068        assert_eq!(backfill.pulled, 1400);
1069        assert_eq!(backfill.bytes, 73_400_320);
1070    }
1071
1072    #[test]
1073    fn backfill_job_fixture_parses_and_is_not_terminal() {
1074        let result: BackfillJobResult = serde_json::from_str(BACKFILL_JOB).unwrap();
1075        assert_eq!(result.bucket, "photos");
1076        let job = result.job.unwrap();
1077        assert_eq!(job.state, "running");
1078        assert!(!job.is_terminal());
1079        assert_eq!(job.skip_existing, SkipExisting::Always);
1080        assert_eq!(job.prefix.as_deref(), Some("photos/"));
1081        assert_eq!(job.last_error.unwrap().class, "source_timeout");
1082        assert_eq!(job.owner.unwrap().node, "node-a:9000");
1083        assert_eq!(job.failed_keys, vec!["9f2c3b0a1d4e5f60"]);
1084        for state in [
1085            "cancelled",
1086            "completed",
1087            "completed_with_failures",
1088            "failed",
1089        ] {
1090            assert!(is_terminal_backfill_state(state), "{state}");
1091        }
1092        for state in ["pending", "running", "paused", "something_new"] {
1093            assert!(!is_terminal_backfill_state(state), "{state}");
1094        }
1095    }
1096
1097    #[test]
1098    fn older_servers_that_omit_fields_still_parse() {
1099        let status: OnDemandMigrationStatus = serde_json::from_str("{}").unwrap();
1100        assert!(!status.configured);
1101        assert_eq!(status.served_by_source_ratio, None);
1102        let config: OnDemandMigrationConfigResult =
1103            serde_json::from_str(r#"{"bucket":"b","config":{"source":{"provider":"s3"}}}"#)
1104                .unwrap();
1105        let config = config.config.unwrap();
1106        assert!(config.enabled);
1107        assert_eq!(config.policy.max_concurrent_pulls, 8);
1108        assert!(config.source.credentials.is_none());
1109        let job: BackfillJobResult = serde_json::from_str(r#"{"bucket":"b"}"#).unwrap();
1110        assert!(job.job.is_none());
1111        // Unknown fields from a newer server are ignored rather than fatal.
1112        let newer: OnDemandMigrationStatus =
1113            serde_json::from_value(json!({"configured": true, "future_field": 1})).unwrap();
1114        assert!(newer.configured);
1115    }
1116
1117    #[test]
1118    fn validation_rejects_what_the_server_would() {
1119        let mut request = fixture_request();
1120        request.source.endpoint = None;
1121        assert!(request.validate().is_err());
1122        request.source.provider = SourceProvider::Aws;
1123        assert!(request.validate().is_ok());
1124
1125        for endpoint in [
1126            "source.example.com",
1127            "ftp://source.example.com",
1128            "https://user:pw@source.example.com",
1129            "https://source.example.com/path",
1130            "https://source.example.com/?x=1",
1131            "https://source.example.com/#frag",
1132        ] {
1133            let mut request = fixture_request();
1134            request.source.endpoint = Some(endpoint.into());
1135            assert!(request.validate().is_err(), "{endpoint}");
1136        }
1137        let mut request = fixture_request();
1138        request.source.endpoint = Some("http://127.0.0.1:9000/".into());
1139        assert!(request.validate().is_ok());
1140
1141        let mut request = fixture_request();
1142        request.source.region = " ".into();
1143        assert!(request.validate().is_err());
1144        let mut request = fixture_request();
1145        request.source.bucket = "a/b".into();
1146        assert!(request.validate().is_err());
1147        let mut request = fixture_request();
1148        request.filter.prefix = Some(String::new());
1149        assert!(request.validate().is_err());
1150        let mut request = fixture_request();
1151        request.policy.inline_max_bytes = MAX_INLINE_MAX_BYTES + 1;
1152        assert!(request.validate().is_err());
1153        let mut request = fixture_request();
1154        request.policy.max_concurrent_pulls = 0;
1155        assert!(request.validate().is_err());
1156        let mut request = fixture_request();
1157        request.source.tls.ca_cert_pem = Some("not a certificate".into());
1158        assert!(request.validate().is_err());
1159        let mut request = fixture_request();
1160        request.source.credentials = None;
1161        assert!(request.validate().is_ok());
1162        assert_eq!(
1163            fixture_request().validate().map_err(|e| e.exit_code()),
1164            Ok(())
1165        );
1166        let mut request = fixture_request();
1167        request.source.endpoint = None;
1168        assert_eq!(request.validate().unwrap_err().exit_code(), 2);
1169    }
1170
1171    #[test]
1172    fn local_bucket_names_cannot_change_the_route() {
1173        for bucket in ["photos", "my.bucket", "a-b_c"] {
1174            assert!(validate_local_bucket(bucket).is_ok(), "{bucket}");
1175        }
1176        for bucket in ["", ".", "..", "a/b", "a b", "a%2Fb", "a?x", "a#f", "a\\b"] {
1177            assert_eq!(
1178                validate_local_bucket(bucket).unwrap_err().exit_code(),
1179                2,
1180                "{bucket}"
1181            );
1182        }
1183    }
1184
1185    #[test]
1186    fn backfill_start_request_omits_unset_fields() {
1187        let request = BackfillStartRequest::default();
1188        assert_eq!(
1189            serde_json::to_value(&request).unwrap(),
1190            json!({"dry_run": false})
1191        );
1192        let request = BackfillStartRequest {
1193            prefix: Some("photos/".into()),
1194            skip_existing: Some(SkipExisting::EtagOrSize),
1195            dry_run: true,
1196        };
1197        assert_eq!(
1198            serde_json::to_value(&request).unwrap(),
1199            json!({"prefix": "photos/", "skip_existing": "etag_or_size", "dry_run": true})
1200        );
1201        assert!(
1202            BackfillStartRequest {
1203                prefix: Some(String::new()),
1204                ..Default::default()
1205            }
1206            .validate()
1207            .is_err()
1208        );
1209    }
1210}