Skip to main content

opendal_service_gdrive/
backend.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use std::fmt::Debug;
19use std::sync::Arc;
20
21use bytes::Buf;
22use http::StatusCode;
23
24use super::core::GdriveFile;
25use super::core::GdriveRecentPathState;
26use super::core::normalize_dir_path;
27use super::core::parse_error;
28use super::core::{ErrorContext, GdriveCore};
29use super::deleter::GdriveDeleter;
30use super::lister::GdriveFlatLister;
31use super::lister::GdriveLister;
32use super::reader::*;
33use super::writer::GdriveWriter;
34use opendal_core::raw::*;
35use opendal_core::*;
36
37use asyncband::mutex::Mutex;
38use log::debug;
39
40use super::GDRIVE_SCHEME;
41use super::config::GdriveConfig;
42use super::core::GdrivePathQuery;
43use super::core::GdriveSigner;
44use super::path_index::GdrivePathIndex;
45
46/// [GoogleDrive](https://drive.google.com/) backend support.
47#[derive(Default)]
48#[doc = include_str!("docs.md")]
49pub struct GdriveBuilder {
50    pub(super) config: GdriveConfig,
51}
52
53impl Debug for GdriveBuilder {
54    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
55        f.debug_struct("GdriveBuilder")
56            .field("config", &self.config)
57            .finish_non_exhaustive()
58    }
59}
60
61impl GdriveBuilder {
62    /// Set root path of GoogleDrive folder.
63    pub fn root(mut self, root: &str) -> Self {
64        self.config.root = if root.is_empty() {
65            None
66        } else {
67            Some(root.to_string())
68        };
69
70        self
71    }
72
73    /// Access token is used for temporary access to the GoogleDrive API.
74    ///
75    /// You can get the access token from [GoogleDrive App Console](https://console.cloud.google.com/apis/credentials)
76    /// or [GoogleDrive OAuth2 Playground](https://developers.google.com/oauthplayground/)
77    ///
78    /// # Note
79    ///
80    /// - An access token is valid for 1 hour.
81    /// - If you want to use the access token for a long time,
82    ///   you can use the refresh token to get a new access token.
83    pub fn access_token(mut self, access_token: &str) -> Self {
84        self.config.access_token = Some(access_token.to_string());
85        self
86    }
87
88    /// Refresh token is used for long term access to the GoogleDrive API.
89    ///
90    /// You can get the refresh token via OAuth 2.0 Flow of GoogleDrive API.
91    ///
92    /// OpenDAL will use this refresh token to get a new access token when the old one is expired.
93    pub fn refresh_token(mut self, refresh_token: &str) -> Self {
94        self.config.refresh_token = Some(refresh_token.to_string());
95        self
96    }
97
98    /// Set the client id for GoogleDrive.
99    ///
100    /// This is required for OAuth 2.0 Flow to refresh the access token.
101    pub fn client_id(mut self, client_id: &str) -> Self {
102        self.config.client_id = Some(client_id.to_string());
103        self
104    }
105
106    /// Set the client secret for GoogleDrive.
107    ///
108    /// This is required for OAuth 2.0 Flow with refresh the access token.
109    pub fn client_secret(mut self, client_secret: &str) -> Self {
110        self.config.client_secret = Some(client_secret.to_string());
111        self
112    }
113}
114
115impl Builder for GdriveBuilder {
116    type Config = GdriveConfig;
117
118    fn build(self) -> Result<impl Service> {
119        let root = normalize_root(&self.config.root.unwrap_or_default());
120        debug!("backend use root {root}");
121
122        let info = ServiceInfo::new(GDRIVE_SCHEME, &root, "");
123        let capability = Capability {
124            stat: true,
125
126            read: true,
127            read_with_suffix: true,
128
129            list: true,
130            list_with_recursive: true,
131
132            write: true,
133
134            create_dir: true,
135            delete: true,
136            delete_with_recursive: true,
137            rename: true,
138            copy: true,
139
140            shared: true,
141
142            ..Default::default()
143        };
144
145        let accessor_info = info;
146        let mut signer = GdriveSigner::new();
147        match (self.config.access_token, self.config.refresh_token) {
148            (Some(access_token), None) => {
149                signer.access_token = access_token;
150                // We will never expire user specified access token.
151                signer.expires_in = Timestamp::MAX;
152            }
153            (None, Some(refresh_token)) => {
154                let client_id = self.config.client_id.ok_or_else(|| {
155                    Error::new(
156                        ErrorKind::ConfigInvalid,
157                        "client_id must be set when refresh_token is set",
158                    )
159                    .with_context("service", GDRIVE_SCHEME)
160                })?;
161                let client_secret = self.config.client_secret.ok_or_else(|| {
162                    Error::new(
163                        ErrorKind::ConfigInvalid,
164                        "client_secret must be set when refresh_token is set",
165                    )
166                    .with_context("service", GDRIVE_SCHEME)
167                })?;
168
169                signer.refresh_token = refresh_token;
170                signer.client_id = client_id;
171                signer.client_secret = client_secret;
172            }
173            (Some(_), Some(_)) => {
174                return Err(Error::new(
175                    ErrorKind::ConfigInvalid,
176                    "access_token and refresh_token cannot be set at the same time",
177                )
178                .with_context("service", GDRIVE_SCHEME));
179            }
180            (None, None) => {
181                return Err(Error::new(
182                    ErrorKind::ConfigInvalid,
183                    "access_token or refresh_token must be set",
184                )
185                .with_context("service", GDRIVE_SCHEME));
186            }
187        };
188
189        let signer = Arc::new(Mutex::new(signer));
190
191        Ok(GdriveBackend {
192            core: Arc::new(GdriveCore {
193                info: accessor_info.clone(),
194                capability,
195                root,
196                signer: signer.clone(),
197                path_index: GdrivePathIndex::new(GdrivePathQuery::new(signer)),
198                recent_entries: Mutex::default(),
199            }),
200        })
201    }
202}
203
204#[derive(Clone, Debug)]
205pub struct GdriveBackend {
206    pub core: Arc<GdriveCore>,
207}
208
209/// Lister type that supports both recursive and non-recursive listing
210pub type GdriveListers = TwoWays<oio::PageLister<GdriveLister>, GdriveFlatLister>;
211
212impl Service for GdriveBackend {
213    type Reader = oio::StreamReader<GdriveReader>;
214    type Writer = oio::OneShotWriter<GdriveWriter>;
215    type Lister = GdriveListers;
216    type Deleter = oio::OneShotDeleter<GdriveDeleter>;
217    type Copier = oio::OneShotCopier;
218    type Composer = ();
219
220    fn info(&self) -> ServiceInfo {
221        self.core.info.clone()
222    }
223
224    fn capability(&self) -> Capability {
225        self.core.capability
226    }
227
228    async fn create_dir(
229        &self,
230        ctx: &OperationContext,
231        path: &str,
232        _args: OpCreateDir,
233    ) -> Result<RpCreateDir> {
234        let path = build_abs_path(&self.core.root, path);
235        let dir_id = self.core.ensure_dir(ctx, &path).await?;
236        let metadata = MetadataBuilder::dir().build();
237
238        self.core.cache_dir_id(&path, &dir_id).await;
239        self.core.record_recent_upsert(&path, metadata).await;
240
241        Ok(RpCreateDir::default())
242    }
243
244    async fn stat(&self, ctx: &OperationContext, path: &str, _args: OpStat) -> Result<RpStat> {
245        let path = build_abs_path(&self.core.root, path);
246
247        match self.core.recent_entry_for_path(&path).await {
248            GdriveRecentPathState::Present(metadata) => {
249                if metadata.mode().is_dir() && path.ends_with('/') {
250                    return Ok(RpStat::new(*metadata));
251                }
252            }
253            GdriveRecentPathState::Deleted => {
254                return Err(Error::new(
255                    ErrorKind::NotFound,
256                    format!("path not found: {path}"),
257                ));
258            }
259            GdriveRecentPathState::Missing => {}
260        }
261
262        let mut file_id = match self.core.resolve_path(ctx, &path).await? {
263            Some(id) => id,
264            None => match self.core.resolve_path_after_refresh(ctx, &path).await? {
265                Some(id) => id,
266                None => {
267                    return Err(Error::new(
268                        ErrorKind::NotFound,
269                        format!("path not found: {path}"),
270                    ));
271                }
272            },
273        };
274        let mut resp = self.core.gdrive_stat_by_id(ctx, &file_id).await?;
275
276        if resp.status() == StatusCode::NOT_FOUND {
277            file_id = self
278                .core
279                .resolve_path_after_refresh(ctx, &path)
280                .await?
281                .ok_or(Error::new(
282                    ErrorKind::NotFound,
283                    format!("path not found: {path}"),
284                ))?;
285            resp = self.core.gdrive_stat_by_id(ctx, &file_id).await?;
286        }
287
288        if resp.status() != StatusCode::OK {
289            return Err(parse_error(
290                ErrorContext::new(ServiceOperation("GetFile")),
291                resp,
292            ));
293        }
294
295        let bs = resp.into_body();
296        let gdrive_file: GdriveFile =
297            serde_json::from_reader(bs.reader()).map_err(new_json_deserialize_error)?;
298
299        let file_type = if gdrive_file.mime_type == "application/vnd.google-apps.folder" {
300            EntryMode::DIR
301        } else {
302            EntryMode::FILE
303        };
304        let mut meta = if file_type == EntryMode::FILE {
305            let size = gdrive_file.size.as_deref().ok_or_else(|| {
306                Error::new(
307                    ErrorKind::Unexpected,
308                    "gdrive stat response does not contain file size",
309                )
310            })?;
311            MetadataBuilder::file(size.parse::<u64>().map_err(|e| {
312                Error::new(ErrorKind::Unexpected, "parse content length").set_source(e)
313            })?)
314        } else {
315            MetadataBuilder::dir()
316        };
317        meta.content_type(gdrive_file.mime_type);
318        if let Some(v) = gdrive_file.modified_time {
319            meta.last_modified(v.parse::<Timestamp>().map_err(|e| {
320                Error::new(ErrorKind::Unexpected, "parse last modified time").set_source(e)
321            })?);
322        }
323        Ok(RpStat::new(meta.build()))
324    }
325    fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
326        let output: oio::StreamReader<GdriveReader> = {
327            Ok(oio::StreamReader::new(GdriveReader::new(
328                self.clone(),
329                ctx.clone(),
330                path,
331                args,
332            )))
333        }?;
334
335        Ok(output)
336    }
337
338    fn write(&self, ctx: &OperationContext, path: &str, _: OpWrite) -> Result<Self::Writer> {
339        let output: oio::OneShotWriter<GdriveWriter> = {
340            let path = build_abs_path(&self.core.root, path);
341
342            Ok(oio::OneShotWriter::new(GdriveWriter::new(
343                self.core.clone(),
344                ctx.clone(),
345                path,
346                None,
347            )))
348        }?;
349
350        Ok(output)
351    }
352
353    fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
354        let output: oio::OneShotDeleter<GdriveDeleter> = {
355            Ok(oio::OneShotDeleter::new(GdriveDeleter::new(
356                self.core.clone(),
357                ctx.clone(),
358            )))
359        }?;
360
361        Ok(output)
362    }
363
364    fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
365        let output: GdriveListers = {
366            let path = build_abs_path(&self.core.root, path);
367
368            if args.recursive() {
369                // Use optimized batch-query recursive lister
370                let l = GdriveFlatLister::new(path, self.core.clone(), ctx.clone());
371                Ok(TwoWays::Two(l))
372            } else {
373                // Use standard page-based lister for non-recursive
374                let l = GdriveLister::new(path, self.core.clone(), ctx.clone());
375                Ok(TwoWays::One(oio::PageLister::new(l)))
376            }
377        }?;
378
379        Ok(output)
380    }
381
382    fn copy(
383        &self,
384        ctx: &OperationContext,
385        from: &str,
386        to: &str,
387        _args: OpCopy,
388    ) -> Result<Self::Copier> {
389        let core = self.core.clone();
390        let ctx = ctx.clone();
391        let from = from.to_string();
392        let to = to.to_string();
393
394        Ok(oio::OneShotCopier::new(async move {
395            let source = build_abs_path(&core.root, &from);
396            let target = build_abs_path(&core.root, &to);
397            let resp = core.gdrive_copy(&ctx, &from, &to).await?;
398
399            match resp.status() {
400                StatusCode::OK => {
401                    let body = resp.into_body();
402                    let meta: GdriveFile = serde_json::from_reader(body.reader())
403                        .map_err(new_json_deserialize_error)?;
404
405                    let to_path = build_abs_path(&core.root, &to);
406                    let is_dir = meta.mime_type == "application/vnd.google-apps.folder";
407                    let metadata = if is_dir {
408                        MetadataBuilder::dir()
409                    } else {
410                        let size = meta.size.as_deref().ok_or_else(|| {
411                            Error::new(
412                                ErrorKind::Unexpected,
413                                "gdrive copy response does not contain file size",
414                            )
415                        })?;
416                        MetadataBuilder::file(size.parse::<u64>().map_err(|e| {
417                            Error::new(ErrorKind::Unexpected, "parse content length").set_source(e)
418                        })?)
419                    };
420
421                    if is_dir {
422                        core.cache_dir_id(&to_path, &meta.id).await;
423                    } else {
424                        core.cache_file_id(&to_path, &meta.id).await;
425                    }
426                    let metadata = metadata.build();
427                    core.record_recent_upsert(&to_path, metadata.clone()).await;
428
429                    Ok(metadata)
430                }
431                StatusCode::NOT_FOUND => {
432                    core.refresh_path(&source).await;
433                    core.refresh_path(&target).await;
434                    let resp = core.gdrive_copy(&ctx, &from, &to).await?;
435                    match resp.status() {
436                        StatusCode::OK => {
437                            let body = resp.into_body();
438                            let meta: GdriveFile = serde_json::from_reader(body.reader())
439                                .map_err(new_json_deserialize_error)?;
440
441                            let to_path = build_abs_path(&core.root, &to);
442                            let is_dir = meta.mime_type == "application/vnd.google-apps.folder";
443                            let metadata = if is_dir {
444                                MetadataBuilder::dir()
445                            } else {
446                                let size = meta.size.as_deref().ok_or_else(|| {
447                                    Error::new(
448                                        ErrorKind::Unexpected,
449                                        "gdrive copy response does not contain file size",
450                                    )
451                                })?;
452                                MetadataBuilder::file(size.parse::<u64>().map_err(|e| {
453                                    Error::new(ErrorKind::Unexpected, "parse content length")
454                                        .set_source(e)
455                                })?)
456                            };
457
458                            if is_dir {
459                                core.cache_dir_id(&to_path, &meta.id).await;
460                            } else {
461                                core.cache_file_id(&to_path, &meta.id).await;
462                            }
463                            let metadata = metadata.build();
464                            core.record_recent_upsert(&to_path, metadata.clone()).await;
465
466                            Ok(metadata)
467                        }
468                        _ => Err(parse_error(
469                            ErrorContext::new(ServiceOperation("CopyFile")),
470                            resp,
471                        )),
472                    }
473                }
474                _ => Err(parse_error(
475                    ErrorContext::new(ServiceOperation("CopyFile")),
476                    resp,
477                )),
478            }
479        }))
480    }
481
482    async fn rename(
483        &self,
484        ctx: &OperationContext,
485        from: &str,
486        to: &str,
487        _args: OpRename,
488    ) -> Result<RpRename> {
489        let source = build_abs_path(&self.core.root, from);
490        let target = build_abs_path(&self.core.root, to);
491
492        // rename will overwrite `to`, delete it if exist
493        self.core.trash_path_if_exists(ctx, &target).await?;
494
495        let resp = self
496            .core
497            .gdrive_patch_metadata_request(ctx, &source, &target)
498            .await?;
499
500        let status = resp.status();
501
502        match status {
503            StatusCode::OK => {
504                let body = resp.into_body();
505                let meta: GdriveFile =
506                    serde_json::from_reader(body.reader()).map_err(new_json_deserialize_error)?;
507
508                let source_path = if meta.mime_type == "application/vnd.google-apps.folder" {
509                    normalize_dir_path(&build_abs_path(&self.core.root, from))
510                } else {
511                    build_abs_path(&self.core.root, from)
512                };
513                let target_path = if meta.mime_type == "application/vnd.google-apps.folder" {
514                    normalize_dir_path(&build_abs_path(&self.core.root, to))
515                } else {
516                    build_abs_path(&self.core.root, to)
517                };
518                let is_dir = meta.mime_type == "application/vnd.google-apps.folder";
519                let metadata = if is_dir {
520                    MetadataBuilder::dir()
521                } else {
522                    let size = meta.size.as_deref().ok_or_else(|| {
523                        Error::new(
524                            ErrorKind::Unexpected,
525                            "gdrive move response does not contain file size",
526                        )
527                    })?;
528                    MetadataBuilder::file(size.parse::<u64>().map_err(|e| {
529                        Error::new(ErrorKind::Unexpected, "parse content length").set_source(e)
530                    })?)
531                };
532
533                if is_dir {
534                    self.core.invalidate_dir_id(&source_path).await;
535                    self.core.cache_dir_id(&target_path, &meta.id).await;
536                } else {
537                    self.core.invalidate_file_id(&source_path).await;
538                    self.core.cache_file_id(&target_path, &meta.id).await;
539                }
540                self.core
541                    .record_recent_delete(
542                        &source_path,
543                        if is_dir {
544                            EntryMode::DIR
545                        } else {
546                            EntryMode::FILE
547                        },
548                    )
549                    .await;
550                self.core
551                    .record_recent_upsert(&target_path, metadata.build())
552                    .await;
553
554                Ok(RpRename::default())
555            }
556            StatusCode::NOT_FOUND => {
557                self.core.refresh_path(&source).await;
558                self.core.refresh_path(&target).await;
559
560                let resp = self
561                    .core
562                    .gdrive_patch_metadata_request(ctx, &source, &target)
563                    .await?;
564
565                if resp.status() != StatusCode::OK {
566                    return Err(parse_error(
567                        ErrorContext::new(ServiceOperation("MoveFile")),
568                        resp,
569                    ));
570                }
571
572                let body = resp.into_body();
573                let meta: GdriveFile =
574                    serde_json::from_reader(body.reader()).map_err(new_json_deserialize_error)?;
575
576                let source_path = if meta.mime_type == "application/vnd.google-apps.folder" {
577                    normalize_dir_path(&build_abs_path(&self.core.root, from))
578                } else {
579                    build_abs_path(&self.core.root, from)
580                };
581                let target_path = if meta.mime_type == "application/vnd.google-apps.folder" {
582                    normalize_dir_path(&build_abs_path(&self.core.root, to))
583                } else {
584                    build_abs_path(&self.core.root, to)
585                };
586                let is_dir = meta.mime_type == "application/vnd.google-apps.folder";
587                let metadata = if is_dir {
588                    MetadataBuilder::dir()
589                } else {
590                    let size = meta.size.as_deref().ok_or_else(|| {
591                        Error::new(
592                            ErrorKind::Unexpected,
593                            "gdrive move response does not contain file size",
594                        )
595                    })?;
596                    MetadataBuilder::file(size.parse::<u64>().map_err(|e| {
597                        Error::new(ErrorKind::Unexpected, "parse content length").set_source(e)
598                    })?)
599                };
600
601                if is_dir {
602                    self.core.invalidate_dir_id(&source_path).await;
603                    self.core.cache_dir_id(&target_path, &meta.id).await;
604                } else {
605                    self.core.invalidate_file_id(&source_path).await;
606                    self.core.cache_file_id(&target_path, &meta.id).await;
607                }
608                self.core
609                    .record_recent_delete(
610                        &source_path,
611                        if is_dir {
612                            EntryMode::DIR
613                        } else {
614                            EntryMode::FILE
615                        },
616                    )
617                    .await;
618                self.core
619                    .record_recent_upsert(&target_path, metadata.build())
620                    .await;
621
622                Ok(RpRename::default())
623            }
624            _ => Err(parse_error(
625                ErrorContext::new(ServiceOperation("MoveFile")),
626                resp,
627            )),
628        }
629    }
630
631    async fn presign(
632        &self,
633        _ctx: &OperationContext,
634        _path: &str,
635        _args: OpPresign,
636    ) -> Result<RpPresign> {
637        Err(Error::new(
638            ErrorKind::Unsupported,
639            "operation is not supported",
640        ))
641    }
642}