1use std::path::{Component, Path, PathBuf};
18use std::sync::{Arc, Mutex};
19use std::vec::IntoIter;
20use std::{fmt, io};
21
22use cap_std::ambient_authority;
23use cap_std::fs::{Dir, DirEntry, OpenOptions, ReadDir};
24use opendal::layers::SimulateLayer;
25use opendal::raw::*;
26use opendal::{
27 Buffer, BytesRange, Capability, EntryMode, Error, ErrorKind, Metadata, OperationContext,
28 Operator, Result,
29};
30use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
31
32const LIST_BATCH_SIZE: usize = 128;
33
34#[derive(Clone)]
36pub struct SecureFsRoot {
37 dir: Arc<Dir>,
38 path: Arc<PathBuf>,
39}
40
41impl fmt::Debug for SecureFsRoot {
42 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
43 f.debug_struct("SecureFsRoot")
44 .field("path", &self.path)
45 .finish_non_exhaustive()
46 }
47}
48
49impl SecureFsRoot {
50 pub fn open(path: impl AsRef<Path>) -> io::Result<Self> {
54 let path = path.as_ref();
55 std::fs::create_dir_all(path)?;
56 if std::fs::symlink_metadata(path)?.file_type().is_symlink() {
57 return Err(io::Error::new(
58 io::ErrorKind::InvalidInput,
59 "filesystem sandbox root must not be a symbolic link",
60 ));
61 }
62
63 let path = path.canonicalize()?;
64 let dir = Dir::open_ambient_dir(&path, ambient_authority())?;
65 Ok(Self {
66 dir: Arc::new(dir),
67 path: Arc::new(path),
68 })
69 }
70
71 pub fn path(&self) -> &Path {
73 &self.path
74 }
75
76 pub fn open_subdir(&self, path: impl AsRef<Path>) -> io::Result<Self> {
78 let path = normalize_relative_path(path.as_ref())?;
79 if path.as_os_str().is_empty() {
80 return Ok(self.clone());
81 }
82
83 let dir = self.dir.open_dir(&path)?;
84 Ok(Self {
85 dir: Arc::new(dir),
86 path: Arc::new(self.path.join(path)),
87 })
88 }
89
90 pub fn create_subdir(&self, path: impl AsRef<Path>) -> io::Result<Self> {
92 let path = normalize_relative_path(path.as_ref())?;
93 if path.as_os_str().is_empty() {
94 return Ok(self.clone());
95 }
96
97 self.dir.create_dir_all(&path)?;
98 let dir = self.dir.open_dir(&path)?;
99 Ok(Self {
100 dir: Arc::new(dir),
101 path: Arc::new(self.path.join(path)),
102 })
103 }
104
105 pub fn build_operator(&self) -> Operator {
107 Operator::from_parts(
108 OperationContext::default(),
109 Arc::new(SecureFsBackend::new(self.clone())) as Servicer,
110 )
111 .layer(SimulateLayer::default())
112 }
113}
114
115fn normalize_relative_path(path: &Path) -> io::Result<PathBuf> {
116 let mut normalized = PathBuf::new();
117 for component in path.components() {
118 match component {
119 Component::CurDir => {}
120 Component::Normal(value) => normalized.push(value),
121 Component::ParentDir | Component::RootDir | Component::Prefix(_) => {
122 return Err(io::Error::new(
123 io::ErrorKind::PermissionDenied,
124 "path escapes the filesystem sandbox",
125 ));
126 }
127 }
128 }
129 Ok(normalized)
130}
131
132fn backend_path(path: &str) -> io::Result<PathBuf> {
133 let path = path.trim_matches('/');
134 if path.is_empty() {
135 Ok(PathBuf::new())
136 } else {
137 normalize_relative_path(Path::new(path))
138 }
139}
140
141fn parse_write_error(error: io::Error, if_not_exists: bool) -> Error {
142 if if_not_exists && error.kind() == io::ErrorKind::AlreadyExists {
143 Error::new(
144 ErrorKind::ConditionNotMatch,
145 "the file already exists in the filesystem",
146 )
147 .set_source(error)
148 } else {
149 new_std_io_error(error)
150 }
151}
152
153fn metadata_from_fs(metadata: cap_std::fs::Metadata) -> Result<Metadata> {
154 let mode = if metadata.is_dir() {
155 EntryMode::DIR
156 } else if metadata.is_file() {
157 EntryMode::FILE
158 } else {
159 EntryMode::Unknown
160 };
161
162 Ok(Metadata::new(mode)
163 .with_content_length(metadata.len())
164 .with_last_modified(Timestamp::try_from(
165 metadata.modified().map_err(new_std_io_error)?.into_std(),
166 )?))
167}
168
169#[derive(Clone, Debug)]
170struct SecureFsBackend {
171 root: SecureFsRoot,
172 info: ServiceInfo,
173 capability: Capability,
174}
175
176impl SecureFsBackend {
177 fn new(root: SecureFsRoot) -> Self {
178 let info = ServiceInfo::new("fs", root.path().to_string_lossy(), "");
179 let capability = Capability {
180 stat: true,
181 read: true,
182 write: true,
183 write_can_empty: true,
184 write_can_append: true,
185 write_can_multi: true,
186 write_with_if_not_exists: true,
187 create_dir: true,
188 delete: true,
189 delete_with_recursive: true,
190 list: true,
191 shared: true,
192 ..Default::default()
193 };
194 Self {
195 root,
196 info,
197 capability,
198 }
199 }
200}
201
202impl Service for SecureFsBackend {
203 type Reader = oio::StreamReader<SecureFsReader>;
204 type Writer = SecureFsWriter;
205 type Lister = SecureFsLister;
206 type Deleter = oio::OneShotDeleter<SecureFsDeleter>;
207 type Copier = ();
208
209 fn info(&self) -> ServiceInfo {
210 self.info.clone()
211 }
212
213 fn capability(&self) -> Capability {
214 self.capability
215 }
216
217 async fn create_dir(
218 &self,
219 _: &OperationContext,
220 path: &str,
221 _: OpCreateDir,
222 ) -> Result<RpCreateDir> {
223 let path = backend_path(path).map_err(new_std_io_error)?;
224 let root = self.root.clone();
225 common_runtime::spawn_blocking_global(move || root.dir.create_dir_all(path))
226 .await
227 .map_err(new_task_join_error)?
228 .map_err(new_std_io_error)?;
229 Ok(RpCreateDir::default())
230 }
231
232 async fn stat(&self, _: &OperationContext, path: &str, _: OpStat) -> Result<RpStat> {
233 let path = backend_path(path).map_err(new_std_io_error)?;
234 let root = self.root.clone();
235 let metadata = common_runtime::spawn_blocking_global(move || {
236 if path.as_os_str().is_empty() {
237 root.dir.dir_metadata()
238 } else {
239 root.dir.metadata(path)
240 }
241 })
242 .await
243 .map_err(new_task_join_error)?
244 .map_err(new_std_io_error)?;
245 Ok(RpStat::new(metadata_from_fs(metadata)?))
246 }
247
248 fn read(&self, _: &OperationContext, path: &str, _: OpRead) -> Result<Self::Reader> {
249 let path = backend_path(path).map_err(new_std_io_error)?;
250 Ok(oio::StreamReader::new(SecureFsReader {
251 root: self.root.clone(),
252 path,
253 }))
254 }
255
256 fn write(&self, _: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
257 let path = backend_path(path).map_err(new_std_io_error)?;
258 Ok(SecureFsWriter {
259 root: self.root.clone(),
260 path,
261 args,
262 file: None,
263 })
264 }
265
266 fn delete(&self, _: &OperationContext) -> Result<Self::Deleter> {
267 Ok(oio::OneShotDeleter::new(SecureFsDeleter {
268 root: self.root.clone(),
269 }))
270 }
271
272 fn list(&self, _: &OperationContext, path: &str, _: OpList) -> Result<Self::Lister> {
273 let path = backend_path(path).map_err(new_std_io_error)?;
274 let display_prefix = if path.as_os_str().is_empty() {
275 String::new()
276 } else {
277 format!("{}/", path.to_string_lossy().replace('\\', "/"))
278 };
279 Ok(SecureFsLister {
280 root: self.root.clone(),
281 path,
282 display_prefix,
283 read_dir: None,
284 entries: vec![].into_iter(),
285 seeded: false,
286 done: false,
287 })
288 }
289
290 fn copy(
291 &self,
292 _: &OperationContext,
293 _: &str,
294 _: &str,
295 _: OpCopy,
296 _: OpCopier,
297 ) -> Result<Self::Copier> {
298 Err(Error::new(
299 ErrorKind::Unsupported,
300 "operation is not supported",
301 ))
302 }
303
304 async fn rename(
305 &self,
306 _: &OperationContext,
307 _: &str,
308 _: &str,
309 _: OpRename,
310 ) -> Result<RpRename> {
311 Err(Error::new(
312 ErrorKind::Unsupported,
313 "operation is not supported",
314 ))
315 }
316
317 async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> {
318 Err(Error::new(
319 ErrorKind::Unsupported,
320 "operation is not supported",
321 ))
322 }
323}
324
325struct SecureFsReader {
326 root: SecureFsRoot,
327 path: PathBuf,
328}
329
330impl oio::StreamRead for SecureFsReader {
331 async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
332 let path = self.path.clone();
333 let root = self.root.clone();
334 let file = common_runtime::spawn_blocking_global(move || root.dir.open(path))
335 .await
336 .map_err(new_task_join_error)?
337 .map_err(new_std_io_error)?;
338 let mut file = tokio::fs::File::from_std(file.into_std());
339 if range.offset() != 0 {
340 file.seek(io::SeekFrom::Start(range.offset()))
341 .await
342 .map_err(new_std_io_error)?;
343 }
344 Ok((
345 RpRead::default(),
346 Box::new(SecureFsReadStream {
347 file,
348 remaining: range.size().unwrap_or(u64::MAX),
349 }),
350 ))
351 }
352}
353
354struct SecureFsReadStream {
355 file: tokio::fs::File,
356 remaining: u64,
357}
358
359impl oio::ReadStream for SecureFsReadStream {
360 async fn read(&mut self) -> Result<Buffer> {
361 if self.remaining == 0 {
362 return Ok(Buffer::new());
363 }
364
365 let size = self.remaining.min(2 * 1024 * 1024) as usize;
366 let mut buffer = vec![0; size];
367 let read = self
368 .file
369 .read(&mut buffer)
370 .await
371 .map_err(new_std_io_error)?;
372 self.remaining = self.remaining.saturating_sub(read as u64);
373 buffer.truncate(read);
374 Ok(Buffer::from(buffer))
375 }
376}
377
378struct SecureFsWriter {
379 root: SecureFsRoot,
380 path: PathBuf,
381 args: OpWrite,
382 file: Option<tokio::fs::File>,
383}
384
385impl SecureFsWriter {
386 async fn ensure_file(&mut self) -> Result<&mut tokio::fs::File> {
387 if self.file.is_none() {
388 let path = self.path.clone();
389 let root = self.root.clone();
390 let if_not_exists = self.args.if_not_exists();
391 let append = self.args.append();
392 let file = common_runtime::spawn_blocking_global(move || {
393 if let Some(parent) = path.parent()
394 && !parent.as_os_str().is_empty()
395 {
396 root.dir.create_dir_all(parent).map_err(new_std_io_error)?;
397 }
398
399 let mut options = OpenOptions::new();
400 options.write(true);
401 if if_not_exists {
402 options.create_new(true);
403 } else {
404 options.create(true);
405 }
406 if append {
407 options.append(true);
408 } else {
409 options.truncate(true);
410 }
411 root.dir
412 .open_with(path, &options)
413 .map_err(|error| parse_write_error(error, if_not_exists))
414 })
415 .await
416 .map_err(new_task_join_error)??;
417
418 self.file = Some(tokio::fs::File::from_std(file.into_std()));
419 }
420 Ok(self.file.as_mut().expect("file must be initialized"))
421 }
422}
423
424impl oio::Write for SecureFsWriter {
425 async fn write(&mut self, buffer: Buffer) -> Result<()> {
426 self.ensure_file()
427 .await?
428 .write_all(&buffer.to_bytes())
429 .await
430 .map_err(new_std_io_error)
431 }
432
433 async fn close(&mut self) -> Result<Metadata> {
434 let file = self.ensure_file().await?;
435 file.flush().await.map_err(new_std_io_error)?;
436 file.sync_all().await.map_err(new_std_io_error)?;
437 let metadata = file.metadata().await.map_err(new_std_io_error)?;
438 Ok(Metadata::new(EntryMode::FILE)
439 .with_content_length(metadata.len())
440 .with_last_modified(Timestamp::try_from(
441 metadata.modified().map_err(new_std_io_error)?,
442 )?))
443 }
444
445 async fn abort(&mut self) -> Result<()> {
446 Err(Error::new(
447 ErrorKind::Unsupported,
448 "filesystem writes cannot be aborted without atomic writes",
449 ))
450 }
451}
452
453struct SecureFsLister {
454 root: SecureFsRoot,
455 path: PathBuf,
456 display_prefix: String,
457 read_dir: Option<Arc<Mutex<ReadDir>>>,
458 entries: IntoIter<oio::Entry>,
459 seeded: bool,
460 done: bool,
461}
462
463impl oio::List for SecureFsLister {
464 async fn next(&mut self) -> Result<Option<oio::Entry>> {
465 if !self.seeded {
466 self.seeded = true;
467 let path = self.path.clone();
468 let root = self.root.clone();
469 let display_prefix = self.display_prefix.clone();
470 let read_dir = common_runtime::spawn_blocking_global(move || {
471 let result = (|| {
472 let dir = if path.as_os_str().is_empty() {
473 root.dir.open_dir(".")?
474 } else {
475 root.dir.open_dir(&path)?
476 };
477 dir.entries()
478 })();
479
480 match result {
481 Ok(read_dir) => Ok(Some(read_dir)),
482 Err(error)
483 if matches!(
484 error.kind(),
485 io::ErrorKind::NotFound | io::ErrorKind::NotADirectory
486 ) =>
487 {
488 Ok(None)
489 }
490 Err(error) => Err(error),
491 }
492 })
493 .await
494 .map_err(new_task_join_error)?
495 .map_err(new_std_io_error)?;
496
497 let Some(read_dir) = read_dir else {
498 self.done = true;
499 return Ok(None);
500 };
501 self.read_dir = Some(Arc::new(Mutex::new(read_dir)));
502 let current_path = oio::Entry::new(
503 if display_prefix.is_empty() {
504 "/"
505 } else {
506 &display_prefix
507 },
508 Metadata::new(EntryMode::DIR),
509 );
510 self.entries = vec![current_path].into_iter();
511 }
512
513 if let Some(entry) = self.entries.next() {
514 return Ok(Some(entry));
515 }
516 if self.done {
517 return Ok(None);
518 }
519
520 let Some(read_dir) = self.read_dir.clone() else {
521 self.done = true;
522 return Ok(None);
523 };
524 let display_prefix = self.display_prefix.clone();
525 let (entries, done) = common_runtime::spawn_blocking_global(move || {
526 let mut read_dir = read_dir
527 .lock()
528 .map_err(|_| io::Error::other("filesystem directory iterator lock is poisoned"))?;
529 read_list_batch(&mut read_dir, &display_prefix)
530 })
531 .await
532 .map_err(new_task_join_error)?
533 .map_err(new_std_io_error)?;
534
535 self.entries = entries.into_iter();
536 self.done = done;
537 Ok(self.entries.next())
538 }
539}
540
541fn read_list_batch(
542 read_dir: &mut ReadDir,
543 display_prefix: &str,
544) -> io::Result<(Vec<oio::Entry>, bool)> {
545 let mut entries = Vec::with_capacity(LIST_BATCH_SIZE);
546 while entries.len() < LIST_BATCH_SIZE {
547 let entry = match read_dir.next() {
548 Some(Ok(entry)) => entry,
549 Some(Err(error)) if error.kind() == io::ErrorKind::NotFound => {
550 return Ok((entries, true));
551 }
552 Some(Err(error)) => return Err(error),
553 None => return Ok((entries, true)),
554 };
555
556 if let Some(entry) = read_list_entry(entry, display_prefix)? {
557 entries.push(entry);
558 }
559 }
560 Ok((entries, false))
561}
562
563fn read_list_entry(entry: DirEntry, display_prefix: &str) -> io::Result<Option<oio::Entry>> {
564 let file_type = match entry.file_type() {
565 Ok(file_type) => file_type,
566 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
567 Err(error) => return Err(error),
568 };
569 let name = entry.file_name().to_string_lossy().to_string();
570 let (path, mode) = if file_type.is_dir() {
571 (format!("{display_prefix}{name}/"), EntryMode::DIR)
572 } else if file_type.is_file() {
573 (format!("{display_prefix}{name}"), EntryMode::FILE)
574 } else {
575 (format!("{display_prefix}{name}"), EntryMode::Unknown)
576 };
577 let metadata = if mode == EntryMode::Unknown {
578 Metadata::new(mode)
579 } else {
580 match entry.metadata() {
581 Ok(metadata) => match metadata_from_fs(metadata) {
582 Ok(metadata) => metadata,
583 Err(error) if error.kind() == ErrorKind::NotFound => return Ok(None),
584 Err(error) => return Err(io::Error::other(error.to_string())),
585 },
586 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
587 Err(error) => return Err(error),
588 }
589 };
590 Ok(Some(oio::Entry::new(&path, metadata)))
591}
592
593struct SecureFsDeleter {
594 root: SecureFsRoot,
595}
596
597impl oio::OneShotDelete for SecureFsDeleter {
598 async fn delete_once(&self, path: String, args: OpDelete) -> Result<()> {
599 let path = backend_path(&path).map_err(new_std_io_error)?;
600 if path.as_os_str().is_empty() {
601 return Err(Error::new(
602 ErrorKind::Unsupported,
603 "deleting the filesystem sandbox root is not supported",
604 ));
605 }
606 let root = self.root.clone();
607 common_runtime::spawn_blocking_global(move || {
608 let metadata = match root.dir.symlink_metadata(&path) {
609 Ok(metadata) => metadata,
610 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
611 Err(error) => return Err(error),
612 };
613
614 if metadata.is_dir() {
615 if args.recursive() {
616 root.dir.remove_dir_all(path)
617 } else {
618 root.dir.remove_dir(path)
619 }
620 } else {
621 root.dir.remove_file(path)
622 }
623 })
624 .await
625 .map_err(new_task_join_error)?
626 .map_err(new_std_io_error)
627 }
628}
629
630#[cfg(test)]
631mod tests {
632 use bytes::Bytes;
633 use common_test_util::temp_dir::create_temp_dir;
634 use opendal::raw::oio::List;
635 use opendal::raw::{OpList, Service};
636 use opendal::{BytesRange, ErrorKind, OperationContext};
637
638 use super::{LIST_BATCH_SIZE, SecureFsBackend, SecureFsRoot, read_list_entry};
639
640 #[tokio::test]
641 async fn test_operator_suffix_reads_final_bytes() {
642 let temp_dir = create_temp_dir("secure_fs_operator_suffix");
643 std::fs::write(temp_dir.path().join("file"), b"0123456789").unwrap();
644 let operator = SecureFsRoot::open(temp_dir.path())
645 .unwrap()
646 .build_operator();
647
648 assert_eq!(
649 Bytes::from_static(b"789"),
650 operator
651 .read_with("file")
652 .range(7..10)
653 .await
654 .unwrap()
655 .to_bytes()
656 );
657 assert_eq!(
658 Bytes::from_static(b"789"),
659 operator
660 .read_with("file")
661 .range(BytesRange::Suffix { size: 3 })
662 .await
663 .unwrap()
664 .to_bytes()
665 );
666 }
667
668 #[tokio::test]
669 async fn test_lister_streams_entries() {
670 let temp_dir = create_temp_dir("secure_fs_lister_streams_entries");
671 for index in 0..129 {
672 std::fs::write(temp_dir.path().join(format!("{index}.parquet")), []).unwrap();
673 }
674
675 let root = SecureFsRoot::open(temp_dir.path()).unwrap();
676 let backend = SecureFsBackend::new(root);
677 let ctx = OperationContext::default();
678 let mut lister = backend.list(&ctx, "/", OpList::new()).unwrap();
679
680 let mut paths = Vec::new();
681 while let Some(entry) = lister.next().await.unwrap() {
682 paths.push(entry.path().to_string());
683 assert!(lister.entries.len() <= LIST_BATCH_SIZE);
684 }
685 assert_eq!(130, paths.len());
686 assert!(paths.iter().any(|path| path == "/"));
687 assert!(paths.iter().any(|path| path == "128.parquet"));
688 }
689
690 #[tokio::test]
691 async fn test_if_not_exists_returns_condition_not_match() {
692 let temp_dir = create_temp_dir("secure_fs_if_not_exists");
693 let operator = SecureFsRoot::open(temp_dir.path())
694 .unwrap()
695 .build_operator();
696 operator
697 .write("existing", Bytes::from_static(b"original"))
698 .await
699 .unwrap();
700
701 let error = operator
702 .write_with("existing", Bytes::from_static(b"replacement"))
703 .if_not_exists(true)
704 .await
705 .unwrap_err();
706
707 assert_eq!(ErrorKind::ConditionNotMatch, error.kind());
708 assert_eq!(
709 Bytes::from_static(b"original"),
710 operator.read("existing").await.unwrap().to_bytes()
711 );
712 }
713
714 #[tokio::test]
715 async fn test_if_not_exists_does_not_remap_parent_directory_error() {
716 let temp_dir = create_temp_dir("secure_fs_if_not_exists_parent_error");
717 std::fs::write(temp_dir.path().join("parent"), []).unwrap();
718 let operator = SecureFsRoot::open(temp_dir.path())
719 .unwrap()
720 .build_operator();
721
722 let error = operator
723 .write_with("parent/file", Bytes::new())
724 .if_not_exists(true)
725 .await
726 .unwrap_err();
727
728 assert_eq!(ErrorKind::AlreadyExists, error.kind());
729 }
730
731 #[tokio::test]
732 async fn test_list_missing_or_non_directory_is_empty() {
733 let temp_dir = create_temp_dir("secure_fs_list_missing_or_non_directory");
734 std::fs::write(temp_dir.path().join("file"), []).unwrap();
735 let operator = SecureFsRoot::open(temp_dir.path())
736 .unwrap()
737 .build_operator();
738
739 assert!(operator.list("missing/").await.unwrap().is_empty());
740 assert!(operator.list("file/").await.unwrap().is_empty());
741 }
742
743 #[test]
752 fn test_lister_skips_entry_removed_during_iteration() {
753 let temp_dir = create_temp_dir("secure_fs_lister_removed_entry");
754 let path = temp_dir.path().join("removed");
755 std::fs::write(&path, []).unwrap();
756 let root = SecureFsRoot::open(temp_dir.path()).unwrap();
757 let mut read_dir = root.dir.entries().unwrap();
758 let entry = read_dir.next().unwrap().unwrap();
759 std::fs::remove_file(path).unwrap();
760
761 let result = read_list_entry(entry, "").unwrap();
762
763 if let Some(entry) = &result {
764 assert_eq!("removed", entry.path());
765 }
766
767 #[cfg(not(windows))]
768 assert!(result.is_none());
769 }
770
771 #[tokio::test]
772 async fn test_delete_root_is_unsupported() {
773 let temp_dir = create_temp_dir("secure_fs_delete_root");
774 let operator = SecureFsRoot::open(temp_dir.path())
775 .unwrap()
776 .build_operator();
777 operator
778 .write("nested/file", Bytes::from_static(b"data"))
779 .await
780 .unwrap();
781
782 let error = operator.delete_with("/").recursive(true).await.unwrap_err();
783
784 assert_eq!(ErrorKind::Unsupported, error.kind());
785 assert!(temp_dir.path().join("nested/file").exists());
786
787 operator
788 .delete_with("nested/")
789 .recursive(true)
790 .await
791 .unwrap();
792 assert!(!temp_dir.path().join("nested").exists());
793 }
794
795 #[tokio::test]
796 async fn test_writer_abort_is_unsupported_without_atomic_write() {
797 let temp_dir = create_temp_dir("secure_fs_writer_abort");
798 let operator = SecureFsRoot::open(temp_dir.path())
799 .unwrap()
800 .build_operator();
801 let mut writer = operator.writer("partial").await.unwrap();
802 writer.write(Bytes::from_static(b"partial")).await.unwrap();
803
804 let error = writer.abort().await.unwrap_err();
805
806 assert_eq!(ErrorKind::Unsupported, error.kind());
807 }
808}