Files
Felis/internal/fileedit/session_test.go
T

601 lines
21 KiB
Go

package fileedit
import (
"errors"
"io"
"os"
"path/filepath"
"reflect"
"strings"
"testing"
"time"
)
// diskStage is a stage on a disk of total bytes whose free space is what the
// files in it leave of free: statfs sees parts land, as a real disk would.
func diskStage(t *testing.T, free, total uint64, minFree float64) *Stage {
t.Helper()
s := &Stage{Dir: filepath.Join(t.TempDir(), "stage"), MinFree: minFree}
prev := statfs
statfs = func(string) (uint64, uint64, error) {
var used uint64
des, _ := os.ReadDir(s.Dir)
for _, de := range des {
if info, err := de.Info(); err == nil {
used += uint64(info.Size())
}
}
return free - used, total, nil
}
t.Cleanup(func() { statfs = prev })
return s
}
func appendString(s *Stage, user, server, id string, offset int64, part string) (Session, error) {
return s.Append(user, server, id, offset, strings.NewReader(part), int64(len(part)), sumOf(part))
}
func readStaged(t *testing.T, s *Stage, id, token string) string {
t.Helper()
f, size, err := s.Open(id, token)
if err != nil {
t.Fatalf("Open: %v", err)
}
defer f.Close()
b, err := io.ReadAll(f)
if err != nil {
t.Fatal(err)
}
if size != int64(len(b)) {
t.Fatalf("Open said %d bytes and served %d", size, len(b))
}
return string(b)
}
// TestSessionArrivesInParts: parts land in order and are listed with their
// digests, Seal hands the Job a token for exactly those bytes, and they stay
// until that Job reports them landed.
func TestSessionArrivesInParts(t *testing.T) {
s := roomyStage(t)
const whole = "PK\x03\x04 first part, second part"
sess, err := s.Begin("u1", "survival", "plugins/big.jar", int64(len(whole)))
if err != nil {
t.Fatalf("Begin: %v", err)
}
if !hexID.MatchString(sess.ID) || !reflect.DeepEqual(sess, Session{ID: sess.ID, Path: "plugins/big.jar", Size: int64(len(whole))}) {
t.Fatalf("Begin = %+v", sess)
}
got, err := appendString(s, "u1", "survival", sess.ID, 0, whole[:16])
if err != nil || got.Received != 16 || got.Size != int64(len(whole)) {
t.Fatalf("first part: %+v, %v", got, err)
}
first := Part{Size: 16, SHA256: digest([]byte(whole[:16]))}
if at, err := s.Status("u1", "survival", sess.ID); err != nil || at.Received != 16 || at.Path != "plugins/big.jar" ||
!reflect.DeepEqual(at.Parts, []Part{first}) {
t.Fatalf("Status = %+v, %v", at, err)
}
if got, err = appendString(s, "u1", "survival", sess.ID, 16, whole[16:]); err != nil || got.Received != int64(len(whole)) {
t.Fatalf("second part: %+v, %v", got, err)
}
if want := []Part{first, {Size: int64(len(whole) - 16), SHA256: digest([]byte(whole[16:]))}}; !reflect.DeepEqual(got.Parts, want) {
t.Fatalf("parts = %+v, want %+v", got.Parts, want)
}
st, err := s.Seal("u1", "survival", sess.ID)
if err != nil {
t.Fatalf("Seal: %v", err)
}
if st.ID != sess.ID || !hexToken.MatchString(st.Token) || st.Size != int64(len(whole)) || st.SHA256 != digest([]byte(whole)) {
t.Fatalf("Seal = %+v, want the digest of %q", st, whole)
}
if body := readStaged(t, s, st.ID, st.Token); body != whole {
t.Fatalf("served %q, want %q", body, whole)
}
if _, _, err := s.Open(st.ID, st.Token); !errors.Is(err, ErrNotStaged) {
t.Fatalf("second Open with the same token: err = %v, want ErrNotStaged", err)
}
// Served whole, and still here: the Job has yet to check the bytes and put
// them in place, and a Job that fails at either is started again on them.
if at, err := s.Status("u1", "survival", sess.ID); err != nil || at.Received != int64(len(whole)) {
t.Fatalf("Status once served: %+v, %v", at, err)
}
if err := s.Landed(st.ID, strings.Repeat("0", 64)); !errors.Is(err, ErrNotStaged) {
t.Fatalf("Landed with a wrong token: err = %v, want ErrNotStaged", err)
}
if names := stagedNames(t, s); len(names) != 1 {
t.Fatalf("after a wrong token: %v", names)
}
if err := s.Landed(st.ID, st.Token); err != nil {
t.Fatalf("Landed: %v", err)
}
if names := stagedNames(t, s); len(names) != 0 {
t.Fatalf("after Landed: %v", names)
}
if _, err := s.Status("u1", "survival", sess.ID); !errors.Is(err, ErrNotStaged) {
t.Fatalf("Status after Landed: err = %v, want ErrNotStaged", err)
}
if err := s.Landed(st.ID, st.Token); !errors.Is(err, ErrNotStaged) {
t.Fatalf("Landed twice: err = %v, want ErrNotStaged", err)
}
}
// A session answers only the user and the server it was begun for.
func TestSessionAnswersItsOwnerOnly(t *testing.T) {
s := roomyStage(t)
sess, err := s.Begin("u1", "survival", "a.zip", 4)
if err != nil {
t.Fatal(err)
}
for name, who := range map[string][2]string{
"another user": {"u2", "survival"},
"another server": {"u1", "creative"},
} {
if _, err := s.Status(who[0], who[1], sess.ID); !errors.Is(err, ErrNotStaged) {
t.Errorf("%s: Status err = %v", name, err)
}
if _, err := appendString(s, who[0], who[1], sess.ID, 0, "abcd"); !errors.Is(err, ErrNotStaged) {
t.Errorf("%s: Append err = %v", name, err)
}
if _, err := s.Seal(who[0], who[1], sess.ID); !errors.Is(err, ErrNotStaged) {
t.Errorf("%s: Seal err = %v", name, err)
}
if err := s.Drop(who[0], who[1], sess.ID); !errors.Is(err, ErrNotStaged) {
t.Errorf("%s: Drop err = %v", name, err)
}
}
if _, err := s.Status("u1", "survival", strings.Repeat("0", 32)); !errors.Is(err, ErrNotStaged) {
t.Errorf("unknown id: err = %v", err)
}
if at, err := s.Status("u1", "survival", sess.ID); err != nil || at.Received != 0 {
t.Fatalf("the owner's session after the others tried: %+v, %v", at, err)
}
}
func TestSessionRefusesAPartThatDoesNotFit(t *testing.T) {
s := roomyStage(t)
sess, err := s.Begin("u1", "survival", "a.zip", 6)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, "abc"); err != nil {
t.Fatal(err)
}
var off *OffsetError
for _, offset := range []int64{0, 2, 4} {
at, err := appendString(s, "u1", "survival", sess.ID, offset, "d")
if !errors.As(err, &off) || off.Received != 3 || at.Received != 3 {
t.Fatalf("offset %d: %+v, err = %v; want an OffsetError at 3", offset, at, err)
}
}
if _, err := appendString(s, "u1", "survival", sess.ID, 3, "defg"); !errors.Is(err, ErrPartTooLarge) {
t.Fatalf("past the declared size: err = %v, want ErrPartTooLarge", err)
}
if _, err := s.Append("u1", "survival", sess.ID, 3, strings.NewReader(""), -1, sumOf("")); !errors.Is(err, ErrPartTooLarge) {
t.Fatalf("negative length: err = %v, want ErrPartTooLarge", err)
}
if at, err := appendString(s, "u1", "survival", sess.ID, 3, "def"); err != nil || at.Received != 6 {
t.Fatalf("the part that fits exactly: %+v, %v", at, err)
}
big, err := s.Begin("u1", "survival", "b.zip", PartBytes+2)
if err != nil {
t.Fatal(err)
}
if _, err := s.Append("u1", "survival", big.ID, 0, strings.NewReader(""), PartBytes+1, sumOf("")); !errors.Is(err, ErrPartTooLarge) {
t.Fatalf("a part over PartBytes: err = %v, want ErrPartTooLarge", err)
}
}
// A part that breaks, runs long or does not match its digest leaves the session
// as it was: the file is cut back and the digest forgets it, so the resent part
// makes the right file.
func TestSessionRollsBackAFailedPart(t *testing.T) {
for name, tc := range map[string]struct {
body io.Reader
want error // nil: any other failure
}{
"breaks": {io.MultiReader(strings.NewReader("XY"), errReader{io.ErrUnexpectedEOF}), ErrShortUpload},
"ends": {strings.NewReader("XY"), ErrShortUpload},
"runs long": {strings.NewReader("XYZWV"), nil},
// Four bytes, as declared, that are not the four the digest was made of.
"changed on the way": {strings.NewReader("dXfg"), ErrDigestMismatch},
} {
t.Run(name, func(t *testing.T) {
s := roomyStage(t)
sess, err := s.Begin("u1", "survival", "a.zip", 7)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, "abc"); err != nil {
t.Fatal(err)
}
at, err := s.Append("u1", "survival", sess.ID, 3, tc.body, 4, sumOf("defg"))
kind := tc.want
if kind == nil {
kind = ErrShortUpload // must not be it
}
if err == nil || errors.Is(err, kind) != (tc.want != nil) || at.Received != 3 || len(at.Parts) != 1 {
t.Fatalf("%+v, err = %v; want a failure at 3 (%v)", at, err, tc.want)
}
info, err := os.Stat(filepath.Join(s.Dir, stagedNames(t, s)[0]))
if err != nil || info.Size() != 3 {
t.Fatalf("staged file is %v bytes (%v), want it cut back to 3", info.Size(), err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 3, "defg"); err != nil {
t.Fatalf("resent part: %v", err)
}
st, err := s.Seal("u1", "survival", sess.ID)
if err != nil || st.SHA256 != digest([]byte("abcdefg")) {
t.Fatalf("Seal = %+v, %v; want the digest of abcdefg", st, err)
}
if body := readStaged(t, s, st.ID, st.Token); body != "abcdefg" {
t.Fatalf("served %q", body)
}
})
}
}
// While a part is arriving nothing else may touch the session, and it is never
// idle.
func TestSessionIsBusyWhileAPartArrives(t *testing.T) {
s := roomyStage(t)
now := time.Date(2026, 9, 28, 12, 0, 0, 0, time.UTC)
s.Now = func() time.Time { return now }
sess, err := s.Begin("u1", "survival", "a.zip", 4)
if err != nil {
t.Fatal(err)
}
pr, pw := io.Pipe()
done := make(chan error, 1)
go func() {
_, err := s.Append("u1", "survival", sess.ID, 0, pr, 4, sumOf("abcd"))
done <- err
}()
// The write returns once Append is copying, which is after it marked busy.
if _, err := pw.Write([]byte("ab")); err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, "abcd"); !errors.Is(err, ErrUploadBusy) {
t.Errorf("a second part: err = %v, want ErrUploadBusy", err)
}
if _, err := s.Seal("u1", "survival", sess.ID); !errors.Is(err, ErrUploadBusy) {
t.Errorf("Seal: err = %v, want ErrUploadBusy", err)
}
if err := s.Drop("u1", "survival", sess.ID); !errors.Is(err, ErrUploadBusy) {
t.Errorf("Drop: err = %v, want ErrUploadBusy", err)
}
now = now.Add(SessionIdle + time.Hour)
if n := s.Expire(); n != 0 {
t.Errorf("Expire dropped %d sessions with a part arriving", n)
}
if _, err := pw.Write([]byte("cd")); err != nil {
t.Fatal(err)
}
pw.Close()
if err := <-done; err != nil {
t.Fatalf("the part in flight: %v", err)
}
if _, err := s.Seal("u1", "survival", sess.ID); err != nil {
t.Fatalf("Seal once the part is in: %v", err)
}
}
// Each Seal arms one fetch with a fresh token, so a Job that failed can be
// started again on the same bytes.
func TestSessionSealArmsOneFetch(t *testing.T) {
s := roomyStage(t)
sess, err := s.Begin("u1", "survival", "a.zip", 4)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, "abc"); err != nil {
t.Fatal(err)
}
if _, err := s.Seal("u1", "survival", sess.ID); !errors.Is(err, ErrUploadIncomplete) {
t.Fatalf("Seal at 3 of 4: err = %v, want ErrUploadIncomplete", err)
}
if _, _, err := s.Open(sess.ID, ""); !errors.Is(err, ErrNotStaged) {
t.Fatalf("Open before any Seal: err = %v, want ErrNotStaged", err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 3, "d"); err != nil {
t.Fatal(err)
}
first, err := s.Seal("u1", "survival", sess.ID)
if err != nil {
t.Fatal(err)
}
second, err := s.Seal("u1", "survival", sess.ID)
if err != nil {
t.Fatal(err)
}
if first.Token == second.Token || first.SHA256 != second.SHA256 {
t.Fatalf("two Seals: %+v and %+v; want fresh tokens for the same bytes", first, second)
}
if _, _, err := s.Open(sess.ID, first.Token); !errors.Is(err, ErrNotStaged) {
t.Fatalf("the replaced token: err = %v, want ErrNotStaged", err)
}
if _, _, err := s.Open(sess.ID, strings.Repeat("0", 64)); !errors.Is(err, ErrNotStaged) {
t.Fatalf("a wrong token: err = %v, want ErrNotStaged", err)
}
// Neither wrong token spent the armed one.
if body := readStaged(t, s, sess.ID, second.Token); body != "abcd" {
t.Fatalf("served %q", body)
}
if _, _, err := s.Open(sess.ID, second.Token); !errors.Is(err, ErrNotStaged) {
t.Fatalf("the spent token: err = %v, want ErrNotStaged", err)
}
// Opened but not served whole: the Job broke midway, and a new Seal serves
// the same bytes again.
third, err := s.Seal("u1", "survival", sess.ID)
if err != nil {
t.Fatal(err)
}
if body := readStaged(t, s, sess.ID, third.Token); body != "abcd" {
t.Fatalf("served %q after a new Seal", body)
}
}
// Only the token the latest Seal minted reports a session landed: a Job an
// earlier commit started changes nothing, and neither does anyone before the
// first Seal.
func TestLandedTakesTheLatestToken(t *testing.T) {
s := roomyStage(t)
sess, err := s.Begin("u1", "survival", "a.zip", 4)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, "abcd"); err != nil {
t.Fatal(err)
}
if err := s.Landed(sess.ID, ""); !errors.Is(err, ErrNotStaged) {
t.Fatalf("Landed before any Seal: err = %v, want ErrNotStaged", err)
}
first, err := s.Seal("u1", "survival", sess.ID)
if err != nil {
t.Fatal(err)
}
readStaged(t, s, sess.ID, first.Token)
second, err := s.Seal("u1", "survival", sess.ID)
if err != nil {
t.Fatal(err)
}
if err := s.Landed(sess.ID, first.Token); !errors.Is(err, ErrNotStaged) {
t.Fatalf("the replaced token: err = %v, want ErrNotStaged", err)
}
if _, err := s.Status("u1", "survival", sess.ID); err != nil {
t.Fatalf("after the replaced token: %v", err)
}
// Not yet fetched with it, and it still says so: the Job holds the token
// whatever became of its fetch.
if err := s.Landed(sess.ID, second.Token); err != nil {
t.Fatalf("the latest token: %v", err)
}
if names := stagedNames(t, s); len(names) != 0 {
t.Fatalf("after Landed: %v", names)
}
}
// An upload staged by Put answers Landed to its own token and is left to the
// release func Put returned.
func TestLandedLeavesPutAlone(t *testing.T) {
s := roomyStage(t)
st, release, err := s.Put(strings.NewReader("abc"), 3, sumOf("abc"))
if err != nil {
t.Fatal(err)
}
defer release()
if err := s.Landed(st.ID, strings.Repeat("0", 64)); !errors.Is(err, ErrNotStaged) {
t.Fatalf("a wrong token: err = %v, want ErrNotStaged", err)
}
if err := s.Landed(st.ID, st.Token); err != nil {
t.Fatalf("Landed: %v", err)
}
if body := readStaged(t, s, st.ID, st.Token); body != "abc" {
t.Fatalf("served %q", body)
}
}
// A session reserves its whole size when it begins, and gives back what it has
// not yet received when it is dropped or expires.
func TestSessionReservesItsSize(t *testing.T) {
t.Run("begin reserves the whole size", func(t *testing.T) {
s := diskStage(t, 1000, 1200, 0.5) // room for 400
if _, err := s.Begin("u1", "survival", "a.zip", 300); err != nil {
t.Fatal(err)
}
if _, err := s.Begin("u2", "survival", "b.zip", 101); !errors.Is(err, ErrStageFull) {
t.Fatalf("101 bytes beside a 300-byte session: err = %v, want ErrStageFull", err)
}
if _, err := s.Begin("u2", "survival", "b.zip", 100); err != nil {
t.Fatalf("100 bytes beside a 300-byte session: %v", err)
}
})
t.Run("a part moves its room from the reservation to the disk", func(t *testing.T) {
s := diskStage(t, 1000, 1200, 0.5)
sess, err := s.Begin("u1", "survival", "a.zip", 300)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, strings.Repeat("x", 200)); err != nil {
t.Fatal(err)
}
if _, err := s.Begin("u2", "survival", "b.zip", 101); !errors.Is(err, ErrStageFull) {
t.Fatalf("after a part landed: err = %v, want ErrStageFull", err)
}
if _, err := s.Begin("u2", "survival", "b.zip", 100); err != nil {
t.Fatalf("after a part landed: %v", err)
}
})
t.Run("drop gives it all back", func(t *testing.T) {
s := diskStage(t, 1000, 1200, 0.5)
sess, err := s.Begin("u1", "survival", "a.zip", 300)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, strings.Repeat("x", 200)); err != nil {
t.Fatal(err)
}
if err := s.Drop("u1", "survival", sess.ID); err != nil {
t.Fatal(err)
}
if names := stagedNames(t, s); len(names) != 0 {
t.Fatalf("after Drop: %v", names)
}
if _, err := s.Begin("u2", "survival", "b.zip", 400); err != nil {
t.Fatalf("after Drop: %v", err)
}
if _, err := s.Status("u1", "survival", sess.ID); !errors.Is(err, ErrNotStaged) {
t.Fatalf("Status after Drop: err = %v", err)
}
})
t.Run("expiry gives it all back", func(t *testing.T) {
s := diskStage(t, 1000, 1200, 0.5)
now := time.Date(2026, 9, 28, 12, 0, 0, 0, time.UTC)
s.Now = func() time.Time { return now }
sess, err := s.Begin("u1", "survival", "a.zip", 300)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, strings.Repeat("x", 200)); err != nil {
t.Fatal(err)
}
now = now.Add(SessionIdle + time.Nanosecond)
if n := s.Expire(); n != 1 {
t.Fatalf("Expire dropped %d, want 1", n)
}
if _, err := s.Begin("u2", "survival", "b.zip", 400); err != nil {
t.Fatalf("after Expire: %v", err)
}
})
}
func TestSessionsPerUserAreBounded(t *testing.T) {
s := diskStage(t, 1000, 1200, 0.5)
var ids []string
for i := range MaxSessionsPerUser {
sess, err := s.Begin("u1", "survival", "a.zip", 10)
if err != nil {
t.Fatalf("session %d: %v", i+1, err)
}
ids = append(ids, sess.ID)
}
if _, err := s.Begin("u1", "creative", "a.zip", 10); !errors.Is(err, ErrTooManySessions) {
t.Fatalf("one more on another server: err = %v, want ErrTooManySessions", err)
}
// The refused Begin kept neither a file nor its reservation: room for 400
// less the four sessions' 40.
if names := stagedNames(t, s); len(names) != MaxSessionsPerUser {
t.Fatalf("staged %d files, want %d", len(names), MaxSessionsPerUser)
}
if _, err := s.Begin("u2", "survival", "b.zip", 360); err != nil {
t.Fatalf("another user: %v", err)
}
if err := s.Drop("u1", "survival", ids[0]); err != nil {
t.Fatal(err)
}
if _, err := s.Begin("u1", "survival", "a.zip", 10); err != nil {
t.Fatalf("after dropping one: %v", err)
}
// With room reserved, a negative size would wrap the reservation sum round
// to a small number and pass the room check.
if _, err := s.Begin("u3", "survival", "a.zip", -1); err == nil || errors.Is(err, ErrStageFull) {
t.Fatalf("a negative size: err = %v, want it refused for being negative", err)
}
}
// Expire drops what has sat untouched for longer than SessionIdle, sealed or
// not, and a part keeps a session alive.
func TestSessionExpiry(t *testing.T) {
s := roomyStage(t)
start := time.Date(2026, 9, 28, 12, 0, 0, 0, time.UTC)
now := start
s.Now = func() time.Time { return now }
idle, err := s.Begin("u1", "survival", "idle.zip", 2)
if err != nil {
t.Fatal(err)
}
sealed, err := s.Begin("u1", "survival", "sealed.zip", 1)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sealed.ID, 0, "x"); err != nil {
t.Fatal(err)
}
if _, err := s.Seal("u1", "survival", sealed.ID); err != nil {
t.Fatal(err)
}
active, err := s.Begin("u1", "survival", "active.zip", 2)
if err != nil {
t.Fatal(err)
}
now = start.Add(time.Hour)
if _, err := appendString(s, "u1", "survival", active.ID, 0, "a"); err != nil {
t.Fatal(err)
}
now = start.Add(SessionIdle)
if n := s.Expire(); n != 0 {
t.Fatalf("at exactly SessionIdle Expire dropped %d", n)
}
now = start.Add(SessionIdle + time.Minute)
if n := s.Expire(); n != 2 {
t.Fatalf("Expire dropped %d, want the idle and the sealed one", n)
}
for _, id := range []string{idle.ID, sealed.ID} {
if _, err := s.Status("u1", "survival", id); !errors.Is(err, ErrNotStaged) {
t.Errorf("%s survived Expire: %v", id, err)
}
}
if at, err := s.Status("u1", "survival", active.ID); err != nil || at.Received != 1 {
t.Fatalf("the session a part touched: %+v, %v", at, err)
}
if names := stagedNames(t, s); len(names) != 1 {
t.Fatalf("files left: %v, want the active session's", names)
}
// A Seal, and the Job's Open, each start the idle clock again: a Job begun
// on an upload that sat for hours still finds it there.
for _, tc := range []struct {
name string
// sealAt and openAt are how long after the last part the Seal and the
// Job's Open come; a negative openAt is no Open.
sealAt, openAt time.Duration
}{
{"seal", 4 * time.Hour, -1},
{"open", 0, 4 * time.Hour},
} {
begun := now
sess, err := s.Begin("u1", "survival", tc.name+".zip", 1)
if err != nil {
t.Fatal(err)
}
if _, err := appendString(s, "u1", "survival", sess.ID, 0, "x"); err != nil {
t.Fatal(err)
}
now = begun.Add(tc.sealAt)
st, err := s.Seal("u1", "survival", sess.ID)
if err != nil {
t.Fatal(err)
}
if tc.openAt >= 0 {
now = begun.Add(tc.openAt)
readStaged(t, s, sess.ID, st.Token)
}
now = begun.Add(SessionIdle + time.Minute)
s.Expire()
if _, err := s.Status("u1", "survival", sess.ID); err != nil {
t.Errorf("%s: a session touched 4h after its last part expired 6h after it: %v", tc.name, err)
}
if err := s.Drop("u1", "survival", sess.ID); err != nil {
t.Fatal(err)
}
}
}