This is an automated email from the ASF dual-hosted git repository. erickguan pushed a commit to branch fix-reader in repository https://gitbox.apache.org/repos/asf/opendal.git
commit 960530ddb9c8fd9f573352506ebe4bb08aae33c0 Author: Erick Guan <[email protected]> AuthorDate: Sun Jul 5 21:49:42 2026 +0800 Progress --- core/core/src/types/context/read.rs | 2 +- core/services/sftp/src/backend.rs | 12 +- core/services/sftp/src/core.rs | 4 + core/services/sftp/src/reader.rs | 321 ++++++++++++++++++++++++++++++++++-- 4 files changed, 313 insertions(+), 26 deletions(-) diff --git a/core/core/src/types/context/read.rs b/core/core/src/types/context/read.rs index 0f8e64fb5..18cb028ca 100644 --- a/core/core/src/types/context/read.rs +++ b/core/core/src/types/context/read.rs @@ -35,7 +35,7 @@ pub struct ReadContext { args: OpRead, /// Options for the reader. options: OpReader, - /// Raw reader returned by [`Access::read`]. + /// Raw reader returned by [`Service::read`]. reader: oio::Reader, /// Complete object metadata observed from successful read opens. metadata: OnceLock<Metadata>, diff --git a/core/services/sftp/src/backend.rs b/core/services/sftp/src/backend.rs index 5b8cb8ab9..abceca072 100644 --- a/core/services/sftp/src/backend.rs +++ b/core/services/sftp/src/backend.rs @@ -203,7 +203,7 @@ pub struct SftpBackend { } impl Service for SftpBackend { - type Reader = oio::StreamReader<SftpReader>; + type Reader = SftpReader; type Writer = SftpLazyWriter; type Lister = SftpLazyLister; type Deleter = oio::OneShotDeleter<SftpDeleter>; @@ -255,15 +255,7 @@ impl Service for SftpBackend { Ok(RpStat::new(meta)) } fn read(&self, _ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> { - let output: oio::StreamReader<SftpReader> = { - Ok(oio::StreamReader::new(SftpReader::new( - self.clone(), - path, - args, - ))) - }?; - - Ok(output) + Ok(SftpReader::new(self.clone(), path, args)) } fn write(&self, ctx: &OperationContext, path: &str, op: OpWrite) -> Result<Self::Writer> { diff --git a/core/services/sftp/src/core.rs b/core/services/sftp/src/core.rs index 84ec65004..169ff76ec 100644 --- a/core/services/sftp/src/core.rs +++ b/core/services/sftp/src/core.rs @@ -27,12 +27,15 @@ use std::fmt::Debug; use std::path::Path; use std::path::PathBuf; use std::sync::Arc; +use std::sync::OnceLock; pub struct SftpCore { pub info: ServiceInfo, pub capability: Capability, pub endpoint: String, pub root: String, + /// Service-wide result of the lazy positioned-read probe. + pub positioned_read_support: OnceLock<bool>, client: Arc<bounded::Pool<Manager>>, } @@ -71,6 +74,7 @@ impl SftpCore { capability, endpoint, root, + positioned_read_support: OnceLock::new(), client, } } diff --git a/core/services/sftp/src/reader.rs b/core/services/sftp/src/reader.rs index a2cf46881..ed90a9b09 100644 --- a/core/services/sftp/src/reader.rs +++ b/core/services/sftp/src/reader.rs @@ -24,18 +24,25 @@ use super::lister::SftpLister; use super::writer::SftpWriter; use bytes::BytesMut; use fastpool::bounded; +use opendal_core::raw::oio::ReadStream as _; use opendal_core::raw::*; use opendal_core::*; +use openssh_sftp_client::Error as SftpClientError; +use openssh_sftp_client::error::SftpErrorKind; use openssh_sftp_client::file::File; +use std::error::Error as StdError; use std::io::SeekFrom; +use std::sync::Arc; +use std::sync::Mutex; use tokio::io::AsyncSeekExt; +const SFTP_READ_CHUNK_SIZE: usize = 2 * 1024 * 1024; + pub struct SftpReadStream { /// Keep the connection alive while data stream is alive. _conn: bounded::Object<Manager>, file: File, - chunk: usize, size: Option<usize>, read: usize, buf: BytesMut, @@ -47,7 +54,6 @@ impl SftpReadStream { _conn: conn, file, size: size.map(|v| v as usize), - chunk: 2 * 1024 * 1024, read: 0, buf: BytesMut::new(), } @@ -61,9 +67,9 @@ impl oio::ReadStream for SftpReadStream { } let size = if let Some(size) = self.size { - (size - self.read).min(self.chunk) + (size - self.read).min(SFTP_READ_CHUNK_SIZE) } else { - self.chunk + SFTP_READ_CHUNK_SIZE }; self.buf.reserve(size); @@ -84,9 +90,36 @@ impl oio::ReadStream for SftpReadStream { } /// Reader returned by this backend. +/// +/// On the first read, the reader opens a file handle and tries a positioned +/// read (seek + read). If the server supports it, the handle is cached for the +/// reader's lifetime and all subsequent reads go through the positioned path. +/// If the server does not support positioned reads, the result is cached +/// service-wide and the reader falls back to streaming. pub struct SftpReader { backend: SftpBackend, path: String, + /// Lazily resolved read strategy for this reader. + /// + /// This intentionally uses std synchronization instead of an async once + /// cell. A race can duplicate the first open/probe, but that does not + /// affect correctness, and we prefer std tools when they are sufficient. + positioned: Arc<Mutex<Option<PositionedReadStrategy>>>, +} + +#[derive(Clone)] +enum PositionedReadStrategy { + Supported(Arc<SftpReaderHandle>), + Unsupported, +} + +impl PositionedReadStrategy { + fn handle(&self) -> Option<Arc<SftpReaderHandle>> { + match self { + Self::Supported(handle) => Some(handle.clone()), + Self::Unsupported => None, + } + } } impl SftpReader { @@ -94,12 +127,30 @@ impl SftpReader { Self { backend, path: path.to_string(), + positioned: Arc::new(Mutex::new(None)), } } -} -impl oio::StreamRead for SftpReader { - async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> { + async fn open_stream( + &self, + range: BytesRange, + ) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> { + let (client, mut file) = self.open_file().await?; + + if range.offset() != 0 { + file.seek(SeekFrom::Start(range.offset())) + .await + .map_err(new_std_io_error)?; + } + + let stream = SftpReadStream::new(client, file, range.size()); + Ok(( + RpRead::default(), + Box::new(stream) as Box<dyn oio::ReadStreamDyn>, + )) + } + + async fn open_file(&self) -> Result<(bounded::Object<Manager>, File)> { let backend = &self.backend; let path = self.path.as_str(); @@ -110,21 +161,261 @@ impl oio::StreamRead for SftpReader { let path = fs.canonicalize(path).await.map_err(parse_sftp_error)?; - let mut f = client + let file = client .open(path.as_path()) .await .map_err(parse_sftp_error)?; - if range.offset() != 0 { - f.seek(SeekFrom::Start(range.offset())) - .await - .map_err(new_std_io_error)?; + Ok((client, file)) + } + + /// Return a cached positioned-read handle for this reader, or `None` + /// if the server does not support positioned reads. + /// + /// On the very first call across the entire service, this opens a file + /// and attempts a positioned read to detect support. The file handle is + /// reused for subsequent reads when supported. The detection result is + /// cached in `SftpCore::positioned_read_support` so all future readers + /// skip the probe entirely. + async fn positioned_handle(&self) -> Result<Option<Arc<SftpReaderHandle>>> { + // Fast path: a previous reader already determined positioned reads + // are not supported, so no file needs to be opened. + if self + .backend + .core + .positioned_read_support + .get() + .is_some_and(|&supported| !supported) + { + return Ok(None); + } + + let cached = { self.positioned.lock().expect("mutex poisoned").clone() }; + if let Some(strategy) = cached { + return Ok(strategy.handle()); } - let rp = RpRead::default(); - let stream = SftpReadStream::new(client, f, range.size()); + let (client, file) = self.open_file().await?; + let handle = Arc::new(SftpReaderHandle::new(client, file)); + + // If a previous reader already resolved support, respect that. + let strategy = if let Some(&supported) = self.backend.core.positioned_read_support.get() { + if supported { + PositionedReadStrategy::Supported(handle) + } else { + PositionedReadStrategy::Unsupported + } + } else { + // First reader across the service: probe by trying a real + // positioned read on the file we already opened. + match probe_positioned_read(&handle).await { + Ok(()) => { + // Probe succeeded: positioned reads work. Cache the + // result service-wide so other readers skip the probe. + let _ = self.backend.core.positioned_read_support.set(true); + PositionedReadStrategy::Supported(handle) + } + Err(err) if is_positioned_read_unsupported(&err) => { + // Positioned reads are not supported. Cache the result + // service-wide so other readers skip the probe. + let _ = self.backend.core.positioned_read_support.set(false); + PositionedReadStrategy::Unsupported + } + Err(err) => return Err(err), + } + }; + + *self.positioned.lock().expect("mutex poisoned") = Some(strategy.clone()); + Ok(strategy.handle()) + } +} + +pub struct SftpReaderHandle { + /// Keep the connection alive while the file handle is cached by `SftpReader`. + _conn: bounded::Object<Manager>, + file: File, +} + +impl SftpReaderHandle { + pub(super) fn new(conn: bounded::Object<Manager>, file: File) -> Self { + Self { _conn: conn, file } + } +} + +async fn read_positioned(handle: &SftpReaderHandle, offset: u64, size: usize) -> Result<Buffer> { + if size == 0 { + return Ok(Buffer::new()); + } + + let mut file = handle.file.clone(); + file.seek(SeekFrom::Start(offset)) + .await + .map_err(new_std_io_error)?; + + match file + .read(size as u32, BytesMut::with_capacity(size)) + .await + .map_err(parse_sftp_error)? + { + Some(bytes) => Ok(Buffer::from(bytes.freeze())), + None => Ok(Buffer::new()), + } +} + +/// Probe whether positioned reads work by seeking to the end of the file and +/// attempting a small read. A server that supports positioned reads will +/// succeed (returning empty at EOF); a server that does not support seek +/// will return an error we classify as "unsupported". +async fn probe_positioned_read(handle: &SftpReaderHandle) -> Result<()> { + let mut file = handle.file.clone(); + let len = file.metadata().await.map_err(parse_sftp_error)?.len(); + + let probe_offset = len.unwrap_or(0); + file.seek(SeekFrom::Start(probe_offset)) + .await + .map_err(new_std_io_error)?; + + // A small read to confirm the data path works after seeking. + let _ = file + .read(1u32, BytesMut::with_capacity(1)) + .await + .map_err(parse_sftp_error)?; + + Ok(()) +} + +fn is_positioned_read_unsupported(err: &Error) -> bool { + if err.kind() == ErrorKind::Unsupported { + return true; + } + + err.source() + .and_then(|source| source.downcast_ref::<SftpClientError>()) + .is_some_and(|source| { + matches!( + source, + SftpClientError::UnsupportedSftpProtocol { .. } + | SftpClientError::SftpError(SftpErrorKind::OpUnsupported, _) + ) + }) +} + +impl oio::Read for SftpReader { + async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> { + if let Some(handle) = self.positioned_handle().await? { + let stream = SftpPositionedReadStream { + handle, + offset: range.offset(), + remaining: range.size(), + }; + Ok(( + RpRead::default(), + Box::new(stream) as Box<dyn oio::ReadStreamDyn>, + )) + } else { + self.open_stream(range).await + } + } + + async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> { + if let Some(handle) = self.positioned_handle().await? { + let size = range.size().ok_or_else(|| { + Error::new(ErrorKind::Unsupported, "read requires a bounded range") + })?; + + let mut offset = range.offset(); + let mut remaining = size; + let mut bufs = Vec::new(); + + while remaining > 0 { + let read_size = remaining.min(SFTP_READ_CHUNK_SIZE as u64) as usize; + let buf = read_positioned(&handle, offset, read_size).await?; + if buf.len() > read_size { + return Err(Error::new( + ErrorKind::Unexpected, + "reader got unexpected data size", + ) + .with_context("expect", read_size) + .with_context("actual", buf.len())); + } + if buf.is_empty() { + return Err(Error::new( + ErrorKind::RangeNotSatisfied, + "range exceeds content length", + ) + .with_context("offset", offset) + .with_context("remaining", remaining)); + } + + let n = buf.len() as u64; + offset += n; + remaining -= n; + bufs.push(buf); + } + + Ok((RpRead::default(), bufs.into_iter().flatten().collect())) + } else { + let expected = range.size().ok_or_else(|| { + Error::new(ErrorKind::Unsupported, "read requires a bounded range") + })?; + let (rp, mut stream) = self.open_stream(range).await?; + let buffer = stream.read_all().await?; + if buffer.len() as u64 != expected { + return Err( + Error::new(ErrorKind::Unexpected, "reader got unexpected data size") + .with_context("expect", expected) + .with_context("actual", buffer.len() as u64), + ); + } + Ok((rp, buffer)) + } + } +} + +struct SftpPositionedReadStream { + handle: Arc<SftpReaderHandle>, + offset: u64, + remaining: Option<u64>, +} + +impl oio::ReadStream for SftpPositionedReadStream { + async fn read(&mut self) -> Result<Buffer> { + if self.remaining == Some(0) { + return Ok(Buffer::new()); + } + + let read_size = self + .remaining + .map(|remaining| remaining.min(SFTP_READ_CHUNK_SIZE as u64) as usize) + .unwrap_or(SFTP_READ_CHUNK_SIZE); + + let buf = read_positioned(&self.handle, self.offset, read_size).await?; + if buf.len() > read_size { + return Err( + Error::new(ErrorKind::Unexpected, "reader got unexpected data size") + .with_context("expect", read_size) + .with_context("actual", buf.len()), + ); + } + if buf.is_empty() { + if let Some(remaining) = self.remaining { + return Err(Error::new( + ErrorKind::RangeNotSatisfied, + "range exceeds content length", + ) + .with_context("offset", self.offset) + .with_context("remaining", remaining)); + } + return Ok(Buffer::new()); + } + + let n = buf.len() as u64; + self.offset += n; + if let Some(remaining) = &mut self.remaining { + *remaining -= n; + } - Ok((rp, Box::new(stream) as Box<dyn oio::ReadStreamDyn>)) + Ok(buf) } }
