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
9fn decode_time_flow_string(str: &str, time_now: i64) -> ApiResult<String> {
11 if let Ok(result) = decode_time_flow_string_time(str, time_now) {
13 return Ok(result);
14 }
15
16 if let Ok(result) = decode_time_flow_string_time(str, time_now - 1000) {
18 return Ok(result);
19 }
20
21 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
31fn 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 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 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 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 secret.pop();
90
91 add = !add;
92
93 time_digit_index -= 1;
95 if time_digit_index < 0 {
96 time_digit_index = time_digits.len() as i64 - 1;
97 }
98 }
99
100 let res_str: String = res.iter().collect();
102
103 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 let prefix_len = res_sep[0].len() + 1; if prefix_len <= res_str.len() {
115 Ok(res_str[prefix_len..].to_string())
116 } else {
117 Ok(String::new())
118 }
119}
120
121#[async_trait]
123pub trait ExplorerApi {
124 async fn list_files(&self, params: &ListFileService) -> ApiResult<ListResponse>;
126
127 async fn get_file_thumb(
129 &self,
130 path: &str,
131 context_hint: Option<&str>,
132 ) -> ApiResult<FileThumbResponse>;
133
134 async fn get_file_info(&self, params: &GetFileInfoService) -> ApiResult<FileResponse>;
136
137 async fn create_file(&self, request: &CreateFileService) -> ApiResult<FileResponse>;
139
140 async fn delete_files(&self, request: &DeleteFileService) -> ApiResult<()>;
142
143 async fn rename_file(&self, request: &RenameFileService) -> ApiResult<FileResponse>;
145
146 async fn move_files(&self, request: &MoveFileService) -> ApiResult<()>;
148
149 async fn restore_files(&self, request: &DeleteFileService) -> ApiResult<()>;
151
152 async fn patch_metadata(&self, request: &PatchMetadataService) -> ApiResult<()>;
154
155 async fn get_file_url(&self, request: &FileURLService) -> ApiResult<FileURLResponse>;
157
158 async fn unlock_files(&self, request: &UnlockFileService) -> ApiResult<()>;
160
161 async fn set_current_version(&self, request: &VersionControlService) -> ApiResult<()>;
163
164 async fn delete_version(&self, request: &VersionControlService) -> ApiResult<()>;
166
167 async fn update_file(&self, params: &FileUpdateService, data: Bytes)
169 -> ApiResult<FileResponse>;
170
171 async fn get_storage_policy_options(&self) -> ApiResult<Vec<StoragePolicy>>;
173
174 async fn mount_storage_policy(
176 &self,
177 request: &MountPolicyService,
178 ) -> ApiResult<Vec<StoragePolicy>>;
179
180 async fn set_permissions(&self, request: &SetPermissionService) -> ApiResult<()>;
182
183 async fn create_upload_session(
185 &self,
186 request: &UploadSessionRequest,
187 ) -> ApiResult<UploadCredential>;
188
189 async fn upload_chunk(
191 &self,
192 session_id: &str,
193 chunk_index: usize,
194 data: Bytes,
195 ) -> ApiResult<()>;
196
197 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 async fn delete_upload_session(&self, request: &DeleteUploadSessionService) -> ApiResult<()>;
208
209 async fn complete_s3_upload(
211 &self,
212 policy_type: &str,
213 session_id: &str,
214 session_key: &str,
215 ) -> ApiResult<()>;
216
217 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 let (page, next_token) = if let Some(prev) = previous_response {
243 let prev_pagination = &prev.res.pagination;
244
245 if prev_pagination.next_token.is_some() {
247 (None, prev_pagination.next_token.clone())
249 } else if prev_pagination.total_items.is_some() {
250 let current_page = prev_pagination.page;
252 (Some(current_page + 1), None)
253 } else {
254 (None, None)
256 }
257 } else {
258 (None, None)
260 };
261
262 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(¶ms).await?;
273
274 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 true
284 } else if let Some(total_items) = response.pagination.total_items {
285 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 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 let mut query_params = vec![format!("uri={}", urlencoding::encode(¶ms.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) = ¶ms.order_by {
314 query_params.push(format!("order_by={}", order_by));
315 }
316 if let Some(order_direction) = ¶ms.order_direction {
317 query_params.push(format!("order_direction={}", order_direction));
318 }
319 if let Some(next_page_token) = ¶ms.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 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 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) = ¶ms.uri {
364 query_params.push(format!("uri={}", urlencoding::encode(uri)));
365 }
366 if let Some(id) = ¶ms.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 let mut query_params = vec![format!("uri={}", urlencoding::encode(¶ms.uri))];
501
502 if let Some(previous) = ¶ms.previous {
503 query_params.push(format!("previous={}", previous));
504 }
505
506 let query = format!("?{}", query_params.join("&"));
507
508 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
661pub struct FileEventSubscription {
663 response: reqwest::Response,
664 buffer: String,
665}
666
667impl FileEventSubscription {
668 fn new(response: reqwest::Response) -> Self {
670 Self {
671 response,
672 buffer: String::new(),
673 }
674 }
675
676 pub async fn next_event(&mut self) -> ApiResult<Option<FileEvent>> {
679 loop {
680 if let Some(event) = self.try_parse_event()? {
682 return Ok(Some(event));
683 }
684
685 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 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 fn try_parse_event(&mut self) -> ApiResult<Option<FileEvent>> {
710 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 let event_block = self.buffer[..event_end].to_string();
728 self.buffer = self.buffer[event_end..].to_string();
729
730 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 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 if data_str == "<nil>" || data_str.is_empty() {
752 return Ok(None);
754 }
755 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 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 Ok(None)
774 }
775 }
776 }
777}
778
779#[async_trait]
781pub trait FileEventsApi {
782 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 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 Ok(FileEventSubscription::new(response))
849 } else {
850 let response_text = response.text().await?;
853
854 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 Err(crate::error::ApiError::SseNotUpgraded {
868 code: -1,
869 message: format!("Unexpected response: {}", response_text),
870 })
871 }
872 }
873}