1use async_trait::async_trait;
10use bytes::Bytes;
11use futures_util::StreamExt;
12
13use super::exact_upload::settle_exact_create;
14use super::http::{self, ensure_ok, ok_bytes, ok_json, NotFound};
15use super::key_encoding::{decode_listed_key, encode_key};
16use super::oauth_rest::{
17 response_to_file, rest_delete, rest_list, rest_read, rest_read_range, validated_range_bytes,
18 ListPage, OAuthRestHome, PageTokenTracker,
19};
20use super::oauth_session::OAuthSession;
21use super::resumable::RangePutSink;
22use super::{
23 sharing, BlobBody, BoxPartSink, CloudAccessOutcome, CloudAccessState, CloudHome,
24 CloudHomeError, CloudHomeJoinInfo, ExactCreateOutcome, ExactSlotStorage, ExactUpload,
25 RevokeOutcome, UploadControl,
26};
27use crate::oauth::OAuthConfig;
28use coven_foundation::id_provider::{IdRef, UuidProvider};
29use coven_protocol::objects::{ObjectSlot, PhysicalObjectLocator};
30
31const DRIVE_API: &str = "https://www.googleapis.com/drive/v3";
32const UPLOAD_API: &str = "https://www.googleapis.com/upload/drive/v3";
33const CREATE_TOKEN_PROPERTY: &str = "covenCreateToken";
34const LOGICAL_KEY_PROPERTY: &str = "covenLogicalKey";
35const DRIVE_FOLDER_MIME_TYPE: &str = "application/vnd.google-apps.folder";
36
37mod content_hash;
38mod storage_impl;
39use storage_impl::*;
40
41pub(crate) fn supports_all_drives(request: reqwest::RequestBuilder) -> reqwest::RequestBuilder {
42 request.query(&[("supportsAllDrives", "true")])
43}
44
45fn drive_upload_cancellation_succeeded(status: reqwest::StatusCode) -> bool {
46 status.is_success() || status == reqwest::StatusCode::NOT_FOUND || status.as_u16() == 499
47}
48
49fn escape_drive_query_value(value: &str) -> String {
50 value.replace('\\', "\\\\").replace('\'', "\\'")
51}
52
53enum DriveNameMatch {
54 Equals,
55 Contains,
56}
57
58impl DriveNameMatch {
59 fn operator(&self) -> &'static str {
60 match self {
61 Self::Equals => "=",
62 Self::Contains => "contains",
63 }
64 }
65}
66
67fn drive_file_query(
68 folder_id: Option<&str>,
69 name_match: DriveNameMatch,
70 name_value: &str,
71 extra_predicate: Option<&str>,
72) -> String {
73 let mut predicates = Vec::new();
74 if let Some(folder_id) = folder_id {
75 let folder_id = escape_drive_query_value(folder_id);
76 predicates.push(format!("'{folder_id}' in parents"));
77 }
78
79 let name_value = escape_drive_query_value(name_value);
80 predicates.push(format!("name {} '{name_value}'", name_match.operator()));
81
82 if let Some(extra_predicate) = extra_predicate {
83 predicates.push(extra_predicate.to_string());
84 }
85 predicates.push("trashed = false".to_string());
86 predicates.join(" and ")
87}
88
89fn drive_app_property_predicate(key: &str, value: &str) -> String {
90 let key = escape_drive_query_value(key);
91 let value = escape_drive_query_value(value);
92 format!("appProperties has {{ key='{key}' and value='{value}' }}")
93}
94
95fn find_file_query(folder_id: &str, encoded_name: &str) -> String {
96 drive_file_query(Some(folder_id), DriveNameMatch::Equals, encoded_name, None)
97}
98
99fn find_created_file_query(folder_id: &str, encoded_name: &str, create_token: &str) -> String {
100 let app_property = drive_app_property_predicate(CREATE_TOKEN_PROPERTY, create_token);
101 drive_file_query(
102 Some(folder_id),
103 DriveNameMatch::Equals,
104 encoded_name,
105 Some(&app_property),
106 )
107}
108
109fn list_file_query(folder_id: &str, prefix: &str) -> String {
110 let encoded_prefix = encode_key(prefix);
111 drive_file_query(
112 Some(folder_id),
113 DriveNameMatch::Contains,
114 &encoded_prefix,
115 None,
116 )
117}
118
119pub(crate) fn folder_search_query(folder_name: &str) -> String {
120 drive_file_query(
121 None,
122 DriveNameMatch::Equals,
123 folder_name,
124 Some(&format!("mimeType = '{DRIVE_FOLDER_MIME_TYPE}'")),
125 )
126}
127
128pub struct GoogleDriveCloudHome {
130 folder_id: String,
131 drive_api: String,
132 upload_api: String,
133 ids: IdRef,
134 session: OAuthSession,
135 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
136}
137
138#[derive(Clone, Debug, PartialEq, Eq)]
142struct DriveFileIdentity {
143 id: String,
144 create_token: String,
145}
146
147struct DriveAppendAttempt {
148 file_id: String,
149 create_token: String,
150}
151
152enum DriveAppendAttemptState {
153 Absent,
154 Owned,
155 Foreign,
156}
157
158enum DriveSlotState {
159 Absent,
160 Exact(DriveExactMetadata),
161 Foreign,
162}
163
164#[derive(Clone, Debug, PartialEq, Eq)]
165struct DriveExactMetadata {
166 size: u64,
167 md5_checksum: String,
168}
169
170impl GoogleDriveCloudHome {
171 pub fn new(
172 folder_id: String,
173 session: OAuthSession,
174 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
175 ) -> Self {
176 Self {
177 folder_id,
178 drive_api: DRIVE_API.to_string(),
179 upload_api: UPLOAD_API.to_string(),
180 ids: std::sync::Arc::new(UuidProvider),
181 session,
182 exact_upload_verification,
183 }
184 }
185
186 pub(crate) fn oauth_config(creds: crate::oauth::OAuthClientCreds) -> OAuthConfig {
187 OAuthConfig {
188 client_id: creds.client_id,
189 client_secret: creds.client_secret,
190 auth_url: "https://accounts.google.com/o/oauth2/v2/auth".to_string(),
191 token_url: "https://oauth2.googleapis.com/token".to_string(),
192 scopes: vec![
193 "https://www.googleapis.com/auth/drive.file".to_string(),
194 "https://www.googleapis.com/auth/userinfo.email".to_string(),
196 ],
197 redirect_port: 19284,
198 extra_auth_params: vec![("access_type".to_string(), "offline".to_string())],
199 }
200 }
201
202 async fn find_file_id(&self, encoded_name: &str) -> Result<Option<String>, CloudHomeError> {
204 let files = self.list_file_identities(encoded_name).await?;
205 Ok(select_drive_file(&files).map(|file| file.id.clone()))
206 }
207
208 async fn list_file_identities(
209 &self,
210 encoded_name: &str,
211 ) -> Result<Vec<DriveFileIdentity>, CloudHomeError> {
212 let query = find_file_query(&self.folder_id, encoded_name);
213 let mut page_token: Option<String> = None;
214 let mut page_tokens = PageTokenTracker::new("Google Drive file identity listing");
215 let mut files = Vec::new();
216
217 loop {
218 let page = page_token.clone();
219 let resp =
220 self.session
221 .api_call(|oauth| {
222 let mut req =
223 supports_all_drives(oauth.get(format!("{}/files", self.drive_api)))
224 .query(&[
225 ("q", query.as_str()),
226 ("fields", "nextPageToken,files(id,appProperties)"),
227 ("pageSize", "1000"),
228 ("includeItemsFromAllDrives", "true"),
229 ]);
230 if let Some(ref page) = page {
231 req = req.query(&[("pageToken", page.as_str())]);
232 }
233 req
234 })
235 .await?;
236 let resp = ensure_ok(resp, "list files", NotFound::Status).await?;
237 let json: serde_json::Value = ok_json(resp, "parse list response").await?;
238 files.extend(parse_drive_file_identities(&json)?);
239
240 match json["nextPageToken"].as_str() {
241 Some(next) => page_token = Some(page_tokens.record(next)?),
242 None => break,
243 }
244 }
245
246 Ok(files)
247 }
248
249 async fn create_file_metadata(
250 &self,
251 key: &str,
252 encoded: &str,
253 ) -> Result<DriveFileIdentity, CloudHomeError> {
254 let create_token = self.ids.new_id();
255 let metadata = create_file_metadata_body(encoded, &self.folder_id, &create_token);
256 let resp = self
257 .session
258 .api_call(|oauth| {
259 supports_all_drives(oauth.post(format!("{}/files", self.drive_api)))
260 .query(&[("fields", "id")])
261 .header("Content-Type", "application/json; charset=UTF-8")
262 .body(metadata.clone())
263 })
264 .await?;
265 let status = resp.status();
266 if !status.is_success() {
267 return Err(classify_write_error(
268 status,
269 &http::body_text(resp).await,
270 key,
271 "create",
272 ));
273 }
274 let id_error = match resp.text().await {
275 Ok(body) => match parse_create_file_id(&body, key) {
276 Ok(id) => return Ok(DriveFileIdentity { id, create_token }),
277 Err(error) => error,
278 },
279 Err(error) => CloudHomeError::transport(format!("create {key}: read response"), error),
280 };
281 match self.find_created_file_id(encoded, &create_token).await {
282 Ok(Some(file_id)) => match self.delete_created_file(key, &file_id).await {
283 Ok(()) => Err(id_error),
284 Err(delete_error) => Err(CloudHomeError::Transport(format!(
285 "create {key}: metadata response id failure: {id_error}; rollback delete failed: {delete_error}"
286 ))),
287 },
288 Ok(None) => Err(id_error),
289 Err(lookup_error) => Err(CloudHomeError::Transport(format!(
290 "create {key}: metadata response id failure: {id_error}; rollback lookup failed: {lookup_error}"
291 ))),
292 }
293 }
294
295 async fn create_file_for_key(
296 &self,
297 key: &str,
298 encoded: &str,
299 ) -> Result<String, CloudHomeError> {
300 let created = self.create_file_metadata(key, encoded).await?;
301 self.reconcile_created_file(key, encoded, created).await
302 }
303
304 async fn reconcile_created_file(
305 &self,
306 key: &str,
307 encoded: &str,
308 created: DriveFileIdentity,
309 ) -> Result<String, CloudHomeError> {
310 let files = self.list_file_identities(encoded).await?;
311 if !files.contains(&created) {
312 return Err(CloudHomeError::Transport(format!(
313 "create {key}: created file {} with token {} was not returned by duplicate check",
314 created.id, created.create_token
315 )));
316 }
317 let Some(winner) = select_drive_file(&files) else {
318 return Err(CloudHomeError::Transport(format!(
319 "create {key}: created file {} was not returned by duplicate check",
320 created.id
321 )));
322 };
323 let winner_id = winner.id.clone();
324
325 for file in files {
326 if file.id != winner_id {
327 self.delete_created_file(key, &file.id).await?;
328 }
329 }
330
331 Ok(winner_id)
332 }
333
334 async fn find_created_file_id(
335 &self,
336 encoded_name: &str,
337 create_token: &str,
338 ) -> Result<Option<String>, CloudHomeError> {
339 let query = find_created_file_query(&self.folder_id, encoded_name, create_token);
340 let resp = self
341 .session
342 .api_call(|oauth| {
343 supports_all_drives(oauth.get(format!("{}/files", self.drive_api))).query(&[
344 ("q", query.as_str()),
345 ("fields", "files(id)"),
346 ("pageSize", "1"),
347 ("includeItemsFromAllDrives", "true"),
348 ])
349 })
350 .await?;
351 let resp = ensure_ok(resp, "list created files", NotFound::Status).await?;
352 let json: serde_json::Value = ok_json(resp, "parse created file list response").await?;
353 Ok(json["files"]
354 .as_array()
355 .and_then(|files| files.first())
356 .and_then(|first| first["id"].as_str())
357 .map(String::from))
358 }
359
360 async fn upload_file_media(
361 &self,
362 key: &str,
363 file_id: &str,
364 body: Bytes,
365 op: &str,
366 ) -> Result<(), CloudHomeError> {
367 let resp = self
368 .session
369 .api_call(|oauth| {
370 supports_all_drives(oauth.patch(format!(
371 "{}/files/{}?uploadType=media",
372 self.upload_api, file_id
373 )))
374 .header("Content-Type", "application/octet-stream")
375 .body(body.clone())
376 })
377 .await?;
378 let status = resp.status();
379 if !status.is_success() {
380 return Err(classify_write_error(
381 status,
382 &http::body_text(resp).await,
383 key,
384 op,
385 ));
386 }
387 Ok(())
388 }
389
390 async fn create_file_with_media(
391 &self,
392 key: &str,
393 encoded: &str,
394 media_body: Bytes,
395 ) -> Result<(), CloudHomeError> {
396 let file_id = self.create_file_for_key(key, encoded).await?;
397 let upload_error = match self
398 .upload_file_media(key, &file_id, media_body, "create")
399 .await
400 {
401 Ok(()) => return Ok(()),
402 Err(error) => error,
403 };
404 match self.delete_created_file(key, &file_id).await {
405 Ok(()) => Err(upload_error),
406 Err(delete_error) => Err(CloudHomeError::Transport(format!(
407 "create {key}: media upload failed after metadata create: {upload_error}; rollback delete failed: {delete_error}"
408 ))),
409 }
410 }
411
412 async fn delete_created_file(&self, key: &str, file_id: &str) -> Result<(), CloudHomeError> {
413 let resp = self
414 .session
415 .api_call(|oauth| {
416 supports_all_drives(oauth.delete(format!("{}/files/{}", self.drive_api, file_id)))
417 })
418 .await?;
419 let status = resp.status();
420 if status.is_success() || status == reqwest::StatusCode::NOT_FOUND {
421 return Ok(());
422 }
423 Err(CloudHomeError::Transport(format!(
424 "delete created file {key} (HTTP {status}): {}",
425 http::body_text(resp).await
426 )))
427 }
428
429 async fn generate_file_id(&self, key: &str) -> Result<String, CloudHomeError> {
430 let response = self
431 .session
432 .api_call(|oauth| {
433 oauth
434 .get(format!("{}/files/generateIds", self.drive_api))
435 .query(&[("count", "1"), ("space", "drive"), ("type", "files")])
436 })
437 .await?;
438 let response =
439 ensure_ok(response, "generate Drive append file id", NotFound::Status).await?;
440 let json: serde_json::Value = ok_json(response, "parse generated Drive file id").await?;
441 parse_generated_file_id(&json, key)
442 }
443
444 async fn new_append_attempt(&self, key: &str) -> Result<DriveAppendAttempt, CloudHomeError> {
445 let file_id = self.generate_file_id(key).await?;
446 let create_token = self.ids.new_id();
447 if create_token == file_id {
448 return Err(CloudHomeError::Transport(format!(
449 "generate append id {key}: create token equals the provider file id"
450 )));
451 }
452 Ok(DriveAppendAttempt {
453 file_id,
454 create_token,
455 })
456 }
457
458 async fn inspect_append_attempt(
459 &self,
460 key: &str,
461 attempt: &DriveAppendAttempt,
462 ) -> Result<DriveAppendAttemptState, CloudHomeError> {
463 let response = self
464 .session
465 .api_call(|oauth| {
466 supports_all_drives(
467 oauth.get(format!("{}/files/{}", self.drive_api, attempt.file_id)),
468 )
469 .query(&[("fields", "id,appProperties,trashed")])
470 })
471 .await?;
472 if response.status() == reqwest::StatusCode::NOT_FOUND {
473 return Ok(DriveAppendAttemptState::Absent);
474 }
475 let response = ensure_ok(response, "inspect failed Drive append", NotFound::Status).await?;
476 let json: serde_json::Value = ok_json(response, "parse failed Drive append").await?;
477 if json["id"].as_str() != Some(attempt.file_id.as_str()) {
478 return Err(CloudHomeError::Transport(format!(
479 "inspect append {key}: exact file response did not identify {}",
480 attempt.file_id
481 )));
482 }
483 Ok(
484 if json["appProperties"][CREATE_TOKEN_PROPERTY].as_str()
485 == Some(attempt.create_token.as_str())
486 {
487 DriveAppendAttemptState::Owned
488 } else {
489 DriveAppendAttemptState::Foreign
490 },
491 )
492 }
493
494 async fn resolve_failed_append(
495 &self,
496 key: &str,
497 attempt: &DriveAppendAttempt,
498 operation: CloudHomeError,
499 may_have_committed: bool,
500 ) -> Result<String, CloudHomeError> {
501 match self.inspect_append_attempt(key, attempt).await {
502 Ok(DriveAppendAttemptState::Absent) => Err(operation),
503 Ok(DriveAppendAttemptState::Foreign) => {
504 Err(CloudHomeError::AlreadyExists(key.to_string()))
505 }
506 Ok(DriveAppendAttemptState::Owned) if may_have_committed => Ok(attempt.file_id.clone()),
507 Ok(DriveAppendAttemptState::Owned) => {
508 match self.delete_created_file(key, &attempt.file_id).await {
509 Ok(()) => Err(operation),
510 Err(cleanup) => Err(CloudHomeError::CleanupFailed {
511 operation: Box::new(operation),
512 cleanup: Box::new(cleanup),
513 }),
514 }
515 }
516 Err(verification) => Err(CloudHomeError::CleanupFailed {
517 operation: Box::new(operation),
518 cleanup: Box::new(verification),
519 }),
520 }
521 }
522
523 fn validate_slot<'a>(&self, slot: &'a ObjectSlot) -> Result<&'a str, CloudHomeError> {
524 slot.validate()?;
525 match slot.physical() {
526 PhysicalObjectLocator::Opaque(file_id) => Ok(file_id),
527 PhysicalObjectLocator::LogicalKey => Err(CloudHomeError::Configuration(format!(
528 "Google Drive slot for {} requires an opaque file id",
529 slot.logical_key()
530 ))),
531 }
532 }
533
534 async fn inspect_slot(&self, slot: &ObjectSlot) -> Result<DriveSlotState, CloudHomeError> {
535 let file_id = self.validate_slot(slot)?;
536 let response =
537 self.session
538 .api_call(|oauth| {
539 supports_all_drives(oauth.get(format!("{}/files/{file_id}", self.drive_api)))
540 .query(&[(
541 "fields",
542 "id,name,parents,trashed,appProperties,size,md5Checksum",
543 )])
544 })
545 .await?;
546 if response.status() == reqwest::StatusCode::NOT_FOUND {
547 return Ok(DriveSlotState::Absent);
548 }
549 let response = ensure_ok(
550 response,
551 &format!("inspect exact {}", slot.logical_key()),
552 NotFound::Status,
553 )
554 .await?;
555 let metadata: serde_json::Value =
556 ok_json(response, "parse exact Drive file metadata").await?;
557 let expected_name = encode_key(slot.logical_key());
558 let id_matches = metadata["id"].as_str() == Some(file_id);
559 let name_matches = metadata["name"].as_str() == Some(expected_name.as_str());
560 let parent_matches = metadata["parents"].as_array().is_some_and(|parents| {
561 parents
562 .iter()
563 .any(|parent| parent.as_str() == Some(&self.folder_id))
564 });
565 let logical_key_matches =
566 metadata["appProperties"][LOGICAL_KEY_PROPERTY].as_str() == Some(slot.logical_key());
567 let is_live = metadata["trashed"].as_bool() == Some(false);
568 if id_matches && name_matches && parent_matches && logical_key_matches && is_live {
569 let size = metadata["size"]
570 .as_str()
571 .and_then(|size| size.parse::<u64>().ok())
572 .ok_or_else(|| {
573 CloudHomeError::Transport(format!(
574 "exact Drive metadata for {} omitted size",
575 slot.logical_key()
576 ))
577 })?;
578 let md5_checksum = metadata["md5Checksum"]
579 .as_str()
580 .filter(|hash| !hash.is_empty())
581 .ok_or_else(|| {
582 CloudHomeError::Transport(format!(
583 "exact Drive metadata for {} omitted md5Checksum",
584 slot.logical_key()
585 ))
586 })?
587 .to_string();
588 Ok(DriveSlotState::Exact(DriveExactMetadata {
589 size,
590 md5_checksum,
591 }))
592 } else {
593 Ok(DriveSlotState::Foreign)
594 }
595 }
596
597 async fn verify_slot(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
598 match self.inspect_slot(slot).await? {
599 DriveSlotState::Exact(_) => Ok(()),
600 DriveSlotState::Absent => Err(CloudHomeError::NotFound(slot.logical_key().to_string())),
601 DriveSlotState::Foreign => Err(CloudHomeError::Transport(format!(
602 "exact Drive slot for {} does not identify its allocated file in folder {}",
603 slot.logical_key(),
604 self.folder_id
605 ))),
606 }
607 }
608
609 async fn verify_exact_upload(
610 &self,
611 upload: &super::ExactUpload<'_>,
612 created_response_was_observed: bool,
613 ) -> Result<(), CloudHomeError> {
614 use coven_foundation::config::ExactUploadVerification;
615
616 match self.exact_upload_verification {
617 ExactUploadVerification::UploadChecksum => Err(CloudHomeError::Configuration(
618 "Google Drive does not accept a caller-supplied upload checksum".to_string(),
619 )),
620 ExactUploadVerification::MetadataHash => {
621 let metadata = match self.inspect_slot(upload.object().slot()).await? {
622 DriveSlotState::Absent => {
623 return Err(CloudHomeError::NotFound(
624 upload.object().slot().logical_key().to_string(),
625 ));
626 }
627 DriveSlotState::Foreign => {
628 return Err(CloudHomeError::SlotCollision(
629 upload.object().slot().logical_key().to_string(),
630 ));
631 }
632 DriveSlotState::Exact(metadata) => metadata,
633 };
634 let expected_md5 = content_hash::md5(upload).await?;
635 if metadata.size != upload.object().stored_size()
636 || metadata.md5_checksum != expected_md5
637 {
638 return Err(CloudHomeError::SlotCollision(
639 upload.object().slot().logical_key().to_string(),
640 ));
641 }
642 Ok(())
643 }
644 ExactUploadVerification::Readback => {
645 let bytes = self.read_at_slot(upload.object().slot()).await?;
646 upload.verify_stored_bytes(&bytes)
647 }
648 ExactUploadVerification::Unchecked => {
649 super::exact_upload::accept_unchecked_create_response(
650 created_response_was_observed,
651 upload.object(),
652 )
653 }
654 }
655 }
656
657 async fn create_small_at(
658 &self,
659 slot: &ObjectSlot,
660 data: Vec<u8>,
661 control: &super::UploadControl,
662 ) -> Result<(), CloudHomeError> {
663 use sha2::{Digest, Sha256};
664
665 let file_id = self.validate_slot(slot)?;
666 let boundary = format!(
667 "coven-exact-{}",
668 hex::encode(Sha256::digest(slot.logical_key().as_bytes()))
669 );
670 let metadata = serde_json::json!({
671 "id": file_id,
672 "name": encode_key(slot.logical_key()),
673 "parents": [self.folder_id],
674 "appProperties": { (LOGICAL_KEY_PROPERTY): slot.logical_key() },
675 })
676 .to_string();
677 let prefix = Bytes::from(format!(
678 "--{boundary}\r\nContent-Type: application/json; charset=UTF-8\r\n\r\n{metadata}\r\n--{boundary}\r\nContent-Type: application/octet-stream\r\n\r\n"
679 ));
680 let data = Bytes::from(data);
681 let suffix = Bytes::from(format!("\r\n--{boundary}--\r\n"));
682 let content_length = prefix.len() + data.len() + suffix.len();
683 let response = self
684 .session
685 .api_call(|oauth| {
686 let prefix_control = control.clone();
687 let prefix = prefix.clone();
688 let prefix = futures_util::stream::once(async move {
689 prefix_control.wait_until_resumed().await;
690 Ok::<_, std::io::Error>(prefix)
691 });
692 let payload = control.clone().stream_part(data.clone(), 0);
693 let suffix_control = control.clone();
694 let suffix = suffix.clone();
695 let suffix = futures_util::stream::once(async move {
696 suffix_control.wait_until_resumed().await;
697 Ok::<_, std::io::Error>(suffix)
698 });
699 let body = reqwest::Body::wrap_stream(prefix.chain(payload).chain(suffix));
700 supports_all_drives(oauth.post(format!(
701 "{}/files?uploadType=multipart&fields=id",
702 self.upload_api
703 )))
704 .header(
705 "Content-Type",
706 format!("multipart/related; boundary={boundary}"),
707 )
708 .header("Content-Length", content_length)
709 .body(body)
710 })
711 .await?;
712 let status = response.status();
713 if status.is_success() {
714 return Ok(());
715 }
716 if status == reqwest::StatusCode::CONFLICT {
717 return Err(CloudHomeError::AlreadyExists(
718 slot.logical_key().to_string(),
719 ));
720 }
721 Err(classify_write_error(
722 status,
723 &http::body_text(response).await,
724 slot.logical_key(),
725 "create exact",
726 ))
727 }
728
729 async fn open_resumable_update_session(
732 &self,
733 key: &str,
734 file_id: &str,
735 ) -> Result<String, CloudHomeError> {
736 let url = format!("{}/files/{}?uploadType=resumable", self.upload_api, file_id);
737 let resp = self
738 .session
739 .api_call(|oauth| {
740 supports_all_drives(oauth.patch(&url))
741 .header("Content-Type", "application/json; charset=UTF-8")
742 .body("{}")
743 })
744 .await?;
745
746 let status = resp.status();
747 if !status.is_success() {
748 return Err(classify_write_error(
749 status,
750 &http::body_text(resp).await,
751 key,
752 "update",
753 ));
754 }
755 resp.headers()
756 .get(reqwest::header::LOCATION)
757 .and_then(|v| v.to_str().ok())
758 .map(String::from)
759 .ok_or_else(|| {
760 CloudHomeError::Transport(format!(
761 "resumable session {key}: no Location header returned"
762 ))
763 })
764 }
765
766 async fn open_resumable_create_session(
767 &self,
768 key: &str,
769 attempt: &DriveAppendAttempt,
770 ) -> Result<String, CloudHomeError> {
771 let metadata = serde_json::json!({
772 "id": attempt.file_id,
773 "name": encode_key(key),
774 "parents": [self.folder_id],
775 "appProperties": {
776 (CREATE_TOKEN_PROPERTY): attempt.create_token,
777 (LOGICAL_KEY_PROPERTY): key,
778 },
779 });
780 let response = self
781 .session
782 .api_call(|oauth| {
783 supports_all_drives(oauth.post(format!(
784 "{}/files?uploadType=resumable&fields=id",
785 self.upload_api
786 )))
787 .json(&metadata)
788 })
789 .await?;
790 let status = response.status();
791 if !status.is_success() {
792 return Err(classify_write_error(
793 status,
794 &http::body_text(response).await,
795 key,
796 "append resumable create",
797 ));
798 }
799 let Some(location) = response.headers().get(reqwest::header::LOCATION) else {
800 return Err(CloudHomeError::Transport(format!(
801 "append resumable create {key}: no Location header returned"
802 )));
803 };
804 let location = location.to_str().map_err(|error| {
805 CloudHomeError::transport(
806 format!("read append resumable create Location header for {key}"),
807 error,
808 )
809 })?;
810 if location.is_empty() {
811 return Err(CloudHomeError::Transport(format!(
812 "append resumable create {key}: empty Location header returned"
813 )));
814 }
815 Ok(location.to_string())
816 }
817
818 async fn create_at_slot(
819 &self,
820 slot: &ObjectSlot,
821 body: BlobBody,
822 control: &super::UploadControl,
823 ) -> Result<(), CloudHomeError> {
824 if body.len() <= self.multipart_threshold() {
825 return self
826 .create_small_at(slot, body.collect().await?, control)
827 .await;
828 }
829 let file_id = self.validate_slot(slot)?.to_string();
830 let attempt = DriveAppendAttempt {
831 file_id,
832 create_token: format!("exact:{}", slot.logical_key()),
833 };
834 let session_url = self
835 .open_resumable_create_session(slot.logical_key(), &attempt)
836 .await?;
837 let key = slot.logical_key().to_string();
838 let classify = Box::new(move |status, response: &str| {
839 classify_write_error(status, response, &key, "create exact")
840 });
841 let sink = self.session.range_put_sink(
842 session_url,
843 308,
844 body.len(),
845 GDRIVE_CHUNK_SIZE,
846 slot.logical_key().to_string(),
847 classify,
848 drive_upload_cancellation_succeeded,
849 );
850 super::blob_body::MultipartUpload::new(slot.logical_key(), body, Box::new(sink), control)
851 .run()
852 .await
853 }
854
855 async fn send_exact_read(
860 &self,
861 slot: &ObjectSlot,
862 range: Option<&str>,
863 ) -> Result<reqwest::Response, CloudHomeError> {
864 self.verify_slot(slot).await?;
865 let file_id = self.validate_slot(slot)?.to_string();
866 let response = self
867 .session
868 .api_call(|oauth| {
869 let request =
870 supports_all_drives(oauth.get(format!("{}/files/{file_id}", self.drive_api)))
871 .query(&[("alt", "media")]);
872 match range {
873 Some(range) => request.header("Range", range),
874 None => request,
875 }
876 })
877 .await?;
878 ensure_ok(
879 response,
880 &format!("read exact {}", slot.logical_key()),
881 NotFound::Status,
882 )
883 .await
884 }
885
886 async fn read_at_slot(&self, slot: &ObjectSlot) -> Result<Vec<u8>, CloudHomeError> {
887 let response = self.send_exact_read(slot, None).await?;
888 ok_bytes(
889 response,
890 &format!("read exact body for {}", slot.logical_key()),
891 )
892 .await
893 }
894
895 async fn read_at_slot_to_file(
896 &self,
897 slot: &ObjectSlot,
898 destination: &std::path::Path,
899 progress: super::DownloadProgress,
900 ) -> Result<(), super::CloudFileReadError> {
901 let response = self.send_exact_read(slot, None).await?;
902 response_to_file(
903 response,
904 destination,
905 &format!("read exact body for {}", slot.logical_key()),
906 progress,
907 )
908 .await
909 }
910
911 async fn delete_at_slot(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
912 let file_id = self.validate_slot(slot)?.to_string();
913 match self.verify_slot(slot).await {
914 Ok(()) => self.delete_created_file(slot.logical_key(), &file_id).await,
915 Err(CloudHomeError::NotFound(_)) => Ok(()),
916 Err(error) => Err(error),
917 }
918 }
919
920 #[cfg(test)]
921 fn with_endpoints(mut self, drive_api: String, upload_api: String) -> Self {
922 self.drive_api = drive_api;
923 self.upload_api = upload_api;
924 self
925 }
926
927 #[cfg(test)]
928 fn with_ids(mut self, ids: IdRef) -> Self {
929 self.ids = ids;
930 self
931 }
932}
933
934#[cfg(test)]
935mod tests;