use futures::future::try_join_all;
use tracing::{debug, error, info};
use crate::contract::{Preprocessor, ProcessConfig, ProcessInput, Uploader};
extern crate tokio;
#[derive(Debug)]
pub struct SynchroniseConfig {
pub process: ProcessConfig,
}
#[derive(Debug)]
pub struct SynchroniseReport {
pub sources: Vec<ExternalSourceReport>,
}
#[derive(Debug)]
pub struct ExternalSourceReport {
pub source_id: i64,
pub source_name: String,
pub items: Vec<ExternalItemReport>,
}
#[derive(Debug)]
pub struct ExternalItemReport {
pub item_id: i64,
pub item_name: String,
}
pub async fn synchronise<P, U>(
preprocessor: &P,
uploader: &U,
downloaded_sources: &[crate::contract::DownloadedSource],
) -> Result<SynchroniseReport, String>
where
P: Preprocessor + Sync,
U: Uploader + Sync,
{
info!("[SYNC] Starting full synchronisation pipeline");
if let Err(e) = empty_bucket(uploader).await {
error!(error = ?e, "[SYNC][ERROR] Failed to empty bucket before sync");
return Err(format!("Failed to empty bucket before sync: {e:?}"));
}
info!("[SYNC] Emptied bucket before sync");
let mut sources_report: Vec<ExternalSourceReport> = Vec::new();
for downloaded in downloaded_sources {
let process_input = ProcessInput {
name: downloaded.logical_name.clone(),
repo_path: downloaded.local_path.clone(),
};
info!(repo_name = %downloaded.logical_name, "[SYNC] Invoking processing step (process strategy)");
let source_for_upload = match preprocessor.process(process_input).await {
Ok(src) => {
info!(
items = src.external_items.len(),
"[SYNC] Processing succeeded"
);
src
}
Err(e) => {
error!(error = ?e, "[SYNC][ERROR] Process step failed");
return Err(format!("Process step failed: {:?}", e));
}
};
let bucket_id: i32 = std::env::var("BUCKET_ID")
.expect("BUCKET_ID env var must be set for uploader")
.parse()
.expect("BUCKET_ID must be an integer");
let new_source = crate::contract::NewExternalSource {
name: &source_for_upload.name,
bucket_id,
};
info!(source_name = %source_for_upload.name, "[SYNC][UPLOAD] Creating new external source");
let ext_source = match uploader.create_source(new_source).await {
Ok(src) => {
info!(
external_source_id = src.external_source_id,
"[SYNC][UPLOAD] create_source succeeded"
);
src
}
Err(e) => {
error!(error = ?e, "[SYNC][ERROR][UPLOAD] create_source (external source) failed");
return Err(format!("[UPLOAD fail @ create_source]: {e:?}"));
}
};
let mut uploaded_items_report: Vec<ExternalItemReport> = Vec::new();
for ext_item in &source_for_upload.external_items {
info!(filename = %ext_item.filename, "[SYNC][UPLOAD] Preparing upload for file");
let content = String::from_utf8_lossy(&ext_item.content);
let item_req = crate::contract::NewExternalItem {
content: &content,
url: &ext_item.filename,
bucket_id: bucket_id as i64,
external_source_id: ext_source.external_source_id as i64,
processing_state: None,
};
let uploaded = match uploader.create_item(item_req).await {
Ok(resp) => {
info!(file = %ext_item.filename, state = %resp.processing_state, "[SYNC][UPLOAD] create_item succeeded");
match serde_json::to_string_pretty(&resp) {
Ok(json) => {
debug!(json = %json, file = %ext_item.filename, "[SYNC][UPLOAD][DEBUG] Uploaded ExternalItem as JSON")
}
Err(e) => {
error!(file = %ext_item.filename, error = ?e, "[SYNC][UPLOAD][DEBUG] Failed to serialize ExternalItem as JSON")
}
}
uploaded_items_report.push(ExternalItemReport {
item_id: resp.external_item_id as i64,
item_name: ext_item.filename.clone(),
});
resp
}
Err(e) => {
error!(file = %ext_item.filename, error = ?e, "[SYNC][ERROR][UPLOAD] create_item (external item) failed");
return Err(format!(
"[UPLOAD fail @ create_item for file={}]: {e:?}",
ext_item.filename
));
}
};
if uploaded.processing_state != "Submitted" {
error!(file = %ext_item.filename, state = %uploaded.processing_state, "[SYNC][ERROR][UPLOAD] Uploaded item's processing_state was not 'Submitted'");
return Err(format!(
"[UPLOAD fail @ create_item post-state: file={}] Uploaded item's processing_state was not 'Submitted': {:?}",
ext_item.filename, uploaded.processing_state
));
}
}
sources_report.push(ExternalSourceReport {
source_id: ext_source.external_source_id as i64,
source_name: ext_source.external_source_name.clone(),
items: uploaded_items_report,
});
}
Ok(SynchroniseReport {
sources: sources_report,
})
}
pub async fn empty_bucket<C>(client: &C) -> Result<(), Box<dyn std::error::Error + Send + Sync>>
where
C: Uploader,
{
let sources = client.list_sources().await?;
let deletions = sources
.into_iter()
.map(|src| client.delete_source_by_id(src.external_source_id));
try_join_all(deletions).await?;
Ok(())
}