1use 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
24pub const ON_DEMAND_MIGRATION_CAPABILITY: &str = "admin.on-demand-migration";
26
27pub const MAX_ON_DEMAND_MIGRATION_RESPONSE_BYTES: usize = 1024 * 1024;
30
31pub const MAX_ON_DEMAND_MIGRATION_CA_CERT_BYTES: usize = 64 * 1024;
33
34pub const REDACTED_SECRET: &str = "REDACTED";
36
37pub const MAX_INLINE_MAX_BYTES: u64 = 256 * 1024 * 1024;
39
40pub const MAX_CONCURRENT_PULLS: u32 = 256;
42
43#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
49#[serde(rename_all = "lowercase")]
50pub enum SourceProvider {
51 S3,
53 Aws,
54 Minio,
55 Rustfs,
56 R2,
57 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 pub const fn requires_endpoint(self) -> bool {
75 !matches!(self, Self::Aws)
76 }
77}
78
79#[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#[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#[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#[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#[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#[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#[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#[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 pub credentials: Option<SourceCredentialsRequest>,
229 pub tls: TlsRequest,
230}
231
232#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize)]
234pub struct FilterRequest {
235 pub prefix: Option<String>,
237 pub source_prefix: Option<String>,
239}
240
241#[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#[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#[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 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 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 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
446fn 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
480pub 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
496pub 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#[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 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#[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 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#[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#[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#[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#[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#[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 #[serde(default)]
701 pub updated_at: Option<String>,
702 #[serde(default)]
703 pub probe: Option<ProbeSummary>,
704}
705
706#[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#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
746pub struct RuntimeCounters {
747 #[serde(default)]
749 pub requests_total: BTreeMap<String, BTreeMap<String, u64>>,
750 #[serde(default)]
751 pub pulled_bytes_total: u64,
752 #[serde(default)]
754 pub pulled_objects_total: BTreeMap<String, u64>,
755 #[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#[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#[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 #[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 #[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#[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 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#[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#[async_trait]
925pub trait OnDemandMigrationApi: Send + Sync {
926 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 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 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 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}