use crate::internal::path_escape;
use crate::runtime::{self, Error};
#[derive(Clone, Debug, Default)]
pub struct ChunkedUploadsListFileUploadSessionPartsOptions {
pub offset: Option<i64>,
pub limit: Option<i64>,
}
#[derive(Clone, Debug, Default)]
pub struct ChunkedUploadsCommitFileUploadSessionOptions {
pub if_match: Option<String>,
pub if_none_match: Option<String>,
}
pub struct ChunkedUploadsListFileUploadSessionPartsPaginator {
manager: ChunkedUploadsManager,
upload_session_id: String,
options: ChunkedUploadsListFileUploadSessionPartsOptions,
buffer: std::vec::IntoIter<crate::models::schemas::UploadPart>,
done: bool,
}
impl ChunkedUploadsListFileUploadSessionPartsPaginator {
pub async fn next(&mut self) -> Option<Result<crate::models::schemas::UploadPart, Error>> {
loop {
if let Some(item) = self.buffer.next() {
return Some(Ok(item));
}
if self.done {
return None;
}
let page = match self
.manager
.list_file_upload_session_parts_page(
self.upload_session_id.clone(),
Some(self.options.clone()),
)
.await
{
Ok(page) => page,
Err(err) => {
self.done = true;
return Some(Err(err));
}
};
let items = page.entries.unwrap_or_default();
if items.is_empty() {
self.done = true;
} else {
let next = self.options.offset.unwrap_or(0) + items.len() as i64;
self.options.offset = Some(next);
}
self.buffer = items.into_iter();
}
}
}
pub struct ChunkedUploadsManager {
session: std::sync::Arc<runtime::Client>,
}
impl ChunkedUploadsManager {
pub(crate) fn new(session: std::sync::Arc<runtime::Client>) -> Self {
Self { session }
}
pub async fn create_file_upload_session(
&self,
body: crate::models::schemas::CreateFileUploadSessionRequest,
) -> Result<crate::models::schemas::UploadSession, Error> {
let mut url = self.session.base_url("upload");
url.push_str("/files");
url.push_str("/upload_sessions");
let mut req = self.session.new_request("POST", &url);
let payload = serde_json::to_vec(&body)?;
req = runtime::with_json_body(req, &payload);
let resp = self.session.fetch(req).await?;
let data = runtime::response_bytes(&resp)?;
Ok(serde_json::from_slice(&data)?)
}
pub async fn create_file_version_upload_session(
&self,
file_id: String,
body: crate::models::schemas::CreateFileVersionUploadSessionRequest,
) -> Result<crate::models::schemas::UploadSession, Error> {
let mut url = self.session.base_url("upload");
url.push_str("/files");
url.push('/');
let seg = path_escape(&file_id);
url.push_str(&seg);
url.push_str("/upload_sessions");
let mut req = self.session.new_request("POST", &url);
let payload = serde_json::to_vec(&body)?;
req = runtime::with_json_body(req, &payload);
let resp = self.session.fetch(req).await?;
let data = runtime::response_bytes(&resp)?;
Ok(serde_json::from_slice(&data)?)
}
pub async fn get_file_upload_session(
&self,
upload_session_id: String,
) -> Result<crate::models::schemas::UploadSession, Error> {
let mut url = self.session.base_url("upload_session");
url.push_str("/files");
url.push_str("/upload_sessions");
url.push('/');
let seg = path_escape(&upload_session_id);
url.push_str(&seg);
let req = self.session.new_request("GET", &url);
let resp = self.session.fetch(req).await?;
let data = runtime::response_bytes(&resp)?;
Ok(serde_json::from_slice(&data)?)
}
pub async fn update_file_upload_session(
&self,
upload_session_id: String,
digest: String,
content_range: String,
body: Vec<u8>,
) -> Result<crate::models::schemas::UploadedPart, Error> {
let mut url = self.session.base_url("upload_session");
url.push_str("/files");
url.push_str("/upload_sessions");
url.push('/');
let seg = path_escape(&upload_session_id);
url.push_str(&seg);
let mut req = self.session.new_request("PUT", &url);
req = runtime::with_header(req, "digest", &digest);
req = runtime::with_header(req, "content-range", &content_range);
req = runtime::with_stream_body(
req,
runtime::Stream::from_bytes(body),
"application/octet-stream",
);
let resp = self.session.fetch(req).await?;
let data = runtime::response_bytes(&resp)?;
Ok(serde_json::from_slice(&data)?)
}
pub async fn delete_file_upload_session(&self, upload_session_id: String) -> Result<(), Error> {
let mut url = self.session.base_url("upload_session");
url.push_str("/files");
url.push_str("/upload_sessions");
url.push('/');
let seg = path_escape(&upload_session_id);
url.push_str(&seg);
let req = self.session.new_request("DELETE", &url);
let _ = self.session.fetch(req).await?;
Ok(())
}
async fn list_file_upload_session_parts_page(
&self,
upload_session_id: String,
opts: Option<ChunkedUploadsListFileUploadSessionPartsOptions>,
) -> Result<crate::models::schemas::UploadParts, Error> {
let mut url = self.session.base_url("upload_session");
url.push_str("/files");
url.push_str("/upload_sessions");
url.push('/');
let seg = path_escape(&upload_session_id);
url.push_str(&seg);
url.push_str("/parts");
let mut req = self.session.new_request("GET", &url);
let opts = opts.unwrap_or_default();
if let Some(value) = opts.offset {
req = runtime::with_query(req, "offset", &value.to_string());
}
if let Some(value) = opts.limit {
req = runtime::with_query(req, "limit", &value.to_string());
}
let resp = self.session.fetch(req).await?;
let data = runtime::response_bytes(&resp)?;
Ok(serde_json::from_slice(&data)?)
}
pub fn list_file_upload_session_parts(
&self,
upload_session_id: String,
opts: Option<ChunkedUploadsListFileUploadSessionPartsOptions>,
) -> ChunkedUploadsListFileUploadSessionPartsPaginator {
ChunkedUploadsListFileUploadSessionPartsPaginator {
manager: ChunkedUploadsManager::new(self.session.clone()),
upload_session_id,
options: opts.unwrap_or_default(),
buffer: Vec::new().into_iter(),
done: false,
}
}
pub async fn commit_file_upload_session(
&self,
upload_session_id: String,
digest: String,
body: crate::models::schemas::CommitFileUploadSessionRequest,
opts: Option<ChunkedUploadsCommitFileUploadSessionOptions>,
) -> Result<crate::models::schemas::Files, Error> {
let mut url = self.session.base_url("upload_session");
url.push_str("/files");
url.push_str("/upload_sessions");
url.push('/');
let seg = path_escape(&upload_session_id);
url.push_str(&seg);
url.push_str("/commit");
let mut req = self.session.new_request("POST", &url);
req = runtime::with_header(req, "digest", &digest);
let opts = opts.unwrap_or_default();
if let Some(value) = opts.if_match {
req = runtime::with_header(req, "if-match", &value);
}
if let Some(value) = opts.if_none_match {
req = runtime::with_header(req, "if-none-match", &value);
}
let payload = serde_json::to_vec(&body)?;
req = runtime::with_json_body(req, &payload);
let resp = self.session.fetch(req).await?;
let data = runtime::response_bytes(&resp)?;
Ok(serde_json::from_slice(&data)?)
}
}