1use async_trait::async_trait;
10use bytes::Bytes;
11
12use super::http::{self, ensure_ok, exists_from_response, NotFound};
13use super::key_encoding::{decode_listed_key, encode_key};
14use super::oauth_rest::{
15 rest_delete, rest_list, rest_read, rest_read_range, ListPage, OAuthRestHome,
16};
17use super::oauth_session::OAuthSession;
18use super::{
19 combine_cleanup_failure, sharing, BlobBody, BoxPartSink, CloudAccessOutcome, CloudAccessState,
20 CloudHome, CloudHomeError, CloudHomeJoinInfo, CloudObjectVersion, CloudVersionedObject,
21 ConditionalWriteOutcome, ExactSlotStorage, RevokeOutcome,
22};
23use crate::oauth::OAuthConfig;
24use coven_protocol::objects::ObjectSlot;
25
26#[path = "onedrive/content_hash.rs"]
27mod content_hash;
28
29#[derive(Clone, Debug, PartialEq, Eq)]
30struct OneDriveExactMetadata {
31 size: u64,
32 sha1_hash: String,
33 version: CloudObjectVersion,
34}
35
36const GRAPH_API: &str = "https://graph.microsoft.com/v1.0";
37
38fn onedrive_upload_cancellation_succeeded(status: reqwest::StatusCode) -> bool {
39 status.is_success() || status == reqwest::StatusCode::NOT_FOUND
40}
41
42pub struct OneDriveCloudHome {
44 drive_id: String,
45 folder_id: String,
46 graph_api: String,
47 session: OAuthSession,
48 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
49}
50
51enum UploadSessionCompletion {
52 Automatic,
53 DeferredPersonal,
54}
55
56impl OneDriveCloudHome {
57 pub fn new(
58 drive_id: String,
59 folder_id: String,
60 session: OAuthSession,
61 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
62 ) -> Self {
63 Self {
64 drive_id,
65 folder_id,
66 graph_api: GRAPH_API.to_string(),
67 session,
68 exact_upload_verification,
69 }
70 }
71
72 pub(crate) fn oauth_config(creds: crate::oauth::OAuthClientCreds) -> OAuthConfig {
73 OAuthConfig {
74 client_id: creds.client_id,
75 client_secret: creds.client_secret,
76 auth_url: "https://login.microsoftonline.com/consumers/oauth2/v2.0/authorize"
77 .to_string(),
78 token_url: "https://login.microsoftonline.com/consumers/oauth2/v2.0/token".to_string(),
79 scopes: vec![
80 "Files.ReadWrite".to_string(),
81 "offline_access".to_string(),
82 "User.Read".to_string(),
84 ],
85 redirect_port: 19284,
86 extra_auth_params: vec![],
87 }
88 }
89
90 fn item_path_url(&self, key: &str) -> String {
92 format!(
93 "{}/drives/{}/items/{}:/{}:",
94 self.graph_api,
95 self.drive_id,
96 self.folder_id,
97 encode_key(key)
98 )
99 }
100
101 fn children_url(&self) -> String {
103 format!(
104 "{}/drives/{}/items/{}/children",
105 self.graph_api, self.drive_id, self.folder_id
106 )
107 }
108
109 async fn create_upload_session(
110 &self,
111 key: &str,
112 conflict_behavior: &str,
113 completion: UploadSessionCompletion,
114 ) -> Result<String, CloudHomeError> {
115 let session_url = format!("{}/createUploadSession", self.item_path_url(key));
116 let body = serde_json::json!({
117 "item": { "@microsoft.graph.conflictBehavior": conflict_behavior },
118 "deferCommit": matches!(completion, UploadSessionCompletion::DeferredPersonal),
119 });
120 let response = self
121 .session
122 .api_call(|oauth| oauth.post(&session_url).json(&body))
123 .await?;
124 let status = response.status();
125 let body = http::body_text(response).await;
126 if !status.is_success() {
127 return Err(classify_write_error(status, &body, key));
128 }
129 let json: serde_json::Value = serde_json::from_str(&body).map_err(|error| {
130 CloudHomeError::transport(format!("parse upload session {key}"), error)
131 })?;
132 json["uploadUrl"]
133 .as_str()
134 .filter(|url| !url.is_empty())
135 .map(str::to_string)
136 .ok_or_else(|| {
137 CloudHomeError::Transport(format!("upload session {key}: no uploadUrl returned"))
138 })
139 }
140
141 async fn commit_deferred_append(
142 &self,
143 key: &str,
144 upload_url: &str,
145 ) -> Result<(), CloudHomeError> {
146 let body = serde_json::json!({
147 "name": encode_key(key),
148 "@microsoft.graph.conflictBehavior": "fail",
149 "@microsoft.graph.sourceUrl": upload_url,
150 });
151 let response = match self
152 .session
153 .api_call(|oauth| oauth.put(self.item_path_url(key)).json(&body))
154 .await
155 {
156 Ok(response) => response,
157 Err(operation) => return Err(operation),
158 };
159 let status = response.status();
160 if !status.is_success() {
161 return Err(classify_write_error(
162 status,
163 &http::body_text(response).await,
164 key,
165 ));
166 }
167 let response_body = match response.bytes().await {
168 Ok(body) => body,
169 Err(operation) => {
170 return Err(CloudHomeError::Transport(format!(
171 "commit create {key}: read response: {operation}"
172 )))
173 }
174 };
175 let _: serde_json::Value = serde_json::from_slice(&response_body).map_err(|error| {
176 CloudHomeError::transport(format!("commit create {key}: parse response"), error)
177 })?;
178 Ok(())
179 }
180
181 async fn commit_deferred_replacement(
182 &self,
183 slot: &ObjectSlot,
184 upload_url: &str,
185 expected: &CloudObjectVersion,
186 ) -> Result<ConditionalWriteOutcome, CloudHomeError> {
187 let key = slot.logical_key();
188 let body = serde_json::json!({
189 "name": encode_key(key),
190 "@microsoft.graph.conflictBehavior": "replace",
191 "@microsoft.graph.sourceUrl": upload_url,
192 });
193 let response = self
194 .session
195 .api_call_no_transient_retry(|oauth| {
196 oauth
197 .put(self.item_path_url(key))
198 .header(reqwest::header::IF_MATCH, expected.as_provider())
199 .json(&body)
200 })
201 .await?;
202 let status = response.status();
203 let response_body = http::body_text(response).await;
204 if status == reqwest::StatusCode::PRECONDITION_FAILED {
205 return Ok(ConditionalWriteOutcome::VersionChanged);
206 }
207 if !status.is_success() {
208 return Err(classify_write_error(status, &response_body, key));
209 }
210 let metadata: serde_json::Value =
211 serde_json::from_str(&response_body).map_err(|error| {
212 CloudHomeError::transport(
213 format!("commit conditional OneDrive replacement {key}: parse response"),
214 error,
215 )
216 })?;
217 let version = metadata["eTag"]
218 .as_str()
219 .filter(|etag| !etag.is_empty())
220 .ok_or_else(|| {
221 CloudHomeError::Transport(format!(
222 "commit conditional OneDrive replacement {key}: response omitted eTag"
223 ))
224 })?;
225 Ok(ConditionalWriteOutcome::Replaced(
226 CloudObjectVersion::from_provider(version.to_string())?,
227 ))
228 }
229
230 async fn exact_metadata(
231 &self,
232 slot: &ObjectSlot,
233 ) -> Result<OneDriveExactMetadata, CloudHomeError> {
234 slot.require_logical_key_for("OneDrive")?;
235 let response = self
236 .session
237 .api_call(|oauth| {
238 oauth
239 .get(self.item_path_url(slot.logical_key()))
240 .query(&[("$select", "id,name,parentReference,deleted,file,size,eTag")])
241 })
242 .await?;
243 let response = ensure_ok(response, "verify exact OneDrive item", NotFound::Status).await?;
244 let metadata: serde_json::Value =
245 http::ok_json(response, "parse exact OneDrive item metadata").await?;
246 let expected_name = encode_key(slot.logical_key());
247 let matches = metadata["id"].as_str().is_some_and(|id| !id.is_empty())
248 && metadata["name"].as_str() == Some(expected_name.as_str())
249 && metadata["parentReference"]["id"].as_str() == Some(self.folder_id.as_str())
250 && metadata["deleted"].is_null()
251 && metadata["file"].is_object();
252 if !matches {
253 return Err(CloudHomeError::Transport(format!(
254 "exact OneDrive slot for {} does not identify {expected_name} in folder {}",
255 slot.logical_key(),
256 self.folder_id
257 )));
258 }
259 let size = metadata["size"].as_u64().ok_or_else(|| {
260 CloudHomeError::Transport(format!(
261 "exact OneDrive metadata for {} omitted size",
262 slot.logical_key()
263 ))
264 })?;
265 let sha1_hash = metadata["file"]["hashes"]["sha1Hash"]
266 .as_str()
267 .filter(|hash| !hash.is_empty())
268 .ok_or_else(|| {
269 CloudHomeError::Transport(format!(
270 "exact OneDrive metadata for {} omitted file.hashes.sha1Hash",
271 slot.logical_key()
272 ))
273 })?
274 .to_string();
275 let version = metadata["eTag"]
276 .as_str()
277 .filter(|etag| !etag.is_empty())
278 .ok_or_else(|| {
279 CloudHomeError::Transport(format!(
280 "exact OneDrive metadata for {} omitted eTag",
281 slot.logical_key()
282 ))
283 })?;
284 Ok(OneDriveExactMetadata {
285 size,
286 sha1_hash,
287 version: CloudObjectVersion::from_provider(version.to_string())?,
288 })
289 }
290
291 async fn verify_slot(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
292 self.exact_metadata(slot).await.map(drop)
293 }
294
295 async fn verify_exact_upload(
296 &self,
297 upload: &super::ExactUpload<'_>,
298 created_response_was_observed: bool,
299 ) -> Result<(), CloudHomeError> {
300 use coven_foundation::config::ExactUploadVerification;
301
302 match self.exact_upload_verification {
303 ExactUploadVerification::UploadChecksum => Err(CloudHomeError::Configuration(
304 "OneDrive does not expose upload-checksum enforcement for this endpoint"
305 .to_string(),
306 )),
307 ExactUploadVerification::MetadataHash => {
308 let metadata = self.exact_metadata(upload.object().slot()).await?;
309 let expected_sha1 = content_hash::sha1(upload).await?;
310 if metadata.size != upload.object().stored_size()
311 || !metadata.sha1_hash.eq_ignore_ascii_case(&expected_sha1)
312 {
313 return Err(CloudHomeError::SlotCollision(
314 upload.object().slot().logical_key().to_string(),
315 ));
316 }
317 Ok(())
318 }
319 ExactUploadVerification::Readback => {
320 let bytes = self.read_at(upload.object().slot()).await?;
321 upload.verify_stored_bytes(&bytes)
322 }
323 ExactUploadVerification::Unchecked => {
324 super::exact_upload::accept_unchecked_create_response(
325 created_response_was_observed,
326 upload.object(),
327 )
328 }
329 }
330 }
331
332 async fn create_at_slot(
333 &self,
334 slot: &ObjectSlot,
335 mut body: BlobBody,
336 control: &super::UploadControl,
337 ) -> Result<(), CloudHomeError> {
338 slot.require_logical_key_for("OneDrive")?;
339 let full_logical_key = slot.logical_key();
340 let upload_url = self
341 .create_upload_session(
342 full_logical_key,
343 "fail",
344 UploadSessionCompletion::DeferredPersonal,
345 )
346 .await?;
347 let key = full_logical_key.to_string();
348 let classify =
349 Box::new(move |status, response: &str| classify_write_error(status, response, &key));
350 let mut uploader = self.session.range_put_uploader(
351 upload_url.clone(),
352 202,
353 body.len(),
354 ONEDRIVE_CHUNK_SIZE,
355 full_logical_key.to_string(),
356 classify,
357 onedrive_upload_cancellation_succeeded,
358 );
359 let total = body.len();
360 let mut offset = 0;
361 loop {
362 let part = match body.next_part(uploader.part_size()).await {
363 Ok(Some(part)) => part,
364 Ok(None) if offset == total => break,
365 Ok(None) => {
366 let operation = CloudHomeError::Transport(format!(
367 "append {full_logical_key}: upload body ended before the final part"
368 ));
369 let cleanup = uploader.abort().await;
370 return Err(combine_cleanup_failure(operation, cleanup));
371 }
372 Err(operation) => {
373 let cleanup = uploader.abort().await;
374 return Err(combine_cleanup_failure(operation, cleanup));
375 }
376 };
377 let length = part.len() as u64;
378 let is_last = offset + length >= total;
379 let completion = uploader
380 .send_deferred_part(part, offset, is_last, control)
381 .await?;
382 offset += length;
383 control.report(offset);
384 if let Some(response) = completion {
385 response.bytes().await.map_err(|error| {
386 CloudHomeError::transport(
387 format!("read unexpected append completion for {full_logical_key}"),
388 error,
389 )
390 })?;
391 let operation = CloudHomeError::Transport(format!(
392 "append {full_logical_key}: deferred upload published before explicit commit"
393 ));
394 let cleanup = self.delete_at_slot(slot).await;
395 return Err(combine_cleanup_failure(operation, cleanup));
396 }
397 }
398 let result = self
399 .commit_deferred_append(full_logical_key, &upload_url)
400 .await;
401 match result {
402 Ok(()) => {
403 uploader.mark_completed();
404 Ok(())
405 }
406 Err(operation) => {
407 let cleanup = uploader.abort().await;
408 Err(combine_cleanup_failure(operation, cleanup))
409 }
410 }
411 }
412
413 async fn read_at_to_file(
414 &self,
415 slot: &ObjectSlot,
416 destination: &std::path::Path,
417 progress: super::DownloadProgress,
418 ) -> Result<(), super::CloudFileReadError> {
419 self.verify_slot(slot).await?;
420 let response = self
421 .session
422 .api_call(|oauth| {
423 oauth.get(format!(
424 "{}/content",
425 self.item_path_url(slot.logical_key())
426 ))
427 })
428 .await?;
429 let response = ensure_ok(response, "read exact OneDrive item", NotFound::Status).await?;
430 super::oauth_rest::response_to_file(
431 response,
432 destination,
433 "read exact OneDrive item body",
434 progress,
435 )
436 .await
437 }
438
439 async fn delete_at_slot(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
440 slot.require_logical_key_for("OneDrive")?;
441 match self.verify_slot(slot).await {
442 Ok(()) => {}
443 Err(CloudHomeError::NotFound(_)) => return Ok(()),
444 Err(error) => return Err(error),
445 }
446 let response = self
447 .session
448 .api_call(|oauth| oauth.delete(self.item_path_url(slot.logical_key())))
449 .await?;
450 let status = response.status();
451 if status.is_success() || status == reqwest::StatusCode::NOT_FOUND {
452 Ok(())
453 } else {
454 Err(CloudHomeError::Transport(format!(
455 "delete exact OneDrive item (HTTP {status}): {}",
456 http::body_text(response).await
457 )))
458 }
459 }
460
461 #[cfg(test)]
462 fn with_graph_api(mut self, graph_api: String) -> Self {
463 self.graph_api = graph_api;
464 self
465 }
466}
467
468fn parse_onedrive_error_code(body: &str) -> Option<String> {
471 http::error_reason(body, |v| {
472 v.get("error")?.get("code")?.as_str().map(String::from)
473 })
474}
475
476fn classify_write_error(status: reqwest::StatusCode, body: &str, key: &str) -> CloudHomeError {
480 let code = parse_onedrive_error_code(body);
481 if code.as_deref() == Some("nameAlreadyExists") {
482 return CloudHomeError::AlreadyExists(key.to_string());
483 }
484 if code.as_deref() == Some("quotaLimitReached") {
485 return CloudHomeError::Transport(
486 "Your OneDrive storage is full. Free up space at onedrive.live.com to keep syncing."
487 .to_string(),
488 );
489 }
490 CloudHomeError::Transport(format!("write {key} (HTTP {status}): {body}"))
491}
492
493const ONEDRIVE_SIMPLE_PUT_MAX: usize = 4 * 1024 * 1024;
496
497const ONEDRIVE_CHUNK_SIZE: usize = 24 * 320 * 1024;
500
501#[async_trait]
502impl OAuthRestHome for OneDriveCloudHome {
503 fn not_found(&self) -> NotFound {
504 NotFound::Status
505 }
506
507 async fn send_read(
508 &self,
509 key: &str,
510 range: Option<(u64, u64)>,
511 ) -> Result<reqwest::Response, CloudHomeError> {
512 let url = format!("{}/content", self.item_path_url(key));
513 let range = range.map(|(start, end)| super::range_header(start, end));
514 self.session
515 .api_call(|oauth| {
516 let mut req = oauth.get(&url);
517 if let Some(ref range) = range {
518 req = req.header("Range", range);
519 }
520 req
521 })
522 .await
523 }
524
525 async fn send_delete(&self, key: &str) -> Result<reqwest::Response, CloudHomeError> {
526 let url = self.item_path_url(key);
527 self.session.api_call(|oauth| oauth.delete(&url)).await
528 }
529
530 async fn send_list_page(
531 &self,
532 _prefix: &str,
533 cursor: Option<&str>,
534 ) -> Result<reqwest::Response, CloudHomeError> {
535 let url = match cursor {
538 Some(next) => next.to_string(),
539 None => format!("{}?$select=name", self.children_url()),
540 };
541 self.session.api_call(|oauth| oauth.get(&url)).await
542 }
543
544 fn parse_list_page(&self, body: &str, prefix: &str) -> Result<ListPage, CloudHomeError> {
545 let json: serde_json::Value = serde_json::from_str(body)
546 .map_err(|e| CloudHomeError::transport("parse list".to_string(), e))?;
547 let encoded_prefix = encode_key(prefix);
548 let mut slots = Vec::new();
549 if let Some(items) = json["value"].as_array() {
550 for item in items {
551 if let Some(name) = item["name"].as_str() {
552 if name.starts_with(&encoded_prefix) {
553 let Some(decoded) = decode_listed_key("OneDrive", name) else {
554 continue;
555 };
556 slots.push(ObjectSlot::logical(decoded)?)
557 }
558 }
559 }
560 }
561 Ok(ListPage {
562 slots,
563 next: json["@odata.nextLink"].as_str().map(String::from),
564 })
565 }
566}
567
568#[async_trait]
569impl CloudHome for OneDriveCloudHome {
570 async fn put_object(&self, key: &str, data: Vec<u8>) -> Result<(), CloudHomeError> {
571 let body = Bytes::from(data);
572 let url = format!("{}/content", self.item_path_url(key));
573 let resp = self
574 .session
575 .api_call(|oauth| {
576 oauth
577 .put(&url)
578 .header("Content-Type", "application/octet-stream")
579 .body(body.clone())
580 })
581 .await?;
582 let status = resp.status();
583 if !status.is_success() {
584 return Err(classify_write_error(
585 status,
586 &http::body_text(resp).await,
587 key,
588 ));
589 }
590 Ok(())
591 }
592
593 async fn open_multipart<'a>(
594 &'a self,
595 key: &str,
596 total_len: u64,
597 ) -> Result<BoxPartSink<'a>, CloudHomeError> {
598 let upload_url = self
599 .create_upload_session(key, "replace", UploadSessionCompletion::Automatic)
600 .await?;
601 let key_owned = key.to_string();
602 let classify =
603 Box::new(move |status, body: &str| classify_write_error(status, body, &key_owned));
604 Ok(Box::new(self.session.range_put_sink(
606 upload_url,
607 202,
608 total_len,
609 ONEDRIVE_CHUNK_SIZE,
610 key.to_string(),
611 classify,
612 onedrive_upload_cancellation_succeeded,
613 )))
614 }
615
616 fn multipart_threshold(&self) -> u64 {
617 ONEDRIVE_SIMPLE_PUT_MAX as u64
618 }
619
620 async fn read(&self, key: &str) -> Result<Vec<u8>, CloudHomeError> {
621 rest_read(self, key).await
622 }
623
624 async fn read_range(&self, key: &str, start: u64, end: u64) -> Result<Vec<u8>, CloudHomeError> {
625 rest_read_range(self, key, start, end).await
626 }
627
628 async fn list(&self, prefix: &str) -> Result<Vec<String>, CloudHomeError> {
629 rest_list(self, prefix).await
630 }
631
632 async fn delete(&self, key: &str) -> Result<(), CloudHomeError> {
633 rest_delete(self, key).await
634 }
635
636 async fn exists(&self, key: &str) -> Result<bool, CloudHomeError> {
637 let url = self.item_path_url(key);
638 let resp = self.session.api_call(|oauth| oauth.get(&url)).await?;
639 exists_from_response(resp, &format!("exists {key}"), NotFound::Status).await
640 }
641
642 async fn set_access(
643 &self,
644 desired: CloudAccessState,
645 ) -> Result<CloudAccessOutcome, CloudHomeError> {
646 let email = desired.require_provider_email("OneDrive")?;
647 let perms_url = format!(
648 "{}/drives/{}/items/{}/permissions",
649 self.graph_api, self.drive_id, self.folder_id
650 );
651 let access = sharing::SharedFolderAccess::new(
652 &self.session,
653 perms_url.clone(),
654 "value",
655 |permission: &serde_json::Value| {
656 permission["grantedToV2"]["user"]["email"]
657 .as_str()
658 .or_else(|| permission["grantedTo"]["user"]["email"].as_str())
659 .map(String::from)
660 },
661 |page: &serde_json::Value| {
662 Ok(page["@odata.nextLink"]
663 .as_str()
664 .map(std::string::ToString::to_string))
665 },
666 |permission_id: &str| format!("{perms_url}/{permission_id}"),
667 |permission: &serde_json::Value| {
668 permission["roles"]
669 .as_array()
670 .is_some_and(|roles| roles.iter().any(|role| role.as_str() == Some("write")))
671 },
672 "write",
673 format!(
674 "{}/drives/{}/items/{}/invite",
675 self.graph_api, self.drive_id, self.folder_id
676 ),
677 serde_json::json!({
678 "recipients": [{"email": email}],
679 "roles": ["write"],
680 "requireSignIn": true,
681 }),
682 );
683 match desired {
684 CloudAccessState::Present { .. } => {
685 access.ensure_present(email).await?;
686 Ok(CloudAccessOutcome::Present(CloudHomeJoinInfo::OneDrive {
687 drive_id: self.drive_id.clone(),
688 folder_id: self.folder_id.clone(),
689 }))
690 }
691 CloudAccessState::Absent { .. } => {
692 access.ensure_absent(email).await?;
693 Ok(CloudAccessOutcome::Absent(RevokeOutcome::Revoked))
694 }
695 }
696 }
697}
698
699#[async_trait]
700impl ExactSlotStorage for OneDriveCloudHome {
701 async fn provider_binding(
702 &self,
703 ) -> Result<coven_protocol::objects::ResolvedProviderBinding, CloudHomeError> {
704 use coven_protocol::objects::{
705 ProviderDeviceBinding, ProviderPrincipalId, ResolvedProviderBinding,
706 StoreProviderBinding,
707 };
708
709 if self.drive_id.is_empty() || self.folder_id.is_empty() {
710 return Err(CloudHomeError::Configuration(
711 "OneDrive provider binding has an empty drive or folder id".to_string(),
712 ));
713 }
714 let user_response = self
715 .session
716 .api_call(|oauth| {
717 oauth
718 .get(format!("{}/me", self.graph_api))
719 .query(&[("$select", "id")])
720 })
721 .await?;
722 let user_response = ensure_ok(
723 user_response,
724 "resolve OneDrive principal",
725 NotFound::Status,
726 )
727 .await?;
728 let user: serde_json::Value =
729 http::ok_json(user_response, "parse OneDrive principal").await?;
730 let user_id = user["id"]
731 .as_str()
732 .filter(|value| !value.is_empty())
733 .ok_or_else(|| {
734 CloudHomeError::Transport(
735 "OneDrive /me response omitted the stable user id".to_string(),
736 )
737 })?
738 .to_string();
739
740 let folder_response = self
741 .session
742 .api_call(|oauth| {
743 oauth
744 .get(format!(
745 "{}/drives/{}/items/{}",
746 self.graph_api, self.drive_id, self.folder_id
747 ))
748 .query(&[("$select", "id,parentReference,folder")])
749 })
750 .await?;
751 let folder_response = ensure_ok(
752 folder_response,
753 "resolve OneDrive folder binding",
754 NotFound::Status,
755 )
756 .await?;
757 let folder: serde_json::Value =
758 http::ok_json(folder_response, "parse OneDrive folder binding").await?;
759 if folder["id"].as_str() != Some(self.folder_id.as_str())
760 || folder["parentReference"]["driveId"].as_str() != Some(self.drive_id.as_str())
761 || !folder["folder"].is_object()
762 {
763 return Err(CloudHomeError::Transport(
764 "OneDrive folder lookup returned a different drive/folder binding".to_string(),
765 ));
766 }
767
768 Ok(ResolvedProviderBinding {
769 store: StoreProviderBinding::OneDrive {
770 drive_id: self.drive_id.clone(),
771 folder_id: self.folder_id.clone(),
772 },
773 device: ProviderDeviceBinding {
774 principal: ProviderPrincipalId::OneDrive { user_id },
775 },
776 })
777 }
778
779 async fn create_at(
780 &self,
781 upload: &super::ExactUpload<'_>,
782 control: &super::UploadControl,
783 ) -> Result<super::ExactCreateOutcome, CloudHomeError> {
784 if matches!(
785 self.exact_upload_verification,
786 coven_foundation::config::ExactUploadVerification::UploadChecksum
787 ) {
788 return Err(CloudHomeError::Configuration(
789 "OneDrive does not expose upload-checksum enforcement for this endpoint"
790 .to_string(),
791 ));
792 }
793 let operation = OneDriveCloudHome::create_at_slot(
794 self,
795 upload.object().slot(),
796 upload.body().await?,
797 control,
798 )
799 .await;
800 super::exact_upload::settle_exact_create(operation, |observed| {
801 self.verify_exact_upload(upload, observed)
802 })
803 .await
804 }
805 async fn list_slots(&self, prefix: &str) -> Result<Vec<ObjectSlot>, CloudHomeError> {
806 crate::cloud::logical_slots(CloudHome::list(self, prefix).await?)
807 }
808 async fn read_at(&self, slot: &ObjectSlot) -> Result<Vec<u8>, CloudHomeError> {
809 self.verify_slot(slot).await?;
810 OneDriveCloudHome::read(self, slot.logical_key()).await
811 }
812 async fn read_versioned_at(
813 &self,
814 slot: &ObjectSlot,
815 ) -> Result<CloudVersionedObject, CloudHomeError> {
816 let before = self.exact_metadata(slot).await?;
817 let bytes = OneDriveCloudHome::read(self, slot.logical_key()).await?;
818 let after = self.exact_metadata(slot).await?;
819 if before.version != after.version {
820 return Err(CloudHomeError::Transport(format!(
821 "OneDrive item {} changed while its versioned body was read",
822 slot.logical_key()
823 )));
824 }
825 if bytes.len() as u64 != after.size
826 || !content_hash::sha1_bytes(&bytes).eq_ignore_ascii_case(&after.sha1_hash)
827 {
828 return Err(CloudHomeError::Transport(format!(
829 "OneDrive item {} body differs from its versioned metadata",
830 slot.logical_key()
831 )));
832 }
833 Ok(CloudVersionedObject {
834 bytes,
835 version: after.version,
836 })
837 }
838 async fn replace_at_if_version(
839 &self,
840 slot: &ObjectSlot,
841 expected: &CloudObjectVersion,
842 bytes: Vec<u8>,
843 ) -> Result<ConditionalWriteOutcome, CloudHomeError> {
844 slot.require_logical_key_for("OneDrive")?;
845 let key = slot.logical_key();
846 let upload_url = self
847 .create_upload_session(key, "replace", UploadSessionCompletion::DeferredPersonal)
848 .await?;
849 let total = bytes.len() as u64;
850 let key_owned = key.to_string();
851 let classify =
852 Box::new(move |status, body: &str| classify_write_error(status, body, &key_owned));
853 let mut uploader = self.session.range_put_uploader(
854 upload_url.clone(),
855 202,
856 total,
857 ONEDRIVE_CHUNK_SIZE,
858 key.to_string(),
859 classify,
860 onedrive_upload_cancellation_succeeded,
861 );
862 let control = super::UploadControl::running(super::no_progress());
863 let mut offset = 0_u64;
864 for part in bytes.chunks(uploader.part_size()) {
865 let length = part.len() as u64;
866 let completion = uploader
867 .send_deferred_part(
868 Bytes::copy_from_slice(part),
869 offset,
870 offset + length == total,
871 &control,
872 )
873 .await?;
874 offset += length;
875 if completion.is_some() {
876 let operation = CloudHomeError::Transport(format!(
877 "conditional upload {key}: deferred upload published before explicit commit"
878 ));
879 let cleanup = self.delete_at_slot(slot).await;
880 return Err(combine_cleanup_failure(operation, cleanup));
881 }
882 }
883 if offset != total {
884 let operation = CloudHomeError::Transport(format!(
885 "conditional upload {key}: uploaded {offset} of {total} bytes"
886 ));
887 let cleanup = uploader.abort().await;
888 return Err(combine_cleanup_failure(operation, cleanup));
889 }
890 match self
891 .commit_deferred_replacement(slot, &upload_url, expected)
892 .await
893 {
894 Ok(ConditionalWriteOutcome::Replaced(version)) => {
895 uploader.mark_completed();
896 Ok(ConditionalWriteOutcome::Replaced(version))
897 }
898 Ok(ConditionalWriteOutcome::VersionChanged) => {
899 uploader.abort().await?;
900 Ok(ConditionalWriteOutcome::VersionChanged)
901 }
902 Err(operation) => {
903 let cleanup = uploader.abort().await;
904 Err(combine_cleanup_failure(operation, cleanup))
905 }
906 }
907 }
908 async fn read_range_at(
909 &self,
910 slot: &ObjectSlot,
911 start: u64,
912 end: u64,
913 ) -> Result<Vec<u8>, CloudHomeError> {
914 self.verify_slot(slot).await?;
915 OneDriveCloudHome::read_range(self, slot.logical_key(), start, end).await
916 }
917 async fn read_at_to_file(
918 &self,
919 slot: &ObjectSlot,
920 destination: &std::path::Path,
921 progress: super::DownloadProgress,
922 ) -> Result<(), super::CloudFileReadError> {
923 OneDriveCloudHome::read_at_to_file(self, slot, destination, progress).await
924 }
925 async fn delete_at(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
926 OneDriveCloudHome::delete_at_slot(self, slot).await
927 }
928}
929
930#[cfg(test)]
931#[path = "onedrive_tests.rs"]
932mod tests;