mirror of
https://github.com/restic/restic.git
synced 2026-10-07 11:07:11 +00:00
Update dependencies
This, among others, updates the `go-flags` library, which includes a feature that closes #198.
This commit is contained in:
1 parent
d9a8dcfd67
commit
d9a90f7b89
60 files changed
+1627
-5187
No files matched your search
+505
-122
@@ -1,11 +1,16 @@
|
||||
package sftp
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding"
|
||||
"encoding/binary"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"path"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/kr/fs"
|
||||
@@ -13,8 +18,19 @@ import (
|
||||
"golang.org/x/crypto/ssh"
|
||||
)
|
||||
|
||||
// MaxPacket sets the maximum size of the payload.
|
||||
func MaxPacket(size int) func(*Client) error {
|
||||
return func(c *Client) error {
|
||||
if size < 1<<15 {
|
||||
return fmt.Errorf("size must be greater or equal to 32k")
|
||||
}
|
||||
c.maxPacket = size
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// New creates a new SFTP client on conn.
|
||||
func NewClient(conn *ssh.Client) (*Client, error) {
|
||||
func NewClient(conn *ssh.Client, opts ...func(*Client) error) (*Client, error) {
|
||||
s, err := conn.NewSession()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -31,21 +47,34 @@ func NewClient(conn *ssh.Client) (*Client, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return NewClientPipe(pr, pw)
|
||||
return NewClientPipe(pr, pw, opts...)
|
||||
}
|
||||
|
||||
// NewClientPipe creates a new SFTP client given a Reader and a WriteCloser.
|
||||
// This can be used for connecting to an SFTP server over TCP/TLS or by using
|
||||
// the system's ssh client program (e.g. via exec.Command).
|
||||
func NewClientPipe(rd io.Reader, wr io.WriteCloser) (*Client, error) {
|
||||
func NewClientPipe(rd io.Reader, wr io.WriteCloser, opts ...func(*Client) error) (*Client, error) {
|
||||
sftp := &Client{
|
||||
w: wr,
|
||||
r: rd,
|
||||
w: wr,
|
||||
r: rd,
|
||||
maxPacket: 1 << 15,
|
||||
inflight: make(map[uint32]chan<- result),
|
||||
recvClosed: make(chan struct{}),
|
||||
}
|
||||
if err := sftp.sendInit(); err != nil {
|
||||
if err := sftp.applyOptions(opts...); err != nil {
|
||||
wr.Close()
|
||||
return nil, err
|
||||
}
|
||||
return sftp, sftp.recvVersion()
|
||||
if err := sftp.sendInit(); err != nil {
|
||||
wr.Close()
|
||||
return nil, err
|
||||
}
|
||||
if err := sftp.recvVersion(); err != nil {
|
||||
wr.Close()
|
||||
return nil, err
|
||||
}
|
||||
go sftp.recv()
|
||||
return sftp, nil
|
||||
}
|
||||
|
||||
// Client represents an SFTP session on a *ssh.ClientConn SSH connection.
|
||||
@@ -54,14 +83,23 @@ func NewClientPipe(rd io.Reader, wr io.WriteCloser) (*Client, error) {
|
||||
//
|
||||
// Client implements the github.com/kr/fs.FileSystem interface.
|
||||
type Client struct {
|
||||
w io.WriteCloser
|
||||
r io.Reader
|
||||
mu sync.Mutex // locks mu and seralises commands to the server
|
||||
nextid uint32
|
||||
w io.WriteCloser
|
||||
r io.Reader
|
||||
|
||||
maxPacket int // max packet size read or written.
|
||||
nextid uint32
|
||||
|
||||
mu sync.Mutex // ensures only on request is in flight to the server at once
|
||||
inflight map[uint32]chan<- result // outstanding requests
|
||||
recvClosed chan struct{} // remote end has closed the connection
|
||||
}
|
||||
|
||||
// Close closes the SFTP session.
|
||||
func (c *Client) Close() error { return c.w.Close() }
|
||||
func (c *Client) Close() error {
|
||||
err := c.w.Close()
|
||||
<-c.recvClosed
|
||||
return err
|
||||
}
|
||||
|
||||
// Create creates the named file mode 0666 (before umask), truncating it if
|
||||
// it already exists. If successful, methods on the returned File can be
|
||||
@@ -78,12 +116,9 @@ func (c *Client) sendInit() error {
|
||||
})
|
||||
}
|
||||
|
||||
// returns the current value of c.nextid and increments it
|
||||
// callers is expected to hold c.mu
|
||||
// returns the next value of c.nextid
|
||||
func (c *Client) nextId() uint32 {
|
||||
v := c.nextid
|
||||
c.nextid++
|
||||
return v
|
||||
return atomic.AddUint32(&c.nextid, 1)
|
||||
}
|
||||
|
||||
func (c *Client) recvVersion() error {
|
||||
@@ -103,6 +138,46 @@ func (c *Client) recvVersion() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// broadcastErr sends an error to all goroutines waiting for a response.
|
||||
func (c *Client) broadcastErr(err error) {
|
||||
c.mu.Lock()
|
||||
listeners := make([]chan<- result, 0, len(c.inflight))
|
||||
for _, ch := range c.inflight {
|
||||
listeners = append(listeners, ch)
|
||||
}
|
||||
c.mu.Unlock()
|
||||
for _, ch := range listeners {
|
||||
ch <- result{err: err}
|
||||
}
|
||||
}
|
||||
|
||||
// recv continuously reads from the server and forwards responses to the
|
||||
// appropriate channel.
|
||||
func (c *Client) recv() {
|
||||
defer close(c.recvClosed)
|
||||
for {
|
||||
typ, data, err := recvPacket(c.r)
|
||||
if err != nil {
|
||||
// Return the error to all listeners.
|
||||
c.broadcastErr(err)
|
||||
return
|
||||
}
|
||||
sid, _ := unmarshalUint32(data)
|
||||
c.mu.Lock()
|
||||
ch, ok := c.inflight[sid]
|
||||
delete(c.inflight, sid)
|
||||
c.mu.Unlock()
|
||||
if !ok {
|
||||
// This is an unexpected occurrence. Send the error
|
||||
// back to all listeners so that they terminate
|
||||
// gracefully.
|
||||
c.broadcastErr(fmt.Errorf("sid: %v not fond", sid))
|
||||
return
|
||||
}
|
||||
ch <- result{typ: typ, data: data}
|
||||
}
|
||||
}
|
||||
|
||||
// Walk returns a new Walker rooted at root.
|
||||
func (c *Client) Walk(root string) *fs.Walker {
|
||||
return fs.WalkFS(root, c)
|
||||
@@ -117,8 +192,6 @@ func (c *Client) ReadDir(p string) ([]os.FileInfo, error) {
|
||||
}
|
||||
defer c.close(handle) // this has to defer earlier than the lock below
|
||||
var attrs []os.FileInfo
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
var done = false
|
||||
for !done {
|
||||
id := c.nextId()
|
||||
@@ -163,8 +236,6 @@ func (c *Client) ReadDir(p string) ([]os.FileInfo, error) {
|
||||
return attrs, err
|
||||
}
|
||||
func (c *Client) opendir(path string) (string, error) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpOpendirPacket{
|
||||
Id: id,
|
||||
@@ -189,8 +260,6 @@ func (c *Client) opendir(path string) (string, error) {
|
||||
}
|
||||
|
||||
func (c *Client) Lstat(p string) (os.FileInfo, error) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpLstatPacket{
|
||||
Id: id,
|
||||
@@ -216,8 +285,6 @@ func (c *Client) Lstat(p string) (os.FileInfo, error) {
|
||||
|
||||
// ReadLink reads the target of a symbolic link.
|
||||
func (c *Client) ReadLink(p string) (string, error) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpReadlinkPacket{
|
||||
Id: id,
|
||||
@@ -247,8 +314,6 @@ func (c *Client) ReadLink(p string) (string, error) {
|
||||
|
||||
// setstat is a convience wrapper to allow for changing of various parts of the file descriptor.
|
||||
func (c *Client) setstat(path string, flags uint32, attrs interface{}) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpSetstatPacket{
|
||||
Id: id,
|
||||
@@ -315,8 +380,6 @@ func (c *Client) OpenFile(path string, f int) (*File, error) {
|
||||
}
|
||||
|
||||
func (c *Client) open(path string, pflags uint32) (*File, error) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpOpenPacket{
|
||||
Id: id,
|
||||
@@ -341,43 +404,10 @@ func (c *Client) open(path string, pflags uint32) (*File, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// readAt reads len(buf) bytes from the remote file indicated by handle starting
|
||||
// from offset.
|
||||
func (c *Client) readAt(handle string, offset uint64, buf []byte) (uint32, error) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpReadPacket{
|
||||
Id: id,
|
||||
Handle: handle,
|
||||
Offset: offset,
|
||||
Len: uint32(len(buf)),
|
||||
})
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
switch typ {
|
||||
case ssh_FXP_DATA:
|
||||
sid, data := unmarshalUint32(data)
|
||||
if sid != id {
|
||||
return 0, &unexpectedIdErr{id, sid}
|
||||
}
|
||||
l, data := unmarshalUint32(data)
|
||||
n := copy(buf, data[:l])
|
||||
return uint32(n), nil
|
||||
case ssh_FXP_STATUS:
|
||||
return 0, eofOrErr(unmarshalStatus(id, data))
|
||||
default:
|
||||
return 0, unimplementedPacketErr(typ)
|
||||
}
|
||||
}
|
||||
|
||||
// close closes a handle handle previously returned in the response
|
||||
// to SSH_FXP_OPEN or SSH_FXP_OPENDIR. The handle becomes invalid
|
||||
// immediately after this request has been sent.
|
||||
func (c *Client) close(handle string) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpClosePacket{
|
||||
Id: id,
|
||||
@@ -395,8 +425,6 @@ func (c *Client) close(handle string) error {
|
||||
}
|
||||
|
||||
func (c *Client) fstat(handle string) (*FileStat, error) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpFstatPacket{
|
||||
Id: id,
|
||||
@@ -420,6 +448,40 @@ func (c *Client) fstat(handle string) (*FileStat, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// Get vfs stats from remote host.
|
||||
// Implementing statvfs@openssh.com SSH_FXP_EXTENDED feature
|
||||
// from http://www.opensource.apple.com/source/OpenSSH/OpenSSH-175/openssh/PROTOCOL?txt
|
||||
func (c *Client) StatVFS(path string) (*StatVFS, error) {
|
||||
// send the StatVFS packet to the server
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpStatvfsPacket{
|
||||
Id: id,
|
||||
Path: path,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
switch typ {
|
||||
// server responded with valid data
|
||||
case ssh_FXP_EXTENDED_REPLY:
|
||||
var response StatVFS
|
||||
err = binary.Read(bytes.NewReader(data), binary.BigEndian, &response)
|
||||
if err != nil {
|
||||
return nil, errors.New("can not parse reply")
|
||||
}
|
||||
|
||||
return &response, nil
|
||||
|
||||
// the resquest failed
|
||||
case ssh_FXP_STATUS:
|
||||
return nil, errors.New(fxp(ssh_FXP_STATUS).String())
|
||||
|
||||
default:
|
||||
return nil, unimplementedPacketErr(typ)
|
||||
}
|
||||
}
|
||||
|
||||
// Join joins any number of path elements into a single path, adding a
|
||||
// separating slash if necessary. The result is Cleaned; in particular, all
|
||||
// empty strings are ignored.
|
||||
@@ -437,8 +499,6 @@ func (c *Client) Remove(path string) error {
|
||||
}
|
||||
|
||||
func (c *Client) removeFile(path string) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpRemovePacket{
|
||||
Id: id,
|
||||
@@ -456,8 +516,6 @@ func (c *Client) removeFile(path string) error {
|
||||
}
|
||||
|
||||
func (c *Client) removeDirectory(path string) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpRmdirPacket{
|
||||
Id: id,
|
||||
@@ -476,8 +534,6 @@ func (c *Client) removeDirectory(path string) error {
|
||||
|
||||
// Rename renames a file.
|
||||
func (c *Client) Rename(oldname, newname string) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpRenamePacket{
|
||||
Id: id,
|
||||
@@ -495,46 +551,41 @@ func (c *Client) Rename(oldname, newname string) error {
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Client) sendRequest(p encoding.BinaryMarshaler) (byte, []byte, error) {
|
||||
if err := sendPacket(c.w, p); err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
return recvPacket(c.r)
|
||||
// result captures the result of receiving the a packet from the server
|
||||
type result struct {
|
||||
typ byte
|
||||
data []byte
|
||||
err error
|
||||
}
|
||||
|
||||
// writeAt writes len(buf) bytes from the remote file indicated by handle starting
|
||||
// from offset.
|
||||
func (c *Client) writeAt(handle string, offset uint64, buf []byte) (uint32, error) {
|
||||
type idmarshaler interface {
|
||||
id() uint32
|
||||
encoding.BinaryMarshaler
|
||||
}
|
||||
|
||||
func (c *Client) sendRequest(p idmarshaler) (byte, []byte, error) {
|
||||
ch := make(chan result, 1)
|
||||
c.dispatchRequest(ch, p)
|
||||
s := <-ch
|
||||
return s.typ, s.data, s.err
|
||||
}
|
||||
|
||||
func (c *Client) dispatchRequest(ch chan<- result, p idmarshaler) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpWritePacket{
|
||||
Id: id,
|
||||
Handle: handle,
|
||||
Offset: offset,
|
||||
Length: uint32(len(buf)),
|
||||
Data: buf,
|
||||
})
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
switch typ {
|
||||
case ssh_FXP_STATUS:
|
||||
if err := okOrErr(unmarshalStatus(id, data)); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return uint32(len(buf)), nil
|
||||
default:
|
||||
return 0, unimplementedPacketErr(typ)
|
||||
c.inflight[p.id()] = ch
|
||||
if err := sendPacket(c.w, p); err != nil {
|
||||
delete(c.inflight, p.id())
|
||||
c.mu.Unlock()
|
||||
ch <- result{err: err}
|
||||
return
|
||||
}
|
||||
c.mu.Unlock()
|
||||
}
|
||||
|
||||
// Creates the specified directory. An error will be returned if a file or
|
||||
// directory with the specified path already exists, or if the directory's
|
||||
// parent folder does not exist (the method cannot create complete paths).
|
||||
func (c *Client) Mkdir(path string) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
id := c.nextId()
|
||||
typ, data, err := c.sendRequest(sshFxpMkdirPacket{
|
||||
Id: id,
|
||||
@@ -551,6 +602,17 @@ func (c *Client) Mkdir(path string) error {
|
||||
}
|
||||
}
|
||||
|
||||
// applyOptions applies options functions to the Client.
|
||||
// If an error is encountered, option processing ceases.
|
||||
func (c *Client) applyOptions(opts ...func(*Client) error) error {
|
||||
for _, f := range opts {
|
||||
if err := f(c); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// File represents a remote file.
|
||||
type File struct {
|
||||
c *Client
|
||||
@@ -565,21 +627,226 @@ func (f *File) Close() error {
|
||||
return f.c.close(f.handle)
|
||||
}
|
||||
|
||||
const maxConcurrentRequests = 64
|
||||
|
||||
// Read reads up to len(b) bytes from the File. It returns the number of
|
||||
// bytes read and an error, if any. EOF is signaled by a zero count with
|
||||
// err set to io.EOF.
|
||||
func (f *File) Read(b []byte) (int, error) {
|
||||
var read int
|
||||
for len(b) > 0 {
|
||||
n, err := f.c.readAt(f.handle, f.offset, b[:min(len(b), maxWritePacket)])
|
||||
f.offset += uint64(n)
|
||||
read += int(n)
|
||||
if err != nil {
|
||||
return read, err
|
||||
}
|
||||
b = b[n:]
|
||||
// Split the read into multiple maxPacket sized concurrent reads
|
||||
// bounded by maxConcurrentRequests. This allows reads with a suitably
|
||||
// large buffer to transfer data at a much faster rate due to
|
||||
// overlapping round trip times.
|
||||
inFlight := 0
|
||||
desiredInFlight := 1
|
||||
offset := f.offset
|
||||
ch := make(chan result)
|
||||
type inflightRead struct {
|
||||
b []byte
|
||||
offset uint64
|
||||
}
|
||||
return read, nil
|
||||
reqs := map[uint32]inflightRead{}
|
||||
type offsetErr struct {
|
||||
offset uint64
|
||||
err error
|
||||
}
|
||||
var firstErr offsetErr
|
||||
|
||||
sendReq := func(b []byte, offset uint64) {
|
||||
reqId := f.c.nextId()
|
||||
f.c.dispatchRequest(ch, sshFxpReadPacket{
|
||||
Id: reqId,
|
||||
Handle: f.handle,
|
||||
Offset: offset,
|
||||
Len: uint32(len(b)),
|
||||
})
|
||||
inFlight++
|
||||
reqs[reqId] = inflightRead{b: b, offset: offset}
|
||||
}
|
||||
|
||||
var read int
|
||||
for len(b) > 0 || inFlight > 0 {
|
||||
for inFlight < desiredInFlight && len(b) > 0 && firstErr.err == nil {
|
||||
l := min(len(b), f.c.maxPacket)
|
||||
rb := b[:l]
|
||||
sendReq(rb, offset)
|
||||
offset += uint64(l)
|
||||
b = b[l:]
|
||||
}
|
||||
|
||||
if inFlight == 0 {
|
||||
break
|
||||
}
|
||||
select {
|
||||
case res := <-ch:
|
||||
inFlight--
|
||||
if res.err != nil {
|
||||
firstErr = offsetErr{offset: 0, err: res.err}
|
||||
break
|
||||
}
|
||||
reqId, data := unmarshalUint32(res.data)
|
||||
req, ok := reqs[reqId]
|
||||
if !ok {
|
||||
firstErr = offsetErr{offset: 0, err: fmt.Errorf("sid: %v not found", reqId)}
|
||||
break
|
||||
}
|
||||
delete(reqs, reqId)
|
||||
switch res.typ {
|
||||
case ssh_FXP_STATUS:
|
||||
if firstErr.err == nil || req.offset < firstErr.offset {
|
||||
firstErr = offsetErr{offset: req.offset, err: eofOrErr(unmarshalStatus(reqId, res.data))}
|
||||
break
|
||||
}
|
||||
case ssh_FXP_DATA:
|
||||
l, data := unmarshalUint32(data)
|
||||
n := copy(req.b, data[:l])
|
||||
read += n
|
||||
if n < len(req.b) {
|
||||
sendReq(req.b[l:], req.offset+uint64(l))
|
||||
}
|
||||
if desiredInFlight < maxConcurrentRequests {
|
||||
desiredInFlight++
|
||||
}
|
||||
default:
|
||||
firstErr = offsetErr{offset: 0, err: unimplementedPacketErr(res.typ)}
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
// If the error is anything other than EOF, then there
|
||||
// may be gaps in the data copied to the buffer so it's
|
||||
// best to return 0 so the caller can't make any
|
||||
// incorrect assumptions about the state of the buffer.
|
||||
if firstErr.err != nil && firstErr.err != io.EOF {
|
||||
read = 0
|
||||
}
|
||||
f.offset += uint64(read)
|
||||
return read, firstErr.err
|
||||
}
|
||||
|
||||
// WriteTo writes the file to w. The return value is the number of bytes
|
||||
// written. Any error encountered during the write is also returned.
|
||||
func (f *File) WriteTo(w io.Writer) (int64, error) {
|
||||
fi, err := f.Stat()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
inFlight := 0
|
||||
desiredInFlight := 1
|
||||
offset := f.offset
|
||||
writeOffset := offset
|
||||
fileSize := uint64(fi.Size())
|
||||
ch := make(chan result)
|
||||
type inflightRead struct {
|
||||
b []byte
|
||||
offset uint64
|
||||
}
|
||||
reqs := map[uint32]inflightRead{}
|
||||
pendingWrites := map[uint64][]byte{}
|
||||
type offsetErr struct {
|
||||
offset uint64
|
||||
err error
|
||||
}
|
||||
var firstErr offsetErr
|
||||
|
||||
sendReq := func(b []byte, offset uint64) {
|
||||
reqId := f.c.nextId()
|
||||
f.c.dispatchRequest(ch, sshFxpReadPacket{
|
||||
Id: reqId,
|
||||
Handle: f.handle,
|
||||
Offset: offset,
|
||||
Len: uint32(len(b)),
|
||||
})
|
||||
inFlight++
|
||||
reqs[reqId] = inflightRead{b: b, offset: offset}
|
||||
}
|
||||
|
||||
var copied int64
|
||||
for firstErr.err == nil || inFlight > 0 {
|
||||
for inFlight < desiredInFlight && firstErr.err == nil {
|
||||
b := make([]byte, f.c.maxPacket)
|
||||
sendReq(b, offset)
|
||||
offset += uint64(f.c.maxPacket)
|
||||
if offset > fileSize {
|
||||
desiredInFlight = 1
|
||||
}
|
||||
}
|
||||
|
||||
if inFlight == 0 {
|
||||
break
|
||||
}
|
||||
select {
|
||||
case res := <-ch:
|
||||
inFlight--
|
||||
if res.err != nil {
|
||||
firstErr = offsetErr{offset: 0, err: res.err}
|
||||
break
|
||||
}
|
||||
reqId, data := unmarshalUint32(res.data)
|
||||
req, ok := reqs[reqId]
|
||||
if !ok {
|
||||
firstErr = offsetErr{offset: 0, err: fmt.Errorf("sid: %v not found", reqId)}
|
||||
break
|
||||
}
|
||||
delete(reqs, reqId)
|
||||
switch res.typ {
|
||||
case ssh_FXP_STATUS:
|
||||
if firstErr.err == nil || req.offset < firstErr.offset {
|
||||
firstErr = offsetErr{offset: req.offset, err: eofOrErr(unmarshalStatus(reqId, res.data))}
|
||||
break
|
||||
}
|
||||
case ssh_FXP_DATA:
|
||||
l, data := unmarshalUint32(data)
|
||||
if req.offset == writeOffset {
|
||||
nbytes, err := w.Write(data)
|
||||
copied += int64(nbytes)
|
||||
if err != nil {
|
||||
firstErr = offsetErr{offset: req.offset + uint64(nbytes), err: err}
|
||||
break
|
||||
}
|
||||
if nbytes < int(l) {
|
||||
firstErr = offsetErr{offset: req.offset + uint64(nbytes), err: io.ErrShortWrite}
|
||||
break
|
||||
}
|
||||
switch {
|
||||
case offset > fileSize:
|
||||
desiredInFlight = 1
|
||||
case desiredInFlight < maxConcurrentRequests:
|
||||
desiredInFlight++
|
||||
}
|
||||
writeOffset += uint64(nbytes)
|
||||
for pendingData, ok := pendingWrites[writeOffset]; ok; pendingData, ok = pendingWrites[writeOffset] {
|
||||
nbytes, err := w.Write(pendingData)
|
||||
if err != nil {
|
||||
firstErr = offsetErr{offset: writeOffset + uint64(nbytes), err: err}
|
||||
break
|
||||
}
|
||||
if nbytes < len(pendingData) {
|
||||
firstErr = offsetErr{offset: writeOffset + uint64(nbytes), err: io.ErrShortWrite}
|
||||
break
|
||||
}
|
||||
writeOffset += uint64(nbytes)
|
||||
inFlight--
|
||||
}
|
||||
} else {
|
||||
// Don't write the data yet because
|
||||
// this response came in out of order
|
||||
// and we need to wait for responses
|
||||
// for earlier segments of the file.
|
||||
inFlight++ // Pending writes should still be considered inFlight.
|
||||
pendingWrites[req.offset] = data
|
||||
}
|
||||
default:
|
||||
firstErr = offsetErr{offset: 0, err: unimplementedPacketErr(res.typ)}
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
if firstErr.err != io.EOF {
|
||||
return copied, firstErr.err
|
||||
}
|
||||
return copied, nil
|
||||
|
||||
}
|
||||
|
||||
// Stat returns the FileInfo structure describing file. If there is an
|
||||
@@ -592,24 +859,140 @@ func (f *File) Stat() (os.FileInfo, error) {
|
||||
return fileInfoFromStat(fs, path.Base(f.path)), nil
|
||||
}
|
||||
|
||||
// clamp writes to less than 32k
|
||||
const maxWritePacket = 1 << 15
|
||||
|
||||
// Write writes len(b) bytes to the File. It returns the number of bytes
|
||||
// written and an error, if any. Write returns a non-nil error when n !=
|
||||
// len(b).
|
||||
func (f *File) Write(b []byte) (int, error) {
|
||||
var written int
|
||||
for len(b) > 0 {
|
||||
n, err := f.c.writeAt(f.handle, f.offset, b[:min(len(b), maxWritePacket)])
|
||||
f.offset += uint64(n)
|
||||
written += int(n)
|
||||
if err != nil {
|
||||
return written, err
|
||||
// Split the write into multiple maxPacket sized concurrent writes
|
||||
// bounded by maxConcurrentRequests. This allows writes with a suitably
|
||||
// large buffer to transfer data at a much faster rate due to
|
||||
// overlapping round trip times.
|
||||
inFlight := 0
|
||||
desiredInFlight := 1
|
||||
offset := f.offset
|
||||
ch := make(chan result)
|
||||
var firstErr error
|
||||
written := len(b)
|
||||
for len(b) > 0 || inFlight > 0 {
|
||||
for inFlight < desiredInFlight && len(b) > 0 && firstErr == nil {
|
||||
l := min(len(b), f.c.maxPacket)
|
||||
rb := b[:l]
|
||||
f.c.dispatchRequest(ch, sshFxpWritePacket{
|
||||
Id: f.c.nextId(),
|
||||
Handle: f.handle,
|
||||
Offset: offset,
|
||||
Length: uint32(len(rb)),
|
||||
Data: rb,
|
||||
})
|
||||
inFlight++
|
||||
offset += uint64(l)
|
||||
b = b[l:]
|
||||
}
|
||||
|
||||
if inFlight == 0 {
|
||||
break
|
||||
}
|
||||
select {
|
||||
case res := <-ch:
|
||||
inFlight--
|
||||
if res.err != nil {
|
||||
firstErr = res.err
|
||||
break
|
||||
}
|
||||
switch res.typ {
|
||||
case ssh_FXP_STATUS:
|
||||
id, _ := unmarshalUint32(res.data)
|
||||
err := okOrErr(unmarshalStatus(id, res.data))
|
||||
if err != nil && firstErr == nil {
|
||||
firstErr = err
|
||||
break
|
||||
}
|
||||
if desiredInFlight < maxConcurrentRequests {
|
||||
desiredInFlight++
|
||||
}
|
||||
default:
|
||||
firstErr = unimplementedPacketErr(res.typ)
|
||||
break
|
||||
}
|
||||
}
|
||||
b = b[n:]
|
||||
}
|
||||
return written, nil
|
||||
// If error is non-nil, then there may be gaps in the data written to
|
||||
// the file so it's best to return 0 so the caller can't make any
|
||||
// incorrect assumptions about the state of the file.
|
||||
if firstErr != nil {
|
||||
written = 0
|
||||
}
|
||||
f.offset += uint64(written)
|
||||
return written, firstErr
|
||||
}
|
||||
|
||||
// ReadFrom reads data from r until EOF and writes it to the file. The return
|
||||
// value is the number of bytes read. Any error except io.EOF encountered
|
||||
// during the read is also returned.
|
||||
func (f *File) ReadFrom(r io.Reader) (int64, error) {
|
||||
inFlight := 0
|
||||
desiredInFlight := 1
|
||||
offset := f.offset
|
||||
ch := make(chan result)
|
||||
var firstErr error
|
||||
read := int64(0)
|
||||
b := make([]byte, f.c.maxPacket)
|
||||
for inFlight > 0 || firstErr == nil {
|
||||
for inFlight < desiredInFlight && firstErr == nil {
|
||||
n, err := r.Read(b)
|
||||
if err != nil {
|
||||
firstErr = err
|
||||
}
|
||||
f.c.dispatchRequest(ch, sshFxpWritePacket{
|
||||
Id: f.c.nextId(),
|
||||
Handle: f.handle,
|
||||
Offset: offset,
|
||||
Length: uint32(n),
|
||||
Data: b[:n],
|
||||
})
|
||||
inFlight++
|
||||
offset += uint64(n)
|
||||
read += int64(n)
|
||||
}
|
||||
|
||||
if inFlight == 0 {
|
||||
break
|
||||
}
|
||||
select {
|
||||
case res := <-ch:
|
||||
inFlight--
|
||||
if res.err != nil {
|
||||
firstErr = res.err
|
||||
break
|
||||
}
|
||||
switch res.typ {
|
||||
case ssh_FXP_STATUS:
|
||||
id, _ := unmarshalUint32(res.data)
|
||||
err := okOrErr(unmarshalStatus(id, res.data))
|
||||
if err != nil && firstErr == nil {
|
||||
firstErr = err
|
||||
break
|
||||
}
|
||||
if desiredInFlight < maxConcurrentRequests {
|
||||
desiredInFlight++
|
||||
}
|
||||
default:
|
||||
firstErr = unimplementedPacketErr(res.typ)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
if firstErr == io.EOF {
|
||||
firstErr = nil
|
||||
}
|
||||
// If error is non-nil, then there may be gaps in the data written to
|
||||
// the file so it's best to return 0 so the caller can't make any
|
||||
// incorrect assumptions about the state of the file.
|
||||
if firstErr != nil {
|
||||
read = 0
|
||||
}
|
||||
f.offset += uint64(read)
|
||||
return read, firstErr
|
||||
}
|
||||
|
||||
// Seek implements io.Seeker by setting the client offset for the next Read or
|
||||
|
||||
+324
-47
@@ -14,15 +14,18 @@ import (
|
||||
"path"
|
||||
"path/filepath"
|
||||
"reflect"
|
||||
"syscall"
|
||||
"testing"
|
||||
"testing/quick"
|
||||
"time"
|
||||
|
||||
"github.com/kr/fs"
|
||||
)
|
||||
|
||||
const (
|
||||
READONLY = true
|
||||
READWRITE = false
|
||||
READONLY = true
|
||||
READWRITE = false
|
||||
NO_DELAY time.Duration = 0
|
||||
|
||||
debuglevel = "ERROR" // set to "DEBUG" for debugging
|
||||
)
|
||||
@@ -30,9 +33,57 @@ const (
|
||||
var testIntegration = flag.Bool("integration", false, "perform integration tests against sftp server process")
|
||||
var testSftp = flag.String("sftp", "/usr/lib/openssh/sftp-server", "location of the sftp server binary")
|
||||
|
||||
type delayedWrite struct {
|
||||
t time.Time
|
||||
b []byte
|
||||
}
|
||||
|
||||
// delayedWriter wraps a writer and artificially delays the write. This is
|
||||
// meant to mimic connections with various latencies. Error's returned from the
|
||||
// underlying writer will panic so this should only be used over reliable
|
||||
// connections.
|
||||
type delayedWriter struct {
|
||||
w io.WriteCloser
|
||||
ch chan delayedWrite
|
||||
closed chan struct{}
|
||||
}
|
||||
|
||||
func newDelayedWriter(w io.WriteCloser, delay time.Duration) io.WriteCloser {
|
||||
ch := make(chan delayedWrite, 128)
|
||||
closed := make(chan struct{})
|
||||
go func() {
|
||||
for writeMsg := range ch {
|
||||
time.Sleep(writeMsg.t.Add(delay).Sub(time.Now()))
|
||||
n, err := w.Write(writeMsg.b)
|
||||
if err != nil {
|
||||
panic("write error")
|
||||
}
|
||||
if n < len(writeMsg.b) {
|
||||
panic("showrt write")
|
||||
}
|
||||
}
|
||||
w.Close()
|
||||
close(closed)
|
||||
}()
|
||||
return delayedWriter{w: w, ch: ch, closed: closed}
|
||||
}
|
||||
|
||||
func (w delayedWriter) Write(b []byte) (int, error) {
|
||||
bcopy := make([]byte, len(b))
|
||||
copy(bcopy, b)
|
||||
w.ch <- delayedWrite{t: time.Now(), b: bcopy}
|
||||
return len(b), nil
|
||||
}
|
||||
|
||||
func (w delayedWriter) Close() error {
|
||||
close(w.ch)
|
||||
<-w.closed
|
||||
return nil
|
||||
}
|
||||
|
||||
// testClient returns a *Client connected to a localy running sftp-server
|
||||
// the *exec.Cmd returned must be defer Wait'd.
|
||||
func testClient(t testing.TB, readonly bool) (*Client, *exec.Cmd) {
|
||||
func testClient(t testing.TB, readonly bool, delay time.Duration) (*Client, *exec.Cmd) {
|
||||
if !*testIntegration {
|
||||
t.Skip("skipping intergration test")
|
||||
}
|
||||
@@ -45,6 +96,9 @@ func testClient(t testing.TB, readonly bool) (*Client, *exec.Cmd) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if delay > NO_DELAY {
|
||||
pw = newDelayedWriter(pw, delay)
|
||||
}
|
||||
pr, err := cmd.StdoutPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -58,19 +112,11 @@ func testClient(t testing.TB, readonly bool) (*Client, *exec.Cmd) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if err := sftp.sendInit(); err != nil {
|
||||
defer cmd.Wait()
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := sftp.recvVersion(); err != nil {
|
||||
defer cmd.Wait()
|
||||
t.Fatal(err)
|
||||
}
|
||||
return sftp, cmd
|
||||
}
|
||||
|
||||
func TestNewClient(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
|
||||
if err := sftp.Close(); err != nil {
|
||||
@@ -79,7 +125,7 @@ func TestNewClient(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientLstat(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -105,7 +151,7 @@ func TestClientLstat(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientLstatMissing(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -122,7 +168,7 @@ func TestClientLstatMissing(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientMkdir(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE)
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -140,7 +186,7 @@ func TestClientMkdir(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientOpen(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -199,7 +245,7 @@ func (s seek) end(t *testing.T, r io.ReadSeeker) {
|
||||
}
|
||||
|
||||
func TestClientSeek(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -243,7 +289,7 @@ func TestClientSeek(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientCreate(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE)
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -262,7 +308,7 @@ func TestClientCreate(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientAppend(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE)
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -281,7 +327,7 @@ func TestClientAppend(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientCreateFailed(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -302,7 +348,7 @@ func TestClientCreateFailed(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientFileStat(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -333,7 +379,7 @@ func TestClientFileStat(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientRemove(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE)
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -350,7 +396,7 @@ func TestClientRemove(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientRemoveDir(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE)
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -367,7 +413,7 @@ func TestClientRemoveDir(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientRemoveFailed(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -384,7 +430,7 @@ func TestClientRemoveFailed(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientRename(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE)
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -405,7 +451,7 @@ func TestClientRename(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestClientReadLine(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE)
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -449,7 +495,7 @@ var clientReadTests = []struct {
|
||||
}
|
||||
|
||||
func TestClientRead(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -537,7 +583,7 @@ var clientWriteTests = []struct {
|
||||
}
|
||||
|
||||
func TestClientWrite(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE)
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -664,7 +710,7 @@ func mark(path string, info os.FileInfo, err error, errors *[]error, clear bool)
|
||||
}
|
||||
|
||||
func TestClientWalk(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY)
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -747,11 +793,29 @@ func TestClientWalk(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func benchmarkRead(b *testing.B, bufsize int) {
|
||||
// sftp/issue/42, abrupt server hangup would result in client hangs.
|
||||
func TestServerRoughDisconnect(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READONLY, NO_DELAY)
|
||||
|
||||
f, err := sftp.Open("/dev/zero")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer f.Close()
|
||||
go func() {
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
cmd.Process.Kill()
|
||||
}()
|
||||
|
||||
io.Copy(ioutil.Discard, f)
|
||||
sftp.Close()
|
||||
}
|
||||
|
||||
func benchmarkRead(b *testing.B, bufsize int, delay time.Duration) {
|
||||
size := 10*1024*1024 + 123 // ~10MiB
|
||||
|
||||
// open sftp client
|
||||
sftp, cmd := testClient(b, READONLY)
|
||||
sftp, cmd := testClient(b, READONLY, delay)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -786,38 +850,50 @@ func benchmarkRead(b *testing.B, bufsize int) {
|
||||
}
|
||||
|
||||
func BenchmarkRead1k(b *testing.B) {
|
||||
benchmarkRead(b, 1*1024)
|
||||
benchmarkRead(b, 1*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkRead16k(b *testing.B) {
|
||||
benchmarkRead(b, 16*1024)
|
||||
benchmarkRead(b, 16*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkRead32k(b *testing.B) {
|
||||
benchmarkRead(b, 32*1024)
|
||||
benchmarkRead(b, 32*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkRead128k(b *testing.B) {
|
||||
benchmarkRead(b, 128*1024)
|
||||
benchmarkRead(b, 128*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkRead512k(b *testing.B) {
|
||||
benchmarkRead(b, 512*1024)
|
||||
benchmarkRead(b, 512*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkRead1MiB(b *testing.B) {
|
||||
benchmarkRead(b, 1024*1024)
|
||||
benchmarkRead(b, 1024*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkRead4MiB(b *testing.B) {
|
||||
benchmarkRead(b, 4*1024*1024)
|
||||
benchmarkRead(b, 4*1024*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func benchmarkWrite(b *testing.B, bufsize int) {
|
||||
func BenchmarkRead4MiBDelay10Msec(b *testing.B) {
|
||||
benchmarkRead(b, 4*1024*1024, 10*time.Millisecond)
|
||||
}
|
||||
|
||||
func BenchmarkRead4MiBDelay50Msec(b *testing.B) {
|
||||
benchmarkRead(b, 4*1024*1024, 50*time.Millisecond)
|
||||
}
|
||||
|
||||
func BenchmarkRead4MiBDelay150Msec(b *testing.B) {
|
||||
benchmarkRead(b, 4*1024*1024, 150*time.Millisecond)
|
||||
}
|
||||
|
||||
func benchmarkWrite(b *testing.B, bufsize int, delay time.Duration) {
|
||||
size := 10*1024*1024 + 123 // ~10MiB
|
||||
|
||||
// open sftp client
|
||||
sftp, cmd := testClient(b, false)
|
||||
sftp, cmd := testClient(b, false, delay)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
@@ -870,29 +946,230 @@ func benchmarkWrite(b *testing.B, bufsize int) {
|
||||
}
|
||||
|
||||
func BenchmarkWrite1k(b *testing.B) {
|
||||
benchmarkWrite(b, 1*1024)
|
||||
benchmarkWrite(b, 1*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkWrite16k(b *testing.B) {
|
||||
benchmarkWrite(b, 16*1024)
|
||||
benchmarkWrite(b, 16*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkWrite32k(b *testing.B) {
|
||||
benchmarkWrite(b, 32*1024)
|
||||
benchmarkWrite(b, 32*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkWrite128k(b *testing.B) {
|
||||
benchmarkWrite(b, 128*1024)
|
||||
benchmarkWrite(b, 128*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkWrite512k(b *testing.B) {
|
||||
benchmarkWrite(b, 512*1024)
|
||||
benchmarkWrite(b, 512*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkWrite1MiB(b *testing.B) {
|
||||
benchmarkWrite(b, 1024*1024)
|
||||
benchmarkWrite(b, 1024*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkWrite4MiB(b *testing.B) {
|
||||
benchmarkWrite(b, 4*1024*1024)
|
||||
benchmarkWrite(b, 4*1024*1024, NO_DELAY)
|
||||
}
|
||||
|
||||
func BenchmarkWrite4MiBDelay10Msec(b *testing.B) {
|
||||
benchmarkWrite(b, 4*1024*1024, 10*time.Millisecond)
|
||||
}
|
||||
|
||||
func BenchmarkWrite4MiBDelay50Msec(b *testing.B) {
|
||||
benchmarkWrite(b, 4*1024*1024, 50*time.Millisecond)
|
||||
}
|
||||
|
||||
func BenchmarkWrite4MiBDelay150Msec(b *testing.B) {
|
||||
benchmarkWrite(b, 4*1024*1024, 150*time.Millisecond)
|
||||
}
|
||||
|
||||
func TestClientStatVFS(t *testing.T) {
|
||||
sftp, cmd := testClient(t, READWRITE, NO_DELAY)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
vfs, err := sftp.StatVFS("/")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// get system stats
|
||||
s := syscall.Statfs_t{}
|
||||
err = syscall.Statfs("/", &s)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// check some stats
|
||||
if vfs.Frsize != uint64(s.Frsize) {
|
||||
t.Fatal("fr_size does not match")
|
||||
}
|
||||
|
||||
if vfs.Bsize != uint64(s.Bsize) {
|
||||
t.Fatal("f_bsize does not match")
|
||||
}
|
||||
|
||||
if vfs.Namemax != uint64(s.Namelen) {
|
||||
t.Fatal("f_namemax does not match")
|
||||
}
|
||||
|
||||
if vfs.Bavail != s.Bavail {
|
||||
t.Fatal("f_bavail does not match")
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func benchmarkCopyDown(b *testing.B, fileSize int64, delay time.Duration) {
|
||||
// Create a temp file and fill it with zero's.
|
||||
src, err := ioutil.TempFile("", "sftptest")
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer src.Close()
|
||||
srcFilename := src.Name()
|
||||
defer os.Remove(srcFilename)
|
||||
zero, err := os.Open("/dev/zero")
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
n, err := io.Copy(src, io.LimitReader(zero, fileSize))
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
if n < fileSize {
|
||||
b.Fatal("short copy")
|
||||
}
|
||||
zero.Close()
|
||||
src.Close()
|
||||
|
||||
sftp, cmd := testClient(b, READONLY, delay)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
b.ResetTimer()
|
||||
b.SetBytes(fileSize)
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
dst, err := ioutil.TempFile("", "sftptest")
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer os.Remove(dst.Name())
|
||||
|
||||
src, err := sftp.Open(srcFilename)
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer src.Close()
|
||||
n, err := io.Copy(dst, src)
|
||||
if err != nil {
|
||||
b.Fatalf("copy error: %v", err)
|
||||
}
|
||||
if n < fileSize {
|
||||
b.Fatal("unable to copy all bytes")
|
||||
}
|
||||
dst.Close()
|
||||
fi, err := os.Stat(dst.Name())
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
|
||||
if fi.Size() != fileSize {
|
||||
b.Fatalf("wrong file size: want %d, got %d", fileSize, fi.Size())
|
||||
}
|
||||
os.Remove(dst.Name())
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkCopyDown10MiBDelay10Msec(b *testing.B) {
|
||||
benchmarkCopyDown(b, 10*1024*1024, 10*time.Millisecond)
|
||||
}
|
||||
|
||||
func BenchmarkCopyDown10MiBDelay50Msec(b *testing.B) {
|
||||
benchmarkCopyDown(b, 10*1024*1024, 50*time.Millisecond)
|
||||
}
|
||||
|
||||
func BenchmarkCopyDown10MiBDelay150Msec(b *testing.B) {
|
||||
benchmarkCopyDown(b, 10*1024*1024, 150*time.Millisecond)
|
||||
}
|
||||
|
||||
func benchmarkCopyUp(b *testing.B, fileSize int64, delay time.Duration) {
|
||||
// Create a temp file and fill it with zero's.
|
||||
src, err := ioutil.TempFile("", "sftptest")
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer src.Close()
|
||||
srcFilename := src.Name()
|
||||
defer os.Remove(srcFilename)
|
||||
zero, err := os.Open("/dev/zero")
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
n, err := io.Copy(src, io.LimitReader(zero, fileSize))
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
if n < fileSize {
|
||||
b.Fatal("short copy")
|
||||
}
|
||||
zero.Close()
|
||||
src.Close()
|
||||
|
||||
sftp, cmd := testClient(b, false, delay)
|
||||
defer cmd.Wait()
|
||||
defer sftp.Close()
|
||||
|
||||
b.ResetTimer()
|
||||
b.SetBytes(fileSize)
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
tmp, err := ioutil.TempFile("", "sftptest")
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
tmp.Close()
|
||||
defer os.Remove(tmp.Name())
|
||||
|
||||
dst, err := sftp.Create(tmp.Name())
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer dst.Close()
|
||||
src, err := os.Open(srcFilename)
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer src.Close()
|
||||
n, err := io.Copy(dst, src)
|
||||
if err != nil {
|
||||
b.Fatalf("copy error: %v", err)
|
||||
}
|
||||
if n < fileSize {
|
||||
b.Fatal("unable to copy all bytes")
|
||||
}
|
||||
|
||||
fi, err := os.Stat(tmp.Name())
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
|
||||
if fi.Size() != fileSize {
|
||||
b.Fatalf("wrong file size: want %d, got %d", fileSize, fi.Size())
|
||||
}
|
||||
os.Remove(tmp.Name())
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkCopyUp10MiBDelay10Msec(b *testing.B) {
|
||||
benchmarkCopyUp(b, 10*1024*1024, 10*time.Millisecond)
|
||||
}
|
||||
|
||||
func BenchmarkCopyUp10MiBDelay50Msec(b *testing.B) {
|
||||
benchmarkCopyUp(b, 10*1024*1024, 50*time.Millisecond)
|
||||
}
|
||||
|
||||
func BenchmarkCopyUp10MiBDelay150Msec(b *testing.B) {
|
||||
benchmarkCopyUp(b, 10*1024*1024, 150*time.Millisecond)
|
||||
}
|
||||
Generated
Vendored
+2
-1
@@ -22,6 +22,7 @@ var (
|
||||
HOST = flag.String("host", "localhost", "ssh server hostname")
|
||||
PORT = flag.Int("port", 22, "ssh server port")
|
||||
PASS = flag.String("pass", os.Getenv("SOCKSIE_SSH_PASSWORD"), "ssh password")
|
||||
SIZE = flag.Int("s", 1<<15, "set max packet size")
|
||||
)
|
||||
|
||||
func init() {
|
||||
@@ -49,7 +50,7 @@ func main() {
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
c, err := sftp.NewClient(conn)
|
||||
c, err := sftp.NewClient(conn, sftp.MaxPacket(*SIZE))
|
||||
if err != nil {
|
||||
log.Fatalf("unable to start sftp subsytem: %v", err)
|
||||
}
|
||||
|
||||
Generated
Vendored
+2
-1
@@ -22,6 +22,7 @@ var (
|
||||
HOST = flag.String("host", "localhost", "ssh server hostname")
|
||||
PORT = flag.Int("port", 22, "ssh server port")
|
||||
PASS = flag.String("pass", os.Getenv("SOCKSIE_SSH_PASSWORD"), "ssh password")
|
||||
SIZE = flag.Int("s", 1<<15, "set max packet size")
|
||||
)
|
||||
|
||||
func init() {
|
||||
@@ -49,7 +50,7 @@ func main() {
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
c, err := sftp.NewClient(conn)
|
||||
c, err := sftp.NewClient(conn, sftp.MaxPacket(*SIZE))
|
||||
if err != nil {
|
||||
log.Fatalf("unable to start sftp subsytem: %v", err)
|
||||
}
|
||||
|
||||
-147
@@ -1,147 +0,0 @@
|
||||
// gsftp implements a simple sftp client.
|
||||
//
|
||||
// gsftp understands the following commands:
|
||||
//
|
||||
// List a directory (and its subdirectories)
|
||||
// gsftp ls DIR
|
||||
//
|
||||
// Fetch a remote file
|
||||
// gsftp fetch FILE
|
||||
//
|
||||
// Put the contents of stdin to a remote file
|
||||
// cat LOCALFILE | gsftp put REMOTEFILE
|
||||
//
|
||||
// Print the details of a remote file
|
||||
// gsftp stat FILE
|
||||
//
|
||||
// Remove a remote file
|
||||
// gsftp rm FILE
|
||||
//
|
||||
// Rename a file
|
||||
// gsftp mv OLD NEW
|
||||
//
|
||||
package main
|
||||
|
||||
import (
|
||||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net"
|
||||
"os"
|
||||
|
||||
"golang.org/x/crypto/ssh"
|
||||
"golang.org/x/crypto/ssh/agent"
|
||||
|
||||
"github.com/pkg/sftp"
|
||||
)
|
||||
|
||||
var (
|
||||
USER = flag.String("user", os.Getenv("USER"), "ssh username")
|
||||
HOST = flag.String("host", "localhost", "ssh server hostname")
|
||||
PORT = flag.Int("port", 22, "ssh server port")
|
||||
PASS = flag.String("pass", os.Getenv("SOCKSIE_SSH_PASSWORD"), "ssh password")
|
||||
)
|
||||
|
||||
func init() {
|
||||
flag.Parse()
|
||||
if len(flag.Args()) < 1 {
|
||||
log.Fatal("subcommand required")
|
||||
}
|
||||
}
|
||||
|
||||
func main() {
|
||||
var auths []ssh.AuthMethod
|
||||
if aconn, err := net.Dial("unix", os.Getenv("SSH_AUTH_SOCK")); err == nil {
|
||||
auths = append(auths, ssh.PublicKeysCallback(agent.NewClient(aconn).Signers))
|
||||
|
||||
}
|
||||
if *PASS != "" {
|
||||
auths = append(auths, ssh.Password(*PASS))
|
||||
}
|
||||
|
||||
config := ssh.ClientConfig{
|
||||
User: *USER,
|
||||
Auth: auths,
|
||||
}
|
||||
addr := fmt.Sprintf("%s:%d", *HOST, *PORT)
|
||||
conn, err := ssh.Dial("tcp", addr, &config)
|
||||
if err != nil {
|
||||
log.Fatalf("unable to connect to [%s]: %v", addr, err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
client, err := sftp.NewClient(conn)
|
||||
if err != nil {
|
||||
log.Fatalf("unable to start sftp subsytem: %v", err)
|
||||
}
|
||||
defer client.Close()
|
||||
switch cmd := flag.Args()[0]; cmd {
|
||||
case "ls":
|
||||
if len(flag.Args()) < 2 {
|
||||
log.Fatalf("%s %s: remote path required", cmd, os.Args[0])
|
||||
}
|
||||
walker := client.Walk(flag.Args()[1])
|
||||
for walker.Step() {
|
||||
if err := walker.Err(); err != nil {
|
||||
log.Println(err)
|
||||
continue
|
||||
}
|
||||
fmt.Println(walker.Path())
|
||||
}
|
||||
case "fetch":
|
||||
if len(flag.Args()) < 2 {
|
||||
log.Fatalf("%s %s: remote path required", cmd, os.Args[0])
|
||||
}
|
||||
f, err := client.Open(flag.Args()[1])
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
defer f.Close()
|
||||
if _, err := io.Copy(os.Stdout, f); err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
case "put":
|
||||
if len(flag.Args()) < 2 {
|
||||
log.Fatalf("%s %s: remote path required", cmd, os.Args[0])
|
||||
}
|
||||
f, err := client.Create(flag.Args()[1])
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
defer f.Close()
|
||||
if _, err := io.Copy(f, os.Stdin); err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
case "stat":
|
||||
if len(flag.Args()) < 2 {
|
||||
log.Fatalf("%s %s: remote path required", cmd, os.Args[0])
|
||||
}
|
||||
f, err := client.Open(flag.Args()[1])
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
defer f.Close()
|
||||
fi, err := f.Stat()
|
||||
if err != nil {
|
||||
log.Fatalf("unable to stat file: %v", err)
|
||||
}
|
||||
fmt.Printf("%s %d %v\n", fi.Name(), fi.Size(), fi.Mode())
|
||||
case "rm":
|
||||
if len(flag.Args()) < 2 {
|
||||
log.Fatalf("%s %s: remote path required", cmd, os.Args[0])
|
||||
}
|
||||
if err := client.Remove(flag.Args()[1]); err != nil {
|
||||
log.Fatalf("unable to remove file: %v", err)
|
||||
}
|
||||
case "mv":
|
||||
if len(flag.Args()) < 3 {
|
||||
log.Fatalf("%s %s: old and new name required", cmd, os.Args[0])
|
||||
}
|
||||
if err := client.Rename(flag.Args()[1], flag.Args()[2]); err != nil {
|
||||
log.Fatalf("unable to rename file: %v", err)
|
||||
}
|
||||
default:
|
||||
log.Fatalf("unknown subcommand: %v", cmd)
|
||||
}
|
||||
}
|
||||
Generated
Vendored
+2
-1
@@ -23,6 +23,7 @@ var (
|
||||
HOST = flag.String("host", "localhost", "ssh server hostname")
|
||||
PORT = flag.Int("port", 22, "ssh server port")
|
||||
PASS = flag.String("pass", os.Getenv("SOCKSIE_SSH_PASSWORD"), "ssh password")
|
||||
SIZE = flag.Int("s", 1<<15, "set max packet size")
|
||||
)
|
||||
|
||||
func init() {
|
||||
@@ -50,7 +51,7 @@ func main() {
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
c, err := sftp.NewClient(conn)
|
||||
c, err := sftp.NewClient(conn, sftp.MaxPacket(*SIZE))
|
||||
if err != nil {
|
||||
log.Fatalf("unable to start sftp subsytem: %v", err)
|
||||
}
|
||||
|
||||
Generated
Vendored
+2
-1
@@ -23,6 +23,7 @@ var (
|
||||
HOST = flag.String("host", "localhost", "ssh server hostname")
|
||||
PORT = flag.Int("port", 22, "ssh server port")
|
||||
PASS = flag.String("pass", os.Getenv("SOCKSIE_SSH_PASSWORD"), "ssh password")
|
||||
SIZE = flag.Int("s", 1<<15, "set max packet size")
|
||||
)
|
||||
|
||||
func init() {
|
||||
@@ -50,7 +51,7 @@ func main() {
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
c, err := sftp.NewClient(conn)
|
||||
c, err := sftp.NewClient(conn, sftp.MaxPacket(*SIZE))
|
||||
if err != nil {
|
||||
log.Fatalf("unable to start sftp subsytem: %v", err)
|
||||
}
|
||||
|
||||
+71
-1
@@ -64,7 +64,6 @@ func unmarshalString(b []byte) (string, []byte) {
|
||||
}
|
||||
|
||||
// sendPacket marshals p according to RFC 4234.
|
||||
|
||||
func sendPacket(w io.Writer, m encoding.BinaryMarshaler) error {
|
||||
bb, err := m.MarshalBinary()
|
||||
if err != nil {
|
||||
@@ -142,6 +141,8 @@ func (p sshFxpReaddirPacket) MarshalBinary() ([]byte, error) {
|
||||
return marshalIdString(ssh_FXP_READDIR, p.Id, p.Handle)
|
||||
}
|
||||
|
||||
func (p sshFxpReaddirPacket) id() uint32 { return p.Id }
|
||||
|
||||
type sshFxpOpendirPacket struct {
|
||||
Id uint32
|
||||
Path string
|
||||
@@ -151,11 +152,15 @@ func (p sshFxpOpendirPacket) MarshalBinary() ([]byte, error) {
|
||||
return marshalIdString(ssh_FXP_OPENDIR, p.Id, p.Path)
|
||||
}
|
||||
|
||||
func (p sshFxpOpendirPacket) id() uint32 { return p.Id }
|
||||
|
||||
type sshFxpLstatPacket struct {
|
||||
Id uint32
|
||||
Path string
|
||||
}
|
||||
|
||||
func (p sshFxpLstatPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpLstatPacket) MarshalBinary() ([]byte, error) {
|
||||
return marshalIdString(ssh_FXP_LSTAT, p.Id, p.Path)
|
||||
}
|
||||
@@ -165,6 +170,8 @@ type sshFxpFstatPacket struct {
|
||||
Handle string
|
||||
}
|
||||
|
||||
func (p sshFxpFstatPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpFstatPacket) MarshalBinary() ([]byte, error) {
|
||||
return marshalIdString(ssh_FXP_FSTAT, p.Id, p.Handle)
|
||||
}
|
||||
@@ -178,11 +185,15 @@ func (p sshFxpClosePacket) MarshalBinary() ([]byte, error) {
|
||||
return marshalIdString(ssh_FXP_CLOSE, p.Id, p.Handle)
|
||||
}
|
||||
|
||||
func (p sshFxpClosePacket) id() uint32 { return p.Id }
|
||||
|
||||
type sshFxpRemovePacket struct {
|
||||
Id uint32
|
||||
Filename string
|
||||
}
|
||||
|
||||
func (p sshFxpRemovePacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpRemovePacket) MarshalBinary() ([]byte, error) {
|
||||
return marshalIdString(ssh_FXP_REMOVE, p.Id, p.Filename)
|
||||
}
|
||||
@@ -192,6 +203,8 @@ type sshFxpRmdirPacket struct {
|
||||
Path string
|
||||
}
|
||||
|
||||
func (p sshFxpRmdirPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpRmdirPacket) MarshalBinary() ([]byte, error) {
|
||||
return marshalIdString(ssh_FXP_RMDIR, p.Id, p.Path)
|
||||
}
|
||||
@@ -201,6 +214,8 @@ type sshFxpReadlinkPacket struct {
|
||||
Path string
|
||||
}
|
||||
|
||||
func (p sshFxpReadlinkPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpReadlinkPacket) MarshalBinary() ([]byte, error) {
|
||||
return marshalIdString(ssh_FXP_READLINK, p.Id, p.Path)
|
||||
}
|
||||
@@ -212,6 +227,8 @@ type sshFxpOpenPacket struct {
|
||||
Flags uint32 // ignored
|
||||
}
|
||||
|
||||
func (p sshFxpOpenPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpOpenPacket) MarshalBinary() ([]byte, error) {
|
||||
l := 1 + 4 +
|
||||
4 + len(p.Path) +
|
||||
@@ -233,6 +250,8 @@ type sshFxpReadPacket struct {
|
||||
Len uint32
|
||||
}
|
||||
|
||||
func (p sshFxpReadPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpReadPacket) MarshalBinary() ([]byte, error) {
|
||||
l := 1 + 4 + // type(byte) + uint32
|
||||
4 + len(p.Handle) +
|
||||
@@ -253,6 +272,8 @@ type sshFxpRenamePacket struct {
|
||||
Newpath string
|
||||
}
|
||||
|
||||
func (p sshFxpRenamePacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpRenamePacket) MarshalBinary() ([]byte, error) {
|
||||
l := 1 + 4 + // type(byte) + uint32
|
||||
4 + len(p.Oldpath) +
|
||||
@@ -274,6 +295,8 @@ type sshFxpWritePacket struct {
|
||||
Data []byte
|
||||
}
|
||||
|
||||
func (s sshFxpWritePacket) id() uint32 { return s.Id }
|
||||
|
||||
func (s sshFxpWritePacket) MarshalBinary() ([]byte, error) {
|
||||
l := 1 + 4 + // type(byte) + uint32
|
||||
4 + len(s.Handle) +
|
||||
@@ -296,6 +319,8 @@ type sshFxpMkdirPacket struct {
|
||||
Flags uint32 // ignored
|
||||
}
|
||||
|
||||
func (p sshFxpMkdirPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpMkdirPacket) MarshalBinary() ([]byte, error) {
|
||||
l := 1 + 4 + // type(byte) + uint32
|
||||
4 + len(p.Path) +
|
||||
@@ -316,6 +341,8 @@ type sshFxpSetstatPacket struct {
|
||||
Attrs interface{}
|
||||
}
|
||||
|
||||
func (p sshFxpSetstatPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpSetstatPacket) MarshalBinary() ([]byte, error) {
|
||||
l := 1 + 4 + // type(byte) + uint32
|
||||
4 + len(p.Path) +
|
||||
@@ -329,3 +356,46 @@ func (p sshFxpSetstatPacket) MarshalBinary() ([]byte, error) {
|
||||
b = marshal(b, p.Attrs)
|
||||
return b, nil
|
||||
}
|
||||
|
||||
type sshFxpStatvfsPacket struct {
|
||||
Id uint32
|
||||
Path string
|
||||
}
|
||||
|
||||
func (p sshFxpStatvfsPacket) id() uint32 { return p.Id }
|
||||
|
||||
func (p sshFxpStatvfsPacket) MarshalBinary() ([]byte, error) {
|
||||
l := 1 + 4 + // type(byte) + uint32
|
||||
len(p.Path) +
|
||||
len("statvfs@openssh.com")
|
||||
|
||||
b := make([]byte, 0, l)
|
||||
b = append(b, ssh_FXP_EXTENDED)
|
||||
b = marshalUint32(b, p.Id)
|
||||
b = marshalString(b, "statvfs@openssh.com")
|
||||
b = marshalString(b, p.Path)
|
||||
return b, nil
|
||||
}
|
||||
|
||||
type StatVFS struct {
|
||||
Id uint32
|
||||
Bsize uint64 /* file system block size */
|
||||
Frsize uint64 /* fundamental fs block size */
|
||||
Blocks uint64 /* number of blocks (unit f_frsize) */
|
||||
Bfree uint64 /* free blocks in file system */
|
||||
Bavail uint64 /* free blocks for non-root */
|
||||
Files uint64 /* total file inodes */
|
||||
Ffree uint64 /* free file inodes */
|
||||
Favail uint64 /* free file inodes for to non-root */
|
||||
Fsid uint64 /* file system id */
|
||||
Flag uint64 /* bit mask of f_flag values */
|
||||
Namemax uint64 /* maximum filename length */
|
||||
}
|
||||
|
||||
func (p *StatVFS) TotalSpace() uint64 {
|
||||
return p.Frsize * p.Blocks
|
||||
}
|
||||
|
||||
func (p *StatVFS) FreeSpace() uint64 {
|
||||
return p.Frsize * p.Bfree
|
||||
}
|
||||
Reference in new issue
Block a user