1use async_trait::async_trait;
10use bytes::Bytes;
11use reqwest::StatusCode;
12use std::fmt::Write as _;
13use std::sync::OnceLock;
14
15#[path = "dropbox/content_hash.rs"]
16mod content_hash;
17
18use super::http::{self, ensure_ok, exists_from_response, NotFound};
19use super::oauth_rest::{
20 rest_delete, rest_list, rest_read, rest_read_range, ListPage, OAuthRestHome,
21};
22use super::oauth_session::OAuthSession;
23use super::{
24 combine_cleanup_failure, BoxPartSink, CloudAccessOutcome, CloudAccessState, CloudHome,
25 CloudHomeError, CloudHomeJoinInfo, CloudObjectVersion, CloudVersionedObject,
26 ConditionalWriteOutcome, ExactSlotStorage, RevokeOutcome,
27};
28use crate::oauth::OAuthConfig;
29use coven_protocol::objects::ObjectSlot;
30use tracing::warn;
31
32#[derive(Clone, Debug, PartialEq, Eq)]
33struct DropboxExactMetadata {
34 size: u64,
35 content_hash: String,
36 version: CloudObjectVersion,
37}
38
39const API_BASE: &str = "https://api.dropboxapi.com/2";
40const CONTENT_BASE: &str = "https://content.dropboxapi.com/2";
41
42pub struct DropboxCloudHome {
44 folder_path: String,
46 api_base: String,
47 content_base: String,
48 session: OAuthSession,
49 namespace_id: OnceLock<String>,
50 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
51}
52
53impl DropboxCloudHome {
54 pub fn new(
55 folder_path: String,
56 session: OAuthSession,
57 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
58 ) -> Self {
59 Self {
60 folder_path,
61 api_base: API_BASE.to_string(),
62 content_base: CONTENT_BASE.to_string(),
63 session,
64 namespace_id: OnceLock::new(),
65 exact_upload_verification,
66 }
67 }
68
69 fn join_info(&self) -> CloudHomeJoinInfo {
74 CloudHomeJoinInfo::Dropbox {
75 folder_path: self.folder_path.clone(),
76 }
77 }
78
79 pub(crate) fn oauth_config(creds: crate::oauth::OAuthClientCreds) -> OAuthConfig {
80 OAuthConfig {
81 client_id: creds.client_id,
82 client_secret: creds.client_secret,
83 auth_url: "https://www.dropbox.com/oauth2/authorize".to_string(),
84 token_url: "https://api.dropboxapi.com/oauth2/token".to_string(),
85 scopes: vec![],
92 redirect_port: 19284,
93 extra_auth_params: vec![("token_access_type".to_string(), "offline".to_string())],
94 }
95 }
96
97 fn namespace_path(key: &str) -> String {
98 format!("/{key}")
99 }
100
101 fn path_root_header(namespace_id: &str) -> String {
102 serde_json::json!({
103 ".tag": "namespace_id",
104 "namespace_id": namespace_id,
105 })
106 .to_string()
107 }
108
109 fn scoped_request(
110 request: reqwest::RequestBuilder,
111 path_root: &str,
112 ) -> reqwest::RequestBuilder {
113 request.header("Dropbox-API-Path-Root", path_root)
114 }
115
116 fn remember_namespace(&self, namespace_id: String) -> Result<String, CloudHomeError> {
117 if namespace_id.is_empty() {
118 return Err(CloudHomeError::Transport(
119 "Dropbox shared namespace id is empty".to_string(),
120 ));
121 }
122 if let Some(current) = self.namespace_id.get() {
123 if current != &namespace_id {
124 return Err(CloudHomeError::Transport(
125 "Dropbox shared namespace changed during one adapter session".to_string(),
126 ));
127 }
128 return Ok(current.clone());
129 }
130 self.namespace_id
131 .set(namespace_id.clone())
132 .map_err(|actual| {
133 CloudHomeError::Transport(format!(
134 "Dropbox shared namespace changed while being resolved: {actual}"
135 ))
136 })?;
137 Ok(namespace_id)
138 }
139
140 async fn start_upload_session(&self, key: &str) -> Result<String, CloudHomeError> {
141 let namespace_id = self.get_or_create_shared_folder_id().await?;
142 let path_root = Self::path_root_header(&namespace_id);
143 let response = self
144 .session
145 .api_call(|oauth| {
146 Self::scoped_request(
147 oauth.post(format!("{}/files/upload_session/start", self.content_base)),
148 &path_root,
149 )
150 .header("Dropbox-API-Arg", r#"{"close":false}"#)
151 .header("Content-Type", "application/octet-stream")
152 .body(Vec::new())
153 })
154 .await?;
155 let status = response.status();
156 let body = http::body_text(response).await;
157 if !status.is_success() {
158 return Err(classify_write_error(status, &body, key));
159 }
160 let json: serde_json::Value = serde_json::from_str(&body).map_err(|error| {
161 CloudHomeError::transport(format!("parse upload session {key}"), error)
162 })?;
163 json["session_id"]
164 .as_str()
165 .filter(|id| !id.is_empty())
166 .map(str::to_string)
167 .ok_or_else(|| {
168 CloudHomeError::Transport(format!("upload session {key}: no session_id returned"))
169 })
170 }
171
172 async fn append_small(
173 &self,
174 key: &str,
175 data: Vec<u8>,
176 control: &super::UploadControl,
177 ) -> Result<(), CloudHomeError> {
178 let namespace_id = self.get_or_create_shared_folder_id().await?;
179 let path_root = Self::path_root_header(&namespace_id);
180 let body = Bytes::from(data);
181 let length = body.len();
182 let api_arg = dropbox_api_arg(&serde_json::json!({
183 "path": Self::namespace_path(key),
184 "mode": { ".tag": "add" },
185 "autorename": false,
186 "strict_conflict": true,
187 "mute": true,
188 }));
189 let response = self
190 .session
191 .api_call_no_transient_retry(|oauth| {
192 let request_body =
193 reqwest::Body::wrap_stream(control.clone().stream_part(body.clone(), 0));
194 Self::scoped_request(
195 oauth.post(format!("{}/files/upload", self.content_base)),
196 &path_root,
197 )
198 .header("Dropbox-API-Arg", &api_arg)
199 .header("Content-Type", "application/octet-stream")
200 .header("Content-Length", length)
201 .body(request_body)
202 })
203 .await?;
204 self.validate_create_response(key, response).await.map(drop)
205 }
206
207 async fn validate_create_response(
208 &self,
209 key: &str,
210 response: reqwest::Response,
211 ) -> Result<serde_json::Value, CloudHomeError> {
212 let status = response.status();
213 let body = http::body_text(response).await;
214 if !status.is_success() {
215 return Err(classify_write_error(status, &body, key));
216 }
217 let json: serde_json::Value = serde_json::from_str(&body).map_err(|error| {
218 CloudHomeError::transport(format!("append {key}: parse response"), error)
219 })?;
220 json["id"]
221 .as_str()
222 .filter(|id| !id.is_empty())
223 .ok_or_else(|| {
224 CloudHomeError::Transport(format!("append {key}: response missing file id"))
225 })?;
226 if json["path_display"].as_str() != Some(Self::namespace_path(key).as_str()) {
227 return Err(CloudHomeError::Transport(format!(
228 "append {key}: response identifies another Dropbox path"
229 )));
230 }
231 Ok(json)
232 }
233
234 async fn send_exact_read(
235 &self,
236 slot: &ObjectSlot,
237 ) -> Result<(reqwest::Response, DropboxExactMetadata), CloudHomeError> {
238 slot.require_logical_key_for("Dropbox")?;
239 let namespace_id = self.get_or_create_shared_folder_id().await?;
240 let path_root = Self::path_root_header(&namespace_id);
241 let path = Self::namespace_path(slot.logical_key());
242 let api_arg = dropbox_api_arg(&serde_json::json!({"path": path}));
243 let response = self
244 .session
245 .api_call(|oauth| {
246 Self::scoped_request(
247 oauth.post(format!("{}/files/download", self.content_base)),
248 &path_root,
249 )
250 .header("Dropbox-API-Arg", &api_arg)
251 })
252 .await?;
253 let response = ensure_ok(response, "read exact Dropbox file", self.not_found()).await?;
254 let metadata = response
255 .headers()
256 .get("Dropbox-API-Result")
257 .ok_or_else(|| {
258 CloudHomeError::Transport(format!(
259 "exact Dropbox read for {} omitted Dropbox-API-Result",
260 slot.logical_key()
261 ))
262 })?
263 .to_str()
264 .map_err(|error| {
265 CloudHomeError::transport(
266 format!("read exact Dropbox metadata for {}", slot.logical_key()),
267 error,
268 )
269 })?;
270 let metadata: serde_json::Value = serde_json::from_str(metadata).map_err(|error| {
271 CloudHomeError::transport(
272 format!("parse exact Dropbox metadata for {}", slot.logical_key()),
273 error,
274 )
275 })?;
276 let metadata = self.exact_metadata_from_json(slot, &metadata)?;
277 Ok((response, metadata))
278 }
279
280 fn exact_metadata_from_json(
281 &self,
282 slot: &ObjectSlot,
283 metadata: &serde_json::Value,
284 ) -> Result<DropboxExactMetadata, CloudHomeError> {
285 slot.require_logical_key_for("Dropbox")?;
286 let expected_path = Self::namespace_path(slot.logical_key());
287 let matches = metadata[".tag"].as_str() == Some("file")
288 && metadata["id"].as_str().is_some_and(|id| !id.is_empty())
289 && metadata["path_display"].as_str() == Some(expected_path.as_str());
290 if !matches {
291 return Err(CloudHomeError::Transport(format!(
292 "exact Dropbox slot for {} does not identify {expected_path}",
293 slot.logical_key(),
294 )));
295 }
296 let size = metadata["size"].as_u64().ok_or_else(|| {
297 CloudHomeError::Transport(format!(
298 "exact Dropbox metadata for {} omitted size",
299 slot.logical_key()
300 ))
301 })?;
302 let content_hash = metadata["content_hash"]
303 .as_str()
304 .filter(|hash| !hash.is_empty())
305 .ok_or_else(|| {
306 CloudHomeError::Transport(format!(
307 "exact Dropbox metadata for {} omitted content_hash",
308 slot.logical_key()
309 ))
310 })?
311 .to_string();
312 let version = metadata["rev"]
313 .as_str()
314 .filter(|revision| !revision.is_empty())
315 .ok_or_else(|| {
316 CloudHomeError::Transport(format!(
317 "exact Dropbox metadata for {} omitted rev",
318 slot.logical_key()
319 ))
320 })?;
321 Ok(DropboxExactMetadata {
322 size,
323 content_hash,
324 version: CloudObjectVersion::from_provider(version.to_string())?,
325 })
326 }
327
328 async fn exact_metadata(
329 &self,
330 slot: &ObjectSlot,
331 ) -> Result<DropboxExactMetadata, CloudHomeError> {
332 slot.require_logical_key_for("Dropbox")?;
333 let namespace_id = self.get_or_create_shared_folder_id().await?;
334 let path_root = Self::path_root_header(&namespace_id);
335 let body = serde_json::json!({"path": Self::namespace_path(slot.logical_key())});
336 let response = self
337 .session
338 .api_call(|oauth| {
339 Self::scoped_request(
340 oauth.post(format!("{}/files/get_metadata", self.api_base)),
341 &path_root,
342 )
343 .json(&body)
344 })
345 .await?;
346 let response = ensure_ok(response, "verify exact Dropbox file", self.not_found()).await?;
347 let metadata: serde_json::Value =
348 http::ok_json(response, "parse exact Dropbox metadata").await?;
349 self.exact_metadata_from_json(slot, &metadata)
350 }
351
352 async fn verify_slot(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
353 self.exact_metadata(slot).await.map(drop)
354 }
355
356 async fn verify_exact_upload(
357 &self,
358 upload: &super::ExactUpload<'_>,
359 created_response_was_observed: bool,
360 ) -> Result<(), CloudHomeError> {
361 use coven_foundation::config::ExactUploadVerification;
362
363 match self.exact_upload_verification {
364 ExactUploadVerification::UploadChecksum => Err(CloudHomeError::Configuration(
365 "Dropbox does not accept a caller-supplied upload checksum".to_string(),
366 )),
367 ExactUploadVerification::MetadataHash => {
368 let metadata = self.exact_metadata(upload.object().slot()).await?;
369 let expected_hash = content_hash::for_upload(upload).await?;
370 if metadata.size != upload.object().stored_size()
371 || metadata.content_hash != expected_hash
372 {
373 return Err(CloudHomeError::SlotCollision(
374 upload.object().slot().logical_key().to_string(),
375 ));
376 }
377 Ok(())
378 }
379 ExactUploadVerification::Readback => {
380 let bytes = self.read_at(upload.object().slot()).await?;
381 upload.verify_stored_bytes(&bytes)
382 }
383 ExactUploadVerification::Unchecked => {
384 super::exact_upload::accept_unchecked_create_response(
385 created_response_was_observed,
386 upload.object(),
387 )
388 }
389 }
390 }
391
392 async fn get_or_create_shared_folder_id(&self) -> Result<String, CloudHomeError> {
395 if let Some(namespace_id) = self.namespace_id.get() {
396 return Ok(namespace_id.clone());
397 }
398 let share_body = serde_json::json!({ "path": self.folder_path });
399 let resp = self
400 .session
401 .api_call(|oauth| {
402 oauth
403 .post(format!("{}/sharing/share_folder", self.api_base))
404 .json(&share_body)
405 })
406 .await?;
407
408 let status = resp.status();
409 let resp_body = http::body_text(resp).await;
410 let json: serde_json::Value = serde_json::from_str(&resp_body).map_err(|e| {
411 CloudHomeError::transport(
412 format!("parse share-folder response (HTTP {status}): {resp_body}"),
413 e,
414 )
415 })?;
416
417 if let Some(id) = json["shared_folder_id"].as_str() {
419 return self.remember_namespace(id.to_string());
420 }
421 if let Some(id) = json["error"]["shared_folder_metadata"]["shared_folder_id"].as_str() {
423 return self.remember_namespace(id.to_string());
424 }
425 if let Some(job_id) = json["async_job_id"].as_str() {
427 let namespace_id = self.poll_share_job(job_id).await?;
428 return self.remember_namespace(namespace_id);
429 }
430 if !status.is_success() {
431 return Err(CloudHomeError::Transport(format!(
432 "share folder (HTTP {status}): {resp_body}"
433 )));
434 }
435 Err(CloudHomeError::Transport(
436 "could not determine shared_folder_id".to_string(),
437 ))
438 }
439
440 async fn poll_share_job(&self, job_id: &str) -> Result<String, CloudHomeError> {
442 self.poll_dropbox_job(
443 job_id,
444 "sharing/check_share_job_status",
445 "share folder",
446 None,
447 |json| {
448 json["shared_folder_id"]
449 .as_str()
450 .map(String::from)
451 .ok_or_else(|| {
452 CloudHomeError::Transport(
453 "share job completed but no shared_folder_id".to_string(),
454 )
455 })
456 },
457 )
458 .await
459 }
460
461 async fn poll_remove_member_job(
462 &self,
463 job_id: &str,
464 shared_folder_id: &str,
465 ) -> Result<(), CloudHomeError> {
466 let path_root = Self::path_root_header(shared_folder_id);
467 self.poll_dropbox_job(
468 job_id,
469 "sharing/check_remove_member_job_status",
470 "remove folder member",
471 Some(&path_root),
472 |_| Ok(()),
473 )
474 .await
475 }
476
477 async fn poll_dropbox_job<T>(
478 &self,
479 job_id: &str,
480 endpoint: &str,
481 operation: &str,
482 path_root: Option<&str>,
483 complete: impl Fn(&serde_json::Value) -> Result<T, CloudHomeError>,
484 ) -> Result<T, CloudHomeError> {
485 let request_body = serde_json::json!({ "async_job_id": job_id });
486 for _ in 0..30 {
487 tokio::time::sleep(std::time::Duration::from_secs(1)).await;
488 let resp = self
489 .session
490 .api_call(|oauth| {
491 let request = oauth
492 .post(format!("{}/{endpoint}", self.api_base))
493 .json(&request_body);
494 match path_root {
495 Some(path_root) => Self::scoped_request(request, path_root),
496 None => request,
497 }
498 })
499 .await?;
500 let status = resp.status();
501 let resp_body = http::body_text(resp).await;
502 if !status.is_success() {
503 return Err(CloudHomeError::Transport(format!(
504 "{operation} job status (HTTP {status}): {resp_body}"
505 )));
506 }
507 let json: serde_json::Value = serde_json::from_str(&resp_body).map_err(|e| {
508 CloudHomeError::transport(
509 format!("parse {operation} job status response: {resp_body}"),
510 e,
511 )
512 })?;
513 match json[".tag"].as_str() {
514 Some("complete") => return complete(&json),
515 Some("failed") => {
516 return Err(CloudHomeError::Transport(format!(
517 "{operation} job failed: {resp_body}"
518 )));
519 }
520 Some("in_progress") => continue,
521 _ => {
522 return Err(CloudHomeError::Transport(format!(
523 "{operation} job status returned an unexpected tag: {resp_body}"
524 )));
525 }
526 }
527 }
528 Err(CloudHomeError::Transport(format!(
529 "{operation} timed out after 30 seconds"
530 )))
531 }
532
533 async fn folder_member_access(
534 &self,
535 shared_folder_id: &str,
536 email: &str,
537 ) -> Result<Option<String>, CloudHomeError> {
538 let path_root = Self::path_root_header(shared_folder_id);
539 let mut endpoint = "sharing/list_folder_members";
540 let mut request = serde_json::json!({
541 "shared_folder_id": shared_folder_id,
542 "include_inherited": false,
543 "limit": 1000,
544 });
545 loop {
546 let resp = self
547 .session
548 .api_call(|oauth| {
549 Self::scoped_request(
550 oauth.post(format!("{}/{endpoint}", self.api_base)),
551 &path_root,
552 )
553 .json(&request)
554 })
555 .await?;
556 let resp = ensure_ok(resp, "list folder members", self.not_found()).await?;
557 let body: serde_json::Value = http::ok_json(resp, "parse folder members").await?;
558 for (array, identity_field) in [("users", "user"), ("invitees", "invitee")] {
559 if let Some(access) = body[array].as_array().and_then(|members| {
560 members.iter().find_map(|member| {
561 let member_email = member[identity_field]["email"].as_str()?;
562 member_email
563 .eq_ignore_ascii_case(email)
564 .then(|| member["access_type"][".tag"].as_str().map(str::to_string))?
565 })
566 }) {
567 return Ok(Some(access));
568 }
569 }
570 let Some(cursor) = body["cursor"].as_str() else {
571 return Ok(None);
572 };
573 endpoint = "sharing/list_folder_members/continue";
574 request = serde_json::json!({ "cursor": cursor });
575 }
576 }
577
578 async fn remove_folder_member(
579 &self,
580 shared_folder_id: &str,
581 email: &str,
582 ) -> Result<(), CloudHomeError> {
583 let path_root = Self::path_root_header(shared_folder_id);
584 let remove_body = serde_json::json!({
585 "shared_folder_id": shared_folder_id,
586 "member": { ".tag": "email", "email": email },
587 "leave_a_copy": false,
588 });
589 let resp = self
590 .session
591 .api_call(|oauth| {
592 Self::scoped_request(
593 oauth.post(format!("{}/sharing/remove_folder_member", self.api_base)),
594 &path_root,
595 )
596 .json(&remove_body)
597 })
598 .await?;
599 let status = resp.status();
600 let body = http::body_text(resp).await;
601 if !status.is_success() {
602 if dropbox_revoke_error_is_already_absent(&body) {
603 return Ok(());
604 }
605 return Err(CloudHomeError::Transport(format!(
606 "revoke access for {email} (HTTP {status}): {body}"
607 )));
608 }
609 match parse_dropbox_revoke_launch(&body)? {
610 DropboxRevokeLaunch::Complete => Ok(()),
611 DropboxRevokeLaunch::AsyncJob(job_id) => {
612 self.poll_remove_member_job(&job_id, shared_folder_id).await
613 }
614 }
615 }
616
617 #[cfg(test)]
618 fn with_endpoints(mut self, api_base: String, content_base: String) -> Self {
619 self.api_base = api_base;
620 self.content_base = content_base;
621 self
622 }
623}
624
625enum DropboxRevokeLaunch {
626 Complete,
627 AsyncJob(String),
628}
629
630fn parse_dropbox_revoke_launch(body: &str) -> Result<DropboxRevokeLaunch, CloudHomeError> {
631 let json: serde_json::Value = serde_json::from_str(body).map_err(|e| {
632 CloudHomeError::transport(format!("parse revoke-access response: {body}"), e)
633 })?;
634 match json[".tag"].as_str() {
635 Some("complete") => Ok(DropboxRevokeLaunch::Complete),
636 Some("async_job_id") => json["async_job_id"]
637 .as_str()
638 .map(|id| DropboxRevokeLaunch::AsyncJob(id.to_string()))
639 .ok_or_else(|| {
640 CloudHomeError::Transport(format!(
641 "revoke access: async_job_id response missing async_job_id: {body}"
642 ))
643 }),
644 _ => Err(CloudHomeError::Transport(format!(
645 "revoke access: unexpected launch response: {body}"
646 ))),
647 }
648}
649
650fn dropbox_api_arg(value: &serde_json::Value) -> String {
651 let json = value.to_string();
652 let mut arg = String::with_capacity(json.len());
653 for c in json.chars() {
654 if c.is_ascii() {
655 arg.push(c);
656 } else {
657 let mut units = [0u16; 2];
658 for unit in c.encode_utf16(&mut units) {
659 write!(&mut arg, "\\u{unit:04x}").expect("writing to a String cannot fail");
660 }
661 }
662 }
663 arg
664}
665
666fn dropbox_revoke_error_is_already_absent(body: &str) -> bool {
667 parse_dropbox_error_summary(body)
668 .as_deref()
669 .is_some_and(|summary| summary.starts_with("member_error/not_a_member"))
670}
671
672struct DropboxSessionSink<'a> {
682 home: &'a DropboxCloudHome,
683 session_id: String,
684 key: String,
685 completion: DropboxSessionCompletion,
686 confirmed_offset: u64,
687 settled: bool,
688}
689
690#[derive(Clone, Copy)]
691enum DropboxSessionCompletion {
692 Overwrite,
693 CreateOnly,
694}
695
696#[async_trait]
697impl super::PartSink for DropboxSessionSink<'_> {
698 fn part_size(&self) -> usize {
699 DROPBOX_CHUNK_SIZE
700 }
701
702 async fn send_part(
703 &mut self,
704 part: bytes::Bytes,
705 offset: u64,
706 is_last: bool,
707 control: &super::UploadControl,
708 ) -> Result<(), CloudHomeError> {
709 let length = part.len() as u64;
710 let namespace_id = self.home.get_or_create_shared_folder_id().await?;
711 let path_root = DropboxCloudHome::path_root_header(&namespace_id);
712 let resp = if is_last {
713 let path = DropboxCloudHome::namespace_path(&self.key);
715 let mode = match self.completion {
716 DropboxSessionCompletion::Overwrite => serde_json::json!({ ".tag": "overwrite" }),
717 DropboxSessionCompletion::CreateOnly => serde_json::json!({ ".tag": "add" }),
718 };
719 let arg = dropbox_api_arg(&serde_json::json!({
720 "cursor": { "session_id": self.session_id, "offset": offset },
721 "commit": {
722 "path": path,
723 "mode": mode,
724 "autorename": false,
725 "strict_conflict": matches!(self.completion, DropboxSessionCompletion::CreateOnly),
726 "mute": true,
727 },
728 }));
729 self.home
730 .session
731 .api_call_no_transient_retry(|oauth| {
732 let body = reqwest::Body::wrap_stream(
733 control.clone().stream_part(part.clone(), offset),
734 );
735 DropboxCloudHome::scoped_request(
736 oauth.post(format!(
737 "{}/files/upload_session/finish",
738 self.home.content_base
739 )),
740 &path_root,
741 )
742 .header("Dropbox-API-Arg", &arg)
743 .header("Content-Type", "application/octet-stream")
744 .header("Content-Length", length)
745 .body(body)
746 })
747 .await?
748 } else {
749 let arg = dropbox_api_arg(&serde_json::json!({
750 "cursor": { "session_id": self.session_id, "offset": offset },
751 "close": false,
752 }));
753 self.home
754 .session
755 .api_call_no_transient_retry(|oauth| {
756 let body = reqwest::Body::wrap_stream(
757 control.clone().stream_part(part.clone(), offset),
758 );
759 DropboxCloudHome::scoped_request(
760 oauth.post(format!(
761 "{}/files/upload_session/append_v2",
762 self.home.content_base
763 )),
764 &path_root,
765 )
766 .header("Dropbox-API-Arg", &arg)
767 .header("Content-Type", "application/octet-stream")
768 .header("Content-Length", length)
769 .body(body)
770 })
771 .await?
772 };
773 let status = resp.status();
774 if !status.is_success() {
775 return Err(classify_write_error(
776 status,
777 &http::body_text(resp).await,
778 &self.key,
779 ));
780 }
781 self.confirmed_offset = offset.checked_add(length).ok_or_else(|| {
782 CloudHomeError::Transport(format!("Dropbox upload offset overflow for {}", self.key))
783 })?;
784 if is_last {
785 self.settled = true;
786 if matches!(self.completion, DropboxSessionCompletion::CreateOnly) {
787 self.home.validate_create_response(&self.key, resp).await?;
788 }
789 }
790 Ok(())
791 }
792
793 async fn abort(&mut self) -> Result<(), CloudHomeError> {
794 if self.settled {
795 return Ok(());
796 }
797 let arg = dropbox_api_arg(&serde_json::json!({
798 "cursor": {
799 "session_id": self.session_id,
800 "offset": self.confirmed_offset,
801 },
802 "close": true,
803 }));
804 let namespace_id = self.home.get_or_create_shared_folder_id().await?;
805 let path_root = DropboxCloudHome::path_root_header(&namespace_id);
806 let response = self
807 .home
808 .session
809 .api_call_no_transient_retry(|oauth| {
810 DropboxCloudHome::scoped_request(
811 oauth.post(format!(
812 "{}/files/upload_session/append_v2",
813 self.home.content_base
814 )),
815 &path_root,
816 )
817 .header("Dropbox-API-Arg", &arg)
818 .header("Content-Type", "application/octet-stream")
819 .body(Vec::new())
820 })
821 .await?;
822 let status = response.status();
823 if !status.is_success() {
824 return Err(classify_write_error(
825 status,
826 &http::body_text(response).await,
827 &self.key,
828 ));
829 }
830 self.settled = true;
831 Ok(())
832 }
833
834 async fn finish(mut self: Box<Self>) -> Result<(), CloudHomeError> {
835 if self.settled {
836 return Ok(());
837 }
838 let operation = CloudHomeError::Transport(format!(
839 "Dropbox upload session for {} reached finish without committing its final part",
840 self.key
841 ));
842 let cleanup = self.abort().await;
843 Err(combine_cleanup_failure(operation, cleanup))
844 }
845}
846
847fn parse_dropbox_error_summary(body: &str) -> Option<String> {
850 http::error_reason(body, |v| v.get("error_summary")?.as_str().map(String::from))
851}
852
853fn classify_write_error(status: reqwest::StatusCode, body: &str, key: &str) -> CloudHomeError {
857 if let Some(summary) = parse_dropbox_error_summary(body) {
858 if summary.starts_with("path/conflict/file") {
859 return CloudHomeError::AlreadyExists(key.to_string());
860 }
861 if summary.starts_with("path/insufficient_space") {
862 return CloudHomeError::Transport(
863 "Your Dropbox storage is full. Free up space at dropbox.com to keep syncing."
864 .to_string(),
865 );
866 }
867 }
868 CloudHomeError::Transport(format!("write {key} (HTTP {status}): {body}"))
869}
870
871const DROPBOX_SIMPLE_UPLOAD_MAX: usize = 4 * 1024 * 1024;
875
876const DROPBOX_CHUNK_SIZE: usize = 8 * 1024 * 1024;
879
880#[async_trait]
881impl OAuthRestHome for DropboxCloudHome {
882 fn not_found(&self) -> NotFound {
883 NotFound::BodyContains {
885 status: StatusCode::CONFLICT,
886 needle: "not_found",
887 }
888 }
889
890 async fn send_read(
891 &self,
892 key: &str,
893 range: Option<(u64, u64)>,
894 ) -> Result<reqwest::Response, CloudHomeError> {
895 let namespace_id = self.get_or_create_shared_folder_id().await?;
896 let path_root = Self::path_root_header(&namespace_id);
897 let arg = dropbox_api_arg(&serde_json::json!({ "path": Self::namespace_path(key) }));
898 let range = range.map(|(start, end)| super::range_header(start, end));
899 self.session
900 .api_call(|oauth| {
901 let mut req = Self::scoped_request(
902 oauth.post(format!("{}/files/download", self.content_base)),
903 &path_root,
904 )
905 .header("Dropbox-API-Arg", &arg);
906 if let Some(ref range) = range {
907 req = req.header("Range", range);
908 }
909 req
910 })
911 .await
912 }
913
914 async fn send_delete(&self, key: &str) -> Result<reqwest::Response, CloudHomeError> {
915 let namespace_id = self.get_or_create_shared_folder_id().await?;
916 let path_root = Self::path_root_header(&namespace_id);
917 let body = serde_json::json!({ "path": Self::namespace_path(key) });
918 self.session
919 .api_call(|oauth| {
920 Self::scoped_request(
921 oauth.post(format!("{}/files/delete_v2", self.api_base)),
922 &path_root,
923 )
924 .json(&body)
925 })
926 .await
927 }
928
929 async fn send_list_page(
930 &self,
931 _prefix: &str,
932 cursor: Option<&str>,
933 ) -> Result<reqwest::Response, CloudHomeError> {
934 let namespace_id = self.get_or_create_shared_folder_id().await?;
935 let path_root = Self::path_root_header(&namespace_id);
936 match cursor {
939 Some(cur) => {
940 let body = serde_json::json!({ "cursor": cur });
941 self.session
942 .api_call(|oauth| {
943 Self::scoped_request(
944 oauth.post(format!("{}/files/list_folder/continue", self.api_base)),
945 &path_root,
946 )
947 .json(&body)
948 })
949 .await
950 }
951 None => {
952 let body = serde_json::json!({
953 "path": "",
954 "recursive": true,
955 "limit": 2000,
956 });
957 self.session
958 .api_call(|oauth| {
959 Self::scoped_request(
960 oauth.post(format!("{}/files/list_folder", self.api_base)),
961 &path_root,
962 )
963 .json(&body)
964 })
965 .await
966 }
967 }
968 }
969
970 fn parse_list_page(&self, body: &str, prefix: &str) -> Result<ListPage, CloudHomeError> {
971 let json: serde_json::Value = serde_json::from_str(body)
972 .map_err(|e| CloudHomeError::transport("parse list".to_string(), e))?;
973 let mut slots = Vec::new();
974 if let Some(entries) = json["entries"].as_array() {
975 for entry in entries {
976 if entry[".tag"].as_str() != Some("file") {
977 continue;
978 }
979 let Some(path_display) = entry["path_display"].as_str() else {
980 warn!("skipping Dropbox file entry without path_display: {entry}");
981 continue;
982 };
983 let Some(key) = path_display.strip_prefix('/') else {
984 warn!(
985 path_display,
986 "skipping Dropbox file entry outside the shared namespace"
987 );
988 continue;
989 };
990 if key.starts_with(prefix) {
991 slots.push(ObjectSlot::logical(key.to_string())?);
992 }
993 }
994 }
995 let next = if json["has_more"].as_bool() == Some(true) {
996 Some(
997 json["cursor"]
998 .as_str()
999 .ok_or_else(|| {
1000 CloudHomeError::Transport(
1001 "Dropbox list response has_more without cursor".to_string(),
1002 )
1003 })?
1004 .to_string(),
1005 )
1006 } else {
1007 None
1008 };
1009 Ok(ListPage { slots, next })
1010 }
1011}
1012
1013#[async_trait]
1014impl CloudHome for DropboxCloudHome {
1015 async fn put_object(&self, key: &str, data: Vec<u8>) -> Result<(), CloudHomeError> {
1016 let namespace_id = self.get_or_create_shared_folder_id().await?;
1017 let path_root = Self::path_root_header(&namespace_id);
1018 let body = Bytes::from(data);
1019 let api_arg = dropbox_api_arg(&serde_json::json!({
1020 "path": Self::namespace_path(key),
1021 "mode": { ".tag": "overwrite" },
1022 "autorename": false,
1023 "mute": true,
1024 }));
1025 let resp = self
1026 .session
1027 .api_call(|oauth| {
1028 Self::scoped_request(
1029 oauth.post(format!("{}/files/upload", self.content_base)),
1030 &path_root,
1031 )
1032 .header("Dropbox-API-Arg", &api_arg)
1033 .header("Content-Type", "application/octet-stream")
1034 .body(body.clone())
1035 })
1036 .await?;
1037 let status = resp.status();
1038 if !status.is_success() {
1039 return Err(classify_write_error(
1040 status,
1041 &http::body_text(resp).await,
1042 key,
1043 ));
1044 }
1045 Ok(())
1046 }
1047
1048 async fn open_multipart<'a>(
1049 &'a self,
1050 key: &str,
1051 _total_len: u64,
1052 ) -> Result<BoxPartSink<'a>, CloudHomeError> {
1053 let session_id = self.start_upload_session(key).await?;
1054 Ok(Box::new(DropboxSessionSink {
1055 home: self,
1056 session_id,
1057 key: key.to_string(),
1058 completion: DropboxSessionCompletion::Overwrite,
1059 confirmed_offset: 0,
1060 settled: false,
1061 }))
1062 }
1063
1064 fn multipart_threshold(&self) -> u64 {
1065 DROPBOX_SIMPLE_UPLOAD_MAX as u64
1066 }
1067
1068 async fn read(&self, key: &str) -> Result<Vec<u8>, CloudHomeError> {
1069 rest_read(self, key).await
1070 }
1071
1072 async fn read_range(&self, key: &str, start: u64, end: u64) -> Result<Vec<u8>, CloudHomeError> {
1073 rest_read_range(self, key, start, end).await
1074 }
1075
1076 async fn list(&self, prefix: &str) -> Result<Vec<String>, CloudHomeError> {
1077 rest_list(self, prefix).await
1078 }
1079
1080 async fn delete(&self, key: &str) -> Result<(), CloudHomeError> {
1081 rest_delete(self, key).await
1082 }
1083
1084 async fn exists(&self, key: &str) -> Result<bool, CloudHomeError> {
1085 let namespace_id = self.get_or_create_shared_folder_id().await?;
1086 let path_root = Self::path_root_header(&namespace_id);
1087 let body = serde_json::json!({ "path": Self::namespace_path(key) });
1088 let resp = self
1089 .session
1090 .api_call(|oauth| {
1091 Self::scoped_request(
1092 oauth.post(format!("{}/files/get_metadata", self.api_base)),
1093 &path_root,
1094 )
1095 .json(&body)
1096 })
1097 .await?;
1098 exists_from_response(resp, &format!("exists {key}"), self.not_found()).await
1099 }
1100
1101 async fn set_access(
1102 &self,
1103 desired: CloudAccessState,
1104 ) -> Result<CloudAccessOutcome, CloudHomeError> {
1105 let email = desired.require_provider_email("Dropbox")?.to_string();
1106 let shared_folder_id = self.get_or_create_shared_folder_id().await?;
1107 let path_root = Self::path_root_header(&shared_folder_id);
1108 match desired {
1109 CloudAccessState::Present { .. } => {
1110 let current = self.folder_member_access(&shared_folder_id, &email).await?;
1111 if current.as_deref() != Some("editor") {
1112 if current.is_some() {
1113 self.remove_folder_member(&shared_folder_id, &email).await?;
1114 }
1115 let add_body = serde_json::json!({
1116 "shared_folder_id": shared_folder_id,
1117 "members": [{
1118 "member": { ".tag": "email", "email": email },
1119 "access_level": { ".tag": "editor" },
1120 }],
1121 "quiet": false,
1122 });
1123 let resp = self
1124 .session
1125 .api_call(|oauth| {
1126 Self::scoped_request(
1127 oauth.post(format!("{}/sharing/add_folder_member", self.api_base)),
1128 &path_root,
1129 )
1130 .json(&add_body)
1131 })
1132 .await?;
1133 ensure_ok(resp, &format!("grant access to {email}"), self.not_found()).await?;
1134 }
1135 if self
1136 .folder_member_access(&shared_folder_id, &email)
1137 .await?
1138 .as_deref()
1139 != Some("editor")
1140 {
1141 return Err(CloudHomeError::Transport(format!(
1142 "editor access for {email} is not visible after update"
1143 )));
1144 }
1145 Ok(CloudAccessOutcome::Present(self.join_info()))
1146 }
1147 CloudAccessState::Absent { .. } => {
1148 if self
1149 .folder_member_access(&shared_folder_id, &email)
1150 .await?
1151 .is_some()
1152 {
1153 self.remove_folder_member(&shared_folder_id, &email).await?;
1154 }
1155 if self
1156 .folder_member_access(&shared_folder_id, &email)
1157 .await?
1158 .is_some()
1159 {
1160 return Err(CloudHomeError::Transport(format!(
1161 "access for {email} remains after removal"
1162 )));
1163 }
1164 Ok(CloudAccessOutcome::Absent(RevokeOutcome::Revoked))
1165 }
1166 }
1167 }
1168}
1169
1170#[async_trait]
1171impl ExactSlotStorage for DropboxCloudHome {
1172 async fn provider_binding(
1173 &self,
1174 ) -> Result<coven_protocol::objects::ResolvedProviderBinding, CloudHomeError> {
1175 use coven_protocol::objects::{
1176 ProviderDeviceBinding, ProviderPrincipalId, ResolvedProviderBinding,
1177 StoreProviderBinding,
1178 };
1179
1180 if self.folder_path.is_empty() {
1181 return Err(CloudHomeError::Configuration(
1182 "Dropbox provider binding has an empty folder path".to_string(),
1183 ));
1184 }
1185 let metadata_body = serde_json::json!({
1186 "path": self.folder_path,
1187 "include_deleted": false,
1188 });
1189 let metadata_response = self
1190 .session
1191 .api_call(|oauth| {
1192 oauth
1193 .post(format!("{}/files/get_metadata", self.api_base))
1194 .json(&metadata_body)
1195 })
1196 .await?;
1197 let metadata_response = ensure_ok(
1198 metadata_response,
1199 "resolve Dropbox folder binding",
1200 self.not_found(),
1201 )
1202 .await?;
1203 let metadata: serde_json::Value =
1204 http::ok_json(metadata_response, "parse Dropbox folder binding").await?;
1205 let folder_id = metadata["id"]
1206 .as_str()
1207 .filter(|value| !value.is_empty())
1208 .ok_or_else(|| {
1209 CloudHomeError::Transport(
1210 "Dropbox folder metadata omitted the stable folder id".to_string(),
1211 )
1212 })?
1213 .to_string();
1214 let path_lower = metadata["path_lower"]
1215 .as_str()
1216 .filter(|value| !value.is_empty())
1217 .ok_or_else(|| {
1218 CloudHomeError::Transport("Dropbox folder metadata omitted path_lower".to_string())
1219 })?
1220 .to_string();
1221 if metadata[".tag"].as_str() != Some("folder")
1222 || !path_lower.eq_ignore_ascii_case(&self.folder_path)
1223 {
1224 return Err(CloudHomeError::Transport(
1225 "Dropbox folder lookup returned a different folder binding".to_string(),
1226 ));
1227 }
1228
1229 let namespace_id = self.get_or_create_shared_folder_id().await?;
1230 let path_root = Self::path_root_header(&namespace_id);
1231 let account_response = self
1232 .session
1233 .api_call(|oauth| {
1234 Self::scoped_request(
1235 oauth.post(format!("{}/users/get_current_account", self.api_base)),
1236 &path_root,
1237 )
1238 .header("Content-Type", "application/json")
1239 .body("null")
1240 })
1241 .await?;
1242 let account_response = ensure_ok(
1243 account_response,
1244 "resolve Dropbox principal",
1245 NotFound::Status,
1246 )
1247 .await?;
1248 let account: serde_json::Value =
1249 http::ok_json(account_response, "parse Dropbox principal").await?;
1250 let account_id = account["account_id"]
1251 .as_str()
1252 .filter(|value| !value.is_empty())
1253 .ok_or_else(|| {
1254 CloudHomeError::Transport(
1255 "Dropbox current-account response omitted account_id".to_string(),
1256 )
1257 })?
1258 .to_string();
1259
1260 if namespace_id.is_empty() || folder_id.is_empty() || path_lower.is_empty() {
1261 return Err(CloudHomeError::Transport(
1262 "Dropbox shared namespace resolution returned an empty identifier".to_string(),
1263 ));
1264 }
1265 Ok(ResolvedProviderBinding {
1266 store: StoreProviderBinding::Dropbox { namespace_id },
1267 device: ProviderDeviceBinding {
1268 principal: ProviderPrincipalId::Dropbox { account_id },
1269 },
1270 })
1271 }
1272
1273 async fn create_at(
1274 &self,
1275 upload: &super::ExactUpload<'_>,
1276 control: &super::UploadControl,
1277 ) -> Result<super::ExactCreateOutcome, CloudHomeError> {
1278 if matches!(
1279 self.exact_upload_verification,
1280 coven_foundation::config::ExactUploadVerification::UploadChecksum
1281 ) {
1282 return Err(CloudHomeError::Configuration(
1283 "Dropbox does not expose upload-checksum enforcement for this endpoint".to_string(),
1284 ));
1285 }
1286 let slot = upload.object().slot();
1287 let body = upload.body().await?;
1288 slot.require_logical_key_for("Dropbox")?;
1289 let key = slot.logical_key();
1290 if body.len() <= self.multipart_threshold() {
1291 let operation = self.append_small(key, body.collect().await?, control).await;
1292 return super::exact_upload::settle_exact_create(operation, |observed| {
1293 self.verify_exact_upload(upload, observed)
1294 })
1295 .await;
1296 }
1297
1298 let session_id = self.start_upload_session(key).await?;
1299 let sink = DropboxSessionSink {
1300 home: self,
1301 session_id,
1302 key: key.to_string(),
1303 completion: DropboxSessionCompletion::CreateOnly,
1304 confirmed_offset: 0,
1305 settled: false,
1306 };
1307 let operation = super::blob_body::MultipartUpload::new(key, body, Box::new(sink), control)
1308 .run()
1309 .await;
1310 super::exact_upload::settle_exact_create(operation, |observed| {
1311 self.verify_exact_upload(upload, observed)
1312 })
1313 .await
1314 }
1315
1316 async fn list_slots(&self, prefix: &str) -> Result<Vec<ObjectSlot>, CloudHomeError> {
1317 crate::cloud::logical_slots(CloudHome::list(self, prefix).await?)
1318 }
1319
1320 async fn read_at(&self, slot: &ObjectSlot) -> Result<Vec<u8>, CloudHomeError> {
1321 let (response, _) = self.send_exact_read(slot).await?;
1322 http::ok_bytes(response, "read exact Dropbox body").await
1323 }
1324
1325 async fn read_versioned_at(
1326 &self,
1327 slot: &ObjectSlot,
1328 ) -> Result<CloudVersionedObject, CloudHomeError> {
1329 let (response, metadata) = self.send_exact_read(slot).await?;
1330 let bytes = http::ok_bytes(response, "read versioned Dropbox body").await?;
1331 Ok(CloudVersionedObject {
1332 bytes,
1333 version: metadata.version,
1334 })
1335 }
1336
1337 async fn replace_at_if_version(
1338 &self,
1339 slot: &ObjectSlot,
1340 expected: &CloudObjectVersion,
1341 bytes: Vec<u8>,
1342 ) -> Result<ConditionalWriteOutcome, CloudHomeError> {
1343 slot.require_logical_key_for("Dropbox")?;
1344 let namespace_id = self.get_or_create_shared_folder_id().await?;
1345 let path_root = Self::path_root_header(&namespace_id);
1346 let api_arg = dropbox_api_arg(&serde_json::json!({
1347 "path": Self::namespace_path(slot.logical_key()),
1348 "mode": { ".tag": "update", "update": expected.as_provider() },
1349 "autorename": false,
1350 "strict_conflict": true,
1351 "mute": true,
1352 }));
1353 let body = Bytes::from(bytes);
1354 let response = self
1355 .session
1356 .api_call_no_transient_retry(|oauth| {
1357 Self::scoped_request(
1358 oauth.post(format!("{}/files/upload", self.content_base)),
1359 &path_root,
1360 )
1361 .header("Dropbox-API-Arg", &api_arg)
1362 .header("Content-Type", "application/octet-stream")
1363 .header("Content-Length", body.len())
1364 .body(body.clone())
1365 })
1366 .await?;
1367 let status = response.status();
1368 let response_body = http::body_text(response).await;
1369 if status == StatusCode::CONFLICT
1370 && parse_dropbox_error_summary(&response_body)
1371 .is_some_and(|summary| summary.starts_with("path/conflict/file"))
1372 {
1373 return Ok(ConditionalWriteOutcome::VersionChanged);
1374 }
1375 if !status.is_success() {
1376 return Err(classify_write_error(
1377 status,
1378 &response_body,
1379 slot.logical_key(),
1380 ));
1381 }
1382 let metadata: serde_json::Value =
1383 serde_json::from_str(&response_body).map_err(|error| {
1384 CloudHomeError::transport(
1385 format!(
1386 "parse conditional Dropbox response for {}",
1387 slot.logical_key()
1388 ),
1389 error,
1390 )
1391 })?;
1392 let metadata = self.exact_metadata_from_json(slot, &metadata)?;
1393 Ok(ConditionalWriteOutcome::Replaced(metadata.version))
1394 }
1395
1396 async fn read_range_at(
1397 &self,
1398 slot: &ObjectSlot,
1399 start: u64,
1400 end: u64,
1401 ) -> Result<Vec<u8>, CloudHomeError> {
1402 let bytes = self.read_at(slot).await?;
1403 let start = usize::try_from(start)
1404 .map_err(|_| CloudHomeError::Configuration("range start is too large".to_string()))?;
1405 let end = usize::try_from(end)
1406 .map_err(|_| CloudHomeError::Configuration("range end is too large".to_string()))?;
1407 bytes.get(start..end).map(<[u8]>::to_vec).ok_or_else(|| {
1408 CloudHomeError::Configuration(format!(
1409 "invalid range {start}..{end} for {} bytes",
1410 bytes.len()
1411 ))
1412 })
1413 }
1414
1415 async fn read_at_to_file(
1416 &self,
1417 slot: &ObjectSlot,
1418 destination: &std::path::Path,
1419 progress: super::DownloadProgress,
1420 ) -> Result<(), super::CloudFileReadError> {
1421 let (response, _) = self.send_exact_read(slot).await?;
1422 super::oauth_rest::response_to_file(
1423 response,
1424 destination,
1425 "read exact Dropbox body",
1426 progress,
1427 )
1428 .await
1429 }
1430
1431 async fn delete_at(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
1432 match self.verify_slot(slot).await {
1433 Ok(()) => {}
1434 Err(CloudHomeError::NotFound(_)) => return Ok(()),
1435 Err(error) => return Err(error),
1436 }
1437 let namespace_id = self.get_or_create_shared_folder_id().await?;
1438 let path_root = Self::path_root_header(&namespace_id);
1439 let body = serde_json::json!({"path": Self::namespace_path(slot.logical_key())});
1440 let response = self
1441 .session
1442 .api_call(|oauth| {
1443 Self::scoped_request(
1444 oauth.post(format!("{}/files/delete_v2", self.api_base)),
1445 &path_root,
1446 )
1447 .json(&body)
1448 })
1449 .await?;
1450 match ensure_ok(response, "delete exact Dropbox file", self.not_found()).await {
1451 Ok(_) | Err(CloudHomeError::NotFound(_)) => Ok(()),
1452 Err(error) => Err(error),
1453 }
1454 }
1455}
1456
1457#[cfg(test)]
1458#[path = "dropbox_tests.rs"]
1459mod tests;