Skip to main content

coven_storage/cloud/
onedrive.rs

1//! OneDrive `CloudHome` implementation.
2//!
3//! Uses the Microsoft Graph API. Files are stored flat in a single folder — path
4//! separators are escaped by the `key_encoding` helpers. The
5//! `read`/`read_range`/`list`/`delete` methods are the shared `OAuthRestHome`
6//! implementations; this file supplies only the Graph request shapes, the page
7//! parser, the upload paths, and sharing.
8
9use 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
42/// OneDrive cloud home backend.
43pub 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                // Lets the joiner fetch its account email for OAuth folder sharing.
83                "User.Read".to_string(),
84            ],
85            redirect_port: 19284,
86            extra_auth_params: vec![],
87        }
88    }
89
90    /// Build the Graph API URL for a file by encoded name within the app folder.
91    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    /// Build the Graph API URL for the folder's children endpoint.
102    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
468/// `error.code` from a Microsoft Graph error body (e.g. `"quotaLimitReached"`),
469/// or `None` if the body isn't Graph JSON.
470fn 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
476/// Map a OneDrive write failure to a `CloudHomeError`. The `quotaLimitReached`
477/// code gets a message naming the provider and the recovery step; everything else
478/// keeps the raw HTTP status + body for debugging.
479fn 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
493/// Files at or below this size go up as a single PUT; larger files use a resumable
494/// session. Microsoft Graph caps a simple PUT at 250 MiB.
495const ONEDRIVE_SIMPLE_PUT_MAX: usize = 4 * 1024 * 1024;
496
497/// Resumable-session part size. Graph requires every part except the last to be a
498/// multiple of 320 KiB; 7.5 MiB (24 × 320 KiB) keeps the request count low.
499const 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        // `@odata.nextLink` is a full URL with all params; the first page is the
536        // children endpoint selecting only names.
537        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        // OneDrive returns 202 Accepted for every non-final part.
605        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;