Test / test (push) Successful in 2m4s
Release / gates (push) Successful in 2m5s
Release / build (amd64, freebsd) (push) Successful in 1m27s
Release / build (amd64, linux) (push) Successful in 1m22s
Release / build (amd64, netbsd) (push) Successful in 1m19s
Release / build (amd64, openbsd) (push) Successful in 1m20s
Release / build (arm64, darwin) (push) Successful in 1m21s
Release / build (arm64, freebsd) (push) Successful in 1m26s
Release / build (arm64, linux) (push) Successful in 1m25s
Release / build (arm64, netbsd) (push) Successful in 1m31s
Release / build (arm64, openbsd) (push) Successful in 1m27s
Release / build (loong64, linux) (push) Successful in 1m37s
Release / build (riscv64, linux) (push) Successful in 1m21s
Release / release (push) Successful in 40s
Assisted-by: GLM 5.3 Flash
266 lines
6.3 KiB
Go
266 lines
6.3 KiB
Go
// Copyright (c) 2026 Petr Balvín <opensource@petrbalvin.org> (https://petrbalvin.org)
|
|
// SPDX-License-Identifier: MIT
|
|
|
|
package main
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
|
|
"sourcedock.dev/petrbalvin/nfs/internal/nfs4"
|
|
"sourcedock.dev/petrbalvin/nfs/internal/nfsclient"
|
|
"sourcedock.dev/petrbalvin/nfs/internal/xdr"
|
|
)
|
|
|
|
// putChunk is the size of one in flight read or write of get and put,
|
|
// the same megabyte the commands have always transferred per compound.
|
|
const putChunk = 1 << 20
|
|
|
|
// cmdGet mirrors a remote file into a local file through READ compounds,
|
|
// up to workers of them in flight. The local file is created first, so a
|
|
// shorter remote leaves no tail behind.
|
|
func cmdGet(cl *nfsclient.Client, remote, local string, workers int) error {
|
|
out, err := os.Create(local)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer out.Close()
|
|
|
|
var next atomic.Int64
|
|
var stop atomic.Bool
|
|
type chunk struct {
|
|
idx int64
|
|
data []byte
|
|
eof bool
|
|
}
|
|
done := make(chan error, 1)
|
|
pages := make(chan chunk, workers)
|
|
var wg sync.WaitGroup
|
|
for range workers {
|
|
wg.Go(func() {
|
|
for {
|
|
if stop.Load() {
|
|
return
|
|
}
|
|
idx := next.Add(1) - 1
|
|
ops := append(pathOps(remote),
|
|
nfs4.AppendReadArgs(nil, nfs4.Stateid{}, uint64(idx)*putChunk, putChunk))
|
|
res, bodies, err := cl.Compound("get", ops)
|
|
if err != nil {
|
|
done <- err
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
if res.Status != nfs4.ErrOK {
|
|
done <- fmt.Errorf("get: nfs status %d", res.Status)
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
body, err := bodyAt("get", bodies, len(bodies)-1)
|
|
if err != nil {
|
|
done <- err
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
d := xdr.NewDecoder(body)
|
|
eof, err := d.Bool()
|
|
if err != nil {
|
|
done <- err
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
data, err := d.VarOpaque()
|
|
if err != nil {
|
|
done <- err
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
pages <- chunk{idx: idx, data: data, eof: eof}
|
|
if eof {
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
}
|
|
})
|
|
}
|
|
go func() { wg.Wait(); close(pages) }()
|
|
|
|
var total uint64
|
|
for page := range pages {
|
|
if _, err := out.WriteAt(page.data, page.idx*putChunk); err != nil {
|
|
return err
|
|
}
|
|
total += uint64(len(page.data))
|
|
if page.eof {
|
|
stop.Store(true)
|
|
}
|
|
}
|
|
select {
|
|
case err := <-done:
|
|
return err
|
|
default:
|
|
}
|
|
fmt.Printf("wrote %d bytes from %s\n", total, remote)
|
|
return nil
|
|
}
|
|
|
|
// cmdPut writes a local file to the server through OPEN and WRITE,
|
|
// up to workers of them in flight after the truncate.
|
|
func cmdPut(cl *nfsclient.Client, local, remote string, workers int) error {
|
|
in, err := os.Open(local)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer in.Close()
|
|
info, err := in.Stat()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
parts := splitPath(remote)
|
|
if len(parts) == 0 {
|
|
return fmt.Errorf("put: empty remote path")
|
|
}
|
|
name := parts[len(parts)-1]
|
|
dirOps := pathOps(strings.Join(parts[:len(parts)-1], "/"))
|
|
openOps := append(dirOps,
|
|
nfs4.AppendOpenArgs(nil, 0, []byte("nfs-cli"), nfs4.ShareAccessBoth, 0,
|
|
true, 0o644, name),
|
|
nfs4.AppendGetfh(nil))
|
|
res, bodies, err := cl.Compound("put-open", openOps)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if res.Status != nfs4.ErrOK {
|
|
return fmt.Errorf("put: open status %d", res.Status)
|
|
}
|
|
var st nfs4.Stateid
|
|
stateBody, err := bodyAt("put", bodies, len(bodies)-2)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
copy(st[:], stateBody)
|
|
fhBody, err := bodyAt("put", bodies, len(bodies)-1)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fh, err := xdr.NewDecoder(fhBody).VarOpaque()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// PUT replaces the whole file: the size is zeroed through SETATTR
|
|
// before the first write, so a shorter file leaves no tail behind.
|
|
tres, _, err := cl.Compound("put-truncate", [][]byte{
|
|
nfs4.AppendPutfh(nil, fh),
|
|
nfs4.AppendSetattrArgs(nil, nfs4.AllZero, nfs4.OfBits(nfs4.AttrSize), nfs4.Attrs{}),
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tres.Status != nfs4.ErrOK {
|
|
return fmt.Errorf("put: truncate status %d", tres.Status)
|
|
}
|
|
|
|
chunks := (info.Size() + putChunk - 1) / putChunk
|
|
var next atomic.Int64
|
|
var stop atomic.Bool
|
|
done := make(chan error, 1)
|
|
var wg sync.WaitGroup
|
|
for range workers {
|
|
wg.Go(func() {
|
|
buf := make([]byte, putChunk)
|
|
for {
|
|
if stop.Load() {
|
|
return
|
|
}
|
|
idx := next.Add(1) - 1
|
|
if idx >= chunks {
|
|
return
|
|
}
|
|
n, rerr := in.ReadAt(buf, idx*putChunk)
|
|
if rerr != nil && !errors.Is(rerr, os.ErrClosed) {
|
|
// A short final read is the file's end, not a failure.
|
|
if !errors.Is(rerr, io.EOF) {
|
|
done <- rerr
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
}
|
|
wres, _, err := cl.Compound("put-write", [][]byte{
|
|
nfs4.AppendPutfh(nil, fh),
|
|
nfs4.AppendWriteArgs(nil, st, uint64(idx)*putChunk, nfs4.StableFileSync, buf[:n]),
|
|
})
|
|
if err != nil {
|
|
done <- err
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
if wres.Status != nfs4.ErrOK {
|
|
done <- fmt.Errorf("put: write status %d", wres.Status)
|
|
stop.Store(true)
|
|
return
|
|
}
|
|
}
|
|
})
|
|
}
|
|
wg.Wait()
|
|
select {
|
|
case err := <-done:
|
|
return err
|
|
default:
|
|
}
|
|
cres, _, err := cl.Compound("put-close", [][]byte{
|
|
nfs4.AppendPutfh(nil, fh),
|
|
nfs4.AppendCloseArgs(nil, st),
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if cres.Status != nfs4.ErrOK {
|
|
return fmt.Errorf("put: close status %d", cres.Status)
|
|
}
|
|
fmt.Printf("wrote %d bytes to %s\n", info.Size(), remote)
|
|
return nil
|
|
}
|
|
|
|
// cmdRm removes one object from the server through REMOVE.
|
|
func cmdRm(cl *nfsclient.Client, path string) error {
|
|
parts := splitPath(path)
|
|
if len(parts) == 0 {
|
|
return errors.New("rm: empty path")
|
|
}
|
|
dir := strings.Join(parts[:len(parts)-1], "/")
|
|
ops := append(pathOps(dir), nfs4.AppendRemoveArgs(nil, parts[len(parts)-1]))
|
|
res, _, err := cl.Compound("rm", ops)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if res.Status != nfs4.ErrOK {
|
|
return fmt.Errorf("rm: nfs status %d", res.Status)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// cmdMkdir makes one directory on the server through CREATE NF4DIR.
|
|
func cmdMkdir(cl *nfsclient.Client, path string) error {
|
|
parts := splitPath(path)
|
|
if len(parts) == 0 {
|
|
return errors.New("mkdir: empty path")
|
|
}
|
|
dir := strings.Join(parts[:len(parts)-1], "/")
|
|
ops := append(pathOps(dir),
|
|
nfs4.AppendCreateArgs(nil, nfs4.NF4Dir, parts[len(parts)-1], "", 0, 0, 0o755))
|
|
res, _, err := cl.Compound("mkdir", ops)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if res.Status != nfs4.ErrOK {
|
|
return fmt.Errorf("mkdir: nfs status %d", res.Status)
|
|
}
|
|
return nil
|
|
}
|