use crate::TileSourceManager;
use crate::config::file::pmtiles::PmtConfig;
use crate::config::file::process::ProcessConfig;
use crate::config::file::tiles::discovery::{
FsDiscovery, FsSourceBuilder, ObjectStoreDiscovery, ObjectStoreParser, ObjectStoreSourceBuilder,
};
use crate::config::file::tiles::driver::{Baseline, NotifyTrigger, PollTrigger, ReloadDriver};
use crate::config::file::{
CachePolicy, FileConfigEnum, SourceBuildResult, TileSourceConfiguration as _, TileSourceWarning,
};
use crate::config::primitives::IdResolver;
use crate::reload::FileKind;
const PMTILES_EXT: &str = "pmtiles";
pub struct PmtilesReloader {
local: ReloadDriver<FsDiscovery, TileSourceManager>,
remote: ReloadDriver<ObjectStoreDiscovery, TileSourceManager>,
}
impl PmtilesReloader {
#[must_use]
pub fn new(
tsm: TileSourceManager,
id_resolver: IdResolver,
config: &FileConfigEnum<PmtConfig>,
default_cache: CachePolicy,
global_process: &ProcessConfig,
) -> Self {
let default_cache = config.cache_or(default_cache);
#[cfg(feature = "_process")]
let process = {
let source_type = match config {
FileConfigEnum::Config(cfg) => ProcessConfig {
#[cfg(feature = "mlt")]
convert_to_mlt: cfg.custom.convert_to_mlt.clone(),
#[cfg(feature = "mlt")]
convert_to_mvt: cfg.custom.convert_to_mvt.clone(),
..Default::default()
},
FileConfigEnum::None | FileConfigEnum::Path(_) | FileConfigEnum::Paths(_) => {
ProcessConfig::default()
}
};
ProcessConfig::layered(global_process, &source_type, &ProcessConfig::default())
};
#[cfg(not(feature = "_process"))]
let process = {
let _ = global_process;
ProcessConfig::default()
};
let pmt_config = match config {
FileConfigEnum::Config(cfg) => cfg.custom.clone(),
FileConfigEnum::None | FileConfigEnum::Path(_) | FileConfigEnum::Paths(_) => {
PmtConfig::default()
}
};
let build_config = pmt_config.clone();
let build: FsSourceBuilder = Box::new(move |id, path, policy| {
let config = build_config.clone();
Box::pin(async move { config.new_sources(id, path, policy).await })
});
let local = FsDiscovery::from_config(
FileKind::Pmtiles,
config,
pmt_config.recursive.unwrap_or_default(),
&[PMTILES_EXT],
id_resolver.clone(),
default_cache,
&process,
build,
);
let parser_config = pmt_config.clone();
let parser: ObjectStoreParser = Box::new(move |url| parser_config.parse_url_opts(url));
let remote_build = ObjectStoreSourceBuilder::Pmtiles(pmt_config.clone());
let remote = ObjectStoreDiscovery::from_config(
config,
&[PMTILES_EXT],
"PmtilesReloader",
pmt_config.reload_interval,
id_resolver,
default_cache,
&process,
parser,
remote_build,
);
Self {
local: ReloadDriver::new(local, tsm.clone()),
remote: ReloadDriver::new(remote, tsm),
}
}
pub async fn init(&mut self) -> SourceBuildResult<Vec<TileSourceWarning>> {
self.local.init().await
}
pub fn start(self) -> notify::Result<()> {
let Self { local, remote } = self;
let directories = local.discovery().directories();
let recursive = local.discovery().recursive();
let has_remote = !remote.discovery().remote_prefixes().is_empty();
let interval = remote.discovery().reload_interval();
if directories.is_empty() && !has_remote {
return Ok(());
}
if !directories.is_empty() {
let trigger = NotifyTrigger::new(&directories, recursive)?;
local.spawn(trigger, Baseline::Initialized);
}
if has_remote {
if interval.is_zero() {
tracing::info!(
"PmtilesReloader: remote prefix polling disabled (reload_interval = 0s)"
);
} else {
let trigger = PollTrigger::new(interval);
remote.spawn(trigger, Baseline::Empty);
}
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use std::path::PathBuf;
use std::time::Duration;
use insta::assert_yaml_snapshot;
use super::*;
use crate::config::file::pmtiles::DEFAULT_RELOAD_INTERVAL;
use crate::config::file::{
CachePolicy, FileConfig, FileConfigSource, FileConfigSrc, OnInvalid,
};
use crate::config::primitives::OptOneMany;
fn make_reloader(config: &FileConfigEnum<PmtConfig>) -> PmtilesReloader {
let tsm = TileSourceManager::new(None, OnInvalid::Warn);
let resolver = IdResolver::new(&[]);
PmtilesReloader::new(
tsm,
resolver,
config,
CachePolicy::default(),
&ProcessConfig::default(),
)
}
#[derive(serde::Serialize)]
struct ReloaderSnapshot {
local_dir_count: usize,
remote_prefix_count: usize,
remote_prefixes: Vec<String>,
interval_secs: u64,
}
impl From<&PmtilesReloader> for ReloaderSnapshot {
fn from(r: &PmtilesReloader) -> Self {
Self {
local_dir_count: r.local.discovery().directories().len(),
remote_prefix_count: r.remote.discovery().remote_prefixes().len(),
remote_prefixes: r
.remote
.discovery()
.remote_prefixes()
.iter()
.map(ToString::to_string)
.collect(),
interval_secs: r.remote.discovery().reload_interval().as_secs(),
}
}
}
#[test]
fn new_with_none_config_yields_default_interval() {
let reloader = make_reloader(&FileConfigEnum::None);
assert!(reloader.local.discovery().directories().is_empty());
assert!(reloader.remote.discovery().remote_prefixes().is_empty());
assert_eq!(
reloader.remote.discovery().reload_interval(),
DEFAULT_RELOAD_INTERVAL
);
}
#[test]
fn new_partitions_local_and_remote_paths() {
let cfg = FileConfigEnum::Config(FileConfig {
collections: OptOneMany::NoVals,
paths: OptOneMany::Many(vec![
PathBuf::from("s3://bucket-a/"),
PathBuf::from("s3://bucket-b/folder/"),
PathBuf::from("https://example.com/tiles/"),
]),
sources: None,
custom: PmtConfig {
reload_interval: Duration::from_secs(30),
..PmtConfig::default()
},
});
assert_yaml_snapshot!(ReloaderSnapshot::from(&make_reloader(&cfg)), @r#"
local_dir_count: 0
remote_prefix_count: 3
remote_prefixes:
- "https://example.com/tiles/"
- "s3://bucket-a/"
- "s3://bucket-b/folder/"
interval_secs: 30
"#);
}
#[test]
fn new_dedups_remote_prefixes() {
let cfg = FileConfigEnum::Config(FileConfig {
collections: OptOneMany::NoVals,
paths: OptOneMany::Many(vec![
PathBuf::from("s3://bucket/"),
PathBuf::from("s3://bucket/"),
]),
sources: None,
custom: PmtConfig::default(),
});
let r = make_reloader(&cfg);
assert_eq!(r.remote.discovery().remote_prefixes().len(), 1);
}
#[test]
fn new_skips_remote_individually_configured_sources() {
let mut sources: BTreeMap<String, FileConfigSrc> = BTreeMap::new();
sources.insert(
"remote_a".to_owned(),
FileConfigSrc::Obj(Box::new(FileConfigSource {
path: PathBuf::from("s3://bucket/file.pmtiles"),
cache: CachePolicy::default(),
#[cfg(feature = "mlt")]
convert_to_mlt: None,
#[cfg(feature = "mlt")]
convert_to_mvt: None,
cache_control: None,
#[cfg(all(feature = "hillshade", feature = "_tiles"))]
convert_to_hillshade: None,
#[cfg(all(feature = "contour", feature = "_tiles"))]
convert_to_contour: None,
})),
);
let cfg = FileConfigEnum::Config(FileConfig {
collections: OptOneMany::NoVals,
paths: OptOneMany::NoVals,
sources: Some(sources),
custom: PmtConfig::default(),
});
let r = make_reloader(&cfg);
assert!(r.local.discovery().directories().is_empty());
assert!(r.remote.discovery().remote_prefixes().is_empty());
}
}