Skip to main content

coven_storage/cloud/
google_drive.rs

1//! Google Drive `CloudHome` implementation.
2//!
3//! Uses the Google Drive REST API v3 with OAuth 2.0 tokens. Files are stored flat
4//! in a single folder — path separators are escaped by
5//! the `key_encoding` helpers). The `read`/`read_range`/`list`/`delete` methods are
6//! the shared `OAuthRestHome` implementations; this file supplies only the Drive
7//! request shapes, the page parser, the upload paths, and sharing.
8
9use 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
128/// Google Drive cloud home backend.
129pub 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/// One Drive file named by the provider id it was given and the create token
139/// this device stamped on it, whether the name came back from a create response
140/// or from listing the folder.
141#[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                // Lets the joiner fetch its account email for OAuth folder sharing.
195                "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    /// Find a file's Google Drive ID by name within our folder.
203    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    /// Open a resumable upload session for an existing Drive file and return its
730    /// session URL (the `Location` header Google returns).
731    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    /// Verify the slot, issue the exact-read GET (`alt=media`, optionally with a
856    /// `Range` header), and check its status — the shared preamble of the three
857    /// exact-read paths, which diverge only in what they do with the response
858    /// body. The Dropbox backend factors its equivalent the same way.
859    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;