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)
     }
 }
 

Reply via email to