osm-git

A WIP POC based on the idea from https://blog.andygol.co.ua/en/2023/05/07/osm-2-0-api-using-git/
git clone git://archive.git.mtrnord.blog/MTRNord/osm-git.git
Log | Files | Refs | LICENSE

commit 40575c5f475acba139471178ad805a7cfba815c3
parent 345ee1e632502e2956af92dff02d5c2955eb421c
Author: MTRNord <mtrnord1@gmail.com>
Date:   Thu, 25 May 2023 01:30:25 +0200

I dont know what I fixed but it works now

Diffstat:
M.gitignore | 7+++++--
MCargo.lock | 68+++++++++++++++++++++++++++++++++++---------------------------------
MCargo.toml | 5+++--
Msrc/git/mod.rs | 39++++++++++++++++++++++++++-------------
Msrc/main.rs | 238+++++++++++++++++++++++--------------------------------------------------------
Msrc/osm/changesets.rs | 166+++++++++++++++++++++++++++++++++++++++++++++----------------------------------
Msrc/osm/osm_data.rs | 224++++++++++++++++++++++++++++++++++++-------------------------------------------
Mtemplates/README.md | 3+--
8 files changed, 338 insertions(+), 412 deletions(-)

diff --git a/.gitignore b/.gitignore @@ -1,4 +1,7 @@ /target /osm-git-repo /debug.xml -/cache -\ No newline at end of file +/cache +/perf.data +/perf.data.old +/flamegraph.svg +\ No newline at end of file diff --git a/Cargo.lock b/Cargo.lock @@ -113,9 +113,9 @@ dependencies = [ [[package]] name = "base64" -version = "0.21.0" +version = "0.21.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a4a4ddaa51a5bc52a6948f74c06d20aaaddb71924eab79b8c97a8c556e942d6a" +checksum = "3f1e31e207a6b8fb791a38ea3105e6cb541f55e4d029902d3039a4ad07cc4105" [[package]] name = "bitflags" @@ -125,9 +125,9 @@ checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" [[package]] name = "bumpalo" -version = "3.12.2" +version = "3.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3c6ed94e98ecff0c12dd1b04c15ec0d7d9458ca8fe806cea6f12954efe74c63b" +checksum = "a3e2c3daef883ecc1b5d58c15adae93470a91d425f3532ba1695849656af3fc1" [[package]] name = "bytes" @@ -866,7 +866,7 @@ dependencies = [ "tokio", "tracing", "tracing-subscriber", - "walkdir", + "zstd", ] [[package]] @@ -1134,15 +1134,6 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f91339c0467de62360649f8d3e185ca8de4224ff281f66000de5eb2a77a79041" [[package]] -name = "same-file" -version = "1.0.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "93fc1dc3aaa9bfed95e02e6eadabb4baf7e3078b0bd1b4d7b6b0b68378900502" -dependencies = [ - "winapi-util", -] - -[[package]] name = "scopeguard" version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1605,16 +1596,6 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" [[package]] -name = "walkdir" -version = "2.3.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "36df944cda56c7d8d8b7496af378e6b16de9284591917d307c9b4d313c44e698" -dependencies = [ - "same-file", - "winapi-util", -] - -[[package]] name = "want" version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1761,15 +1742,6 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" [[package]] -name = "winapi-util" -version = "0.1.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "70ec6ce85bb158151cae5e5c87f95a8e97d2c0c4b001223f33a334e3ce5de178" -dependencies = [ - "winapi", -] - -[[package]] name = "winapi-x86_64-pc-windows-gnu" version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1915,3 +1887,33 @@ checksum = "80d0f4e272c85def139476380b12f9ac60926689dd2e01d4923222f40580869d" dependencies = [ "winapi", ] + +[[package]] +name = "zstd" +version = "0.12.3+zstd.1.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76eea132fb024e0e13fd9c2f5d5d595d8a967aa72382ac2f9d39fcc95afd0806" +dependencies = [ + "zstd-safe", +] + +[[package]] +name = "zstd-safe" +version = "6.0.5+zstd.1.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d56d9e60b4b1758206c238a10165fbcae3ca37b01744e394c463463f6529d23b" +dependencies = [ + "libc", + "zstd-sys", +] + +[[package]] +name = "zstd-sys" +version = "2.0.8+zstd.1.5.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5556e6ee25d32df2586c098bbfa278803692a20d0ab9565e049480d52707ec8c" +dependencies = [ + "cc", + "libc", + "pkg-config", +] diff --git a/Cargo.toml b/Cargo.toml @@ -9,7 +9,7 @@ edition = "2021" bytes = "1.4.0" clap = { version = "4.3.0", features = ["derive"] } color-eyre = "0.6.2" -flate2 = "1.0.26" +flate2 = { version = "1.0.26" } git2 = "0.17.1" memmap2 = "0.6.1" quick-xml = { version = "0.28.2", features = ["async-tokio", "encoding", "escape-html", "overlapped-lists"] } @@ -20,4 +20,4 @@ time = { version = "0.3.21", features = ["formatting", "parsing"] } tokio = { version = "1.28.1", features = ["full"] } tracing = "0.1.37" tracing-subscriber = "0.3.17" -walkdir = "2.3.3" +zstd = { version = "0.12.3", features = ["zstdmt"] } +\ No newline at end of file diff --git a/src/git/mod.rs b/src/git/mod.rs @@ -1,4 +1,4 @@ -use std::io::Write; +use std::{io::Write, path::Path}; use color_eyre::eyre::Result; use git2::{Oid, Repository, Signature}; @@ -23,7 +23,6 @@ pub fn init_git_repository( git_repo_path: &str, data_url: &str, author: &Signature, - changeset_url: &str, ) -> Result<Repository> { // Check if the git repo already exists if std::path::Path::new(git_repo_path).exists() { @@ -39,7 +38,7 @@ pub fn init_git_repository( // Create the git repo if it doesn't exist let repository = Repository::init(git_repo_path)?; - generate_readme_from_template(&repository, data_url, changeset_url)?; + generate_readme_from_template(&repository, data_url)?; // Commit the README.md file commit( @@ -54,16 +53,11 @@ pub fn init_git_repository( } /// Generate the README.md file from the template and write it to the git repo -pub fn generate_readme_from_template( - repository: &Repository, - data_url: &str, - changeset_url: &str, -) -> Result<()> { +pub fn generate_readme_from_template(repository: &Repository, data_url: &str) -> Result<()> { let template_file = include_str!("../../templates/README.md"); // Replace the template variables with the actual values let template_file = template_file.replace("$server_url", data_url); - let template_file = template_file.replace("$changeset_server_url", changeset_url); // Get the version of this binary let version = env!("CARGO_PKG_VERSION"); @@ -102,16 +96,35 @@ pub fn commit( let tree_id = { let mut index = repository.index()?; for file in added_or_changed_files { - index.add_path(std::path::Path::new(&file))?; + let file_path = Path::new(&file); + let path = if file_path.starts_with(repository.path().parent().unwrap()) { + Path::new(&file).strip_prefix(repository.path().parent().unwrap())? + } else { + Path::new(&file) + }; + index.add_path(path)?; } for file in removed_files { - index.remove_path(std::path::Path::new(&file))?; + let file_path = Path::new(&file); + let path = if file_path.starts_with(repository.path().parent().unwrap()) { + Path::new(&file).strip_prefix(repository.path().parent().unwrap())? + } else { + Path::new(&file) + }; + index.remove_path(path)?; } index.write()?; index.write_tree()? }; let tree = repository.find_tree(tree_id)?; + let head_id = repository.refname_to_id("HEAD"); + if let Ok(head_id) = head_id { + let parent = repository.find_commit(head_id)?; - let oid = repository.commit(Some("HEAD"), author, committer, message, &tree, &[])?; - Ok(oid) + let oid = repository.commit(Some("HEAD"), author, committer, message, &tree, &[&parent])?; + Ok(oid) + } else { + let oid = repository.commit(Some("HEAD"), author, committer, message, &tree, &[])?; + Ok(oid) + } } diff --git a/src/main.rs b/src/main.rs @@ -24,30 +24,19 @@ struct Cli { default_value = "https://planet.openstreetmap.org/replication/day" )] replication_server: String, - /// The server to get changeset data from - #[arg( - long, - default_value = "https://planet.openstreetmap.org/replication/changesets" - )] - changeset_server: String, /// Where to write cache files #[arg(long, default_value = "./cache")] cache_path: String, /// If the git repo should be removed and recreated #[arg(short, long)] clean: bool, - /// Where to start downloading changesets from + /// Where to start downloading data from #[arg(long, default_value = "000/000/000")] - start_changesets: String, - /// The time to wait between downloading changesets + start_data: String, + /// The time to wait between downloading data /// This is to avoid causing a lot of load on the OSM servers #[arg(long, default_value = "500")] wait_time: u64, - /// If we should use torrent files from the cache instead for changesets - /// - /// These need to be downloaded manually and be put into the cache_folder/changesets/torrents folder - #[arg(long)] - use_torrents: bool, } #[tokio::main] @@ -76,122 +65,14 @@ async fn main() -> Result<()> { let author = Signature::now("osm-git-replay", "osm-git-replay@localhost")?; - let repository = init_git_repository( - &cli.git_repo_path, - &cli.replication_server, - &author, - &cli.changeset_server, - )?; + let repository = init_git_repository(&cli.git_repo_path, &cli.replication_server, &author)?; info!("Git repository initialized"); - if !cli.use_torrents { - // Main download loop - let mut changeset_position_top = cli.start_changesets[0..3].parse::<u64>()?; - let mut changeset_position_middle = cli.start_changesets[4..7].parse::<u64>()?; - let mut changeset_position_bottom = cli.start_changesets[8..11].parse::<u64>()?; - let mut changeset_position_middle_incremented = false; - let mut changeset_position_top_incremented = false; - - loop { - // Check for cache and use it if it exists - let cache_file_path = format!( - "{}/changesets/{:03}/{:03}/{:03}.osm.gz", - cli.cache_path, - changeset_position_top, - changeset_position_middle, - changeset_position_bottom - ); - - if std::path::Path::new(&cache_file_path).exists() { - info!( - "We already got the changeset file at {}. Skipping", - cache_file_path - ); - - // Increment the changeset position - changeset_position_bottom += 1; - changeset_position_middle_incremented = false; - changeset_position_top_incremented = false; - } else { - { - // First we download the changeset files - let changeset_url = format!( - "{}/{:03}/{:03}/{:03}.osm.gz", - cli.changeset_server, - changeset_position_top, - changeset_position_middle, - changeset_position_bottom - ); - info!("Downloading changeset file from {}", changeset_url); - let changeset_response: reqwest::Response = - client.get(&changeset_url).send().await?; - if changeset_response.status() == reqwest::StatusCode::NOT_FOUND { - warn!("Changeset file not found at {}", changeset_url); - // We've reached the end of the changesets for this bottom position. - // If we incremented top and failed again, we're done. - if changeset_position_top_incremented { - info!("Finished or failed downloading changesets"); - info!( - "Changeset position: {} {} {}", - changeset_position_top, - changeset_position_middle, - changeset_position_bottom - ); - warn!("Response body: {:?}", changeset_response.text().await?); - // TODO: We want to have an endless loop here optionally so that we can keep trying to download changesets. - break; - } - // We reset bottom to 0 and increment middle. - // We also mark middle as incremented so that we increment top on the next failure. - if !changeset_position_middle_incremented && changeset_position_bottom != 0 - { - changeset_position_bottom = 0; - changeset_position_middle += 1; - changeset_position_middle_incremented = true; - changeset_position_top_incremented = false; - } else { - changeset_position_middle_incremented = false; - changeset_position_top += 1; - changeset_position_top_incremented = true; - changeset_position_middle = 0; - changeset_position_bottom = 0; - } - continue; - } - let changeset_data = changeset_response.bytes().await?; - info!("Caching changeset file to disk"); - std::fs::create_dir_all( - std::path::Path::new(&cache_file_path).parent().unwrap(), - )?; - std::fs::write(&cache_file_path, &changeset_data)?; - info!("Changeset file downloaded"); - }; - - // TODO: We need to dynamically do this based on the data instead. Otherwise we dont have ram larrge enough - // let file = File::open(cache_file_path)?; - // let changeset_data = unsafe { Mmap::map(&file)? }; - - // let parsed_changeset = parse_changeset(&changeset_data)?; - // info!("Changeset file parsed"); - // changesets.extend(parsed_changeset); - - // Increment the changeset position - changeset_position_bottom += 1; - changeset_position_middle_incremented = false; - changeset_position_top_incremented = false; - - // Wait a few seconds before downloading the next changeset file - tokio::time::sleep(Duration::from_millis(cli.wait_time)).await; - } - } - } - // Data download metadata - let mut data_position_top = 0; - let mut data_position_middle = 0; - let mut data_position_bottom = 0; - let mut data_position_middle_incremented = false; - let mut data_position_top_incremented = false; + // TODO: We should probably detect where to resume from + let mut data_position_top = cli.start_data[0..3].parse::<u16>()?; + let mut data_position_middle = cli.start_data[4..7].parse::<u16>()?; + let mut data_position_bottom = cli.start_data[8..11].parse::<u16>()?; // Parse the changesets and convert them to git objects loop { @@ -205,19 +86,29 @@ async fn main() -> Result<()> { info!("Using cached data file at {}", cache_file_path); let file = File::open(&cache_file_path)?; let data = unsafe { Mmap::map(&file)? }; - convert_objects_to_git( - &repository, - &author, - &data, - cli.cache_path.clone(), - cli.use_torrents, - )?; + let changeset_location = format!("{}/changesets/torrents", cli.cache_path); + convert_objects_to_git(&repository, &author, &data, &changeset_location)?; info!("Data file parsed"); // Increment the data position - data_position_bottom += 1; - data_position_middle_incremented = false; - data_position_top_incremented = false; + if data_position_top == 999 + && data_position_middle == 999 + && data_position_bottom == 999 + { + // Uhhhhhh?! + break; + } + + if data_position_middle == 999 && data_position_bottom == 999 { + data_position_middle = 0; + data_position_bottom = 0; + data_position_top += 1; + } + + if data_position_bottom == 999 { + data_position_bottom = 0; + data_position_middle += 1; + } } else { { // Download minute replication files and find the changesets that were modified in that minute @@ -233,33 +124,30 @@ async fn main() -> Result<()> { if data_response.status() == reqwest::StatusCode::NOT_FOUND { warn!("data file not found at {}", data_url); - // We've reached the end of the data for this bottom position. - // If we incremented top and failed again, we're done. - if data_position_top_incremented { - info!("Finished or failed downloading data"); - info!( - "Data position: {} {} {}", - data_position_top, data_position_middle, data_position_bottom - ); - warn!("Response body: {:?}", data_response.text().await?); - - // TODO: We want to have an endless loop here optionally so that we can keep trying to download changesets. + // Increment the data position + if data_position_top == 999 + && data_position_middle == 999 + && data_position_bottom == 999 + { + // Uhhhhhh?! break; } - // We reset bottom to 0 and increment middle. - // We also mark middle as incremented so that we increment top on the next failure. - if !data_position_middle_incremented && data_position_bottom != 0 { + + if data_position_middle == 999 && data_position_bottom == 999 { + data_position_middle = 0; data_position_bottom = 0; - data_position_middle += 1; - data_position_middle_incremented = true; - data_position_top_incremented = false; - } else { - data_position_middle_incremented = false; data_position_top += 1; - data_position_top_incremented = true; - data_position_middle = 0; + } + + if data_position_bottom == 999 { data_position_bottom = 0; + data_position_middle += 1; + } + + if data_position_bottom < 999 { + data_position_bottom += 1; } + continue; } @@ -273,18 +161,32 @@ async fn main() -> Result<()> { let file = File::open(cache_file_path)?; let data = unsafe { Mmap::map(&file)? }; - convert_objects_to_git( - &repository, - &author, - &data, - cli.cache_path.clone(), - cli.use_torrents, - )?; + let changeset_location = format!("{}/changesets/torrents", cli.cache_path); + convert_objects_to_git(&repository, &author, &data, &changeset_location)?; // Increment the data position - data_position_bottom += 1; - data_position_middle_incremented = false; - data_position_top_incremented = false; + if data_position_top == 999 + && data_position_middle == 999 + && data_position_bottom == 999 + { + // Uhhhhhh?! + break; + } + + if data_position_middle == 999 && data_position_bottom == 999 { + data_position_middle = 0; + data_position_bottom = 0; + data_position_top += 1; + } + + if data_position_bottom == 999 { + data_position_bottom = 0; + data_position_middle += 1; + } + + if data_position_bottom < 999 { + data_position_bottom += 1; + } // Wait a few seconds before downloading the next data file tokio::time::sleep(Duration::from_millis(cli.wait_time)).await; diff --git a/src/osm/changesets.rs b/src/osm/changesets.rs @@ -1,5 +1,4 @@ use color_eyre::eyre::Result; -use flate2::bufread::GzDecoder; use quick_xml::{ events::{BytesStart, Event}, name::QName, @@ -9,9 +8,11 @@ use std::{ borrow::Cow, collections::HashMap, convert::Infallible, - io::{Read, Write}, + fs::File, + io::{BufReader, Write}, }; use tracing::{debug, error, info, warn}; +use zstd::stream::Decoder; #[derive(Debug, Clone, PartialEq)] pub struct Changeset { @@ -29,7 +30,11 @@ pub struct Changeset { } impl Changeset { - fn new_from_element(reader: &mut Reader<&[u8]>, element: BytesStart) -> Result<Self> { + fn new_from_element( + reader: &mut Reader<BufReader<Decoder<'_, BufReader<File>>>>, + element: BytesStart, + read_tags: bool, + ) -> Result<Self> { let changeset_attributes: HashMap<String, String> = element .attributes() .filter_map(|attr_result| attr_result.ok()) @@ -38,10 +43,9 @@ impl Changeset { .decoder() .decode(attr.key.local_name().as_ref()) .or_else(|err| { - dbg!( + debug!( "unable to read key in DefaultSettings attribute {:?}, utf8 error {:?}", - &attr, - err + &attr, err ); Ok::<Cow<'_, str>, Infallible>(std::borrow::Cow::from("")) }) @@ -50,10 +54,9 @@ impl Changeset { let value = attr .decode_and_unescape_value(reader) .or_else(|err| { - dbg!( + debug!( "unable to read key in DefaultSettings attribute {:?}, utf8 error {:?}", - &attr, - err + &attr, err ); Ok::<Cow<'_, str>, Infallible>(std::borrow::Cow::from("")) }) @@ -63,13 +66,22 @@ impl Changeset { }) .collect(); + //debug!("changeset_attributes: {:?}", changeset_attributes); + let mut changeset = Changeset { id: changeset_attributes.get("id").unwrap().parse().unwrap(), created_at: changeset_attributes.get("created_at").unwrap().to_string(), closed_at: changeset_attributes.get("closed_at").map(|s| s.to_string()), open: changeset_attributes.get("open").unwrap().parse().unwrap(), - user: changeset_attributes.get("user").unwrap().to_string(), - uid: changeset_attributes.get("uid").unwrap().parse().unwrap(), + user: changeset_attributes + .get("user") + .map(|s| s.to_string()) + .unwrap_or_else(|| "Unknown".to_string()), + uid: changeset_attributes + .get("uid") + .unwrap_or(&"0".to_string()) + .parse() + .unwrap(), min_lat: changeset_attributes .get("min_lat") .map(|s| s.parse().unwrap()), @@ -85,78 +97,73 @@ impl Changeset { tags: HashMap::new(), }; - let mut element_buf = Vec::new(); - loop { - let event = reader.read_event_into(&mut element_buf)?; + let mut new_buf = Vec::new(); + if read_tags { + loop { + let event = reader.read_event_into(&mut new_buf)?; - if let Event::End(ref e) = event { - if e.name() == element.name() { - break; + if let Event::End(ref e) = event { + if e.name() == element.name() { + break; + } } - } - if let Event::Start(ref e) = event { - let name = e.name(); - if name == QName(b"tag") { - let mut key = Cow::Borrowed(""); - let mut value = Cow::Borrowed(""); - - for attr_result in element.attributes() { - let a = attr_result?; - match a.key.as_ref() { - b"k" => key = a.decode_and_unescape_value(reader)?, - b"v" => value = a.decode_and_unescape_value(reader)?, - _ => (), + if let Event::Start(ref e) = event { + let name = e.name(); + if name == QName(b"tag") { + let mut key = Cow::Borrowed(""); + let mut value = Cow::Borrowed(""); + + for attr_result in element.attributes() { + let a = attr_result?; + match a.key.as_ref() { + b"k" => key = a.decode_and_unescape_value(reader)?, + b"v" => value = a.decode_and_unescape_value(reader)?, + _ => (), + } } - } - changeset.tags.insert(key.to_string(), value.to_string()); - } else { - warn!("Unexpected tag: {:?}", name); - reader.read_to_end(name)?; - } - } else { - if let Event::Text(ref text) = event { - if text.borrow().starts_with(b"\n") { - continue; + changeset.tags.insert(key.to_string(), value.to_string()); + } else { + warn!("Unexpected tag: {:?}", name); } - } else if let Event::End(ref e) = event { - if e.name() == QName(b"tag") { - continue; + } else { + if let Event::Text(ref text) = event { + if text.borrow().starts_with(b"\n") { + continue; + } + } else if let Event::End(ref e) = event { + if e.name() == QName(b"tag") { + continue; + } } - } - warn!("Unexpected event in changeset: {:?}", event); - // Write the data to file for debugging + warn!("Unexpected event in changeset: {:?}", event); + // Write the data to file for debugging - let mut file = std::fs::File::create("debug.xml")?; - file.write_all(&element_buf)?; - file.sync_all()?; + let mut file = std::fs::File::create("debug.xml")?; + file.write_all(&new_buf)?; + file.sync_all()?; + } + new_buf = Vec::new(); } } + Ok(changeset) } } -pub fn parse_changeset(changeset_data: &[u8]) -> Result<Vec<Changeset>> { - // If file is empty, return an empty vector - if changeset_data.is_empty() { - return Ok(Vec::new()); - } +pub fn uncompress_changeset_file<'a>( + file: File, +) -> Reader<BufReader<Decoder<'a, BufReader<File>>>> { // Decompress the changeset file - let mut changeset_data_reader = GzDecoder::new(changeset_data); - let mut changeset_data = String::new(); - changeset_data_reader.read_to_string(&mut changeset_data)?; - debug!( - "Changeset file decompressed. Size: {}", - changeset_data.len() - ); - - // If the changeset file is empty, return an empty vector - if changeset_data.is_empty() { - return Ok(Vec::new()); - } - - let mut changeset_data = Reader::from_str(&changeset_data); + info!("Decompressing changeset file"); + let reader: BufReader<Decoder<BufReader<File>>> = BufReader::new(Decoder::new(file).unwrap()); + Reader::from_reader(reader) +} +pub fn parse_changeset( + changeset_data: &mut Reader<BufReader<Decoder<'_, BufReader<File>>>>, + changeset_id: Option<u64>, +) -> Result<Vec<Changeset>> { // == Handling empty elements == // To simply our processing code // we want the same events for empty elements, like: @@ -175,10 +182,26 @@ pub fn parse_changeset(changeset_data: &[u8]) -> Result<Vec<Changeset>> { Event::Start(element) => { if let b"changeset" = element.name().as_ref() { // TODO: What do we do in case of an error? - let changeset = - Changeset::new_from_element(&mut changeset_data, element.clone()); + let changeset = Changeset::new_from_element( + changeset_data, + element.clone(), + changeset_id.is_some(), + ); + match changeset { - Ok(changeset) => changesets.push(changeset), + Ok(changeset) => { + if changeset_id.is_none() { + changesets.push(changeset); + continue; + } + + if let Some(changeset_id) = changeset_id { + if changeset_id == changeset.id { + changesets.push(changeset); + break; + } + } + } Err(err) => { error!( "unable to read changeset element {:?}, utf8 error {:?}", @@ -191,6 +214,7 @@ pub fn parse_changeset(changeset_data: &[u8]) -> Result<Vec<Changeset>> { Event::Eof => break, // exits the loop when reaching end of file _ => (), // There are `Event` types not considered here } + buf = Vec::new(); } Ok(changesets) } diff --git a/src/osm/osm_data.rs b/src/osm/osm_data.rs @@ -11,17 +11,15 @@ use std::{ borrow::Cow, collections::BTreeMap, convert::Infallible, - fs::File, + fs::{File, OpenOptions}, io::{Read, Write}, - path::Path, }; use time::{format_description::well_known::Iso8601, OffsetDateTime}; use tracing::{debug, error, info, warn}; -use walkdir::WalkDir; use crate::git::commit; -use super::changesets::{parse_changeset, Changeset}; +use super::changesets::{parse_changeset, uncompress_changeset_file, Changeset}; const FILE_VERSION: &str = "0.1.0"; @@ -39,7 +37,7 @@ pub struct Node { pub legacy_object_version: Option<String>, pub lat: f64, pub lon: f64, - #[serde(skip_serializing_if = "BTreeMap::is_empty")] + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] pub tags: BTreeMap<String, String>, } impl Node { @@ -104,9 +102,9 @@ impl Node { file_version: FILE_VERSION.to_string(), }; - let mut element_buf = Vec::new(); + let mut buf = Vec::new(); loop { - let event = reader.read_event_into(&mut element_buf)?; + let event = reader.read_event_into(&mut buf)?; if let Event::End(ref e) = event { if e.name() == element.name() { @@ -132,8 +130,8 @@ impl Node { node.tags.insert(key.to_string(), value.to_string()); } else { warn!("Unexpected tag: {:?}", name); - reader.read_to_end(name)?; } + reader.read_to_end(name)?; } else { if let Event::Text(ref text) = event { if text.borrow().starts_with(b"\n") { @@ -148,9 +146,10 @@ impl Node { // Write the data to file for debugging let mut file = std::fs::File::create("debug.xml")?; - file.write_all(&element_buf)?; + file.write_all(&buf)?; file.sync_all()?; } + buf = Vec::new(); } Ok(node) @@ -169,9 +168,9 @@ pub struct Way { pub file_version: String, #[serde(skip_serializing_if = "Option::is_none")] pub legacy_object_version: Option<String>, - #[serde(skip_serializing_if = "BTreeMap::is_empty")] + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] pub tags: BTreeMap<String, String>, - #[serde(skip_serializing_if = "Vec::is_empty")] + #[serde(default, skip_serializing_if = "Vec::is_empty")] pub nodes: Vec<u64>, } @@ -228,9 +227,9 @@ impl Way { file_version: FILE_VERSION.to_string(), }; - let mut element_buf = Vec::new(); + let mut buf = Vec::new(); loop { - let event = reader.read_event_into(&mut element_buf)?; + let event = reader.read_event_into(&mut buf)?; if let Event::End(ref e) = event { if e.name() == element.name() { @@ -272,8 +271,8 @@ impl Way { ); } else { warn!("Unexpected tag: {:?}", name); - reader.read_to_end(name)?; } + reader.read_to_end(name)?; } else { if let Event::Text(ref text) = event { if text.borrow().starts_with(b"\n") { @@ -288,9 +287,10 @@ impl Way { // Write the data to file for debugging let mut file = std::fs::File::create("debug.xml")?; - file.write_all(&element_buf)?; + file.write_all(&buf)?; file.sync_all()?; } + buf = Vec::new(); } Ok(way) @@ -319,9 +319,9 @@ pub struct Relation { pub file_version: String, #[serde(skip_serializing_if = "Option::is_none")] pub legacy_object_version: Option<String>, - #[serde(skip_serializing_if = "BTreeMap::is_empty")] + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] pub tags: BTreeMap<String, String>, - #[serde(skip_serializing_if = "Vec::is_empty")] + #[serde(default, skip_serializing_if = "Vec::is_empty")] pub member: Vec<RelationMember>, } @@ -378,9 +378,9 @@ impl Relation { file_version: FILE_VERSION.to_string(), }; - let mut element_buf = Vec::new(); + let mut buf = Vec::new(); loop { - let event = reader.read_event_into(&mut element_buf)?; + let event = reader.read_event_into(&mut buf)?; if let Event::End(ref e) = event { if e.name() == element.name() { @@ -435,8 +435,8 @@ impl Relation { }); } else { warn!("Unexpected tag: {:?}", name); - reader.read_to_end(name)?; } + reader.read_to_end(name)?; } else { if let Event::Text(ref text) = event { if text.borrow().starts_with(b"\n") { @@ -451,9 +451,10 @@ impl Relation { // Write the data to file for debugging let mut file = std::fs::File::create("debug.xml")?; - file.write_all(&element_buf)?; + file.write_all(&buf)?; file.sync_all()?; } + buf = Vec::new(); } Ok(relation) @@ -472,8 +473,7 @@ pub fn convert_objects_to_git( repository: &Repository, committer: &Signature, data: &[u8], - cache_folder: String, - use_torrents: bool, + changesets_location: &str, ) -> Result<()> { // If the file is empty we skip it if data.is_empty() { @@ -483,7 +483,10 @@ pub fn convert_objects_to_git( // Decompress the changeset file let mut data_reader = GzDecoder::new(data); let mut file_data = String::new(); - data_reader.read_to_string(&mut file_data)?; + if let Err(e) = data_reader.read_to_string(&mut file_data) { + error!("Unable to decompress data file: {:?}. Moving on", e); + return Ok(()); + } debug!("Data file decompressed. Size: {}", file_data.len()); // If the file is empty we skip it @@ -503,6 +506,7 @@ pub fn convert_objects_to_git( data.expand_empty_elements(true); let mut buf = Vec::new(); + let mut skip_buf = Vec::new(); let mut created_or_modified_objects_for_changeset = BTreeMap::new(); let mut deleted_objects_for_changeset = BTreeMap::new(); @@ -515,9 +519,8 @@ pub fn convert_objects_to_git( let mut created_objects = Vec::new(); - let mut element_buf = Vec::new(); loop { - let event = data.read_event_into(&mut element_buf)?; + let event = data.read_event_into(&mut skip_buf)?; if let Event::End(ref e) = event { if e.name() == element.name() { @@ -579,10 +582,12 @@ pub fn convert_objects_to_git( file.write_all(file_data.as_bytes())?; file.sync_all()?; } + skip_buf = Vec::new(); } // write the objects to the git repo as yaml files let repository_folder = repository.path().parent().unwrap(); + // TODO: We should chunk the world and split it into folders... Otherwise good luck for object in created_objects { let object_file_name = match object { OSMObject::Node(ref node) => format!("{}.yaml", node.id), @@ -610,9 +615,8 @@ pub fn convert_objects_to_git( let mut deleted_objects = Vec::new(); - let mut element_buf = Vec::new(); loop { - let event = data.read_event_into(&mut element_buf)?; + let event = data.read_event_into(&mut skip_buf)?; if let Event::End(ref e) = event { if e.name() == element.name() { @@ -674,6 +678,7 @@ pub fn convert_objects_to_git( file.write_all(file_data.as_bytes())?; file.sync_all()?; } + skip_buf = Vec::new(); } // write the objects to the git repo as yaml files @@ -685,9 +690,22 @@ pub fn convert_objects_to_git( OSMObject::Relation(ref relation) => format!("{}.yaml", relation.id), }; let object_file_path = repository_folder.join(object_file_name); - // Change the file according to the changeset - let mut object_file = std::fs::File::open(object_file_path)?; + + // If we got the file we open it otherwise we create a new object + if !object_file_path.exists() { + // We need to create the file + let object_file = OpenOptions::new() + .read(true) + .write(true) + .create(true) + .open(object_file_path.clone())?; + serde_yaml::to_writer(object_file, &object)?; + } + let mut object_file = OpenOptions::new() + .read(true) + .open(object_file_path.clone())?; + let mut file_object: OSMObject = serde_yaml::from_reader(&mut object_file)?; match object { @@ -726,12 +744,19 @@ pub fn convert_objects_to_git( } } } + let object_file = OpenOptions::new() + .read(true) + .write(true) + .truncate(true) + .open(object_file_path)?; + serde_yaml::to_writer(object_file, &object)?; // Add the object to the list of created objects for the changeset based on the changeset id let changeset = match object { OSMObject::Node(ref node) => node.changeset, OSMObject::Way(ref way) => way.changeset, OSMObject::Relation(ref relation) => relation.changeset, }; + created_or_modified_objects_for_changeset .entry(changeset) .or_insert_with(Vec::new) @@ -743,9 +768,8 @@ pub fn convert_objects_to_git( let mut deleted_objects = Vec::new(); - let mut element_buf = Vec::new(); loop { - let event = data.read_event_into(&mut element_buf)?; + let event = data.read_event_into(&mut skip_buf)?; if let Event::End(ref e) = event { if e.name() == element.name() { @@ -807,6 +831,7 @@ pub fn convert_objects_to_git( file.write_all(file_data.as_bytes())?; file.sync_all()?; } + skip_buf = Vec::new(); } // write the objects to the git repo as yaml files @@ -819,8 +844,10 @@ pub fn convert_objects_to_git( }; let object_file_path = repository_folder.join(object_file_name); - // Delete the file - std::fs::remove_file(object_file_path)?; + // Delete the file if it exists + if object_file_path.exists() { + std::fs::remove_file(object_file_path)?; + } // Add the object to the list of created objects for the changeset based on the changeset id let changeset = match object { @@ -839,6 +866,7 @@ pub fn convert_objects_to_git( Event::Eof => break, // exits the loop when reaching end of file _ => (), // There are `Event` types not considered here } + buf = Vec::new(); } // For all the objects changed apply the changesets as commits @@ -848,16 +876,14 @@ pub fn convert_objects_to_git( .chain(deleted_objects_for_changeset.keys()) .collect(); - for changeset in changeset_list { + for changeset_id in changeset_list { // Find the changeset within the files of the cache - let changeset = find_changesets_in_cache( - cache_folder.clone(), - use_torrents, - *changeset, - None, - None, - None, - )?; + let changeset = find_changesets_in_cache(changesets_location, *changeset_id)?; + + if changeset.is_none() { + warn!("Unable to find changeset {:?}", changeset_id); + continue; + } if let Some(changeset) = changeset { // Get comment tag if it exists and trim it @@ -929,6 +955,13 @@ pub fn convert_objects_to_git( let note = changeset .tags .iter() + .filter_map(|(key, value)| { + if key.trim().is_empty() { + None + } else { + Some((key, value)) + } + }) .map(|(key, value)| format!("{}: {}", key, value)) .collect::<Vec<String>>() .join("\n"); @@ -954,91 +987,40 @@ pub fn convert_objects_to_git( /// /// The changeset if found fn find_changesets_in_cache( - cache_folder: String, - use_torrents: bool, + changesets_location: &str, changeset_id: u64, - changeset_position_top: Option<u64>, - changeset_position_middle: Option<u64>, - changeset_position_bottom: Option<u64>, ) -> Result<Option<Changeset>> { - // Changesets are of the format "{cache_folder}/changesets/{:03}/{:03}/{:03}.osm.gz". - // - // Each file has to be parsed using "parse_changeset(&changeset_data)?;" - - let changeset_folder = if use_torrents { - format!("{}/changesets/torrents", cache_folder) - } else { - format!("{}/changesets", cache_folder) - }; - - if !use_torrents { - let mut changeset_position_top = changeset_position_top.unwrap_or(0); - let mut changeset_position_middle = changeset_position_middle.unwrap_or(0); - let mut changeset_position_bottom = changeset_position_bottom.unwrap_or(0); - let changeset_file = format!( - "{:03}/{:03}/{:03}.osm.gz", - changeset_position_top, changeset_position_middle, changeset_position_bottom - ); - let changeset_path = format!("{}/{}", changeset_folder, changeset_file); - - if !Path::new(&changeset_path).exists() { - return Ok(None); + let mut changeset = None; + + // Find latest changeset file (highest number in filename after "changesets-" and before ".osm.zst") + let changeset_files = std::fs::read_dir(changesets_location)?; + let mut last_highest_id = 0; + let mut changeset_path = String::new(); + for changeset_file in changeset_files { + let changeset_file = changeset_file?; + let changeset_file_path = changeset_file.path(); + let changeset_file_name = changeset_file_path.file_name().unwrap().to_str().unwrap(); + let changeset_file_name = changeset_file_name.trim_end_matches(".osm.zst"); + let changeset_file_name = changeset_file_name.trim_start_matches("changesets-"); + let changeset_file_name = changeset_file_name.parse::<u64>(); + if let Ok(changeset_file_name) = changeset_file_name { + if changeset_file_name > last_highest_id { + last_highest_id = changeset_file_name; + changeset_path = changeset_file_path.to_str().unwrap().to_string(); + } } + } - let mut changeset_file = File::open(changeset_path)?; - let mut changeset_data = Vec::new(); - changeset_file.read_to_end(&mut changeset_data)?; + let changeset_file = File::open(changeset_path)?; + let mut uncompressed_data = uncompress_changeset_file(changeset_file); - let changesets = parse_changeset(&changeset_data)?; + let changesets = parse_changeset(&mut uncompressed_data, Some(changeset_id)); - // Check if the file has the correct changeset id in vector otherwise we recurse to the next file + if let Ok(changesets) = changesets { if changesets.iter().any(|c| c.id == changeset_id) { - return Ok(changesets.into_iter().find(|c| c.id == changeset_id)); - } - - // We recurse to the next file since we found no changeset with the correct id - - if changeset_position_top == 999 - && changeset_position_middle == 999 - && changeset_position_bottom == 999 - { - // Uhhhhhh?! - return Ok(None); - } - - if changeset_position_middle == 999 && changeset_position_bottom == 999 { - changeset_position_middle = 0; - changeset_position_bottom = 0; - changeset_position_top += 1; + changeset = changesets.into_iter().find(|c| c.id == changeset_id); } - - if changeset_position_bottom == 999 { - changeset_position_bottom = 0; - changeset_position_middle += 1; - } - - find_changesets_in_cache( - cache_folder, - use_torrents, - changeset_id, - Some(changeset_position_top), - Some(changeset_position_middle), - Some(changeset_position_bottom), - ) - } else { - let mut changeset: Option<Changeset> = None; - - for entry in WalkDir::new(changeset_folder) { - let mut changeset_file = File::open(entry?.path())?; - let mut changeset_data = Vec::new(); - changeset_file.read_to_end(&mut changeset_data)?; - - let changesets = parse_changeset(&changeset_data)?; - if changesets.iter().any(|c| c.id == changeset_id) { - changeset = changesets.into_iter().find(|c| c.id == changeset_id); - } - } - - Ok(changeset) } + + Ok(changeset) } diff --git a/templates/README.md b/templates/README.md @@ -1,8 +1,7 @@ # OSM-Mirror This is a mirror of the OpenStreetMap database. It is rebuild from the -<$server_url> server data files and the <$changeset_server_url> -changeset files. +<$server_url> server data files and changeset files. This mirror is autogenerated currently.