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
The following commit(s) were added to refs/heads/fix-reader by this push:
new de26b1963 fix(services/sftp): replace PositionRead with adaptive
probe-based reader
de26b1963 is described below
commit de26b1963f2ef1b3262501e0b38be4c1c89d86f2
Author: Erick Guan <[email protected]>
AuthorDate: Sun Jul 5 13:50:37 2026 +0800
fix(services/sftp): replace PositionRead with adaptive probe-based reader
Replace the static PositionRead implementation with an adaptive reader
that probes the server for positioned-read support at runtime.
Previously, SftpReader always used oio::PositionRead, assuming all SFTP
servers support seek-based positioned reads. Some servers (e.g. those
with older SFTP protocol versions or limited extensions) do not support
seek + read, causing failures.
Now SftpReader implements oio::Read directly and:
- Probes positioned-read support on first use (reading one byte at EOF
or offset 0), caching the result service-wide in SftpCore via
OnceCell<bool> so all subsequent readers skip the probe.
- If positioned reads are supported, uses a cached SftpReaderHandle with
SftpPositionedReadStream for efficient range reads (seek + read).
- If not supported, falls back to SftpReadStream which uses sequential
reads with an optional initial seek (open_stream path).
Add the mea dependency for OnceCell used by the probe cache.
---
core/Cargo.lock | 1 +
core/services/sftp/Cargo.toml | 1 +
core/services/sftp/src/backend.rs | 8 +-
core/services/sftp/src/core.rs | 4 +
core/services/sftp/src/reader.rs | 326 ++++++++++++++++++++++++++++++++++----
5 files changed, 306 insertions(+), 34 deletions(-)
diff --git a/core/Cargo.lock b/core/Cargo.lock
index 1a34d396d..aec48713f 100644
--- a/core/Cargo.lock
+++ b/core/Cargo.lock
@@ -7507,6 +7507,7 @@ dependencies = [
"fastpool",
"futures",
"log",
+ "mea",
"opendal-core",
"openssh",
"openssh-sftp-client",
diff --git a/core/services/sftp/Cargo.toml b/core/services/sftp/Cargo.toml
index 705cb28e2..fc3f59d2e 100644
--- a/core/services/sftp/Cargo.toml
+++ b/core/services/sftp/Cargo.toml
@@ -37,6 +37,7 @@ bytes = { workspace = true }
fastpool = "1.0.2"
futures = { workspace = true }
log = { workspace = true }
+mea = { workspace = true }
openssh = "0.11.0"
openssh-sftp-client = { version = "0.15.3", features = ["openssh", "tracing"] }
serde = { workspace = true, features = ["derive"] }
diff --git a/core/services/sftp/src/backend.rs
b/core/services/sftp/src/backend.rs
index 256cbd1da..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::PositionReader<SftpReader>;
+ type Reader = SftpReader;
type Writer = SftpLazyWriter;
type Lister = SftpLazyLister;
type Deleter = oio::OneShotDeleter<SftpDeleter>;
@@ -255,11 +255,7 @@ impl Service for SftpBackend {
Ok(RpStat::new(meta))
}
fn read(&self, _ctx: &OperationContext, path: &str, args: OpRead) ->
Result<Self::Reader> {
- Ok(oio::PositionReader::new(SftpReader::new(
- self.clone(),
- path,
- args,
- )))
+ 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..853944510 100644
--- a/core/services/sftp/src/core.rs
+++ b/core/services/sftp/src/core.rs
@@ -17,6 +17,7 @@
use fastpool::{ManageObject, ObjectStatus, bounded};
use log::debug;
+use mea::once::OnceCell;
use opendal_core::raw::*;
use opendal_core::*;
use openssh::KnownHosts;
@@ -33,6 +34,8 @@ pub struct SftpCore {
pub capability: Capability,
pub endpoint: String,
pub root: String,
+ /// Service-wide result of the lazy positioned-read probe.
+ pub positioned_read_support: OnceCell<bool>,
client: Arc<bounded::Pool<Manager>>,
}
@@ -71,6 +74,7 @@ impl SftpCore {
capability,
endpoint,
root,
+ positioned_read_support: OnceCell::new(),
client,
}
}
diff --git a/core/services/sftp/src/reader.rs b/core/services/sftp/src/reader.rs
index 39e97a053..227fa55b8 100644
--- a/core/services/sftp/src/reader.rs
+++ b/core/services/sftp/src/reader.rs
@@ -22,17 +22,80 @@ use super::core::is_sftp_failure;
use super::core::parse_sftp_error;
use super::lister::SftpLister;
use super::writer::SftpWriter;
+use bytes::BytesMut;
use fastpool::bounded;
+use mea::once::OnceCell;
+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 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,
+ size: Option<usize>,
+ read: usize,
+ buf: BytesMut,
+}
+
+impl SftpReadStream {
+ pub fn new(conn: bounded::Object<Manager>, file: File, size: Option<u64>)
-> Self {
+ Self {
+ _conn: conn,
+ file,
+ size: size.map(|v| v as usize),
+ read: 0,
+ buf: BytesMut::new(),
+ }
+ }
+}
+
+impl oio::ReadStream for SftpReadStream {
+ async fn read(&mut self) -> Result<Buffer> {
+ if self.read >= self.size.unwrap_or(usize::MAX) {
+ return Ok(Buffer::new());
+ }
+
+ let size = if let Some(size) = self.size {
+ (size - self.read).min(SFTP_READ_CHUNK_SIZE)
+ } else {
+ SFTP_READ_CHUNK_SIZE
+ };
+ self.buf.reserve(size);
+
+ let Some(bytes) = self
+ .file
+ .read(size as u32, self.buf.split_off(0))
+ .await
+ .map_err(parse_sftp_error)?
+ else {
+ return Ok(Buffer::new());
+ };
+
+ self.read += bytes.len();
+ self.buf = bytes;
+ let bs = self.buf.split();
+ Ok(Buffer::from(bs.freeze()))
+ }
+}
+
/// Reader returned by this backend.
pub struct SftpReader {
backend: SftpBackend,
path: String,
+ /// Lazily opened positioned-read handle, cached per reader for the file's
+ /// lifetime. Only set when the server supports positioned reads.
+ positioned: Arc<OnceCell<Arc<SftpReaderHandle>>>,
}
impl SftpReader {
@@ -40,12 +103,91 @@ impl SftpReader {
Self {
backend,
path: path.to_string(),
+ positioned: Arc::new(OnceCell::new()),
+ }
+ }
+
+ 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();
+
+ let client = backend.core.connect().await?;
+
+ let mut fs = client.fs();
+ fs.set_cwd(&backend.core.root);
+
+ let path = fs.canonicalize(path).await.map_err(parse_sftp_error)?;
+
+ let file = client
+ .open(path.as_path())
+ .await
+ .map_err(parse_sftp_error)?;
+
+ Ok((client, file))
+ }
+
+ /// Probe whether the server supports positioned reads (seek + read),
+ /// caching the result service-wide so subsequent readers skip the probe.
+ async fn probe_positioned_read_support(&self) -> Result<bool> {
+ self.backend
+ .core
+ .positioned_read_support
+ .get_or_try_init(|| async {
+ // Open a temporary file just for the probe.
+ let (client, file) = self.open_file().await?;
+ let handle = SftpReaderHandle::new(client, file);
+
+ match probe_positioned_read(&handle).await {
+ Ok(supported) => Ok(supported),
+ Err(err) if is_positioned_read_unsupported(&err) =>
Ok(false),
+ Err(err) => Err(err),
+ }
+ })
+ .await
+ .copied()
+ }
+
+ /// Return a cached positioned-read handle for this reader, or `None`
+ /// if the server does not support positioned reads.
+ async fn positioned_handle(&self) -> Result<Option<Arc<SftpReaderHandle>>>
{
+ if !self.probe_positioned_read_support().await? {
+ return Ok(None);
}
+
+ // Positioned read is supported — lazily open or reuse the handle.
+ let handle = self
+ .positioned
+ .get_or_try_init(|| async {
+ let (client, file) = self.open_file().await?;
+ Ok(Arc::new(SftpReaderHandle::new(client, file)))
+ })
+ .await?;
+
+ Ok(Some(handle.clone()))
}
}
pub struct SftpReaderHandle {
- /// Keep the connection alive while the file handle is cached by
`PositionReader`.
+ /// Keep the connection alive while the file handle is cached by
`SftpReader`.
_conn: bounded::Object<Manager>,
file: File,
}
@@ -56,46 +198,174 @@ impl SftpReaderHandle {
}
}
-impl oio::PositionRead for SftpReader {
- type Handle = SftpReaderHandle;
+async fn read_positioned(handle: &SftpReaderHandle, offset: u64, size: usize)
-> Result<Buffer> {
+ if size == 0 {
+ return Ok(Buffer::new());
+ }
- async fn open(&self) -> Result<Self::Handle> {
- let backend = &self.backend;
- let path = self.path.as_str();
+ 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()),
+ }
+}
- let client = backend.core.connect().await?;
+/// Probe whether positioned reads work by reading one byte.
+///
+/// If the file has a known length, read at EOF (expecting an empty result).
+/// Otherwise, read at offset 0 (expecting at most one byte).
+async fn probe_positioned_read(handle: &SftpReaderHandle) -> Result<bool> {
+ let mut file = handle.file.clone();
+ let len = file.metadata().await.map_err(parse_sftp_error)?.len();
+
+ if let Some(len) = len {
+ let buf = read_positioned(handle, len, 1).await?;
+ Ok(buf.is_empty())
+ } else {
+ read_positioned(handle, 0, 1).await.map(|_| true)
+ }
+}
- let mut fs = client.fs();
- fs.set_cwd(&backend.core.root);
+fn is_positioned_read_unsupported(err: &Error) -> bool {
+ if err.kind() == ErrorKind::Unsupported {
+ return true;
+ }
- let path = fs.canonicalize(path).await.map_err(parse_sftp_error)?;
+ err.source()
+ .and_then(|source| source.downcast_ref::<SftpClientError>())
+ .is_some_and(|source| {
+ matches!(
+ source,
+ SftpClientError::UnsupportedSftpProtocol { .. }
+ | SftpClientError::SftpError(SftpErrorKind::OpUnsupported,
_)
+ )
+ })
+}
- let f = client
- .open(path.as_path())
- .await
- .map_err(parse_sftp_error)?;
+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
+ }
+ }
- Ok(SftpReaderHandle::new(client, f))
+ 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>,
+}
- async fn read_at(handle: &Self::Handle, offset: u64, size: usize) ->
Result<Buffer> {
- if size == 0 {
+impl oio::ReadStream for SftpPositionedReadStream {
+ async fn read(&mut self) -> Result<Buffer> {
+ if self.remaining == Some(0) {
return Ok(Buffer::new());
}
- let mut file = handle.file.clone();
- file.seek(SeekFrom::Start(offset))
- .await
- .map_err(new_std_io_error)?;
+ 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());
+ }
- match file
- .read(size as u32, bytes::BytesMut::with_capacity(size))
- .await
- .map_err(parse_sftp_error)?
- {
- Some(bytes) => Ok(Buffer::from(bytes.freeze())),
- None => Ok(Buffer::new()),
+ let n = buf.len() as u64;
+ self.offset += n;
+ if let Some(remaining) = &mut self.remaining {
+ *remaining -= n;
}
+
+ Ok(buf)
}
}