use csv::{Writer, WriterBuilder};
use futures::stream::{self, StreamExt};
use indexmap::IndexMap;
use reqwest::{self, Response};
use std::{fs::File, io::Write};
use tokio::sync::mpsc;
use crate::types::{ExportFormat, VtigerQueryResponse, VtigerResponse};
const LINK_ENDPOINT: &str = "restapi/v1/vtiger/default";
#[derive(Debug)]
pub struct Vtiger {
url: String,
username: String,
access_key: String,
client: reqwest::Client,
}
impl Vtiger {
pub fn new(url: &str, username: &str, access_key: &str) -> Self {
Vtiger {
url: url.to_string(),
username: username.to_string(),
access_key: access_key.to_string(),
client: reqwest::Client::new(),
}
}
async fn get(
&self,
endpoint: &str,
query: &[(&str, &str)],
) -> Result<Response, reqwest::Error> {
let url = format!("{}{}{}", self.url, LINK_ENDPOINT, endpoint);
let mut request = self.client.get(&url);
if !query.is_empty() {
request = request.query(query);
}
request
.basic_auth(&self.username, Some(&self.access_key))
.header("Content-Type", "application/json")
.header("Accept", "application/json")
.send()
.await
}
async fn post(
&self,
endpoint: &str,
query: &[(&str, &str)],
) -> Result<Response, reqwest::Error> {
let url = format!("{}{}{}", self.url, LINK_ENDPOINT, endpoint);
let mut request = self.client.post(&url);
if !query.is_empty() {
request = request.query(query);
}
request
.basic_auth(&self.username, Some(&self.access_key))
.header("Content-Type", "application/json")
.header("Accept", "application/json")
.send()
.await
}
pub async fn me(&self) -> Result<VtigerResponse, reqwest::Error> {
let response = self.get("/me", &[]).await?;
let vtiger_response = response.json::<VtigerResponse>().await?;
Ok(vtiger_response)
}
pub async fn list_types(
&self,
query: &[(&str, &str)],
) -> Result<VtigerResponse, reqwest::Error> {
let response = self.get("/listtypes", query).await?;
let vtiger_response = response.json::<VtigerResponse>().await?;
Ok(vtiger_response)
}
pub async fn describe(&self, module_name: &str) -> Result<VtigerResponse, reqwest::Error> {
let response = self
.get(&format!("/describe"), &[("elementType", module_name)])
.await?;
let vtiger_response = response.json::<VtigerResponse>().await?;
Ok(vtiger_response)
}
pub async fn export(
&self,
module: &str,
query_filter: Option<(&str, Vec<String>)>,
batch_size: usize,
concurrency: usize,
format: Vec<ExportFormat>,
) -> Result<(), Box<dyn std::error::Error>> {
let (column, query_filters) = match query_filter {
Some((column, filters)) => (column, filters),
None => ("", vec![]),
};
let query_filter_counts: Vec<(String, i32)> = stream::iter(query_filters)
.map(|query_filter| async move {
let query = format!(
"SELECT count(*) from {} WHERE {} LIKE '{}-%';",
module, column, query_filter
);
println!("Executing query: {}", query);
let count = match self.query(&query).await {
Ok(result) => match result.result {
Some(records) => {
if let Some(first_record) = records.get(0) {
if let Some(count_val) = first_record.get("count") {
if let Some(count_str) = count_val.as_str() {
count_str.parse::<i32>().unwrap_or(0)
} else {
0
}
} else {
0
}
} else {
0
}
}
None => 0,
},
Err(_) => 0,
};
if count != 0 {
println!("Received {} records for {}", count, query_filter);
}
(query_filter, count)
})
.buffer_unordered(concurrency)
.collect()
.await;
let work_items = query_filter_counts
.into_iter()
.flat_map(|(query_filter, count)| {
(0..count)
.step_by(batch_size.clone())
.map(|offset| (query_filter.clone(), offset, batch_size))
.collect::<Vec<_>>()
})
.collect::<Vec<_>>();
let (record_sender, mut record_receiver) =
mpsc::unbounded_channel::<IndexMap<String, serde_json::Value>>();
let writer_task = tokio::spawn(async move {
let mut file_json_lines: Option<File> = None;
let mut file_json: Option<File> = None;
let mut file_csv: Option<Writer<File>> = None;
if format.contains(&ExportFormat::JsonLines) {
file_json_lines = Some(File::create("output.jsonl").unwrap());
}
if format.contains(&ExportFormat::Json) {
file_json = Some(File::create("output.json").unwrap());
file_json
.as_mut()
.unwrap()
.write_all(b"[")
.expect("Could not write JSON array start");
}
if format.contains(&ExportFormat::CSV) {
file_csv = Some(
WriterBuilder::new()
.quote(b'"')
.quote_style(csv::QuoteStyle::Always)
.escape(b'\\')
.terminator(csv::Terminator::CRLF)
.from_path("output.csv")
.expect("CSV file creation failed"),
);
}
let mut count = 0;
while let Some(record) = record_receiver.recv().await {
if let Some(ref mut file) = file_json_lines {
writeln!(file, "{}", serde_json::to_string(&record).unwrap()).unwrap();
}
if let Some(ref mut file) = file_json {
if count != 0 {
write!(file, ",\n").unwrap();
}
writeln!(file, "{}", serde_json::to_string_pretty(&record).unwrap()).unwrap();
}
if let Some(ref mut file) = file_csv {
if count == 0 {
let header: Vec<&str> = record.keys().map(|k| k.as_str()).collect();
file.write_record(header)
.expect("Could not write CSV header record!");
}
let values: Vec<String> = record
.values()
.map(|v| match v {
serde_json::Value::String(s) => s
.replace('\n', "\\n")
.replace('\r', "\\r")
.replace('\t', "\\t"),
serde_json::Value::Number(s) => s.to_string(),
serde_json::Value::Bool(b) => b.to_string(),
serde_json::Value::Array(_) | serde_json::Value::Object(_) => {
serde_json::to_string(v).unwrap_or_else(|_| String::new())
}
_ => "".to_string(),
})
.collect();
file.write_record(values)
.expect("Could not write CSV values record@");
}
count += 1;
if count % 10000 == 0 {
println!("Processed {} records", count);
}
}
println!("Finished writing {} total records", count);
if let Some(mut file) = file_json_lines {
file.flush().unwrap();
}
if let Some(mut file) = file_json {
write!(file, "\n]").unwrap(); file.flush().unwrap();
}
if let Some(mut file) = file_csv {
file.flush().unwrap();
}
println!("All files flushed and closed");
});
stream::iter(work_items)
.map(|(query_filter, offset, batch_size)| {
let sender = record_sender.clone();
async move {
let query = format!(
"SELECT * FROM {} WHERE {} LIKE '{}%' LIMIT {}, {};",
module, column, query_filter, offset, batch_size
);
println!("Executing query: {}", query);
if let Ok(result) = self.query(&query).await {
if let Some(records) = result.result {
let record_count = records.len();
for record in records {
if let Err(_) = sender.send(record) {
eprintln!("Failed to send record, writer may have stopped");
break;
}
}
println!(
"Sent {} records from {} offset {}",
record_count, query_filter, offset
);
}
}
Ok::<(), Box<dyn std::error::Error + Send + Sync>>(())
}
})
.buffer_unordered(concurrency)
.collect::<Vec<_>>()
.await;
drop(record_sender);
writer_task.await.unwrap();
Ok(())
}
pub async fn create(
&self,
module_name: &str,
fields: &[(&str, &str)],
) -> Result<VtigerResponse, reqwest::Error> {
let fields_map: IndexMap<&str, &str> = fields.iter().cloned().collect();
let element_json = serde_json::to_string(&fields_map)
.expect("Failed to serialize string IndexMap to JSON");
let response = self
.post(
"/create",
&[("elementType", module_name), ("element", &element_json)],
)
.await?;
let vtiger_response = response.json::<VtigerResponse>().await?;
Ok(vtiger_response)
}
pub async fn retrieve(&self, record_id: &str) -> Result<VtigerResponse, reqwest::Error> {
let response = self
.get(&format!("/retrieve"), &[("id", record_id)])
.await?;
let vtiger_response = response.json::<VtigerResponse>().await?;
Ok(vtiger_response)
}
pub async fn query(&self, query: &str) -> Result<VtigerQueryResponse, reqwest::Error> {
let response = self.get(&format!("/query"), &[("query", query)]).await?;
let vtiger_response = response.json::<VtigerQueryResponse>().await?;
Ok(vtiger_response)
}
}