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
409 lines
13 KiB
Go
409 lines
13 KiB
Go
// Copyright (c) 2026 Petr Balvín <opensource@petrbalvin.org> (https://petrbalvin.org)
|
|
// SPDX-License-Identifier: MIT
|
|
|
|
package nfs4server
|
|
|
|
import (
|
|
"fmt"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"sourcedock.dev/petrbalvin/nfs/internal/nfs4"
|
|
)
|
|
|
|
// A slot is one entry of the fore channel slot table: the sequence number
|
|
// the slot is at and the cached bytes of the last answer, which is what a
|
|
// retry after a lost reply replays.
|
|
type slot struct {
|
|
sequence uint32
|
|
status uint32 // the COMPOUND status the cached answer carries
|
|
cached []byte // the result ops of the cached answer
|
|
used bool
|
|
}
|
|
|
|
// A session is one CREATE_SESSION product: its identifier, its slot
|
|
// table, the callback program the client picked and the connection the
|
|
// back channel rides on.
|
|
type session struct {
|
|
id nfs4.SessionID
|
|
|
|
// slMu guards the slot table: slot sequences and their cached
|
|
// replies are the per session state, contended only by the
|
|
// compounds of the session that owns them.
|
|
slMu sync.Mutex
|
|
slots []slot
|
|
cb *connCB
|
|
cbProg uint32
|
|
cbMu sync.Mutex
|
|
cbSeq uint32
|
|
}
|
|
|
|
// A client is one EXCHANGE_ID identity: the verifier it rebooted with,
|
|
// the sessions it created and the CREATE_SESSION replay state.
|
|
type client struct {
|
|
id uint64
|
|
verifier [8]byte
|
|
sessions map[nfs4.SessionID]*session
|
|
lastCSSeq uint32
|
|
lastCSID nfs4.SessionID
|
|
csHas bool
|
|
|
|
// renewNS carries the last renewal as unix nanoseconds, read and
|
|
// written atomically: the lease check runs per operation and must
|
|
// not serialise the clients against each other.
|
|
renewNS atomic.Int64
|
|
}
|
|
|
|
// sessionStore keeps the clients and sessions of the server. Every method
|
|
// is safe for concurrent use.
|
|
type sessionStore struct {
|
|
mu sync.RWMutex
|
|
prefix [4]byte
|
|
nextSess uint32
|
|
nextID uint64
|
|
byOwner map[string]*ownerEntry // the owner id is the client identity
|
|
byID map[uint64]*client // clientid, assigned by the server
|
|
sessions map[nfs4.SessionID]*session
|
|
}
|
|
|
|
// An ownerEntry pairs a client identity with the verifier of the life it
|
|
// was registered under.
|
|
type ownerEntry struct {
|
|
verifier [8]byte
|
|
client *client
|
|
}
|
|
|
|
func newSessionStore(prefix [4]byte) *sessionStore {
|
|
return &sessionStore{
|
|
prefix: prefix,
|
|
// Client ids start from a random base and count up: the ids stay
|
|
// unique, and a foreign client cannot walk another's id by
|
|
// guessing a small counter.
|
|
nextID: randCounter(),
|
|
byOwner: make(map[string]*ownerEntry),
|
|
byID: make(map[uint64]*client),
|
|
sessions: make(map[nfs4.SessionID]*session),
|
|
}
|
|
}
|
|
|
|
// exchangeID resolves the owner to a client id. The same verifier and
|
|
// owner id confirm the client it already assigned; a new verifier with a
|
|
// known owner id means the client rebooted and takes everything with it.
|
|
// The flags of the request are refused by design, RFC 8881 section 13.1:
|
|
// the reply carries the server's own roles, never an echo.
|
|
func (s *sessionStore) exchangeID(verifier [8]byte, ownerID []byte, now time.Time) (clientid uint64, sequence uint32, outFlags uint32, rebooted uint64) {
|
|
// The reply flags carry this server's own roles, RFC 8881 section
|
|
// 13.1: it never echoes the request. The server is a metadata server
|
|
// with itself as the data server, and it serves referrals. USE_PNFS
|
|
// MDS and USE_NON_PNFS are mutually exclusive roles and the Linux
|
|
// client rejects a reply that claims both; this server speaks pNFS,
|
|
// so it claims MDS and DS. CONFIRMED_R is added only when a session
|
|
// already exists; BIND_PRINC_STATEID and SUPP_FENCE_OPS stay off
|
|
// because this build binds no stateids to principals and fences
|
|
// nothing.
|
|
serverFlags := uint32(nfs4.ExchgIDUsePnfsMds | nfs4.ExchgIDUsePnfsDs |
|
|
nfs4.ExchgIDSuppMovedRefer)
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if e, ok := s.byOwner[string(ownerID)]; ok {
|
|
if e.verifier == verifier {
|
|
// The same life of the same client: confirm it. The exchange
|
|
// itself proves liveness, so it renews the lease.
|
|
e.client.renewNS.Store(now.UnixNano())
|
|
return e.client.id, e.client.lastCSSeq, serverFlags | nfs4.ExchgIDConfirmedR, 0
|
|
}
|
|
// A new verifier for a known owner is a reboot: everything the
|
|
// client had is gone with the old life. The caller releases the
|
|
// old client's state everywhere; this store drops its sessions.
|
|
old := e.client.id
|
|
s.dropClient(e.client)
|
|
rebooted = old
|
|
}
|
|
s.nextID++
|
|
c := &client{id: s.nextID, verifier: verifier, sessions: make(map[nfs4.SessionID]*session)}
|
|
c.renewNS.Store(now.UnixNano())
|
|
s.byOwner[string(ownerID)] = &ownerEntry{verifier: verifier, client: c}
|
|
s.byID[c.id] = c
|
|
return c.id, 0, serverFlags, 0
|
|
}
|
|
|
|
// createSession makes a session for the client, or replays the cached
|
|
// answer when the sequence repeats. Every session carries a distinct
|
|
// id: four random bytes of the server prefix and a per-session number,
|
|
// so a second CREATE_SESSION never reuses the first session's id and
|
|
// its slot table, RFC 8881 section 18.36. The boolean reports a replay.
|
|
func (s *sessionStore) createSession(clientid uint64, sequence, cbProgram uint32) (id nfs4.SessionID, replay bool, status uint32) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
c, ok := s.byID[clientid]
|
|
if !ok {
|
|
return id, false, nfs4.ErrStaleClientID
|
|
}
|
|
if c.csHas {
|
|
if sequence == c.lastCSSeq {
|
|
// A retry of the same CREATE_SESSION: answer with the id the
|
|
// first attempt minted.
|
|
return c.lastCSID, true, nfs4.ErrOK
|
|
}
|
|
if sequence < c.lastCSSeq {
|
|
return id, false, nfs4.ErrSeqMisordered
|
|
}
|
|
}
|
|
s.nextSess++
|
|
id = nfs4.MakeNumberedSessionID(s.prefix, s.nextSess, clientid)
|
|
sess := &session{id: id, slots: make([]slot, defaultSlots), cbProg: cbProgram}
|
|
c.sessions[id] = sess
|
|
c.lastCSSeq = sequence
|
|
c.lastCSID = id
|
|
c.csHas = true
|
|
c.renewNS.Store(time.Now().UnixNano())
|
|
s.sessions[id] = sess
|
|
return id, false, nfs4.ErrOK
|
|
}
|
|
|
|
// destroySession removes the session.
|
|
func (s *sessionStore) destroySession(id nfs4.SessionID) uint32 {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if _, ok := s.sessions[id]; !ok {
|
|
return nfs4.ErrBadSession
|
|
}
|
|
delete(s.sessions, id)
|
|
if c, ok := s.byID[id.ClientIDOf()]; ok {
|
|
delete(c.sessions, id)
|
|
}
|
|
return nfs4.ErrOK
|
|
}
|
|
|
|
// sequence advances the slot. The reply cached on the slot answers a
|
|
// repeated sequence; a sequence that is neither the cached one nor the
|
|
// next one is misordered. The boolean reports a replay.
|
|
func (s *sessionStore) sequence(id nfs4.SessionID, sequence, slotID uint32) (*session, bool, uint32) {
|
|
s.mu.RLock()
|
|
sess, ok := s.sessions[id]
|
|
s.mu.RUnlock()
|
|
if !ok {
|
|
return nil, false, nfs4.ErrBadSession
|
|
}
|
|
if slotID >= uint32(len(sess.slots)) {
|
|
return nil, false, nfs4.ErrBadSlot
|
|
}
|
|
// The slot belongs to this session alone; other sessions of other
|
|
// clients proceed beside it.
|
|
sess.slMu.Lock()
|
|
defer sess.slMu.Unlock()
|
|
sl := &sess.slots[slotID]
|
|
switch {
|
|
case !sl.used:
|
|
sl.used = true
|
|
sl.sequence = sequence
|
|
return sess, false, nfs4.ErrOK
|
|
case sl.sequence == sequence:
|
|
return sess, true, nfs4.ErrOK
|
|
case sl.sequence+1 == sequence:
|
|
sl.sequence = sequence
|
|
return sess, false, nfs4.ErrOK
|
|
default:
|
|
return nil, false, nfs4.ErrSeqMisordered
|
|
}
|
|
}
|
|
|
|
// cacheReply stores the answer bytes of one slot for its replay.
|
|
func (s *sessionStore) cacheReply(id nfs4.SessionID, slotID uint32, status uint32, ops []byte) {
|
|
s.mu.RLock()
|
|
sess, ok := s.sessions[id]
|
|
s.mu.RUnlock()
|
|
if !ok || slotID >= uint32(len(sess.slots)) {
|
|
return
|
|
}
|
|
sess.slMu.Lock()
|
|
defer sess.slMu.Unlock()
|
|
sess.slots[slotID].status = status
|
|
sess.slots[slotID].cached = ops
|
|
}
|
|
|
|
// replay returns the cached answer of the slot.
|
|
func (s *sessionStore) replay(id nfs4.SessionID, slotID uint32) (status uint32, ops []byte, ok bool) {
|
|
s.mu.RLock()
|
|
sess, ok := s.sessions[id]
|
|
s.mu.RUnlock()
|
|
if !ok || slotID >= uint32(len(sess.slots)) {
|
|
return 0, nil, false
|
|
}
|
|
sess.slMu.Lock()
|
|
defer sess.slMu.Unlock()
|
|
sl := &sess.slots[slotID]
|
|
return sl.status, sl.cached, sl.used
|
|
}
|
|
|
|
// dropClient removes a client and every session it made.
|
|
func (s *sessionStore) dropClient(c *client) {
|
|
for id := range c.sessions {
|
|
delete(s.sessions, id)
|
|
}
|
|
delete(s.byID, c.id)
|
|
}
|
|
|
|
// defaultSlots is the fore channel slot table the server grants.
|
|
const defaultSlots = 8
|
|
|
|
// renew marks the client's lease as refreshed. A SEQUENCE from any session
|
|
// of the client renews it, as do the stateful operations.
|
|
func (s *sessionStore) renew(clientid uint64, now time.Time) {
|
|
s.mu.RLock()
|
|
c, ok := s.byID[clientid]
|
|
s.mu.RUnlock()
|
|
if ok {
|
|
c.renewNS.Store(now.UnixNano())
|
|
}
|
|
}
|
|
|
|
// leaseExpired reports whether the client's lease has lapsed under the
|
|
// given period. A period of zero or less disables lease enforcement, and
|
|
// a client the store does not know is not this store's business.
|
|
func (s *sessionStore) leaseExpired(clientid uint64, period time.Duration, now time.Time) bool {
|
|
if period <= 0 {
|
|
return false
|
|
}
|
|
s.mu.RLock()
|
|
c, ok := s.byID[clientid]
|
|
s.mu.RUnlock()
|
|
if !ok {
|
|
return false
|
|
}
|
|
return now.UnixNano()-c.renewNS.Load() > period.Nanoseconds()
|
|
}
|
|
|
|
// destroyClientID drops the client and its sessions and reports whether
|
|
// the client id was known.
|
|
func (s *sessionStore) destroyClientID(clientid uint64) bool {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
c, ok := s.byID[clientid]
|
|
if !ok {
|
|
return false
|
|
}
|
|
for id := range c.sessions {
|
|
delete(s.sessions, id)
|
|
}
|
|
delete(s.byID, clientid)
|
|
for owner, e := range s.byOwner {
|
|
if e.client == c {
|
|
delete(s.byOwner, owner)
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
// knownClient reports whether the client id is live.
|
|
func (s *sessionStore) knownClient(clientid uint64) bool {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
_, ok := s.byID[clientid]
|
|
return ok
|
|
}
|
|
|
|
// attachCB binds a negotiated back channel to the session. The callback
|
|
// program the client named in CREATE_SESSION travels with the session,
|
|
// so a later BIND_CONN_TO_SESSION binds the connection under it.
|
|
func (s *sessionStore) attachCB(id nfs4.SessionID, cb *connCB) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if sess, ok := s.sessions[id]; ok {
|
|
sess.cb = cb
|
|
cb.setProgram(sess.cbProg)
|
|
}
|
|
}
|
|
|
|
// queueCB posts one fire and forget CB_COMPOUND behind a CB_SEQUENCE
|
|
// on the session's back channel. The sequence number is drawn and the
|
|
// work enqueued under one lock, so the order the connection's worker
|
|
// sends in is the order the numbers were drawn in. The done action runs
|
|
// on the worker when delivery ends.
|
|
func (s *sessionStore) queueCB(id nfs4.SessionID, tag string, ops [][]byte, done func(cbResult)) error {
|
|
sess, ok := s.lookupSessionPtr(id)
|
|
if !ok {
|
|
return errNoBackChannel("no such session")
|
|
}
|
|
sess.cbMu.Lock()
|
|
defer sess.cbMu.Unlock()
|
|
if sess.cb == nil {
|
|
return errNoBackChannel("session has no back channel")
|
|
}
|
|
sess.cbSeq++
|
|
seqArgs := nfs4.AppendCBSequenceArgs(nil, id, sess.cbSeq, 0, 0, true)
|
|
all := append([][]byte{seqArgs}, ops...)
|
|
return sess.cb.post(tag, id.ClientIDOf(), all, done)
|
|
}
|
|
|
|
// callCB posts one CB_COMPOUND the same way and waits for its delivery.
|
|
func (s *sessionStore) callCB(id nfs4.SessionID, tag string, ops [][]byte) (nfs4.CompoundRes, [][]byte, error) {
|
|
sess, ok := s.lookupSessionPtr(id)
|
|
if !ok {
|
|
return nfs4.CompoundRes{}, nil, errNoBackChannel("no such session")
|
|
}
|
|
sess.cbMu.Lock()
|
|
if sess.cb == nil {
|
|
sess.cbMu.Unlock()
|
|
return nfs4.CompoundRes{}, nil, errNoBackChannel("session has no back channel")
|
|
}
|
|
sess.cbSeq++
|
|
seqArgs := nfs4.AppendCBSequenceArgs(nil, id, sess.cbSeq, 0, 0, true)
|
|
w := cbWork{tag: tag, clientid: id.ClientIDOf(),
|
|
ops: append([][]byte{seqArgs}, ops...),
|
|
result: make(chan cbResult, 1)}
|
|
err := sess.cb.tryQueue(w)
|
|
sess.cbMu.Unlock()
|
|
if err != nil {
|
|
return nfs4.CompoundRes{}, nil, err
|
|
}
|
|
r := <-w.result
|
|
return r.res, r.bodies, r.err
|
|
}
|
|
|
|
// errNoBackChannel marks a callback that has no channel to travel on.
|
|
func errNoBackChannel(why string) error {
|
|
return fmt.Errorf("nfs4server: %s", why)
|
|
}
|
|
|
|
// lookupSessionPtr resolves a session id to the session pointer.
|
|
func (s *sessionStore) lookupSessionPtr(id nfs4.SessionID) (*session, bool) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
sess, ok := s.sessions[id]
|
|
return sess, ok
|
|
}
|
|
|
|
// lookupSession resolves a session id to the session.
|
|
func (s *sessionStore) lookupSession(id nfs4.SessionID) (*session, uint32) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
sess, ok := s.sessions[id]
|
|
if !ok {
|
|
return nil, nfs4.ErrBadSession
|
|
}
|
|
return sess, nfs4.ErrOK
|
|
}
|
|
|
|
// sessionOfClient resolves the session of the client its callbacks
|
|
// travel on, preferring one with a live back channel, so a recall never
|
|
// fails while another session of the client could carry it.
|
|
func (s *sessionStore) sessionOfClient(clientid uint64) (nfs4.SessionID, bool) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
c, ok := s.byID[clientid]
|
|
if !ok {
|
|
return nfs4.SessionID{}, false
|
|
}
|
|
fallback := nfs4.SessionID{}
|
|
have := false
|
|
for id, sess := range c.sessions {
|
|
if sess.cb != nil {
|
|
return id, true
|
|
}
|
|
fallback, have = id, true
|
|
}
|
|
return fallback, have
|
|
}
|