#![cfg(feature = "cloudflare_tasks")]
#![allow(clippy::too_many_arguments, clippy::type_complexity)]
#![allow(clippy::missing_errors_doc, clippy::doc_markdown, clippy::useless_format)]
#![allow(unused_imports)]
use foundation_netio::{DynNetClient, PreparedRequestBuilder};
use foundation_netio::shared::client::http_client::HttpClient;
use serde::{Deserialize, Serialize};
use foundation_macros::JsonHash;
use super::shared::ApiResponse;
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperApiV4Error {
#[serde(flatten)]
pub data: std::collections::HashMap<String, serde_json::Value>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperApiV4Message {
#[serde(flatten)]
pub data: std::collections::HashMap<String, serde_json::Value>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperApiV4Success {
pub errors: Option<Vec<std::collections::HashMap<String, serde_json::Value>>>,
pub messages: Option<Vec<String>>,
pub success: Option<bool>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperCreateJobRequest {
pub overwrite: Option<bool>,
pub source: Option<R2SlurperSourceJobSchema>,
pub target: Option<R2SlurperR2TargetSchema>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperGCSLikeCredsSchema {
#[serde(rename = "clientEmail")]
pub client_email: String,
#[serde(rename = "privateKey")]
pub private_key: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperGCSSourceSchema {
pub bucket: String,
pub keys: Option<Vec<String>>,
#[serde(rename = "pathPrefix")]
pub path_prefix: Option<String>,
pub secret: R2SlurperGCSLikeCredsSchema,
pub vendor: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperJobResponse {
#[serde(rename = "createdAt")]
pub created_at: Option<String>,
#[serde(rename = "finishedAt")]
pub finished_at: Option<String>,
pub id: Option<String>,
pub overwrite: Option<bool>,
pub source: Option<serde_json::Value>,
pub status: Option<R2SlurperJobStatus>,
pub target: Option<std::collections::HashMap<String, serde_json::Value>>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperJobStatus {
#[serde(flatten)]
pub data: std::collections::HashMap<String, serde_json::Value>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperJurisdiction {
#[serde(flatten)]
pub data: std::collections::HashMap<String, serde_json::Value>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperR2SourceSchema {
pub bucket: String,
pub jurisdiction: Option<R2SlurperJurisdiction>,
pub keys: Option<Vec<String>>,
#[serde(rename = "pathPrefix")]
pub path_prefix: Option<String>,
pub secret: R2SlurperS3LikeCredsSchema,
pub vendor: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperR2TargetSchema {
pub bucket: String,
pub jurisdiction: Option<R2SlurperJurisdiction>,
pub secret: R2SlurperS3LikeCredsSchema,
pub vendor: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperS3LikeCredsSchema {
#[serde(rename = "accessKeyId")]
pub access_key_id: String,
#[serde(rename = "secretAccessKey")]
pub secret_access_key: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperS3SourceSchema {
pub bucket: String,
pub endpoint: Option<String>,
pub keys: Option<Vec<String>>,
#[serde(rename = "pathPrefix")]
pub path_prefix: Option<String>,
pub region: Option<String>,
pub secret: R2SlurperS3LikeCredsSchema,
pub vendor: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct R2SlurperSourceJobSchema {
#[serde(flatten)]
pub data: std::collections::HashMap<String, serde_json::Value>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct SlurperAbortAllJobsResponse {
pub result: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct SlurperAbortJobResponse {
pub result: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct SlurperCreateJobResponse {
pub result: Option<std::collections::HashMap<String, serde_json::Value>>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct SlurperGetJobResponse {
pub result: Option<R2SlurperJobResponse>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct SlurperListJobsResponse {
pub result: Option<Vec<R2SlurperJobResponse>>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct SlurperPauseJobResponse {
pub result: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonHash)]
pub struct SlurperResumeJobResponse {
pub result: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize, JsonHash)]
pub struct SlurperListJobsArgs {
pub account_id: String,
pub limit: Option<String>,
pub offset: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize, JsonHash)]
pub struct SlurperCreateJobArgs {
pub account_id: String,
pub body: R2SlurperCreateJobRequest,
}
#[derive(Debug, Clone, Default, Serialize, JsonHash)]
pub struct SlurperAbortAllJobsArgs {
pub account_id: String,
}
#[derive(Debug, Clone, Default, Serialize, JsonHash)]
pub struct SlurperGetJobArgs {
pub account_id: String,
pub job_id: String,
}
#[derive(Debug, Clone, Default, Serialize, JsonHash)]
pub struct SlurperAbortJobArgs {
pub account_id: String,
pub job_id: String,
}
#[derive(Debug, Clone, Default, Serialize, JsonHash)]
pub struct SlurperPauseJobArgs {
pub account_id: String,
pub job_id: String,
}
#[derive(Debug, Clone, Default, Serialize, JsonHash)]
pub struct SlurperResumeJobArgs {
pub account_id: String,
pub job_id: String,
}
pub async fn slurper_list_jobs_request<F>(
client: DynNetClient,
args: &SlurperListJobsArgs,
base_url: &str,
builder_mod: Option<F>,
) -> Result<ApiResponse<SlurperListJobsResponse>, super::shared::ApiError>
where
F: FnOnce(&mut PreparedRequestBuilder),
{
let path = format!("/accounts/{}/slurper/jobs",
args.account_id,
);
let endpoint_url = format!("{}{}", base_url, path);
let mut builder = PreparedRequestBuilder::get(&endpoint_url)
.map_err(|e| super::shared::ApiError::RequestBuildFailed(e.to_string()))?;
builder = builder.query("limit", args.limit.as_deref());
builder = builder.query("offset", args.offset.as_deref());
if let Some(f) = builder_mod {
f(&mut builder);
}
let response = client.send_async(builder.build()).await
.map_err(|e| super::shared::ApiError::RequestSendFailed(e.to_string()))?;
let status: usize = response.get_status().into();
let headers = response.get_headers_ref().clone();
if status < 200 || status >= 300 {
let error_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let body = (!error_bytes.is_empty())
.then(|| String::from_utf8_lossy(&error_bytes).into_owned());
return Err(super::shared::ApiError::HttpStatus { code: status as u16, headers, body });
}
let body_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let parsed: SlurperListJobsResponse = serde_json::from_slice(&body_bytes).map_err(|e: serde_json::Error| super::shared::ApiError::ParseFailed(e.to_string()))?;
Ok(ApiResponse { status: status as u16, headers, body: parsed })
}
pub async fn slurper_create_job_request<F>(
client: DynNetClient,
args: &SlurperCreateJobArgs,
base_url: &str,
builder_mod: Option<F>,
) -> Result<ApiResponse<SlurperCreateJobResponse>, super::shared::ApiError>
where
F: FnOnce(&mut PreparedRequestBuilder),
{
let path = format!("/accounts/{}/slurper/jobs",
args.account_id,
);
let endpoint_url = format!("{}{}", base_url, path);
let mut builder = PreparedRequestBuilder::post(&endpoint_url)
.map_err(|e| super::shared::ApiError::RequestBuildFailed(e.to_string()))?;
builder = builder.body_json(&args.body)
.map_err(|e| super::shared::ApiError::RequestBuildFailed(e.to_string()))?;
if let Some(f) = builder_mod {
f(&mut builder);
}
let response = client.send_async(builder.build()).await
.map_err(|e| super::shared::ApiError::RequestSendFailed(e.to_string()))?;
let status: usize = response.get_status().into();
let headers = response.get_headers_ref().clone();
if status < 200 || status >= 300 {
let error_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let body = (!error_bytes.is_empty())
.then(|| String::from_utf8_lossy(&error_bytes).into_owned());
return Err(super::shared::ApiError::HttpStatus { code: status as u16, headers, body });
}
let body_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let parsed: SlurperCreateJobResponse = serde_json::from_slice(&body_bytes).map_err(|e: serde_json::Error| super::shared::ApiError::ParseFailed(e.to_string()))?;
Ok(ApiResponse { status: status as u16, headers, body: parsed })
}
pub async fn slurper_abort_all_jobs_request<F>(
client: DynNetClient,
args: &SlurperAbortAllJobsArgs,
base_url: &str,
builder_mod: Option<F>,
) -> Result<ApiResponse<SlurperAbortAllJobsResponse>, super::shared::ApiError>
where
F: FnOnce(&mut PreparedRequestBuilder),
{
let path = format!("/accounts/{}/slurper/jobs/abortAll",
args.account_id,
);
let endpoint_url = format!("{}{}", base_url, path);
let mut builder = PreparedRequestBuilder::put(&endpoint_url)
.map_err(|e| super::shared::ApiError::RequestBuildFailed(e.to_string()))?;
if let Some(f) = builder_mod {
f(&mut builder);
}
let response = client.send_async(builder.build()).await
.map_err(|e| super::shared::ApiError::RequestSendFailed(e.to_string()))?;
let status: usize = response.get_status().into();
let headers = response.get_headers_ref().clone();
if status < 200 || status >= 300 {
let error_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let body = (!error_bytes.is_empty())
.then(|| String::from_utf8_lossy(&error_bytes).into_owned());
return Err(super::shared::ApiError::HttpStatus { code: status as u16, headers, body });
}
let body_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let parsed: SlurperAbortAllJobsResponse = serde_json::from_slice(&body_bytes).map_err(|e: serde_json::Error| super::shared::ApiError::ParseFailed(e.to_string()))?;
Ok(ApiResponse { status: status as u16, headers, body: parsed })
}
pub async fn slurper_get_job_request<F>(
client: DynNetClient,
args: &SlurperGetJobArgs,
base_url: &str,
builder_mod: Option<F>,
) -> Result<ApiResponse<SlurperGetJobResponse>, super::shared::ApiError>
where
F: FnOnce(&mut PreparedRequestBuilder),
{
let path = format!("/accounts/{}/slurper/jobs/{}",
args.account_id,
args.job_id,
);
let endpoint_url = format!("{}{}", base_url, path);
let mut builder = PreparedRequestBuilder::get(&endpoint_url)
.map_err(|e| super::shared::ApiError::RequestBuildFailed(e.to_string()))?;
if let Some(f) = builder_mod {
f(&mut builder);
}
let response = client.send_async(builder.build()).await
.map_err(|e| super::shared::ApiError::RequestSendFailed(e.to_string()))?;
let status: usize = response.get_status().into();
let headers = response.get_headers_ref().clone();
if status < 200 || status >= 300 {
let error_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let body = (!error_bytes.is_empty())
.then(|| String::from_utf8_lossy(&error_bytes).into_owned());
return Err(super::shared::ApiError::HttpStatus { code: status as u16, headers, body });
}
let body_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let parsed: SlurperGetJobResponse = serde_json::from_slice(&body_bytes).map_err(|e: serde_json::Error| super::shared::ApiError::ParseFailed(e.to_string()))?;
Ok(ApiResponse { status: status as u16, headers, body: parsed })
}
pub async fn slurper_abort_job_request<F>(
client: DynNetClient,
args: &SlurperAbortJobArgs,
base_url: &str,
builder_mod: Option<F>,
) -> Result<ApiResponse<SlurperAbortJobResponse>, super::shared::ApiError>
where
F: FnOnce(&mut PreparedRequestBuilder),
{
let path = format!("/accounts/{}/slurper/jobs/{}/abort",
args.account_id,
args.job_id,
);
let endpoint_url = format!("{}{}", base_url, path);
let mut builder = PreparedRequestBuilder::put(&endpoint_url)
.map_err(|e| super::shared::ApiError::RequestBuildFailed(e.to_string()))?;
if let Some(f) = builder_mod {
f(&mut builder);
}
let response = client.send_async(builder.build()).await
.map_err(|e| super::shared::ApiError::RequestSendFailed(e.to_string()))?;
let status: usize = response.get_status().into();
let headers = response.get_headers_ref().clone();
if status < 200 || status >= 300 {
let error_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let body = (!error_bytes.is_empty())
.then(|| String::from_utf8_lossy(&error_bytes).into_owned());
return Err(super::shared::ApiError::HttpStatus { code: status as u16, headers, body });
}
let body_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let parsed: SlurperAbortJobResponse = serde_json::from_slice(&body_bytes).map_err(|e: serde_json::Error| super::shared::ApiError::ParseFailed(e.to_string()))?;
Ok(ApiResponse { status: status as u16, headers, body: parsed })
}
pub async fn slurper_pause_job_request<F>(
client: DynNetClient,
args: &SlurperPauseJobArgs,
base_url: &str,
builder_mod: Option<F>,
) -> Result<ApiResponse<SlurperPauseJobResponse>, super::shared::ApiError>
where
F: FnOnce(&mut PreparedRequestBuilder),
{
let path = format!("/accounts/{}/slurper/jobs/{}/pause",
args.account_id,
args.job_id,
);
let endpoint_url = format!("{}{}", base_url, path);
let mut builder = PreparedRequestBuilder::put(&endpoint_url)
.map_err(|e| super::shared::ApiError::RequestBuildFailed(e.to_string()))?;
if let Some(f) = builder_mod {
f(&mut builder);
}
let response = client.send_async(builder.build()).await
.map_err(|e| super::shared::ApiError::RequestSendFailed(e.to_string()))?;
let status: usize = response.get_status().into();
let headers = response.get_headers_ref().clone();
if status < 200 || status >= 300 {
let error_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let body = (!error_bytes.is_empty())
.then(|| String::from_utf8_lossy(&error_bytes).into_owned());
return Err(super::shared::ApiError::HttpStatus { code: status as u16, headers, body });
}
let body_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let parsed: SlurperPauseJobResponse = serde_json::from_slice(&body_bytes).map_err(|e: serde_json::Error| super::shared::ApiError::ParseFailed(e.to_string()))?;
Ok(ApiResponse { status: status as u16, headers, body: parsed })
}
pub async fn slurper_resume_job_request<F>(
client: DynNetClient,
args: &SlurperResumeJobArgs,
base_url: &str,
builder_mod: Option<F>,
) -> Result<ApiResponse<SlurperResumeJobResponse>, super::shared::ApiError>
where
F: FnOnce(&mut PreparedRequestBuilder),
{
let path = format!("/accounts/{}/slurper/jobs/{}/resume",
args.account_id,
args.job_id,
);
let endpoint_url = format!("{}{}", base_url, path);
let mut builder = PreparedRequestBuilder::put(&endpoint_url)
.map_err(|e| super::shared::ApiError::RequestBuildFailed(e.to_string()))?;
if let Some(f) = builder_mod {
f(&mut builder);
}
let response = client.send_async(builder.build()).await
.map_err(|e| super::shared::ApiError::RequestSendFailed(e.to_string()))?;
let status: usize = response.get_status().into();
let headers = response.get_headers_ref().clone();
if status < 200 || status >= 300 {
let error_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let body = (!error_bytes.is_empty())
.then(|| String::from_utf8_lossy(&error_bytes).into_owned());
return Err(super::shared::ApiError::HttpStatus { code: status as u16, headers, body });
}
let body_bytes = foundation_netio::shared::client::body_reader::collect_bytes_from_send_safe(response.take_body());
let parsed: SlurperResumeJobResponse = serde_json::from_slice(&body_bytes).map_err(|e: serde_json::Error| super::shared::ApiError::ParseFailed(e.to_string()))?;
Ok(ApiResponse { status: status as u16, headers, body: parsed })
}