Skip to main content

cloudreve_sdk_api/api/
explorer.rs

1use crate::client::{Client, RequestOptions, CR_HEADER_PREFIX};
2use crate::error::ApiResult;
3use crate::models::common::ListAllRes;
4use crate::models::explorer::*;
5use async_trait::async_trait;
6use bytes::Bytes;
7use reqwest::Body;
8
9/// Decode time flow string (for obfuscated thumbnail URLs)
10fn decode_time_flow_string(str: &str, time_now: i64) -> ApiResult<String> {
11    // Try with current time
12    if let Ok(result) = decode_time_flow_string_time(str, time_now) {
13        return Ok(result);
14    }
15
16    // Try with time - 1000
17    if let Ok(result) = decode_time_flow_string_time(str, time_now - 1000) {
18        return Ok(result);
19    }
20
21    // Try with time + 1000
22    if let Ok(result) = decode_time_flow_string_time(str, time_now + 1000) {
23        return Ok(result);
24    }
25
26    Err(crate::error::ApiError::Other(
27        "Failed to decode time flow string".to_string(),
28    ))
29}
30
31/// Decode time flow string time (for obfuscated thumbnail URLs)
32fn decode_time_flow_string_time(str: &str, time_now: i64) -> ApiResult<String> {
33    let mut time_now = time_now / 1000;
34    let time_now_backup = time_now;
35
36    // Extract time digits
37    let mut time_digits: Vec<i64> = Vec::new();
38
39    if str.is_empty() {
40        return Ok(String::new());
41    }
42
43    while time_now > 0 {
44        time_digits.push(time_now % 10);
45        time_now /= 10;
46    }
47
48    if time_digits.is_empty() {
49        return Err(crate::error::ApiError::Other(
50            "Invalid time value".to_string(),
51        ));
52    }
53
54    // Convert string to character array
55    let chars: Vec<char> = str.chars().collect();
56    let mut res: Vec<char> = chars.clone();
57    let mut secret: Vec<char> = chars.clone();
58
59    let mut add = secret.len() % 2 == 0;
60    let mut time_digit_index = ((secret.len() - 1) % time_digits.len()) as i64;
61    let l = secret.len();
62
63    for pos in 0..l {
64        let res_index = l - 1 - pos;
65        let mut new_index = res_index as i64;
66
67        if add {
68            new_index = new_index + time_digits[time_digit_index as usize] * time_digit_index;
69        } else {
70            new_index = 2 * time_digit_index * time_digits[time_digit_index as usize] - new_index;
71        }
72
73        if new_index < 0 {
74            new_index = new_index * -1;
75        }
76
77        new_index = new_index % secret.len() as i64;
78        let new_index_usize = new_index as usize;
79
80        res[res_index] = secret[new_index_usize];
81
82        // Swap elements in secret
83        let a = secret[res_index];
84        let b = secret[new_index_usize];
85        secret[new_index_usize] = a;
86        secret[res_index] = b;
87
88        // Remove last element from secret
89        secret.pop();
90
91        add = !add;
92
93        // Decrement timeDigitIndex
94        time_digit_index -= 1;
95        if time_digit_index < 0 {
96            time_digit_index = time_digits.len() as i64 - 1;
97        }
98    }
99
100    // Convert result back to string
101    let res_str: String = res.iter().collect();
102
103    // Validate the result
104    let res_sep: Vec<&str> = res_str.split('|').collect();
105
106    if res_sep.is_empty() || res_sep[0] != time_now_backup.to_string() {
107        return Err(crate::error::ApiError::Other(
108            "Invalid time flow string".to_string(),
109        ));
110    }
111
112    // Return the part after the first "|"
113    let prefix_len = res_sep[0].len() + 1; // +1 for the "|"
114    if prefix_len <= res_str.len() {
115        Ok(res_str[prefix_len..].to_string())
116    } else {
117        Ok(String::new())
118    }
119}
120
121/// File explorer API methods
122#[async_trait]
123pub trait ExplorerApi {
124    /// List files in a directory
125    async fn list_files(&self, params: &ListFileService) -> ApiResult<ListResponse>;
126
127    /// Get file thumbnail
128    async fn get_file_thumb(
129        &self,
130        path: &str,
131        context_hint: Option<&str>,
132    ) -> ApiResult<FileThumbResponse>;
133
134    /// Get file information
135    async fn get_file_info(&self, params: &GetFileInfoService) -> ApiResult<FileResponse>;
136
137    /// Create a new file or folder
138    async fn create_file(&self, request: &CreateFileService) -> ApiResult<FileResponse>;
139
140    /// Delete files
141    async fn delete_files(&self, request: &DeleteFileService) -> ApiResult<()>;
142
143    /// Rename a file
144    async fn rename_file(&self, request: &RenameFileService) -> ApiResult<FileResponse>;
145
146    /// Move files
147    async fn move_files(&self, request: &MoveFileService) -> ApiResult<()>;
148
149    /// Restore files from trash
150    async fn restore_files(&self, request: &DeleteFileService) -> ApiResult<()>;
151
152    /// Patch file metadata
153    async fn patch_metadata(&self, request: &PatchMetadataService) -> ApiResult<()>;
154
155    /// Get file entity URL
156    async fn get_file_url(&self, request: &FileURLService) -> ApiResult<FileURLResponse>;
157
158    /// Unlock files
159    async fn unlock_files(&self, request: &UnlockFileService) -> ApiResult<()>;
160
161    /// Set current version
162    async fn set_current_version(&self, request: &VersionControlService) -> ApiResult<()>;
163
164    /// Delete version
165    async fn delete_version(&self, request: &VersionControlService) -> ApiResult<()>;
166
167    /// Update file content
168    async fn update_file(&self, params: &FileUpdateService, data: Bytes)
169        -> ApiResult<FileResponse>;
170
171    /// Get storage policy options
172    async fn get_storage_policy_options(&self) -> ApiResult<Vec<StoragePolicy>>;
173
174    /// Mount storage policy
175    async fn mount_storage_policy(
176        &self,
177        request: &MountPolicyService,
178    ) -> ApiResult<Vec<StoragePolicy>>;
179
180    /// Set file permissions
181    async fn set_permissions(&self, request: &SetPermissionService) -> ApiResult<()>;
182
183    /// Create upload session
184    async fn create_upload_session(
185        &self,
186        request: &UploadSessionRequest,
187    ) -> ApiResult<UploadCredential>;
188
189    /// Upload chunk
190    async fn upload_chunk(
191        &self,
192        session_id: &str,
193        chunk_index: usize,
194        data: Bytes,
195    ) -> ApiResult<()>;
196
197    /// Upload chunk using streaming body (for large chunks)
198    async fn upload_chunk_stream(
199        &self,
200        session_id: &str,
201        chunk_index: usize,
202        content_length: u64,
203        body: Body,
204    ) -> ApiResult<()>;
205
206    /// Delete upload session
207    async fn delete_upload_session(&self, request: &DeleteUploadSessionService) -> ApiResult<()>;
208
209    /// Complete S3-like upload
210    async fn complete_s3_upload(
211        &self,
212        policy_type: &str,
213        session_id: &str,
214        session_key: &str,
215    ) -> ApiResult<()>;
216
217    /// Complete OneDrive upload
218    async fn complete_onedrive_upload(&self, session_id: &str, session_key: &str) -> ApiResult<()>;
219}
220
221#[async_trait]
222pub trait ExplorerApiExt {
223    async fn list_files_all(
224        &self,
225        previous_response: Option<&ListAllRes<ListResponse>>,
226        uri: &str,
227        page_size: i32,
228    ) -> ApiResult<ListAllRes<ListResponse>>;
229}
230
231#[async_trait]
232impl ExplorerApiExt for Client {
233    async fn list_files_all(
234        &self,
235        previous_response: Option<&ListAllRes<ListResponse>>,
236        uri: &str,
237        page_size: i32,
238    ) -> ApiResult<ListAllRes<ListResponse>> {
239        const MIN_PAGE_SIZE: i32 = 1;
240
241        // Extract pagination info from previous response
242        let (page, next_token) = if let Some(prev) = previous_response {
243            let prev_pagination = &prev.res.pagination;
244
245            // Determine next page parameters based on pagination type
246            if prev_pagination.next_token.is_some() {
247                // Token-based pagination
248                (None, prev_pagination.next_token.clone())
249            } else if prev_pagination.total_items.is_some() {
250                // Page-based pagination
251                let current_page = prev_pagination.page;
252                (Some(current_page + 1), None)
253            } else {
254                // No pagination info, start fresh
255                (None, None)
256            }
257        } else {
258            // First page
259            (None, None)
260        };
261
262        // Call list_files with current pagination state
263        let params = ListFileService {
264            uri: uri.to_string(),
265            page,
266            page_size: Some(page_size),
267            order_by: None,
268            order_direction: None,
269            next_page_token: next_token,
270        };
271
272        let response = self.list_files(&params).await?;
273
274        // Determine if there's more data to load
275        let page_size_val = if response.pagination.page_size > 0 {
276            response.pagination.page_size
277        } else {
278            MIN_PAGE_SIZE
279        };
280
281        let has_more = if response.pagination.next_token.is_some() {
282            // Token-based: more data if next_token exists
283            true
284        } else if let Some(total_items) = response.pagination.total_items {
285            // Page-based: calculate if there are more pages
286            let total_pages = (total_items as f64 / page_size_val as f64).ceil() as i32;
287            let current_page = response.pagination.page;
288            current_page + 1 < total_pages
289        } else {
290            // No pagination info, assume no more data
291            false
292        };
293
294        Ok(ListAllRes {
295            res: response,
296            more: has_more,
297        })
298    }
299}
300
301#[async_trait]
302impl ExplorerApi for Client {
303    async fn list_files(&self, params: &ListFileService) -> ApiResult<ListResponse> {
304        // Build query string
305        let mut query_params = vec![format!("uri={}", urlencoding::encode(&params.uri))];
306
307        if let Some(page) = params.page {
308            query_params.push(format!("page={}", page));
309        }
310        if let Some(page_size) = params.page_size {
311            query_params.push(format!("page_size={}", page_size));
312        }
313        if let Some(order_by) = &params.order_by {
314            query_params.push(format!("order_by={}", order_by));
315        }
316        if let Some(order_direction) = &params.order_direction {
317            query_params.push(format!("order_direction={}", order_direction));
318        }
319        if let Some(next_page_token) = &params.next_page_token {
320            query_params.push(format!("next_page_token={}", next_page_token));
321        }
322
323        let query = format!("?{}", query_params.join("&"));
324
325        self.get(
326            &format!("/file{}", query),
327            RequestOptions::new().with_purchase_ticket(),
328        )
329        .await
330    }
331
332    async fn get_file_thumb(
333        &self,
334        path: &str,
335        _context_hint: Option<&str>,
336    ) -> ApiResult<FileThumbResponse> {
337        let query = format!("?uri={}", urlencoding::encode(path));
338
339        // TODO: Add context hint header support if needed
340        let mut response: FileThumbResponse = self
341            .get(
342                &format!("/file/thumb{}", query),
343                RequestOptions::new().with_purchase_ticket(),
344            )
345            .await?;
346
347        if response.obfuscated {
348            // Decode the obfuscated URL
349            let time_now_sec = std::time::SystemTime::now()
350                .duration_since(std::time::UNIX_EPOCH)
351                .map_err(|e| crate::error::ApiError::Other(format!("System time error: {}", e)))?
352                .as_secs() as i64;
353
354            response.url = decode_time_flow_string(&response.url, time_now_sec)?;
355        }
356
357        Ok(response)
358    }
359
360    async fn get_file_info(&self, params: &GetFileInfoService) -> ApiResult<FileResponse> {
361        let mut query_params = vec![];
362
363        if let Some(uri) = &params.uri {
364            query_params.push(format!("uri={}", urlencoding::encode(uri)));
365        }
366        if let Some(id) = &params.id {
367            query_params.push(format!("id={}", id));
368        }
369        if let Some(extended) = params.extended {
370            query_params.push(format!("extended={}", extended));
371        }
372        if let Some(folder_summary) = params.folder_summary {
373            query_params.push(format!("folder_summary={}", folder_summary));
374        }
375
376        let query = if query_params.is_empty() {
377            String::new()
378        } else {
379            format!("?{}", query_params.join("&"))
380        };
381
382        self.get(
383            &format!("/file/info{}", query),
384            RequestOptions::new().with_purchase_ticket(),
385        )
386        .await
387    }
388
389    async fn create_file(&self, request: &CreateFileService) -> ApiResult<FileResponse> {
390        self.post(
391            "/file/create",
392            request,
393            RequestOptions::new().with_purchase_ticket(),
394        )
395        .await
396    }
397
398    async fn delete_files(&self, request: &DeleteFileService) -> ApiResult<()> {
399        let opts = if request.uris.len() == 1 {
400            RequestOptions::new()
401                .with_purchase_ticket()
402                .skip_batch_error()
403        } else {
404            RequestOptions::new().with_purchase_ticket()
405        };
406
407        self.delete_with_body("/file", request, opts).await
408    }
409
410    async fn rename_file(&self, request: &RenameFileService) -> ApiResult<FileResponse> {
411        self.post(
412            "/file/rename",
413            request,
414            RequestOptions::new().with_purchase_ticket(),
415        )
416        .await
417    }
418
419    async fn move_files(&self, request: &MoveFileService) -> ApiResult<()> {
420        let opts = if request.uris.len() == 1 {
421            RequestOptions::new()
422                .with_purchase_ticket()
423                .skip_batch_error()
424        } else {
425            RequestOptions::new().with_purchase_ticket()
426        };
427
428        self.post("/file/move", request, opts).await
429    }
430
431    async fn restore_files(&self, request: &DeleteFileService) -> ApiResult<()> {
432        let opts = if request.uris.len() == 1 {
433            RequestOptions::new()
434                .with_purchase_ticket()
435                .skip_batch_error()
436        } else {
437            RequestOptions::new().with_purchase_ticket()
438        };
439
440        self.post("/file/restore", request, opts).await
441    }
442
443    async fn patch_metadata(&self, request: &PatchMetadataService) -> ApiResult<()> {
444        let opts = if request.uris.len() == 1 {
445            RequestOptions::new()
446                .with_purchase_ticket()
447                .skip_batch_error()
448        } else {
449            RequestOptions::new().with_purchase_ticket()
450        };
451
452        self.patch("/file/metadata", request, opts).await
453    }
454
455    async fn get_file_url(&self, request: &FileURLService) -> ApiResult<FileURLResponse> {
456        let opts = if request.uris.len() == 1 {
457            RequestOptions::new()
458                .with_purchase_ticket()
459                .skip_batch_error()
460        } else {
461            RequestOptions::new().with_purchase_ticket()
462        };
463
464        self.post("/file/url", request, opts).await
465    }
466
467    async fn unlock_files(&self, request: &UnlockFileService) -> ApiResult<()> {
468        self.delete_with_body(
469            "/file/lock",
470            request,
471            RequestOptions::new().skip_lock_conflict(),
472        )
473        .await
474    }
475
476    async fn set_current_version(&self, request: &VersionControlService) -> ApiResult<()> {
477        self.post(
478            "/file/version/current",
479            request,
480            RequestOptions::new().with_purchase_ticket(),
481        )
482        .await
483    }
484
485    async fn delete_version(&self, request: &VersionControlService) -> ApiResult<()> {
486        self.delete_with_body(
487            "/file/version",
488            request,
489            RequestOptions::new().with_purchase_ticket(),
490        )
491        .await
492    }
493
494    async fn update_file(
495        &self,
496        params: &FileUpdateService,
497        data: Bytes,
498    ) -> ApiResult<FileResponse> {
499        // Build query string
500        let mut query_params = vec![format!("uri={}", urlencoding::encode(&params.uri))];
501
502        if let Some(previous) = &params.previous {
503            query_params.push(format!("previous={}", previous));
504        }
505
506        let query = format!("?{}", query_params.join("&"));
507
508        // We need to use a custom request here since we're sending binary data
509        let url = self.build_url(&format!("/file/content{}", query));
510        let token = self.get_access_token().await?;
511
512        let response = self
513            .http_client
514            .put(&url)
515            .header("Authorization", format!("Bearer {}", token))
516            .header("Content-Type", "application/octet-stream")
517            .body(data)
518            .send()
519            .await?;
520
521        let api_response: crate::error::ApiResponse<FileResponse> = response.json().await?;
522
523        if api_response.code != 0 {
524            return Err(crate::error::ApiError::from_response(api_response));
525        }
526
527        api_response.data.ok_or_else(|| {
528            crate::error::ApiError::Other("API returned success but no data".to_string())
529        })
530    }
531
532    async fn get_storage_policy_options(&self) -> ApiResult<Vec<StoragePolicy>> {
533        self.get("/user/setting/policies", RequestOptions::new())
534            .await
535    }
536
537    async fn mount_storage_policy(
538        &self,
539        request: &MountPolicyService,
540    ) -> ApiResult<Vec<StoragePolicy>> {
541        self.patch(
542            "/file/policy",
543            request,
544            RequestOptions::new().with_purchase_ticket(),
545        )
546        .await
547    }
548
549    async fn set_permissions(&self, request: &SetPermissionService) -> ApiResult<()> {
550        let opts = if request.uris.len() == 1 {
551            RequestOptions::new()
552                .with_purchase_ticket()
553                .skip_batch_error()
554        } else {
555            RequestOptions::new().with_purchase_ticket()
556        };
557
558        self.post("/file/permission", request, opts).await
559    }
560
561    async fn create_upload_session(
562        &self,
563        request: &UploadSessionRequest,
564    ) -> ApiResult<UploadCredential> {
565        self.put(
566            "/file/upload",
567            request,
568            RequestOptions::new().with_purchase_ticket(),
569        )
570        .await
571    }
572
573    async fn upload_chunk(
574        &self,
575        session_id: &str,
576        chunk_index: usize,
577        data: Bytes,
578    ) -> ApiResult<()> {
579        let url = self.build_url(&format!("/file/upload/{}/{}", session_id, chunk_index));
580        let token = self.get_access_token().await?;
581
582        let response = self
583            .http_client
584            .post(&url)
585            .header("Authorization", format!("Bearer {}", token))
586            .header("Content-Type", "application/octet-stream")
587            .body(data)
588            .send()
589            .await?;
590
591        let api_response: crate::error::ApiResponse<UploadCredential> = response.json().await?;
592
593        if api_response.code != 0 {
594            return Err(crate::error::ApiError::from_response(api_response));
595        }
596
597        Ok(())
598    }
599
600    async fn upload_chunk_stream(
601        &self,
602        session_id: &str,
603        chunk_index: usize,
604        content_length: u64,
605        body: Body,
606    ) -> ApiResult<()> {
607        let url = self.build_url(&format!("/file/upload/{}/{}", session_id, chunk_index));
608        let token = self.get_access_token().await?;
609
610        let response = self
611            .http_client
612            .post(&url)
613            .header("Authorization", format!("Bearer {}", token))
614            .header("Content-Type", "application/octet-stream")
615            .header("Content-Length", content_length)
616            .body(body)
617            .send()
618            .await?;
619
620        let api_response: crate::error::ApiResponse<()> = response.json().await?;
621
622        if api_response.code != 0 {
623            return Err(crate::error::ApiError::from_response(api_response));
624        }
625
626        Ok(())
627    }
628
629    async fn delete_upload_session(&self, request: &DeleteUploadSessionService) -> ApiResult<()> {
630        self.delete_with_body(
631            "/file/upload",
632            request,
633            RequestOptions::new().with_purchase_ticket(),
634        )
635        .await
636    }
637
638    async fn complete_s3_upload(
639        &self,
640        policy_type: &str,
641        session_id: &str,
642        session_key: &str,
643    ) -> ApiResult<()> {
644        self.get(
645            &format!("/callback/{}/{}/{}", policy_type, session_id, session_key),
646            RequestOptions::new(),
647        )
648        .await
649    }
650
651    async fn complete_onedrive_upload(&self, session_id: &str, session_key: &str) -> ApiResult<()> {
652        self.post::<(), ()>(
653            &format!("/callback/onedrive/{}/{}", session_id, session_key),
654            &(),
655            RequestOptions::new(),
656        )
657        .await
658    }
659}
660
661/// A subscription handle for file events SSE stream
662pub struct FileEventSubscription {
663    response: reqwest::Response,
664    buffer: String,
665}
666
667impl FileEventSubscription {
668    /// Create a new subscription from a response
669    fn new(response: reqwest::Response) -> Self {
670        Self {
671            response,
672            buffer: String::new(),
673        }
674    }
675
676    /// Receive the next file event from the stream.
677    /// Returns None when the stream ends.
678    pub async fn next_event(&mut self) -> ApiResult<Option<FileEvent>> {
679        loop {
680            // Try to parse a complete event from the buffer
681            if let Some(event) = self.try_parse_event()? {
682                return Ok(Some(event));
683            }
684
685            // Need more data from the stream
686            match self.response.chunk().await {
687                Ok(Some(chunk)) => {
688                    let text = String::from_utf8_lossy(&chunk);
689                    self.buffer.push_str(&text);
690                }
691                Ok(None) => {
692                    // Stream ended
693                    // Try to parse any remaining data
694                    if !self.buffer.is_empty() {
695                        if let Some(event) = self.try_parse_event()? {
696                            return Ok(Some(event));
697                        }
698                    }
699                    return Ok(None);
700                }
701                Err(e) => {
702                    return Err(crate::error::ApiError::SseStreamError(e.to_string()));
703                }
704            }
705        }
706    }
707
708    /// Try to parse a complete SSE event from the buffer
709    fn try_parse_event(&mut self) -> ApiResult<Option<FileEvent>> {
710        // SSE events are separated by double newlines
711        // Format:
712        // event:eventname
713        // data:payload
714        //
715        // (blank line)
716
717        // Find the end of an event (double newline)
718        let event_end = if let Some(pos) = self.buffer.find("\n\n") {
719            pos + 2
720        } else if let Some(pos) = self.buffer.find("\r\n\r\n") {
721            pos + 4
722        } else {
723            return Ok(None);
724        };
725
726        // Extract the event block
727        let event_block = self.buffer[..event_end].to_string();
728        self.buffer = self.buffer[event_end..].to_string();
729
730        // Parse the event
731        let mut event_type: Option<&str> = None;
732        let mut data: Option<&str> = None;
733
734        for line in event_block.lines() {
735            if let Some(rest) = line.strip_prefix("event:") {
736                event_type = Some(rest.trim());
737            } else if let Some(rest) = line.strip_prefix("data:") {
738                data = Some(rest.trim());
739            }
740        }
741
742        // Match on event type
743        match event_type {
744            Some("resumed") => Ok(Some(FileEvent::Resumed)),
745            Some("subscribed") => Ok(Some(FileEvent::Subscribed)),
746            Some("keep-alive") | Some("keepalive") => Ok(Some(FileEvent::KeepAlive)),
747            Some("reconnect-required") => Ok(Some(FileEvent::ReconnectRequired)),
748            Some("event") => {
749                if let Some(data_str) = data {
750                    // Skip nil data
751                    if data_str == "<nil>" || data_str.is_empty() {
752                        // This shouldn't happen for "event" type, but handle gracefully
753                        return Ok(None);
754                    }
755                    // Try to parse as array first (batch of events)
756                    if let Ok(event_data_list) =
757                        serde_json::from_str::<Vec<FileEventData>>(data_str)
758                    {
759                        if event_data_list.is_empty() {
760                            return Ok(None);
761                        }
762                        return Ok(Some(FileEvent::Event(event_data_list)));
763                    }
764                    // Fall back to parsing as single event for backwards compatibility
765                    let event_data: FileEventData = serde_json::from_str(data_str)?;
766                    Ok(Some(FileEvent::Event(vec![event_data])))
767                } else {
768                    Ok(None)
769                }
770            }
771            _ => {
772                // Unknown event type, skip it
773                Ok(None)
774            }
775        }
776    }
777}
778
779/// File events SSE API methods
780#[async_trait]
781pub trait FileEventsApi {
782    /// Subscribe to file events for a given URI.
783    ///
784    /// This connects to the SSE endpoint at /v4/file/events with the provided URI.
785    /// Returns a subscription handle that can be used to receive events.
786    ///
787    /// # Arguments
788    /// * `uri` - The filesystem URI to watch for events (e.g., "cloudreve://my-drive/")
789    ///
790    /// # Returns
791    /// * `Ok(FileEventSubscription)` - A handle to receive events from
792    /// * `Err(ApiError::SseNotUpgraded)` - If the server returned an error instead of SSE stream
793    /// * `Err(ApiError::RequestError)` - If the HTTP request failed
794    ///
795    /// # Example
796    /// ```no_run
797    /// use cloudreve_sdk_api::{Client, ClientConfig};
798    /// use cloudreve_sdk_api::api::explorer::FileEventsApi;
799    ///
800    /// async fn watch_events(client: &Client) -> Result<(), Box<dyn std::error::Error>> {
801    ///     let mut subscription = client.subscribe_file_events("cloudreve://my-drive/").await?;
802    ///
803    ///     while let Some(event) = subscription.next_event().await? {
804    ///         match event {
805    ///             cloudreve_sdk_api::models::explorer::FileEvent::Event(events) => {
806    ///                 for data in events {
807    ///                     println!("File event: {:?} on {}", data.event_type, data.from);
808    ///                 }
809    ///             }
810    ///             _ => {}
811    ///         }
812    ///     }
813    ///
814    ///     Ok(())
815    /// }
816    /// ```
817    async fn subscribe_file_events(&self, uri: &str) -> ApiResult<FileEventSubscription>;
818}
819
820#[async_trait]
821impl FileEventsApi for Client {
822    async fn subscribe_file_events(&self, uri: &str) -> ApiResult<FileEventSubscription> {
823        let query = format!("?uri={}", urlencoding::encode(uri));
824        let url = self.build_url(&format!("/file/events{}", query));
825        let token = self.get_access_token().await?;
826
827        let response = self
828            .http_client
829            .get(&url)
830            .header(
831                format!("{}Client-Id", CR_HEADER_PREFIX),
832                self.config.client_id.clone(),
833            )
834            .header("Authorization", format!("Bearer {}", token))
835            .header("Accept", "text/event-stream")
836            .send()
837            .await?;
838
839        // Check if we got an SSE response by looking at content-type
840        let content_type = response
841            .headers()
842            .get(reqwest::header::CONTENT_TYPE)
843            .and_then(|v| v.to_str().ok())
844            .unwrap_or("");
845
846        if content_type.contains("text/event-stream") {
847            // Successfully upgraded to SSE
848            Ok(FileEventSubscription::new(response))
849        } else {
850            // Server returned a regular response (likely an error)
851            // Try to parse it as an API error response
852            let response_text = response.text().await?;
853
854            // Try to parse as API response
855            if let Ok(api_response) =
856                serde_json::from_str::<crate::error::ApiResponse<()>>(&response_text)
857            {
858                if api_response.code != 0 {
859                    return Err(crate::error::ApiError::SseNotUpgraded {
860                        code: api_response.code,
861                        message: api_response.msg,
862                    });
863                }
864            }
865
866            // If we couldn't parse it, return a generic error
867            Err(crate::error::ApiError::SseNotUpgraded {
868                code: -1,
869                message: format!("Unexpected response: {}", response_text),
870            })
871        }
872    }
873}