601 lines
21 KiB
Go
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)
|
|
}
|
|
}
|
|
}
|