1use 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#[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 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 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 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 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 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 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
209pub 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 let l = GdriveFlatLister::new(path, self.core.clone(), ctx.clone());
371 Ok(TwoWays::Two(l))
372 } else {
373 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 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}