Skip to main content

coven_storage/cloud/google_drive/
storage_impl.rs

1use super::*;
2
3pub(crate) fn parse_drive_file_identities(
4    page: &serde_json::Value,
5) -> Result<Vec<DriveFileIdentity>, CloudHomeError> {
6    let files = page["files"].as_array().ok_or_else(|| {
7        CloudHomeError::Transport("Drive file identity response omitted files".to_string())
8    })?;
9    files
10        .iter()
11        .map(|file| {
12            let id = file["id"]
13                .as_str()
14                .filter(|id| !id.is_empty())
15                .ok_or_else(|| {
16                    CloudHomeError::Transport(
17                        "Drive file identity response omitted a file id".to_string(),
18                    )
19                })?
20                .to_string();
21            let create_token = file["appProperties"][CREATE_TOKEN_PROPERTY]
22                .as_str()
23                .filter(|token| !token.is_empty())
24                .ok_or_else(|| {
25                    CloudHomeError::Transport(format!(
26                        "Drive file {id} omitted its Coven create token"
27                    ))
28                })?
29                .to_string();
30            Ok(DriveFileIdentity { id, create_token })
31        })
32        .collect()
33}
34
35pub(crate) fn select_drive_file(files: &[DriveFileIdentity]) -> Option<&DriveFileIdentity> {
36    files.iter().min_by(|left, right| {
37        left.create_token
38            .cmp(&right.create_token)
39            .then_with(|| left.id.cmp(&right.id))
40    })
41}
42
43pub(crate) fn parse_create_file_id(body: &str, key: &str) -> Result<String, CloudHomeError> {
44    let json: serde_json::Value = serde_json::from_str(body)
45        .map_err(|e| CloudHomeError::transport(format!("create {key}: parse response"), e))?;
46    match json.get("id").and_then(|id| id.as_str()) {
47        Some(id) if !id.is_empty() => Ok(id.to_string()),
48        _ => Err(CloudHomeError::Transport(format!(
49            "create {key}: response missing id"
50        ))),
51    }
52}
53
54pub(crate) fn parse_generated_file_id(
55    response: &serde_json::Value,
56    key: &str,
57) -> Result<String, CloudHomeError> {
58    match response["ids"].as_array() {
59        Some(ids) if ids.len() == 1 => match ids[0].as_str() {
60            Some(id) if !id.is_empty() => Ok(id.to_string()),
61            _ => Err(CloudHomeError::Transport(format!(
62                "generate append id {key}: generated id is empty"
63            ))),
64        },
65        _ => Err(CloudHomeError::Transport(format!(
66            "generate append id {key}: expected exactly one generated id"
67        ))),
68    }
69}
70
71pub(crate) fn create_file_metadata_body(
72    encoded_name: &str,
73    folder_id: &str,
74    create_token: &str,
75) -> String {
76    let mut app_properties = serde_json::Map::new();
77    app_properties.insert(
78        CREATE_TOKEN_PROPERTY.to_string(),
79        serde_json::Value::String(create_token.to_string()),
80    );
81    serde_json::json!({
82        "name": encoded_name,
83        "parents": [folder_id],
84        "appProperties": app_properties,
85    })
86    .to_string()
87}
88
89/// A Drive resumable sink. New objects use a resumable-create session, which
90/// keeps the file absent until the final part commits; `finish` then resolves
91/// concurrent same-name creates by the create token. Existing objects use a
92/// resumable-update session and require no post-commit reconciliation.
93pub(crate) struct DriveMultipartSink<'a> {
94    home: &'a GoogleDriveCloudHome,
95    inner: RangePutSink,
96    key: String,
97    encoded: String,
98    created: Option<DriveFileIdentity>,
99}
100
101#[async_trait]
102impl crate::cloud::PartSink for DriveMultipartSink<'_> {
103    fn part_size(&self) -> usize {
104        self.inner.part_size()
105    }
106
107    async fn send_part(
108        &mut self,
109        part: Bytes,
110        offset: u64,
111        is_last: bool,
112        control: &crate::cloud::UploadControl,
113    ) -> Result<(), CloudHomeError> {
114        self.inner.send_part(part, offset, is_last, control).await
115    }
116
117    async fn abort(&mut self) -> Result<(), CloudHomeError> {
118        self.inner.abort().await
119    }
120
121    async fn finish(mut self: Box<Self>) -> Result<(), CloudHomeError> {
122        Box::new(self.inner).finish().await?;
123        if let Some(created) = self.created.take() {
124            self.home
125                .reconcile_created_file(&self.key, &self.encoded, created)
126                .await?;
127        }
128        Ok(())
129    }
130}
131
132/// First `error.errors[].reason` in a Google API error body (the shape Drive,
133/// Sheets, and other googleapis.com endpoints share), or `None` if the body isn't
134/// that JSON.
135pub(crate) fn parse_google_api_error_reason(body: &str) -> Option<String> {
136    http::error_reason(body, |v| {
137        v.get("error")?
138            .get("errors")?
139            .as_array()?
140            .first()?
141            .get("reason")?
142            .as_str()
143            .map(String::from)
144    })
145}
146
147/// Map a Drive write failure to a `CloudHomeError`. The `storageQuotaExceeded`
148/// reason gets a message naming the provider and the recovery step; everything
149/// else keeps the raw HTTP status + body so transient failures stay debuggable.
150pub(crate) fn classify_write_error(
151    status: reqwest::StatusCode,
152    body: &str,
153    key: &str,
154    op: &str,
155) -> CloudHomeError {
156    if status == reqwest::StatusCode::FORBIDDEN
157        && parse_google_api_error_reason(body).as_deref() == Some("storageQuotaExceeded")
158    {
159        return CloudHomeError::Transport(
160            "Your Google Drive storage is full. Free up space at drive.google.com to keep syncing."
161                .to_string(),
162        );
163    }
164    CloudHomeError::Transport(format!("{op} {key} (HTTP {status}): {body}"))
165}
166
167/// Files at or below this size go up through simple Drive requests; larger files
168/// use a resumable session. Drive accepts a simple upload up to 5 MB.
169pub(crate) const GDRIVE_SIMPLE_UPLOAD_MAX: usize = 4 * 1024 * 1024;
170
171/// Resumable-session part size. Drive requires every part except the last to be a
172/// multiple of 256 KiB; 8 MiB (32 × 256 KiB) keeps the request count low.
173pub(crate) const GDRIVE_CHUNK_SIZE: usize = 8 * 1024 * 1024;
174
175#[async_trait]
176impl OAuthRestHome for GoogleDriveCloudHome {
177    fn not_found(&self) -> NotFound {
178        NotFound::Status
179    }
180
181    async fn send_read(
182        &self,
183        key: &str,
184        range: Option<(u64, u64)>,
185    ) -> Result<reqwest::Response, CloudHomeError> {
186        let file_id = self
187            .find_file_id(&encode_key(key))
188            .await?
189            .ok_or_else(|| CloudHomeError::NotFound(key.to_string()))?;
190        let range = range.map(|(start, end)| crate::cloud::range_header(start, end));
191        self.session
192            .api_call(|oauth| {
193                let mut req =
194                    supports_all_drives(oauth.get(format!("{}/files/{}", self.drive_api, file_id)))
195                        .query(&[("alt", "media")]);
196                if let Some(ref range) = range {
197                    req = req.header("Range", range);
198                }
199                req
200            })
201            .await
202    }
203
204    async fn send_delete(&self, key: &str) -> Result<reqwest::Response, CloudHomeError> {
205        // No file id ⇒ already absent; surface as not-found so `rest_delete` treats
206        // it as success.
207        let file_id = self
208            .find_file_id(&encode_key(key))
209            .await?
210            .ok_or_else(|| CloudHomeError::NotFound(key.to_string()))?;
211        self.session
212            .api_call(|oauth| {
213                supports_all_drives(oauth.delete(format!("{}/files/{}", self.drive_api, file_id)))
214            })
215            .await
216    }
217
218    async fn send_list_page(
219        &self,
220        prefix: &str,
221        cursor: Option<&str>,
222    ) -> Result<reqwest::Response, CloudHomeError> {
223        let query = list_file_query(&self.folder_id, prefix);
224        let page = cursor.map(str::to_string);
225        self.session
226            .api_call(|oauth| {
227                let mut req = supports_all_drives(oauth.get(format!("{}/files", self.drive_api)))
228                    .query(&[
229                        ("q", query.as_str()),
230                        ("fields", "nextPageToken,files(id,name)"),
231                        ("pageSize", "1000"),
232                        ("includeItemsFromAllDrives", "true"),
233                    ]);
234                if let Some(ref pt) = page {
235                    req = req.query(&[("pageToken", pt.as_str())]);
236                }
237                req
238            })
239            .await
240    }
241
242    fn parse_list_page(&self, body: &str, prefix: &str) -> Result<ListPage, CloudHomeError> {
243        let json: serde_json::Value = serde_json::from_str(body)
244            .map_err(|e| CloudHomeError::transport("parse list".to_string(), e))?;
245        let mut slots = Vec::new();
246        if let Some(files) = json["files"].as_array() {
247            for file in files {
248                let (Some(name), Some(id)) = (file["name"].as_str(), file["id"].as_str()) else {
249                    continue;
250                };
251                let Some(decoded) = decode_listed_key("Google Drive", name) else {
252                    continue;
253                };
254                // The `contains` query may match mid-string, so filter to the
255                // actual prefix.
256                if decoded.starts_with(prefix) {
257                    slots.push(ObjectSlot::opaque(decoded, id.to_string())?);
258                }
259            }
260        }
261        Ok(ListPage {
262            slots,
263            next: json["nextPageToken"].as_str().map(String::from),
264        })
265    }
266}
267
268#[async_trait]
269impl CloudHome for GoogleDriveCloudHome {
270    async fn put_object(&self, key: &str, data: Vec<u8>) -> Result<(), CloudHomeError> {
271        let media_body = Bytes::from(data);
272        let encoded = encode_key(key);
273        if let Some(file_id) = self.find_file_id(&encoded).await? {
274            self.upload_file_media(key, &file_id, media_body, "update")
275                .await?;
276        } else {
277            self.create_file_with_media(key, &encoded, media_body)
278                .await?;
279        }
280        Ok(())
281    }
282
283    async fn open_multipart<'a>(
284        &'a self,
285        key: &str,
286        total_len: u64,
287    ) -> Result<BoxPartSink<'a>, CloudHomeError> {
288        let encoded = encode_key(key);
289        let (session_url, created, op) = match self.find_file_id(&encoded).await? {
290            Some(file_id) => (
291                self.open_resumable_update_session(key, &file_id).await?,
292                None,
293                "update",
294            ),
295            None => {
296                let attempt = self.new_append_attempt(key).await?;
297                let session_url = match self.open_resumable_create_session(key, &attempt).await {
298                    Ok(session_url) => session_url,
299                    Err(operation) => {
300                        return match self
301                            .resolve_failed_append(key, &attempt, operation, false)
302                            .await
303                        {
304                            Err(error) => Err(error),
305                            Ok(file_id) => Err(CloudHomeError::Transport(format!(
306                            "open mutable Drive upload {key}: uncommitted session resolved as {}",
307                            file_id
308                        ))),
309                        }
310                    }
311                };
312                (
313                    session_url,
314                    Some(DriveFileIdentity {
315                        id: attempt.file_id,
316                        create_token: attempt.create_token,
317                    }),
318                    "create",
319                )
320            }
321        };
322        let key_owned = key.to_string();
323        let classify =
324            Box::new(move |status, body: &str| classify_write_error(status, body, &key_owned, op));
325        // Drive returns 308 Resume Incomplete for every non-final part.
326        let inner = self.session.range_put_sink(
327            session_url,
328            308,
329            total_len,
330            GDRIVE_CHUNK_SIZE,
331            key.to_string(),
332            classify,
333            drive_upload_cancellation_succeeded,
334        );
335        Ok(Box::new(DriveMultipartSink {
336            home: self,
337            inner,
338            key: key.to_string(),
339            encoded,
340            created,
341        }))
342    }
343
344    fn multipart_threshold(&self) -> u64 {
345        GDRIVE_SIMPLE_UPLOAD_MAX as u64
346    }
347
348    async fn read(&self, key: &str) -> Result<Vec<u8>, CloudHomeError> {
349        rest_read(self, key).await
350    }
351
352    async fn read_range(&self, key: &str, start: u64, end: u64) -> Result<Vec<u8>, CloudHomeError> {
353        rest_read_range(self, key, start, end).await
354    }
355
356    async fn list(&self, prefix: &str) -> Result<Vec<String>, CloudHomeError> {
357        rest_list(self, prefix).await
358    }
359
360    async fn delete(&self, key: &str) -> Result<(), CloudHomeError> {
361        rest_delete(self, key).await
362    }
363
364    async fn exists(&self, key: &str) -> Result<bool, CloudHomeError> {
365        // A name query confirms existence in one request; the generic 2xx/404 rule
366        // doesn't fit (the query returns 200 with an empty array for an absent key).
367        Ok(self.find_file_id(&encode_key(key)).await?.is_some())
368    }
369
370    async fn set_access(
371        &self,
372        desired: CloudAccessState,
373    ) -> Result<CloudAccessOutcome, CloudHomeError> {
374        let email = desired.require_provider_email("Google Drive")?;
375        let list_url = format!(
376            "{}/files/{}/permissions?fields=permissions(id,emailAddress,role),nextPageToken&supportsAllDrives=true",
377            self.drive_api, self.folder_id
378        );
379        let access = sharing::SharedFolderAccess::new(
380            &self.session,
381            list_url.clone(),
382            "permissions",
383            |permission: &serde_json::Value| permission["emailAddress"].as_str().map(String::from),
384            |page: &serde_json::Value| drive_permissions_next_page_url(&list_url, page),
385            |permission_id: &str| {
386                format!(
387                    "{}/files/{}/permissions/{}?supportsAllDrives=true",
388                    self.drive_api, self.folder_id, permission_id
389                )
390            },
391            |permission: &serde_json::Value| permission["role"].as_str() == Some("writer"),
392            "writer",
393            format!(
394                "{}/files/{}/permissions?supportsAllDrives=true",
395                self.drive_api, self.folder_id
396            ),
397            serde_json::json!({
398                "type": "user",
399                "role": "writer",
400                "emailAddress": email,
401            }),
402        );
403        match desired {
404            CloudAccessState::Present { .. } => {
405                access.ensure_present(email).await?;
406                Ok(CloudAccessOutcome::Present(
407                    CloudHomeJoinInfo::GoogleDrive {
408                        folder_id: self.folder_id.clone(),
409                    },
410                ))
411            }
412            CloudAccessState::Absent { .. } => {
413                access.ensure_absent(email).await?;
414                Ok(CloudAccessOutcome::Absent(RevokeOutcome::Revoked))
415            }
416        }
417    }
418}
419
420#[async_trait]
421impl ExactSlotStorage for GoogleDriveCloudHome {
422    async fn provider_binding(
423        &self,
424    ) -> Result<coven_protocol::objects::ResolvedProviderBinding, CloudHomeError> {
425        use coven_protocol::objects::{
426            GoogleDriveCorpus, ProviderDeviceBinding, ProviderPrincipalId, ResolvedProviderBinding,
427            StoreProviderBinding,
428        };
429
430        if self.folder_id.is_empty() {
431            return Err(CloudHomeError::Configuration(
432                "Google Drive provider binding has an empty folder id".to_string(),
433            ));
434        }
435        let folder_response = self
436            .session
437            .api_call(|oauth| {
438                supports_all_drives(
439                    oauth.get(format!("{}/files/{}", self.drive_api, self.folder_id)),
440                )
441                .query(&[("fields", "id,driveId")])
442            })
443            .await?;
444        let folder_response = ensure_ok(
445            folder_response,
446            "resolve Google Drive corpus",
447            NotFound::Status,
448        )
449        .await?;
450        let folder: serde_json::Value =
451            ok_json(folder_response, "parse Google Drive corpus").await?;
452        if folder["id"].as_str() != Some(self.folder_id.as_str()) {
453            return Err(CloudHomeError::Transport(
454                "Google Drive folder lookup returned a different folder id".to_string(),
455            ));
456        }
457        let corpus = match folder.get("driveId") {
458            None | Some(serde_json::Value::Null) => GoogleDriveCorpus::MyDrive {
459                folder_id: self.folder_id.clone(),
460            },
461            Some(value) => GoogleDriveCorpus::SharedDrive {
462                drive_id: value
463                    .as_str()
464                    .filter(|value| !value.is_empty())
465                    .ok_or_else(|| {
466                        CloudHomeError::Transport(
467                            "Google Drive folder returned a malformed drive id".to_string(),
468                        )
469                    })?
470                    .to_string(),
471                folder_id: self.folder_id.clone(),
472            },
473        };
474
475        let about_response = self
476            .session
477            .api_call(|oauth| {
478                oauth
479                    .get(format!("{}/about", self.drive_api))
480                    .query(&[("fields", "user(permissionId)")])
481            })
482            .await?;
483        let about_response = ensure_ok(
484            about_response,
485            "resolve Google Drive principal",
486            NotFound::Status,
487        )
488        .await?;
489        let about: serde_json::Value =
490            ok_json(about_response, "parse Google Drive principal").await?;
491        let permission_id = about["user"]["permissionId"]
492            .as_str()
493            .filter(|value| !value.is_empty())
494            .ok_or_else(|| {
495                CloudHomeError::Transport(
496                    "Google Drive about response omitted the stable permission id".to_string(),
497                )
498            })?
499            .to_string();
500
501        Ok(ResolvedProviderBinding {
502            store: StoreProviderBinding::GoogleDrive { corpus },
503            device: ProviderDeviceBinding {
504                principal: ProviderPrincipalId::GoogleDrive { permission_id },
505            },
506        })
507    }
508
509    async fn allocate_slot(&self, logical_key: &str) -> Result<ObjectSlot, CloudHomeError> {
510        ObjectSlot::opaque(
511            logical_key.to_string(),
512            self.generate_file_id(logical_key).await?,
513        )
514        .map_err(CloudHomeError::from)
515    }
516
517    async fn create_at(
518        &self,
519        upload: &crate::cloud::ExactUpload<'_>,
520        control: &crate::cloud::UploadControl,
521    ) -> Result<crate::cloud::ExactCreateOutcome, CloudHomeError> {
522        if matches!(
523            self.exact_upload_verification,
524            coven_foundation::config::ExactUploadVerification::UploadChecksum
525        ) {
526            return Err(CloudHomeError::Configuration(
527                "Google Drive does not accept a caller-supplied upload checksum".to_string(),
528            ));
529        }
530        let operation = GoogleDriveCloudHome::create_at_slot(
531            self,
532            upload.object().slot(),
533            upload.body().await?,
534            control,
535        )
536        .await;
537        settle_exact_create(operation, |observed| {
538            self.verify_exact_upload(upload, observed)
539        })
540        .await
541    }
542
543    async fn create_versioned_at(
544        &self,
545        _upload: &ExactUpload<'_>,
546        _control: &UploadControl,
547    ) -> Result<ExactCreateOutcome, CloudHomeError> {
548        Err(CloudHomeError::Configuration(
549            "Google Drive v3 does not expose provider-enforced conditional file replacement"
550                .to_string(),
551        ))
552    }
553    /// Drive mints its own file ids, so the listing reports the id it saw
554    /// beside the name rather than deriving a locator from the key.
555    async fn list_slots(&self, prefix: &str) -> Result<Vec<ObjectSlot>, CloudHomeError> {
556        crate::cloud::oauth_rest::rest_list_slots(self, prefix).await
557    }
558
559    async fn read_at(&self, slot: &ObjectSlot) -> Result<Vec<u8>, CloudHomeError> {
560        GoogleDriveCloudHome::read_at_slot(self, slot).await
561    }
562    async fn read_versioned_at(
563        &self,
564        _slot: &ObjectSlot,
565    ) -> Result<crate::cloud::CloudVersionedObject, CloudHomeError> {
566        Err(CloudHomeError::Configuration(
567            "Google Drive v3 does not expose a conditional file-replacement revision".to_string(),
568        ))
569    }
570    async fn replace_at_if_version(
571        &self,
572        _slot: &ObjectSlot,
573        _expected: &crate::cloud::CloudObjectVersion,
574        _bytes: Vec<u8>,
575    ) -> Result<crate::cloud::ConditionalWriteOutcome, CloudHomeError> {
576        Err(CloudHomeError::Configuration(
577            "Google Drive v3 does not support provider-enforced conditional file replacement"
578                .to_string(),
579        ))
580    }
581    async fn read_range_at(
582        &self,
583        slot: &ObjectSlot,
584        start: u64,
585        end: u64,
586    ) -> Result<Vec<u8>, CloudHomeError> {
587        let range = crate::cloud::range_header(start, end);
588        let response = self.send_exact_read(slot, Some(&range)).await?;
589        validated_range_bytes(response, "read exact Drive range", start, end).await
590    }
591    async fn read_at_to_file(
592        &self,
593        slot: &ObjectSlot,
594        destination: &std::path::Path,
595        progress: crate::cloud::DownloadProgress,
596    ) -> Result<(), crate::cloud::CloudFileReadError> {
597        GoogleDriveCloudHome::read_at_slot_to_file(self, slot, destination, progress).await
598    }
599    async fn delete_at(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
600        GoogleDriveCloudHome::delete_at_slot(self, slot).await
601    }
602}
603
604pub(crate) fn drive_permissions_next_page_url(
605    list_url: &str,
606    page: &serde_json::Value,
607) -> Result<Option<String>, CloudHomeError> {
608    let Some(token) = page["nextPageToken"].as_str() else {
609        return Ok(None);
610    };
611    let query = serde_urlencoded::to_string([("pageToken", token)])
612        .map_err(|e| CloudHomeError::transport("encode Drive page token".to_string(), e))?;
613    let separator = if list_url.contains('?') { '&' } else { '?' };
614    Ok(Some(format!("{list_url}{separator}{query}")))
615}