feat(files): 大文件分片上传、停服解压 zip 先列冲突再覆盖、文件和文件夹可下载;导出和下载不再带出 RCON 密码与转发密钥
This commit is contained in:
76 files changed
+10542
-569
No files matched your search
+15
-2
@@ -582,8 +582,8 @@ func (a *API) externalAPIRoutes() []apiRoute {
|
||||
{Method: "POST", Pattern: "/api/v1/servers/{name}/world/export", h: a.handleExportWorld},
|
||||
{Method: "GET", Pattern: "/api/v1/exports/{ticket}", h: a.handleExportStatus},
|
||||
{Method: "GET", Pattern: "/api/v1/exports/{ticket}/download", h: a.handleExportDownload},
|
||||
// Server file manager: list, read, write, make a folder, delete, rename and
|
||||
// upload in a STOPPED server's world volume (handlers_files.go). App-tier,
|
||||
// Server file manager: list, read, write, make a folder, delete, rename,
|
||||
// upload, unzip and download in a STOPPED server's world volume (handlers_files.go). App-tier,
|
||||
// exactly like the backup pair above and for the same reason — every route
|
||||
// gates on owner-or-admin inside the handler, so an owner repairs their own
|
||||
// broken server without an admin's Zero-Trust path. The path travels as ?path=
|
||||
@@ -605,6 +605,19 @@ func (a *API) externalAPIRoutes() []apiRoute {
|
||||
{Method: "POST", Pattern: "/api/v1/servers/{name}/files/mkdir", h: a.handleMkdir},
|
||||
{Method: "POST", Pattern: "/api/v1/servers/{name}/files/rename", h: a.handleRenameFile},
|
||||
{Method: "PUT", Pattern: "/api/v1/servers/{name}/files/upload", h: a.handleUploadFile},
|
||||
// A file too big for one request goes up in parts as an upload session, and
|
||||
// lands, like an unzip, as a Job the request does not wait on; files/ops
|
||||
// reports how those went (handlers_fileops.go).
|
||||
{Method: "POST", Pattern: "/api/v1/servers/{name}/files/uploads", h: a.handleBeginFileUpload},
|
||||
{Method: "GET", Pattern: "/api/v1/servers/{name}/files/uploads/{id}", h: a.handleFileUploadStatus},
|
||||
{Method: "PUT", Pattern: "/api/v1/servers/{name}/files/uploads/{id}", h: a.handleFileUploadPart},
|
||||
{Method: "DELETE", Pattern: "/api/v1/servers/{name}/files/uploads/{id}", h: a.handleDropFileUpload},
|
||||
{Method: "POST", Pattern: "/api/v1/servers/{name}/files/uploads/{id}/commit", h: a.handleCommitFileUpload},
|
||||
{Method: "POST", Pattern: "/api/v1/servers/{name}/files/unzip", h: a.handleUnzipFile},
|
||||
{Method: "GET", Pattern: "/api/v1/servers/{name}/files/ops", h: a.handleListFileOps},
|
||||
// A file or folder download is an export (exports.go): it answers 202
|
||||
// with a ticket the export routes above serve.
|
||||
{Method: "POST", Pattern: "/api/v1/servers/{name}/files/download", h: a.handleDownloadFile},
|
||||
// Scheduled tasks (handlers_schedules.go): a console command, restart, stop,
|
||||
// start or backup at set times, which felis-api's runner fires. App-tier and
|
||||
// owner-or-admin inside the handler, like the console and power routes they
|
||||
|
||||
+185
-100
@@ -8,30 +8,33 @@ import (
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"hash"
|
||||
"io"
|
||||
"log"
|
||||
"mime"
|
||||
"net/http"
|
||||
pathpkg "path"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||
"felis.lolicon.best/internal/fileedit"
|
||||
"felis.lolicon.best/internal/maintenance"
|
||||
"felis.lolicon.best/internal/naming"
|
||||
"felis.lolicon.best/internal/worldexport"
|
||||
)
|
||||
|
||||
// World export. The owner downloads a tar.gz of their world, either as it is
|
||||
// now (the server stopped) or as one of its backups, straight into the browser:
|
||||
// World export and file download. The owner downloads a tar.gz of their world,
|
||||
// either as it is now (the server stopped) or as one of its backups, or one file
|
||||
// or folder of a stopped world (a folder as a zip), straight into the browser:
|
||||
//
|
||||
// 1. POST /servers/{name}/world/export or /servers/{name}/backups/{id}/export
|
||||
// checks the caller and the server, admits the export against the limits
|
||||
// below and starts a one-shot felis-export Job (internal/worldexport) that
|
||||
// reads the world or the archive read-only. It answers 202 with a ticket:
|
||||
// 256 random bits, good for the caller who started it and nobody else.
|
||||
// 1. POST /servers/{name}/world/export, /servers/{name}/backups/{id}/export or
|
||||
// /servers/{name}/files/download checks the caller and the server, admits
|
||||
// the export against the limits below and starts a one-shot felis-export
|
||||
// Job (internal/worldexport) that reads the world or the archive read-only.
|
||||
// It answers 202 with a ticket: 256 random bits, good for the caller who
|
||||
// started it and nobody else.
|
||||
// 2. The Job PUTs the archive to the internal face (PUT
|
||||
// /api/v1/internal/exports/{id} with the one-time token it was started
|
||||
// with), and that request waits, body unread, for the browser.
|
||||
@@ -41,10 +44,10 @@ import (
|
||||
// archive passes through felis-api's memory once and never touches a disk
|
||||
// it owns, and the Job moves at the browser's pace.
|
||||
//
|
||||
// A backup is checked against the sha256 recorded when it was written as it
|
||||
// streams, and the last read is held back until the digest is known: a mismatch
|
||||
// aborts the response, so the browser reports a failed download and never keeps
|
||||
// a complete-looking corrupt file, and the Job is told backup_corrupt.
|
||||
// A backup is checked by the Job against the sha256 recorded when it was
|
||||
// written (cmd/felis export): on a mismatch the Job aborts its upload short of
|
||||
// the archive's end, the download aborts with it, and the browser reports a
|
||||
// failed download and never keeps a complete-looking corrupt file.
|
||||
//
|
||||
// Tickets live in felis-api's memory. A restart forgets them, and a Job that
|
||||
// then PUTs finds nothing and fails, which the jobs list shows.
|
||||
@@ -57,15 +60,57 @@ type Exporter interface {
|
||||
}
|
||||
|
||||
// Limits on exports. Each one keeps a Job, a connection and a 64 KiB copy
|
||||
// buffer alive for as long as a download takes, and a world export also keeps
|
||||
// its server from starting.
|
||||
// buffer alive for as long as a download takes, and a world export or a file
|
||||
// download also keeps its server from starting.
|
||||
const (
|
||||
exportMaxActive = 2 // admitted and not yet over, install-wide
|
||||
exportMaxPerUser = 1
|
||||
exportPerHour = 6 // started by one user in any hour
|
||||
exportCopyBuffer = 32 << 10
|
||||
|
||||
// File downloads count apart from the world and backup exports, with
|
||||
// their own limits: one is a file or a folder, usually small and over in
|
||||
// seconds, and an owner fetching a few configs one after another must
|
||||
// neither wait on an export nor hold one off.
|
||||
fileExportMaxActive = 4
|
||||
fileExportMaxPerUser = 2
|
||||
fileExportPerHour = 30
|
||||
)
|
||||
|
||||
// exportClass is one set of export limits and how a refusal names them.
|
||||
type exportClass struct {
|
||||
maxActive, perUser, perHour int
|
||||
busyUser, busyActive, busyHour string // each formats its limit
|
||||
}
|
||||
|
||||
var (
|
||||
worldExports = exportClass{
|
||||
maxActive: exportMaxActive, perUser: exportMaxPerUser, perHour: exportPerHour,
|
||||
busyUser: "you already have %d export in progress; download it or let it expire first",
|
||||
busyActive: "%d exports are already in progress; retry in a few minutes",
|
||||
busyHour: "you have started %d exports in the last hour; retry later",
|
||||
}
|
||||
fileExports = exportClass{
|
||||
maxActive: fileExportMaxActive, perUser: fileExportMaxPerUser, perHour: fileExportPerHour,
|
||||
busyUser: "you already have %d file downloads in progress; let one finish first",
|
||||
busyActive: "%d file downloads are already in progress; retry in a minute",
|
||||
busyHour: "you have started %d file downloads in the last hour; retry later",
|
||||
}
|
||||
)
|
||||
|
||||
func (e *exportEntry) class() exportClass {
|
||||
if e.files {
|
||||
return fileExports
|
||||
}
|
||||
return worldExports
|
||||
}
|
||||
|
||||
// exportStarts keys a user's recent starts within one class.
|
||||
type exportStarts struct {
|
||||
userID string
|
||||
files bool
|
||||
}
|
||||
|
||||
// Timings. Vars only so a test can shrink them.
|
||||
var (
|
||||
// exportClaimTTL is how long the Job's upload waits for the browser.
|
||||
@@ -102,8 +147,6 @@ type exportStatusView struct {
|
||||
Message string `json:"message,omitempty"`
|
||||
}
|
||||
|
||||
var errExportDigest = errors.New("the archive does not match the sha256 recorded when it was written")
|
||||
|
||||
func errExportExpired() error {
|
||||
return newError(http.StatusGone, "export_expired", "this export has expired or was already downloaded; start a new one")
|
||||
}
|
||||
@@ -119,13 +162,16 @@ type exportEntry struct {
|
||||
userID string
|
||||
server string
|
||||
mode string
|
||||
sha256 string // what a backup must hash to; empty for a world
|
||||
files bool // a file download, counted in fileExports
|
||||
filename string
|
||||
job string
|
||||
state string
|
||||
message string
|
||||
at time.Time // when it entered its state
|
||||
upload *exportUpload
|
||||
// contentType is what the download is served as: a gzip archive, a zip,
|
||||
// or a single file's raw bytes.
|
||||
contentType string
|
||||
job string
|
||||
state string
|
||||
message string
|
||||
at time.Time // when it entered its state
|
||||
upload *exportUpload
|
||||
}
|
||||
|
||||
// exportUpload is the Job's PUT, parked until a browser claims it.
|
||||
@@ -144,7 +190,7 @@ type exportRegistry struct {
|
||||
mu sync.Mutex
|
||||
byTicket map[string]*exportEntry
|
||||
byID map[string]*exportEntry
|
||||
starts map[string][]time.Time // per user, oldest first
|
||||
starts map[exportStarts][]time.Time // oldest first
|
||||
}
|
||||
|
||||
func (a *API) exportTickets() *exportRegistry {
|
||||
@@ -152,7 +198,7 @@ func (a *API) exportTickets() *exportRegistry {
|
||||
a.exports = &exportRegistry{
|
||||
byTicket: map[string]*exportEntry{},
|
||||
byID: map[string]*exportEntry{},
|
||||
starts: map[string][]time.Time{},
|
||||
starts: map[exportStarts][]time.Time{},
|
||||
}
|
||||
})
|
||||
return a.exports
|
||||
@@ -187,38 +233,37 @@ func randomHex(n int) string {
|
||||
return hex.EncodeToString(b)
|
||||
}
|
||||
|
||||
// admit reserves the export e describes (its user, server, mode, filename and
|
||||
// digest) and returns its upload token, or refuses it with export_busy.
|
||||
// admit reserves the export e describes (its user, server, mode, class,
|
||||
// filename and content type) and returns its upload token, or refuses it with
|
||||
// export_busy. Only exports of e's class count against it.
|
||||
func (g *exportRegistry) admit(e *exportEntry, now time.Time) (string, error) {
|
||||
g.mu.Lock()
|
||||
defer g.mu.Unlock()
|
||||
g.sweepLocked(now)
|
||||
active, mine := 0, 0
|
||||
for _, o := range g.byTicket {
|
||||
if o.active() {
|
||||
if o.active() && o.files == e.files {
|
||||
active++
|
||||
if o.userID == e.userID {
|
||||
mine++
|
||||
}
|
||||
}
|
||||
}
|
||||
switch starts := g.starts[e.userID]; {
|
||||
case mine >= exportMaxPerUser:
|
||||
return "", newError(http.StatusTooManyRequests, "export_busy",
|
||||
"you already have an export in progress; download it or let it expire first").retryAfter(exportClaimTTL)
|
||||
case active >= exportMaxActive:
|
||||
return "", newError(http.StatusTooManyRequests, "export_busy",
|
||||
"%d exports are already in progress; retry in a few minutes", exportMaxActive).retryAfter(time.Minute)
|
||||
case len(starts) >= exportPerHour:
|
||||
return "", newError(http.StatusTooManyRequests, "export_busy",
|
||||
"you have started %d exports in the last hour; retry later", exportPerHour).retryAfter(starts[0].Add(time.Hour).Sub(now))
|
||||
c, key := e.class(), exportStarts{e.userID, e.files}
|
||||
switch starts := g.starts[key]; {
|
||||
case mine >= c.perUser:
|
||||
return "", newError(http.StatusTooManyRequests, "export_busy", c.busyUser, c.perUser).retryAfter(exportClaimTTL)
|
||||
case active >= c.maxActive:
|
||||
return "", newError(http.StatusTooManyRequests, "export_busy", c.busyActive, c.maxActive).retryAfter(time.Minute)
|
||||
case len(starts) >= c.perHour:
|
||||
return "", newError(http.StatusTooManyRequests, "export_busy", c.busyHour, c.perHour).retryAfter(starts[0].Add(time.Hour).Sub(now))
|
||||
}
|
||||
token := randomHex(32)
|
||||
e.ticket, e.id, e.tokenHash = randomHex(32), randomHex(8), sha256.Sum256([]byte(token))
|
||||
e.state, e.at = exportPending, now
|
||||
g.byTicket[e.ticket] = e
|
||||
g.byID[e.id] = e
|
||||
g.starts[e.userID] = append(g.starts[e.userID], now)
|
||||
g.starts[key] = append(g.starts[key], now)
|
||||
return token, nil
|
||||
}
|
||||
|
||||
@@ -228,10 +273,15 @@ func (g *exportRegistry) drop(e *exportEntry) {
|
||||
defer g.mu.Unlock()
|
||||
delete(g.byTicket, e.ticket)
|
||||
delete(g.byID, e.id)
|
||||
ts := g.starts[e.userID]
|
||||
key := exportStarts{e.userID, e.files}
|
||||
ts := g.starts[key]
|
||||
for i := len(ts) - 1; i >= 0; i-- {
|
||||
if ts[i].Equal(e.at) {
|
||||
g.starts[e.userID] = append(ts[:i:i], ts[i+1:]...)
|
||||
if rest := append(ts[:i:i], ts[i+1:]...); len(rest) > 0 {
|
||||
g.starts[key] = rest
|
||||
} else {
|
||||
delete(g.starts, key)
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
@@ -377,13 +427,13 @@ func (a *API) handleExportBackup(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
e := &exportEntry{userID: p.UserID, server: name, mode: worldexport.ModeBackup,
|
||||
filename: fmt.Sprintf("%s-backup-%s.tar.gz", name, backup.ID), sha256: backup.SHA256}
|
||||
filename: fmt.Sprintf("%s-backup-%s.tar.gz", name, backup.ID), contentType: archiveContentType}
|
||||
token, err := a.exportTickets().admit(e, a.now())
|
||||
if err != nil {
|
||||
writeError(w, r, err)
|
||||
return
|
||||
}
|
||||
if !a.startExport(w, r, e, token, backup.BackupRef) {
|
||||
if !a.startExport(w, r, e, token, worldexport.Request{BackupRef: backup.BackupRef, BackupSHA256: backup.SHA256}) {
|
||||
return
|
||||
}
|
||||
a.auditEntry(r, AuditEntry{Actor: auditActor(p), ActorUserID: p.UserID, Action: "backup.export", ServerName: rec.Name,
|
||||
@@ -425,7 +475,8 @@ func (a *API) handleExportWorld(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
reg := a.exportTickets()
|
||||
e := &exportEntry{userID: p.UserID, server: name, mode: worldexport.ModeWorld,
|
||||
filename: fmt.Sprintf("%s-world-%s.tar.gz", name, a.now().UTC().Format("20060102-150405"))}
|
||||
filename: fmt.Sprintf("%s-world-%s.tar.gz", name, a.now().UTC().Format("20060102-150405")),
|
||||
contentType: archiveContentType}
|
||||
token, err := reg.admit(e, a.now())
|
||||
if err != nil {
|
||||
writeError(w, r, err)
|
||||
@@ -437,15 +488,79 @@ func (a *API) handleExportWorld(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
defer release()
|
||||
if !a.startExport(w, r, e, token, "") {
|
||||
if !a.startExport(w, r, e, token, worldexport.Request{}) {
|
||||
return
|
||||
}
|
||||
a.auditEntry(r, AuditEntry{Actor: auditActor(p), ActorUserID: p.UserID, Action: "world.export", ServerName: rec.Name})
|
||||
writeJSON(w, http.StatusAccepted, exportTicketView{Ticket: e.ticket, State: exportPending, Filename: e.filename})
|
||||
}
|
||||
|
||||
// archiveContentType is how a world or a backup export is served.
|
||||
const archiveContentType = "application/gzip"
|
||||
|
||||
// handleDownloadFile starts the download of one file or folder of a stopped
|
||||
// server's world (POST /api/v1/servers/{name}/files/download?path=…&dir=true
|
||||
// for a folder). The gate is the file manager's (authorizeFileOp), plus an
|
||||
// account to bind the ticket to. The download is an export: a felis-export Job
|
||||
// in files mode reads the file, or zips the folder, from the world volume and
|
||||
// hands it over through a ticket like any other, under the world-volume lock
|
||||
// (as KindExport) so the server cannot start mid-zip.
|
||||
//
|
||||
// The path is passed to the Job as it came, as every file route does (see
|
||||
// handleListFiles): the Job's os.Root is the containment, and the guards run
|
||||
// there. Only the world root is refused here, since a whole world is what the
|
||||
// world export is for.
|
||||
func (a *API) handleDownloadFile(w http.ResponseWriter, r *http.Request) {
|
||||
name, ok := a.authorizeFileOp(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
p := principalFromContext(r.Context())
|
||||
if p.UserID == "" {
|
||||
writeError(w, r, errForbidden)
|
||||
return
|
||||
}
|
||||
path, ok := requirePath(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
dir := r.URL.Query().Get("dir") == "true"
|
||||
base := pathpkg.Base(pathpkg.Clean("/" + path))
|
||||
if base == "/" {
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_path",
|
||||
"the whole world is not a file download; export the world from the backups page instead"))
|
||||
return
|
||||
}
|
||||
if a.Exporter == nil || a.InternalBaseURL == "" {
|
||||
writeError(w, r, errExportUnavailable())
|
||||
return
|
||||
}
|
||||
e := &exportEntry{userID: p.UserID, server: name, mode: worldexport.ModeFiles, files: true,
|
||||
filename: base, contentType: fileedit.DownloadFileType}
|
||||
if dir {
|
||||
e.filename, e.contentType = base+".zip", fileedit.DownloadZipType
|
||||
}
|
||||
reg := a.exportTickets()
|
||||
token, err := reg.admit(e, a.now())
|
||||
if err != nil {
|
||||
writeError(w, r, err)
|
||||
return
|
||||
}
|
||||
release, ok := a.acquireWorld(w, r, name, maintenance.KindExport, "stop the server before editing its files")
|
||||
if !ok {
|
||||
reg.drop(e)
|
||||
return
|
||||
}
|
||||
defer release()
|
||||
if !a.startExport(w, r, e, token, worldexport.Request{Path: path, Dir: dir}) {
|
||||
return
|
||||
}
|
||||
a.auditFile(r, "file.download", name, path, map[string]any{"dir": dir})
|
||||
writeJSON(w, http.StatusAccepted, exportTicketView{Ticket: e.ticket, State: exportPending, Filename: e.filename})
|
||||
}
|
||||
|
||||
func errExportUnavailable() error {
|
||||
return newError(http.StatusServiceUnavailable, "export_unavailable", "world export is not configured")
|
||||
return newError(http.StatusServiceUnavailable, "export_unavailable", "world export and downloads are not configured")
|
||||
}
|
||||
|
||||
// exportGate is the front half both export routes share: a valid name, a known
|
||||
@@ -471,12 +586,12 @@ func (a *API) exportGate(w http.ResponseWriter, r *http.Request) (string, *Serve
|
||||
}
|
||||
|
||||
// startExport creates the admitted export's Job, or forgets the export and
|
||||
// writes the error.
|
||||
func (a *API) startExport(w http.ResponseWriter, r *http.Request, e *exportEntry, token, backupRef string) bool {
|
||||
job, err := a.Exporter.Start(r.Context(), worldexport.Request{
|
||||
Server: e.server, Mode: e.mode, BackupRef: backupRef, ID: e.id, Token: token,
|
||||
TargetURL: a.InternalBaseURL + "/api/v1/internal/exports/" + e.id,
|
||||
})
|
||||
// writes the error. what carries the mode's own fields (the backup and its
|
||||
// digest, or the path); the rest comes from e.
|
||||
func (a *API) startExport(w http.ResponseWriter, r *http.Request, e *exportEntry, token string, what worldexport.Request) bool {
|
||||
what.Server, what.Mode, what.ID, what.Token = e.server, e.mode, e.id, token
|
||||
what.TargetURL = a.InternalBaseURL + "/api/v1/internal/exports/" + e.id
|
||||
job, err := a.Exporter.Start(r.Context(), what)
|
||||
if err != nil {
|
||||
a.exportTickets().drop(e)
|
||||
writeError(w, r, err)
|
||||
@@ -533,7 +648,7 @@ func (a *API) handleExportDownload(w http.ResponseWriter, r *http.Request) {
|
||||
defer reg.finish(e.ticket, a.now())
|
||||
|
||||
h := w.Header()
|
||||
h.Set("Content-Type", "application/gzip")
|
||||
h.Set("Content-Type", e.contentType)
|
||||
if cd := mime.FormatMediaType("attachment", map[string]string{"filename": e.filename}); cd != "" {
|
||||
h.Set("Content-Disposition", cd)
|
||||
} else {
|
||||
@@ -546,7 +661,7 @@ func (a *API) handleExportDownload(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
w.WriteHeader(http.StatusOK)
|
||||
|
||||
err = copyExport(w, e.upload.body, e.sha256)
|
||||
err = copyExport(w, e.upload.body)
|
||||
e.upload.done <- err
|
||||
if err != nil {
|
||||
log.Printf("api: export %s of %s ended early: %v", e.id, e.server, err)
|
||||
@@ -557,53 +672,27 @@ func (a *API) handleExportDownload(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
// copyExport copies body into w through a fixed 32 KiB buffer, restarting w's
|
||||
// write deadline on every write. With want set, the last read is held back
|
||||
// until the whole body hashes to it.
|
||||
func copyExport(w http.ResponseWriter, body io.Reader, want string) error {
|
||||
out := &heldWriter{w: w, rc: http.NewResponseController(w)}
|
||||
var sum hash.Hash
|
||||
if want != "" {
|
||||
sum = sha256.New()
|
||||
body = io.TeeReader(body, sum)
|
||||
}
|
||||
// write deadline on every write.
|
||||
func copyExport(w http.ResponseWriter, body io.Reader) error {
|
||||
out := &stallWriter{w: w, rc: http.NewResponseController(w)}
|
||||
if _, err := io.CopyBuffer(out, body, make([]byte, exportCopyBuffer)); err != nil {
|
||||
return err
|
||||
}
|
||||
if sum != nil && !strings.EqualFold(hex.EncodeToString(sum.Sum(nil)), want) {
|
||||
return errExportDigest
|
||||
}
|
||||
if err := out.flush(); err != nil {
|
||||
return err
|
||||
}
|
||||
_ = out.rc.SetWriteDeadline(time.Time{}) // the connection may serve another request
|
||||
return nil
|
||||
}
|
||||
|
||||
// heldWriter passes each write on one behind, keeping the latest back until
|
||||
// flush, so the end of an archive reaches the browser only once it is checked.
|
||||
// It has no ReadFrom, so io.CopyBuffer uses the buffer it is given.
|
||||
type heldWriter struct {
|
||||
w io.Writer
|
||||
rc *http.ResponseController
|
||||
held []byte
|
||||
// stallWriter restarts the connection's write deadline before every write, so
|
||||
// a write fails only once the browser has taken nothing for exportStall. It
|
||||
// has no ReadFrom, so io.CopyBuffer uses the buffer it is given.
|
||||
type stallWriter struct {
|
||||
w io.Writer
|
||||
rc *http.ResponseController
|
||||
}
|
||||
|
||||
func (h *heldWriter) Write(p []byte) (int, error) {
|
||||
if err := h.flush(); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
h.held = append(h.held[:0], p...)
|
||||
return len(p), nil
|
||||
}
|
||||
|
||||
func (h *heldWriter) flush() error {
|
||||
if len(h.held) == 0 {
|
||||
return nil
|
||||
}
|
||||
_ = h.rc.SetWriteDeadline(time.Now().Add(exportStall))
|
||||
_, err := h.w.Write(h.held)
|
||||
h.held = h.held[:0]
|
||||
return err
|
||||
func (s *stallWriter) Write(p []byte) (int, error) {
|
||||
_ = s.rc.SetWriteDeadline(time.Now().Add(exportStall))
|
||||
return s.w.Write(p)
|
||||
}
|
||||
|
||||
// handleInternalExportUpload takes an export Job's archive (PUT
|
||||
@@ -611,9 +700,8 @@ func (h *heldWriter) flush() error {
|
||||
// export is the check, as for file uploads: an unknown id, a wrong token and a
|
||||
// token already used all read the same 404. The request then waits for the
|
||||
// browser, up to exportClaimTTL, and answers once the download has ended: 204
|
||||
// when it got the whole archive, 409 backup_corrupt when the archive did not
|
||||
// match its recorded digest, 410 export_expired when nobody came for it or the
|
||||
// browser left early.
|
||||
// when it got the whole archive, 410 export_expired when nobody came for it,
|
||||
// the browser left early or the Job itself cut the upload short.
|
||||
func (a *API) handleInternalExportUpload(w http.ResponseWriter, r *http.Request) {
|
||||
token, ok := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer ")
|
||||
if !ok {
|
||||
@@ -644,14 +732,11 @@ func (a *API) handleInternalExportUpload(w http.ResponseWriter, r *http.Request)
|
||||
return
|
||||
}
|
||||
}
|
||||
switch err := <-up.done; {
|
||||
case err == nil:
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
case errors.Is(err, errExportDigest):
|
||||
writeError(w, r, newError(http.StatusConflict, "backup_corrupt", "%v", err))
|
||||
default:
|
||||
if err := <-up.done; err != nil {
|
||||
writeError(w, r, newError(http.StatusGone, "export_expired", "the download ended before the archive did: %v", err))
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
|
||||
// stallBody restarts the connection's read deadline on every read, so reading
|
||||
|
||||
+295
-70
@@ -4,13 +4,12 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"regexp"
|
||||
"strconv"
|
||||
"strings"
|
||||
@@ -19,6 +18,7 @@ import (
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||
"felis.lolicon.best/internal/fileedit"
|
||||
"felis.lolicon.best/internal/maintenance"
|
||||
"felis.lolicon.best/internal/worldexport"
|
||||
batchv1 "k8s.io/api/batch/v1"
|
||||
@@ -189,6 +189,7 @@ func randomBytes(n int) []byte {
|
||||
func TestExportBackupGate(t *testing.T) {
|
||||
t.Run("former owner starts a backup export", func(t *testing.T) {
|
||||
a, repo, cl, ex := exportFixture()
|
||||
repo.backups[0].sha256 = strings.Repeat("cd", 32)
|
||||
v := beginExport(t, a.ExternalHandler(), backupPath, "owner1")
|
||||
if !hex64.MatchString(v.Ticket) || v.State != "pending" || v.Filename != "survival-backup-bk1.tar.gz" {
|
||||
t.Fatalf("ticket = %+v", v)
|
||||
@@ -198,6 +199,7 @@ func TestExportBackupGate(t *testing.T) {
|
||||
}
|
||||
r := ex.reqs[0]
|
||||
if r.Server != "survival" || r.Mode != worldexport.ModeBackup || r.BackupRef != "/backups/survival-bk1.tar.gz" ||
|
||||
r.BackupSHA256 != strings.Repeat("cd", 32) || r.Path != "" || r.Dir ||
|
||||
!hex16.MatchString(r.ID) || !hex64.MatchString(r.Token) || r.Token == v.Ticket ||
|
||||
r.TargetURL != exportBase+"/api/v1/internal/exports/"+r.ID {
|
||||
t.Fatalf("export request = %+v", r)
|
||||
@@ -271,7 +273,8 @@ func TestExportWorldGate(t *testing.T) {
|
||||
if v.Filename != "survival-world-"+exportStamp+".tar.gz" || v.State != "pending" {
|
||||
t.Fatalf("ticket = %+v", v)
|
||||
}
|
||||
if len(ex.reqs) != 1 || ex.reqs[0].Mode != worldexport.ModeWorld || ex.reqs[0].BackupRef != "" || ex.reqs[0].Server != "survival" {
|
||||
if len(ex.reqs) != 1 || ex.reqs[0].Mode != worldexport.ModeWorld || ex.reqs[0].BackupRef != "" || ex.reqs[0].BackupSHA256 != "" ||
|
||||
ex.reqs[0].Path != "" || ex.reqs[0].Server != "survival" {
|
||||
t.Fatalf("export requests = %+v", ex.reqs)
|
||||
}
|
||||
if strings.Join(cl.acquired, ",") != "survival:"+maintenance.KindExport || strings.Join(cl.released, ",") != "survival" {
|
||||
@@ -363,7 +366,7 @@ func TestExportWorldGate(t *testing.T) {
|
||||
}
|
||||
ex.err = nil
|
||||
beginExport(t, h, worldPath, "owner1")
|
||||
if n := len(a.exportTickets().starts["owner1"]); n != 1 {
|
||||
if n := len(a.exportTickets().starts[exportStarts{userID: "owner1"}]); n != 1 {
|
||||
t.Fatalf("hourly starts = %d, want only the export that got a Job", n)
|
||||
}
|
||||
})
|
||||
@@ -681,82 +684,299 @@ func awaitResponse(t *testing.T, ch <-chan *http.Response) *http.Response {
|
||||
}
|
||||
}
|
||||
|
||||
// TestExportBackupDigest: a backup streams with its length and is checked
|
||||
// against the sha256 recorded when it was written. A match downloads whole; a
|
||||
// mismatch withholds the tail and aborts, so the browser never holds a
|
||||
// complete-looking corrupt file, and the Job hears backup_corrupt.
|
||||
func TestExportBackupDigest(t *testing.T) {
|
||||
// TestExportStreamsWhatTheJobSends: the archive reaches the browser byte for
|
||||
// byte, with the length the Job declared when it declared one. The Job checks a
|
||||
// backup against its recorded digest itself and, on a mismatch, cuts its upload
|
||||
// short of the end; the download then aborts too, so the browser never holds a
|
||||
// complete-looking file.
|
||||
func TestExportStreamsWhatTheJobSends(t *testing.T) {
|
||||
archive := randomBytes(3*exportCopyBuffer + 4321)
|
||||
sum := sha256.Sum256(archive)
|
||||
|
||||
t.Run("whole, with its length", func(t *testing.T) {
|
||||
a, _, _, ex := exportFixture()
|
||||
ext, in := exportServers(t, a)
|
||||
v := beginExport(t, a.ExternalHandler(), backupPath, "owner1")
|
||||
up := realUpload(t, in, ex.reqs[0], bytes.NewReader(archive), int64(len(archive)))
|
||||
waitExportReady(t, a.ExternalHandler(), v.Ticket, "owner1")
|
||||
resp, got, err := realDownload(ext, v.Ticket)
|
||||
if resp == nil {
|
||||
t.Fatalf("download: %v", err)
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK || resp.Header.Get("Content-Length") != strconv.Itoa(len(archive)) ||
|
||||
err != nil || !bytes.Equal(got, archive) {
|
||||
t.Fatalf("download = %d, Content-Length %q, read %d of %d bytes, err %v",
|
||||
resp.StatusCode, resp.Header.Get("Content-Length"), len(got), len(archive), err)
|
||||
}
|
||||
if upResp := awaitResponse(t, up); upResp.StatusCode != http.StatusNoContent {
|
||||
t.Fatalf("upload answered %d, want 204", upResp.StatusCode)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("an upload the Job cuts short aborts the download", func(t *testing.T) {
|
||||
a, _, _, ex := exportFixture()
|
||||
ext, in := exportServers(t, a)
|
||||
v := beginExport(t, a.ExternalHandler(), backupPath, "owner1")
|
||||
// As the Job's digest check fails: all but the end, then a read error,
|
||||
// which aborts the chunked PUT.
|
||||
pr, pw := io.Pipe()
|
||||
go func() {
|
||||
_, _ = pw.Write(archive[:len(archive)-100])
|
||||
pw.CloseWithError(errors.New("the backup archive does not match the sha256 recorded when it was written"))
|
||||
}()
|
||||
up := realUpload(t, in, ex.reqs[0], pr, -1)
|
||||
waitExportReady(t, a.ExternalHandler(), v.Ticket, "owner1")
|
||||
resp, got, err := realDownload(ext, v.Ticket)
|
||||
if resp == nil {
|
||||
t.Fatalf("download: %v", err)
|
||||
}
|
||||
if err == nil || len(got) > len(archive)-100 || !bytes.Equal(got, archive[:len(got)]) {
|
||||
t.Fatalf("read %d of %d bytes, err %v; want an error short of the end", len(got), len(archive), err)
|
||||
}
|
||||
if upResp := awaitResponse(t, up); upResp.StatusCode != -1 {
|
||||
t.Fatalf("the aborted upload answered %d", upResp.StatusCode)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a browser that leaves early is what the Job hears", func(t *testing.T) {
|
||||
a, _, _, ex := exportFixture()
|
||||
ext, _ := exportServers(t, a)
|
||||
v := beginExport(t, a.ExternalHandler(), worldPath, "owner1")
|
||||
up := uploadExport(a.InternalHandler(), ex.reqs[0].ID, ex.reqs[0].Token, io.LimitReader(zeros{}, 1<<30))
|
||||
waitExportReady(t, a.ExternalHandler(), v.Ticket, "owner1")
|
||||
req, _ := http.NewRequest("GET", ext.URL+"/api/v1/exports/"+v.Ticket+"/download", nil)
|
||||
req.Header.Set("X-Test-User", "owner1")
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := io.ReadFull(resp.Body, make([]byte, 1000)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
resp.Body.Close()
|
||||
if w := awaitUpload(t, up); w.Code != http.StatusGone || decodeErr(t, w) != "export_expired" {
|
||||
t.Fatalf("upload answered %d %s, want 410 export_expired", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
const fileDownloadPath = "/api/v1/servers/survival/files/download?path="
|
||||
|
||||
// fileDownloadFixture is exportFixture with the file manager wired, which the
|
||||
// file routes' gate requires.
|
||||
func fileDownloadFixture() (*API, *fakeRepo, *fakeCluster, *fakeExporter) {
|
||||
a, repo, cl, ex := exportFixture()
|
||||
a.Files = &fakeFileEditor{}
|
||||
return a, repo, cl, ex
|
||||
}
|
||||
|
||||
func TestFileDownloadGate(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
digest string
|
||||
name, query, filename, contentType, payload string
|
||||
want worldexport.Request
|
||||
}{
|
||||
{"recorded digest matches", hex.EncodeToString(sum[:])},
|
||||
{"no digest recorded", ""},
|
||||
{"digest mismatch", strings.Repeat("ab", 32)},
|
||||
{"a file", "plugins/Essentials/config.yml", "config.yml", fileedit.DownloadFileType, `{"dir":false}`,
|
||||
worldexport.Request{Server: "survival", Mode: worldexport.ModeFiles, Path: "plugins/Essentials/config.yml"}},
|
||||
{"a folder, as a zip named after it", "plugins/Essentials/&dir=true", "Essentials.zip", fileedit.DownloadZipType, `{"dir":true}`,
|
||||
worldexport.Request{Server: "survival", Mode: worldexport.ModeFiles, Path: "plugins/Essentials/", Dir: true}},
|
||||
{"dir other than true is a file", "a.yml&dir=false", "a.yml", fileedit.DownloadFileType, `{"dir":false}`,
|
||||
worldexport.Request{Server: "survival", Mode: worldexport.ModeFiles, Path: "a.yml"}},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
a, repo, _, ex := exportFixture()
|
||||
repo.backups[0].sha256 = tc.digest
|
||||
ext, in := exportServers(t, a)
|
||||
v := beginExport(t, a.ExternalHandler(), backupPath, "owner1")
|
||||
up := realUpload(t, in, ex.reqs[0], bytes.NewReader(archive), int64(len(archive)))
|
||||
waitExportReady(t, a.ExternalHandler(), v.Ticket, "owner1")
|
||||
resp, got, err := realDownload(ext, v.Ticket)
|
||||
if resp == nil {
|
||||
t.Fatalf("download: %v", err)
|
||||
a, repo, cl, ex := fileDownloadFixture()
|
||||
v := beginExport(t, a.ExternalHandler(), fileDownloadPath+tc.query, "owner1")
|
||||
if !hex64.MatchString(v.Ticket) || v.State != "pending" || v.Filename != tc.filename {
|
||||
t.Fatalf("ticket = %+v", v)
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK || resp.Header.Get("Content-Length") != strconv.Itoa(len(archive)) {
|
||||
t.Fatalf("download = %d, Content-Length %q", resp.StatusCode, resp.Header.Get("Content-Length"))
|
||||
if len(ex.reqs) != 1 {
|
||||
t.Fatalf("exporter started %d Jobs, want 1", len(ex.reqs))
|
||||
}
|
||||
upResp := awaitResponse(t, up)
|
||||
if tc.digest == strings.Repeat("ab", 32) {
|
||||
if err == nil || len(got) >= len(archive) || !bytes.Equal(got, archive[:len(got)]) {
|
||||
t.Fatalf("mismatch: read %d of %d bytes, err %v; want an error short of the end", len(got), len(archive), err)
|
||||
}
|
||||
if upResp.StatusCode != http.StatusConflict || errCode(mustRead(t, upResp.Body)) != "backup_corrupt" {
|
||||
t.Fatalf("upload answered %d, want 409 backup_corrupt", upResp.StatusCode)
|
||||
}
|
||||
return
|
||||
r := ex.reqs[0]
|
||||
if !hex16.MatchString(r.ID) || !hex64.MatchString(r.Token) || r.TargetURL != exportBase+"/api/v1/internal/exports/"+r.ID {
|
||||
t.Fatalf("export request = %+v", r)
|
||||
}
|
||||
if err != nil || !bytes.Equal(got, archive) {
|
||||
t.Fatalf("read %d of %d bytes, err %v", len(got), len(archive), err)
|
||||
r.ID, r.Token, r.TargetURL = "", "", ""
|
||||
if r != tc.want {
|
||||
t.Fatalf("export request = %+v, want %+v", r, tc.want)
|
||||
}
|
||||
if upResp.StatusCode != http.StatusNoContent {
|
||||
t.Fatalf("upload answered %d, want 204", upResp.StatusCode)
|
||||
if e := a.exportTickets().byTicket[v.Ticket]; e.contentType != tc.contentType || !e.files {
|
||||
t.Fatalf("export served as %q, file download %v", e.contentType, e.files)
|
||||
}
|
||||
if strings.Join(cl.acquired, ",") != "survival:"+maintenance.KindExport || strings.Join(cl.released, ",") != "survival" {
|
||||
t.Fatalf("lock acquired %v, released %v", cl.acquired, cl.released)
|
||||
}
|
||||
if len(repo.audits) != 1 || repo.audits[0].Action != "file.download" || repo.audits[0].ActorUserID != "owner1" ||
|
||||
repo.audits[0].ServerName != "survival:"+tc.want.Path || string(repo.audits[0].Payload) != tc.payload {
|
||||
t.Fatalf("audit = %+v", repo.audits)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// The tail is withheld from the response writer itself, not only from
|
||||
// whatever the connection had yet to send: with every write captured, a
|
||||
// mismatch still ends short of the archive.
|
||||
t.Run("the tail waits for the digest", func(t *testing.T) {
|
||||
a, repo, _, ex := exportFixture()
|
||||
repo.backups[0].sha256 = strings.Repeat("ab", 32)
|
||||
ext := a.ExternalHandler()
|
||||
v := beginExport(t, ext, backupPath, "owner1")
|
||||
up := uploadExport(a.InternalHandler(), ex.reqs[0].ID, ex.reqs[0].Token, bytes.NewReader(archive))
|
||||
waitExportReady(t, ext, v.Ticket, "owner1")
|
||||
w := httptest.NewRecorder()
|
||||
func() {
|
||||
defer func() {
|
||||
if p := recover(); p != http.ErrAbortHandler {
|
||||
t.Errorf("the download ended with %v, want the abort", p)
|
||||
}
|
||||
}()
|
||||
r := httptest.NewRequest("GET", "/api/v1/exports/"+v.Ticket+"/download", nil)
|
||||
r.Header.Set("X-Test-User", "owner1")
|
||||
ext.ServeHTTP(w, r)
|
||||
}()
|
||||
if got := w.Body.Bytes(); len(got) != 3*exportCopyBuffer || !bytes.Equal(got, archive[:len(got)]) {
|
||||
t.Fatalf("wrote %d of %d bytes, want all but the last read (%d)", len(got), len(archive), 3*exportCopyBuffer)
|
||||
}
|
||||
if w := awaitUpload(t, up); w.Code != http.StatusConflict || decodeErr(t, w) != "backup_corrupt" {
|
||||
t.Fatalf("upload answered %d %s, want 409 backup_corrupt", w.Code, w.Body.String())
|
||||
t.Run("admin", func(t *testing.T) {
|
||||
a, _, _, ex := fileDownloadFixture()
|
||||
beginExport(t, a.ExternalHandler(), fileDownloadPath+"server.properties", "admin1")
|
||||
if len(ex.reqs) != 1 {
|
||||
t.Fatalf("exporter started %d Jobs, want 1", len(ex.reqs))
|
||||
}
|
||||
})
|
||||
|
||||
busy := func(_ *API, c *fakeCluster) {
|
||||
c.maintErr["survival"] = &MaintenanceBusyError{Kind: maintenance.KindFileWrite}
|
||||
}
|
||||
for _, tc := range []struct {
|
||||
name, user, query string
|
||||
edit func(*API, *fakeCluster)
|
||||
code int
|
||||
errCode string
|
||||
}{
|
||||
{name: "stranger", user: "stranger", query: "a.yml", code: http.StatusForbidden, errCode: "forbidden"},
|
||||
{name: "principal without an account", user: "nouser", query: "a.yml", code: http.StatusForbidden, errCode: "forbidden"},
|
||||
{name: "running", user: "owner1", query: "a.yml", edit: func(_ *API, c *fakeCluster) { c.byName["survival"].Ready = true },
|
||||
code: http.StatusConflict, errCode: "not_stopped"},
|
||||
{name: "no world volume", user: "owner1", query: "a.yml", edit: func(_ *API, c *fakeCluster) { c.noWorld["survival"] = true },
|
||||
code: http.StatusConflict, errCode: "no_world_volume"},
|
||||
{name: "no file editor", user: "owner1", query: "a.yml", edit: func(a *API, _ *fakeCluster) { a.Files = nil },
|
||||
code: http.StatusServiceUnavailable, errCode: "files_unavailable"},
|
||||
{name: "no path", user: "owner1", query: "", code: http.StatusBadRequest, errCode: "bad_request"},
|
||||
{name: "the root", user: "owner1", query: ".&dir=true", code: http.StatusBadRequest, errCode: "bad_path"},
|
||||
{name: "the root, slashed", user: "owner1", query: "/&dir=true", code: http.StatusBadRequest, errCode: "bad_path"},
|
||||
{name: "the root, dotted", user: "owner1", query: "./", code: http.StatusBadRequest, errCode: "bad_path"},
|
||||
{name: "the root, walked back", user: "owner1", query: "plugins/..", code: http.StatusBadRequest, errCode: "bad_path"},
|
||||
{name: "no exporter", user: "owner1", query: "a.yml", edit: func(a *API, _ *fakeCluster) { a.Exporter = nil },
|
||||
code: http.StatusServiceUnavailable, errCode: "export_unavailable"},
|
||||
{name: "no internal URL", user: "owner1", query: "a.yml", edit: func(a *API, _ *fakeCluster) { a.InternalBaseURL = "" },
|
||||
code: http.StatusServiceUnavailable, errCode: "export_unavailable"},
|
||||
{name: "world busy", user: "owner1", query: "a.yml", edit: busy, code: http.StatusConflict, errCode: "maintenance_in_progress"},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
a, repo, cl, ex := fileDownloadFixture()
|
||||
if tc.edit != nil {
|
||||
tc.edit(a, cl)
|
||||
}
|
||||
w := do(a.ExternalHandler(), "POST", fileDownloadPath+tc.query, "", as(tc.user))
|
||||
if w.Code != tc.code {
|
||||
t.Fatalf("code = %d, want %d (%s)", w.Code, tc.code, w.Body.String())
|
||||
}
|
||||
if got := decodeErr(t, w); got != tc.errCode {
|
||||
t.Errorf("error code = %q, want %q", got, tc.errCode)
|
||||
}
|
||||
// Nothing is left behind: no Job, no audit, no ticket, no start
|
||||
// counted against the hour.
|
||||
reg := a.exportTickets()
|
||||
if len(ex.reqs) != 0 || len(repo.audits) != 0 || len(cl.released) != 0 || len(reg.byTicket) != 0 || len(reg.starts) != 0 {
|
||||
t.Errorf("a refused download started %d Jobs, wrote %d audits, released %v, left %d tickets and %v",
|
||||
len(ex.reqs), len(repo.audits), cl.released, len(reg.byTicket), reg.starts)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
t.Run("a failed Job is refunded and releases the lock", func(t *testing.T) {
|
||||
a, _, cl, ex := fileDownloadFixture()
|
||||
ex.err = errors.New("jobs is forbidden")
|
||||
if w := do(a.ExternalHandler(), "POST", fileDownloadPath+"a.yml", "", as("owner1")); w.Code != http.StatusInternalServerError {
|
||||
t.Fatalf("failed Job: code = %d, want 500 (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if reg := a.exportTickets(); len(reg.byTicket) != 0 || len(reg.starts) != 0 || strings.Join(cl.released, ",") != "survival" {
|
||||
t.Fatalf("after a failed Job: %d tickets, starts %v, released %v", len(reg.byTicket), reg.starts, cl.released)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestFileDownloadLimits(t *testing.T) {
|
||||
busy := func(t *testing.T, w *httptest.ResponseRecorder, message, retry string) {
|
||||
t.Helper()
|
||||
var raw map[string]map[string]string
|
||||
_ = json.Unmarshal(w.Body.Bytes(), &raw)
|
||||
if w.Code != http.StatusTooManyRequests || raw["error"]["code"] != "export_busy" || raw["error"]["message"] != message ||
|
||||
w.Header().Get("Retry-After") != retry {
|
||||
t.Fatalf("refusal = %d %s Retry-After %q; want 429 %q Retry-After %s", w.Code, w.Body.String(), w.Header().Get("Retry-After"), message, retry)
|
||||
}
|
||||
}
|
||||
|
||||
t.Run("two per user, counted apart from exports", func(t *testing.T) {
|
||||
a, _, _, ex := fileDownloadFixture()
|
||||
h := a.ExternalHandler()
|
||||
beginExport(t, h, fileDownloadPath+"a.yml", "owner1")
|
||||
beginExport(t, h, fileDownloadPath+"b.yml", "owner1")
|
||||
beginExport(t, h, worldPath, "owner1") // file downloads do not hold an export off
|
||||
busy(t, do(h, "POST", fileDownloadPath+"c.yml", "", as("owner1")),
|
||||
"you already have 2 file downloads in progress; let one finish first", "90")
|
||||
if len(ex.reqs) != 3 {
|
||||
t.Fatalf("exporter started %d Jobs, want 3", len(ex.reqs))
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("four across the install", func(t *testing.T) {
|
||||
a, _, _, ex := fileDownloadFixture()
|
||||
h := a.ExternalHandler()
|
||||
for _, u := range []string{"owner1", "owner1", "owner3", "owner3"} {
|
||||
server := map[string]string{"owner1": "survival", "owner3": "gamma"}[u]
|
||||
beginExport(t, h, "/api/v1/servers/"+server+"/files/download?path=a.yml", u)
|
||||
}
|
||||
busy(t, do(h, "POST", fileDownloadPath+"a.yml", "", as("admin1")),
|
||||
"4 file downloads are already in progress; retry in a minute", "60")
|
||||
if len(ex.reqs) != 4 {
|
||||
t.Fatalf("exporter started %d Jobs, want 4", len(ex.reqs))
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("thirty per user per hour", func(t *testing.T) {
|
||||
defer func(old time.Duration) { exportPendingTTL = old }(exportPendingTTL)
|
||||
exportPendingTTL = time.Minute
|
||||
a, _, _, _ := fileDownloadFixture()
|
||||
var clock atomic.Int64
|
||||
a.Now = func() time.Time { return time.Unix(clock.Load(), 0) }
|
||||
h := a.ExternalHandler()
|
||||
for i := range fileExportPerHour {
|
||||
clock.Store(1_700_000_000 + int64(i)*100) // each start outlives the last one's pending TTL
|
||||
beginExport(t, h, fileDownloadPath+"a.yml", "owner1")
|
||||
}
|
||||
clock.Store(1_700_000_000 + 2950)
|
||||
busy(t, do(h, "POST", fileDownloadPath+"a.yml", "", as("owner1")),
|
||||
"you have started 30 file downloads in the last hour; retry later", "650")
|
||||
beginExport(t, h, worldPath, "owner1") // exports keep their own hour
|
||||
clock.Store(1_700_000_000 + 3600)
|
||||
beginExport(t, h, fileDownloadPath+"a.yml", "owner1")
|
||||
})
|
||||
}
|
||||
|
||||
// TestFileDownloadServed: a file goes out as its raw bytes under its own name,
|
||||
// with its length; a folder as a zip, streamed without one.
|
||||
func TestFileDownloadServed(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name, query, contentType, disposition string
|
||||
size int64
|
||||
}{
|
||||
{"a file", url.QueryEscape("plugins/配置 file.yml"), fileedit.DownloadFileType,
|
||||
"attachment; filename*=utf-8''%E9%85%8D%E7%BD%AE%20file.yml", 64 << 10},
|
||||
{"a folder", "plugins&dir=true", fileedit.DownloadZipType, "attachment; filename=plugins.zip", -1},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
a, _, _, ex := fileDownloadFixture()
|
||||
ext, in := exportServers(t, a)
|
||||
v := beginExport(t, a.ExternalHandler(), fileDownloadPath+tc.query, "owner1")
|
||||
// Past what the server would buffer and measure itself when the
|
||||
// handler sets no length.
|
||||
body := randomBytes(64 << 10)
|
||||
up := realUpload(t, in, ex.reqs[0], bytes.NewReader(body), tc.size)
|
||||
waitExportReady(t, a.ExternalHandler(), v.Ticket, "owner1")
|
||||
resp, got, err := realDownload(ext, v.Ticket)
|
||||
if resp == nil || err != nil || !bytes.Equal(got, body) {
|
||||
t.Fatalf("download: read %d bytes, %v", len(got), err)
|
||||
}
|
||||
wantLength := ""
|
||||
if tc.size >= 0 {
|
||||
wantLength = strconv.FormatInt(tc.size, 10)
|
||||
}
|
||||
h := resp.Header
|
||||
if h.Get("Content-Type") != tc.contentType || h.Get("Content-Disposition") != tc.disposition ||
|
||||
h.Get("Content-Length") != wantLength || h.Get("X-Content-Type-Options") != "nosniff" {
|
||||
t.Fatalf("download headers = %v", h)
|
||||
}
|
||||
if upResp := awaitResponse(t, up); upResp.StatusCode != http.StatusNoContent {
|
||||
t.Fatalf("upload answered %d, want 204", upResp.StatusCode)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// zeros reads as an endless run of zero bytes.
|
||||
@@ -873,7 +1093,7 @@ func TestK8sExportJobs(t *testing.T) {
|
||||
job := func(id, mode string) *batchv1.Job {
|
||||
j, err := worldexport.ExportJob(worldexport.JobParams{
|
||||
Server: "survival", ID: id, Mode: mode,
|
||||
WorldPVC: "world-survival-0", BackupPVC: "felis-backups", BackupRef: "/backups/a.tar.gz",
|
||||
WorldPVC: "world-survival-0", BackupPVC: "felis-backups", BackupRef: "/backups/a.tar.gz", Path: "plugins",
|
||||
TargetURL: exportBase + "/x", Token: "t", Namespace: "minecraft", Image: "felis:1",
|
||||
})
|
||||
if err != nil {
|
||||
@@ -882,9 +1102,10 @@ func TestK8sExportJobs(t *testing.T) {
|
||||
return j
|
||||
}
|
||||
world, backupJob := job("1111111111111111", worldexport.ModeWorld), job("2222222222222222", worldexport.ModeBackup)
|
||||
files := job("3333333333333333", worldexport.ModeFiles)
|
||||
restoring := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Namespace: "minecraft", Name: "restore-survival-cc",
|
||||
Labels: map[string]string{jobServerLabel: "survival", jobManagedByLabel: jobManagedByRestore}}}
|
||||
c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(world, backupJob, restoring).
|
||||
c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(world, backupJob, files, restoring).
|
||||
WithStatusSubresource(&batchv1.Job{}).Build()
|
||||
k := NewK8sJobStatus(c, "minecraft")
|
||||
ctx := context.Background()
|
||||
@@ -900,7 +1121,8 @@ func TestK8sExportJobs(t *testing.T) {
|
||||
for _, j := range jobs {
|
||||
kinds[j.Name] = j.Kind + "/" + j.State
|
||||
}
|
||||
want := map[string]string{world.Name: "export_world/running", backupJob.Name: "export_backup/running", restoring.Name: "restore/running"}
|
||||
want := map[string]string{world.Name: "export_world/running", backupJob.Name: "export_backup/running",
|
||||
files.Name: "export_files/running", restoring.Name: "restore/running"}
|
||||
if len(kinds) != len(want) {
|
||||
t.Fatalf("jobs = %v, want %v", kinds, want)
|
||||
}
|
||||
@@ -910,11 +1132,14 @@ func TestK8sExportJobs(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// The kinds agree with maintenance.JobKind: a world export holds the world,
|
||||
// a backup export does not.
|
||||
// The kinds agree with maintenance.JobKind: a world export and a file
|
||||
// download hold the world, a backup export does not.
|
||||
if kind, holds := maintenance.JobKind(world); kind != maintenance.KindExport || !holds {
|
||||
t.Errorf("JobKind(world export) = %q, %v", kind, holds)
|
||||
}
|
||||
if kind, holds := maintenance.JobKind(files); kind != maintenance.KindExport || !holds {
|
||||
t.Errorf("JobKind(file download) = %q, %v", kind, holds)
|
||||
}
|
||||
if _, holds := maintenance.JobKind(backupJob); holds {
|
||||
t.Error("a backup export holds the world")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,446 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/fileedit"
|
||||
"felis.lolicon.best/internal/maintenance"
|
||||
"felis.lolicon.best/internal/naming"
|
||||
)
|
||||
|
||||
// A file too big for the one-request upload (handleUploadFile) arrives as an
|
||||
// upload session instead: POST …/files/uploads begins one for a path and a
|
||||
// size, PUT …/files/uploads/{id}?offset= sends it in parts of at most
|
||||
// fileedit.PartBytes (each part fits the Cloudflare edge's body limit), and POST
|
||||
// …/files/uploads/{id}/commit lands it. The session lives on felis-api's staging
|
||||
// disk (fileedit.Stage, session.go), bound to the account and the server it was
|
||||
// begun for; the room for the whole file is reserved when it begins, so there is
|
||||
// no product ceiling on the size, only the disk.
|
||||
//
|
||||
// Landing a file that size, like unzipping an archive, can outlast a request,
|
||||
// so both answer 202 with the op, and GET …/files/ops reports how far it has got
|
||||
// and how it ended (fileedit.Editor.StartUpload, StartUnzip, Ops). The Job holds
|
||||
// the world volume while it runs, as any file write does, and the server cannot
|
||||
// start until it ends.
|
||||
|
||||
// fileSessionView is where an upload session stands.
|
||||
type fileSessionView struct {
|
||||
ID string `json:"id"`
|
||||
Path string `json:"path"`
|
||||
Size int64 `json:"size"`
|
||||
Received int64 `json:"received"`
|
||||
PartMaxBytes int64 `json:"part_max_bytes"`
|
||||
}
|
||||
|
||||
func sessionView(s fileedit.Session) fileSessionView {
|
||||
return fileSessionView{ID: s.ID, Path: s.Path, Size: s.Size, Received: s.Received, PartMaxBytes: fileedit.PartBytes}
|
||||
}
|
||||
|
||||
// beginFileUploadRequest is the POST …/files/uploads body.
|
||||
type beginFileUploadRequest struct {
|
||||
Size *int64 `json:"size"`
|
||||
}
|
||||
|
||||
// handleBeginFileUpload serves POST /api/v1/servers/{name}/files/uploads?path=…
|
||||
// — begin an upload session for a file of body.size bytes that will land at
|
||||
// path. The gate is the file manager's, so a server that is running is refused
|
||||
// before any byte is sent; the parts that follow need only the account and the
|
||||
// server, so starting the server midway costs the upload nothing but the
|
||||
// commit's refusal until it is stopped again.
|
||||
func (a *API) handleBeginFileUpload(w http.ResponseWriter, r *http.Request) {
|
||||
name, ok := a.authorizeFileOp(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
user, ok := a.requireFileStage(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
path, ok := requirePath(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
// The Job checks the path against the volume when the file lands, hours of
|
||||
// upload later for a big one; one that could never land is refused now.
|
||||
if !filepath.IsLocal(path) || filepath.Clean(path) == "." {
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_path",
|
||||
"the path must name a file inside the world folder"))
|
||||
return
|
||||
}
|
||||
var body beginFileUploadRequest
|
||||
if err := decodeJSON(w, r, &body); err != nil {
|
||||
writeError(w, r, err)
|
||||
return
|
||||
}
|
||||
if body.Size == nil || *body.Size < 0 {
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_request",
|
||||
"size must be the file's length in bytes"))
|
||||
return
|
||||
}
|
||||
s, err := a.FileStage.Begin(user, name, path, *body.Size)
|
||||
if err != nil {
|
||||
writeFileSessionError(w, r, err)
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, sessionView(s))
|
||||
}
|
||||
|
||||
// handleFileUploadStatus serves GET /api/v1/servers/{name}/files/uploads/{id}
|
||||
// — where the caller's session stands, so a client that lost a part's answer
|
||||
// resumes from received.
|
||||
func (a *API) handleFileUploadStatus(w http.ResponseWriter, r *http.Request) {
|
||||
name, user, ok := a.authorizeFileSession(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
s, err := a.FileStage.Status(user, name, r.PathValue("id"))
|
||||
if err != nil {
|
||||
writeFileSessionError(w, r, err)
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, sessionView(s))
|
||||
}
|
||||
|
||||
// handleFileUploadPart serves PUT
|
||||
// /api/v1/servers/{name}/files/uploads/{id}?offset=… — append the raw body to
|
||||
// the caller's session. offset must be where the session ends (409
|
||||
// upload_offset_mismatch otherwise; the status says where), and Content-Length
|
||||
// is required, as for the one-request upload: the part is taken whole or not at
|
||||
// all, and a part that breaks midway leaves the session where it was. A part
|
||||
// over fileedit.PartBytes is refused before a byte of it is read (Append).
|
||||
func (a *API) handleFileUploadPart(w http.ResponseWriter, r *http.Request) {
|
||||
name, user, ok := a.authorizeFileSession(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
offset, err := strconv.ParseInt(r.URL.Query().Get("offset"), 10, 64)
|
||||
if err != nil {
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_request",
|
||||
"offset must be the byte position the part starts at"))
|
||||
return
|
||||
}
|
||||
if r.ContentLength < 0 {
|
||||
writeError(w, r, newError(http.StatusLengthRequired, "length_required",
|
||||
"a part needs a Content-Length"))
|
||||
return
|
||||
}
|
||||
s, err := a.FileStage.Append(user, name, r.PathValue("id"), offset, r.Body, r.ContentLength)
|
||||
if err != nil {
|
||||
writeFileSessionError(w, r, err)
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, sessionView(s))
|
||||
}
|
||||
|
||||
// handleDropFileUpload serves DELETE /api/v1/servers/{name}/files/uploads/{id}
|
||||
// — cancel the caller's session and free the room it holds.
|
||||
func (a *API) handleDropFileUpload(w http.ResponseWriter, r *http.Request) {
|
||||
name, user, ok := a.authorizeFileSession(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if err := a.FileStage.Drop(user, name, r.PathValue("id")); err != nil {
|
||||
writeFileSessionError(w, r, err)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
|
||||
// startFileOpRequest is the body of a commit or an unzip.
|
||||
type startFileOpRequest struct {
|
||||
Overwrite bool `json:"overwrite"`
|
||||
}
|
||||
|
||||
// handleCommitFileUpload serves POST
|
||||
// /api/v1/servers/{name}/files/uploads/{id}/commit — land the caller's
|
||||
// finished session at its path, replacing a file there only with
|
||||
// body.overwrite (the op ends file_exists otherwise). It answers 202 with the
|
||||
// op; GET …/files/ops reports how it ends.
|
||||
//
|
||||
// The world lock is taken BEFORE the session is sealed: a commit made while an
|
||||
// earlier commit's Job is still fetching the file is refused by that Job's hold
|
||||
// on the volume, so it never mints the fresh token that would lock the running
|
||||
// Job out. The session outlives a Job that fails before fetching every byte, so
|
||||
// such a commit is simply made again; one that fetched them all is gone.
|
||||
func (a *API) handleCommitFileUpload(w http.ResponseWriter, r *http.Request) {
|
||||
name, ok := a.authorizeFileOp(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
user, ok := a.requireFileStage(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
var body startFileOpRequest
|
||||
if err := decodeJSON(w, r, &body); err != nil {
|
||||
writeError(w, r, err)
|
||||
return
|
||||
}
|
||||
id := r.PathValue("id")
|
||||
s, err := a.FileStage.Status(user, name, id)
|
||||
// Refused before the world lock is asked for; Seal checks again under the
|
||||
// stage's own lock, for a part that arrives in between.
|
||||
if err == nil && s.Received != s.Size {
|
||||
err = fmt.Errorf("%w: %d of %d bytes are here", fileedit.ErrUploadIncomplete, s.Received, s.Size)
|
||||
}
|
||||
if err != nil {
|
||||
writeFileSessionError(w, r, err)
|
||||
return
|
||||
}
|
||||
release, ok := a.acquireWorld(w, r, name, maintenance.KindFileWrite, "stop the server before editing its files")
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
defer release()
|
||||
staged, err := a.FileStage.Seal(user, name, id)
|
||||
if err != nil {
|
||||
writeFileSessionError(w, r, err)
|
||||
return
|
||||
}
|
||||
op, err := a.Files.StartUpload(r.Context(), name, s.Path, fileedit.UploadSource{
|
||||
URL: a.InternalBaseURL + "/api/v1/internal/file-uploads/" + id,
|
||||
Token: staged.Token,
|
||||
Size: staged.Size,
|
||||
SHA256: staged.SHA256,
|
||||
}, body.Overwrite)
|
||||
if err != nil {
|
||||
writeFileEditError(w, r, err)
|
||||
return
|
||||
}
|
||||
a.auditFile(r, "file.upload", name, s.Path, map[string]any{
|
||||
"size_bytes": staged.Size, "sha256": staged.SHA256, "overwrite": body.Overwrite,
|
||||
})
|
||||
writeJSON(w, http.StatusAccepted, map[string]any{"op": opView(op)})
|
||||
}
|
||||
|
||||
// handleUnzipFile serves POST /api/v1/servers/{name}/files/unzip?path=… —
|
||||
// extract the .zip at path into the folder holding it. Without body.overwrite
|
||||
// an archive that would replace any file changes nothing and the op ends
|
||||
// file_exists with the list of them, for the caller to confirm and run again
|
||||
// with overwrite. It answers 202 with the op, like a commit.
|
||||
//
|
||||
// Only the name is checked here, so a caller who picked the wrong file hears
|
||||
// so at once; whether it is a zip, and whether every entry is safe to extract,
|
||||
// is the Job's to decide (fileedit/unzip.go).
|
||||
func (a *API) handleUnzipFile(w http.ResponseWriter, r *http.Request) {
|
||||
name, ok := a.authorizeFileOp(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
path, ok := requirePath(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if !strings.HasSuffix(strings.ToLower(path), ".zip") {
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_path", "only a .zip file can be extracted"))
|
||||
return
|
||||
}
|
||||
var body startFileOpRequest
|
||||
if err := decodeJSON(w, r, &body); err != nil {
|
||||
writeError(w, r, err)
|
||||
return
|
||||
}
|
||||
release, ok := a.acquireWorld(w, r, name, maintenance.KindFileWrite, "stop the server before editing its files")
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
defer release()
|
||||
op, err := a.Files.StartUnzip(r.Context(), name, path, body.Overwrite)
|
||||
if err != nil {
|
||||
writeFileEditError(w, r, err)
|
||||
return
|
||||
}
|
||||
a.auditFile(r, "file.unzip", name, path, map[string]any{"overwrite": body.Overwrite})
|
||||
writeJSON(w, http.StatusAccepted, map[string]any{"op": opView(op)})
|
||||
}
|
||||
|
||||
// handleListFileOps serves GET /api/v1/servers/{name}/files/ops — the server's
|
||||
// background uploads and unzips, newest first: the one running, if any, and
|
||||
// those that ended within the last half hour. It has no stopped gate: while
|
||||
// one runs the server cannot start, and a finished one is still worth showing
|
||||
// after it has.
|
||||
func (a *API) handleListFileOps(w http.ResponseWriter, r *http.Request) {
|
||||
name, ok := a.authorizeServerFiles(w, r)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if a.Files == nil {
|
||||
writeError(w, r, newError(http.StatusServiceUnavailable, "files_unavailable",
|
||||
"the file editor is not configured"))
|
||||
return
|
||||
}
|
||||
ops, err := a.Files.Ops(r.Context(), name)
|
||||
if err != nil {
|
||||
writeError(w, r, err)
|
||||
return
|
||||
}
|
||||
views := make([]fileOpView, 0, len(ops))
|
||||
for _, op := range ops {
|
||||
views = append(views, opView(op))
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"ops": views})
|
||||
}
|
||||
|
||||
// fileOpView is one background file operation as the API shows it.
|
||||
type fileOpView struct {
|
||||
ID string `json:"id"`
|
||||
Op string `json:"op"`
|
||||
Path string `json:"path"`
|
||||
State string `json:"state"`
|
||||
StartedAt time.Time `json:"started_at"`
|
||||
FinishedAt *time.Time `json:"finished_at,omitempty"`
|
||||
Done int64 `json:"done"`
|
||||
Total int64 `json:"total"`
|
||||
Files int `json:"files,omitempty"`
|
||||
Bytes int64 `json:"bytes,omitempty"`
|
||||
Error *fileOpError `json:"error,omitempty"`
|
||||
}
|
||||
|
||||
// fileOpError is why an op failed. Code is the one the synchronous file routes
|
||||
// answer with for the same refusal (writeFileEditError), or an unzip's own
|
||||
// (archive_invalid, archive_unsafe, archive_symlink, type_conflict), or
|
||||
// job_failed for a Job that ended without saying why; the rest is what the Job
|
||||
// reported about it.
|
||||
type fileOpError struct {
|
||||
Code string `json:"code"`
|
||||
Message string `json:"message"`
|
||||
Entry string `json:"entry,omitempty"`
|
||||
Conflicts []string `json:"conflicts,omitempty"`
|
||||
ConflictCount int `json:"conflict_count,omitempty"`
|
||||
Need int64 `json:"need,omitempty"`
|
||||
Avail int64 `json:"avail,omitempty"`
|
||||
}
|
||||
|
||||
func opView(op fileedit.OpState) fileOpView {
|
||||
v := fileOpView{
|
||||
ID: op.ID, Op: op.Op, Path: op.Path, State: op.State, StartedAt: op.Started,
|
||||
Done: op.Done, Total: op.Total,
|
||||
}
|
||||
if !op.Finished.IsZero() {
|
||||
v.FinishedAt = &op.Finished
|
||||
}
|
||||
if op.Result != nil && op.Result.Code == "" {
|
||||
v.Files, v.Bytes = op.Result.Files, op.Result.Bytes
|
||||
}
|
||||
if op.State == fileedit.OpFailed {
|
||||
v.Error = opError(op)
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
// opError maps a failed op onto the API's codes. A Job that printed no result
|
||||
// carries only its condition's reason (DeadlineExceeded, BackoffLimitExceeded):
|
||||
// the log it left is the world's content and the runtime's, and none of it is
|
||||
// the caller's to read.
|
||||
func opError(op fileedit.OpState) *fileOpError {
|
||||
res := op.Result
|
||||
if res == nil {
|
||||
msg := "the file operation stopped before it could report how it went (%s); run it again"
|
||||
if op.Reason == "DeadlineExceeded" {
|
||||
msg = "the file operation ran out of time (%s); run it again"
|
||||
}
|
||||
return &fileOpError{Code: "job_failed", Message: fmt.Sprintf(msg, op.Reason)}
|
||||
}
|
||||
code := res.Code
|
||||
switch res.Code {
|
||||
case fileedit.CodeExists:
|
||||
code = "file_exists"
|
||||
case fileedit.CodeNoSpace:
|
||||
code = "volume_full"
|
||||
case fileedit.CodeConflict:
|
||||
code = "file_changed"
|
||||
}
|
||||
return &fileOpError{
|
||||
Code: code, Message: res.Error, Entry: res.Entry,
|
||||
Conflicts: res.Conflicts, ConflictCount: res.ConflictCount, Need: res.Need, Avail: res.Avail,
|
||||
}
|
||||
}
|
||||
|
||||
// requireFileStage checks the caller has an account to bind an upload session
|
||||
// to and that sessions are configured, and returns the account.
|
||||
func (a *API) requireFileStage(w http.ResponseWriter, r *http.Request) (string, bool) {
|
||||
p := principalFromContext(r.Context())
|
||||
if p == nil || p.UserID == "" {
|
||||
writeError(w, r, errForbidden)
|
||||
return "", false
|
||||
}
|
||||
if a.FileStage == nil || a.InternalBaseURL == "" {
|
||||
writeError(w, r, newError(http.StatusServiceUnavailable, "files_unavailable",
|
||||
"uploads are not configured"))
|
||||
return "", false
|
||||
}
|
||||
return p.UserID, true
|
||||
}
|
||||
|
||||
// authorizeServerFiles is authorizeFileOp without the stopped and world-volume
|
||||
// gates: the name is valid, the server exists, and the caller owns it or is
|
||||
// staff. It returns the server name.
|
||||
func (a *API) authorizeServerFiles(w http.ResponseWriter, r *http.Request) (string, bool) {
|
||||
name := r.PathValue("name")
|
||||
if err := naming.ValidateServerName(name); err != nil {
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_name", "invalid server name: %v", err))
|
||||
return "", false
|
||||
}
|
||||
rec, err := a.Repo.ServerByName(r.Context(), name)
|
||||
if err != nil {
|
||||
a.writeLookupError(w, r, err)
|
||||
return "", false
|
||||
}
|
||||
if !a.isOwnerOrAdmin(principalFromContext(r.Context()), rec) {
|
||||
writeError(w, r, errForbidden)
|
||||
return "", false
|
||||
}
|
||||
return name, true
|
||||
}
|
||||
|
||||
// authorizeFileSession is the gate of a session's parts, status and cancel:
|
||||
// the caller still owns the server (or is staff) and has the account the
|
||||
// session was begun under. A session answers only that account on that server
|
||||
// (fileedit.Stage), so it returns both.
|
||||
func (a *API) authorizeFileSession(w http.ResponseWriter, r *http.Request) (name, user string, ok bool) {
|
||||
name, ok = a.authorizeServerFiles(w, r)
|
||||
if !ok {
|
||||
return "", "", false
|
||||
}
|
||||
user, ok = a.requireFileStage(w, r)
|
||||
return name, user, ok
|
||||
}
|
||||
|
||||
// writeFileSessionError maps the upload session's errors onto HTTP statuses.
|
||||
func writeFileSessionError(w http.ResponseWriter, r *http.Request, err error) {
|
||||
var offset *fileedit.OffsetError
|
||||
switch {
|
||||
case errors.Is(err, fileedit.ErrNotStaged):
|
||||
writeError(w, r, newError(http.StatusNotFound, "upload_not_found",
|
||||
"no such upload; it was cancelled, landed, or left idle too long, so start it again"))
|
||||
case errors.Is(err, fileedit.ErrStageFull):
|
||||
writeError(w, r, newError(http.StatusInsufficientStorage, "upload_staging_full",
|
||||
"felis has no room to take this upload right now; try again later or ask an admin"))
|
||||
case errors.Is(err, fileedit.ErrTooManySessions):
|
||||
writeError(w, r, newError(http.StatusTooManyRequests, "too_many_uploads",
|
||||
"you have %d uploads in progress; finish or cancel one first", fileedit.MaxSessionsPerUser))
|
||||
case errors.Is(err, fileedit.ErrUploadBusy):
|
||||
writeError(w, r, newError(http.StatusConflict, "upload_busy",
|
||||
"another request is still writing this upload; read where it stands and continue from there"))
|
||||
case errors.As(err, &offset):
|
||||
writeError(w, r, newError(http.StatusConflict, "upload_offset_mismatch",
|
||||
"the upload holds %d bytes; send the part that starts there", offset.Received))
|
||||
case errors.Is(err, fileedit.ErrPartTooLarge):
|
||||
writeError(w, r, newError(http.StatusRequestEntityTooLarge, "part_too_large",
|
||||
"the part is larger than part_max_bytes, or runs past the size the upload began with"))
|
||||
case errors.Is(err, fileedit.ErrShortUpload):
|
||||
writeError(w, r, newError(http.StatusBadRequest, "upload_incomplete",
|
||||
"the part ended before its Content-Length; read where the upload stands and send it again"))
|
||||
case errors.Is(err, fileedit.ErrUploadIncomplete):
|
||||
writeError(w, r, newError(http.StatusConflict, "upload_incomplete",
|
||||
"the upload has not finished arriving; read where it stands and send the rest"))
|
||||
default:
|
||||
writeError(w, r, err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,787 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"reflect"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||
"felis.lolicon.best/internal/fileedit"
|
||||
"felis.lolicon.best/internal/maintenance"
|
||||
)
|
||||
|
||||
const (
|
||||
sessionsRoute = "/api/v1/servers/survival/files/uploads"
|
||||
unzipRoute = "/api/v1/servers/survival/files/unzip"
|
||||
opsRoute = "/api/v1/servers/survival/files/ops"
|
||||
sessionPath = "world/region/r.0.0.mca"
|
||||
)
|
||||
|
||||
var (
|
||||
fileOwner = &Principal{UserID: "owner1", Email: "[email protected]", Role: "user"}
|
||||
fileStranger = &Principal{UserID: "stranger", Email: "[email protected]", Role: "user"}
|
||||
fileAdmin = &Principal{UserID: "admin1", Email: "[email protected]", Role: "admin", ViaAdminAccess: true}
|
||||
)
|
||||
|
||||
// beginSession begins a session for size bytes at sessionPath and returns it.
|
||||
func beginSession(t *testing.T, api *API, size int) fileSessionView {
|
||||
t.Helper()
|
||||
w := do(api.ExternalHandler(), "POST", sessionsRoute+"?path="+sessionPath, `{"size":`+strconv.Itoa(size)+`}`, jsonHeader)
|
||||
if w.Code != http.StatusCreated {
|
||||
t.Fatalf("begin: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
return sessionAnswer(t, w)
|
||||
}
|
||||
|
||||
// sessionAnswer decodes a session answer, refusing a field the view lacks.
|
||||
func sessionAnswer(t *testing.T, w *httptest.ResponseRecorder) fileSessionView {
|
||||
t.Helper()
|
||||
var s fileSessionView
|
||||
dec := json.NewDecoder(strings.NewReader(w.Body.String()))
|
||||
dec.DisallowUnknownFields()
|
||||
if err := dec.Decode(&s); err != nil {
|
||||
t.Fatalf("session answer: %v (%s)", err, w.Body.String())
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
// doPart sends one part with the Content-Length given, whatever the body's own
|
||||
// length: -1 sends none, and one past the body is a part cut short.
|
||||
func doPart(h http.Handler, target, body string, length int64) *httptest.ResponseRecorder {
|
||||
r := httptest.NewRequest("PUT", target, strings.NewReader(body))
|
||||
r.Header.Set("Content-Type", "application/octet-stream")
|
||||
r.ContentLength = length
|
||||
w := httptest.NewRecorder()
|
||||
h.ServeHTTP(w, r)
|
||||
recordContract(r, body, w)
|
||||
return w
|
||||
}
|
||||
|
||||
func partAt(id string, offset int) string {
|
||||
return sessionsRoute + "/" + id + "?offset=" + strconv.Itoa(offset)
|
||||
}
|
||||
|
||||
func errMessage(t *testing.T, w *httptest.ResponseRecorder) string {
|
||||
t.Helper()
|
||||
var raw map[string]map[string]string
|
||||
if err := json.Unmarshal(w.Body.Bytes(), &raw); err != nil {
|
||||
t.Fatalf("error body not JSON: %v (%s)", err, w.Body.String())
|
||||
}
|
||||
return raw["error"]["message"]
|
||||
}
|
||||
|
||||
// fetchStaged is the Job's fetch of what a commit handed it.
|
||||
func fetchStaged(api *API, src fileedit.UploadSource) *httptest.ResponseRecorder {
|
||||
at := strings.TrimPrefix(src.URL, api.InternalBaseURL)
|
||||
return do(api.InternalHandler(), "GET", at, "", map[string]string{"Authorization": "Bearer " + src.Token})
|
||||
}
|
||||
|
||||
var opStarted = time.Date(2026, 9, 28, 10, 0, 0, 0, time.UTC)
|
||||
|
||||
// TestFileUploadSession drives a session from begin to the Job's fetch across
|
||||
// both faces, and each way it can go wrong on the way.
|
||||
func TestFileUploadSession(t *testing.T) {
|
||||
commit := func(api *API, id, body string) *httptest.ResponseRecorder {
|
||||
return do(api.ExternalHandler(), "POST", sessionsRoute+"/"+id+"/commit", body, jsonHeader)
|
||||
}
|
||||
status := func(api *API, id string) *httptest.ResponseRecorder {
|
||||
return do(api.ExternalHandler(), "GET", sessionsRoute+"/"+id, "", nil)
|
||||
}
|
||||
|
||||
t.Run("begin, parts, commit, and the Job fetches the whole file once", func(t *testing.T) {
|
||||
api, repo, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
// Every other internal route wants a service token; the Job has none.
|
||||
api.Internal = CallerTokens{CallerVelocity: "s3cr3t"}
|
||||
files.op = fileedit.OpState{ID: "op1", Op: fileedit.OpUpload, Path: sessionPath,
|
||||
State: fileedit.OpRunning, Started: opStarted}
|
||||
h := api.ExternalHandler()
|
||||
|
||||
s := beginSession(t, api, 10)
|
||||
if len(s.ID) != 32 || s.Path != sessionPath || s.Size != 10 || s.Received != 0 || s.PartMaxBytes != fileedit.PartBytes {
|
||||
t.Fatalf("begin = %+v", s)
|
||||
}
|
||||
|
||||
if w := doPart(h, partAt(s.ID, 0), "hello", 5); w.Code != http.StatusOK || sessionAnswer(t, w).Received != 5 {
|
||||
t.Fatalf("part 1: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
// The same part again, as a client that lost the answer might send it.
|
||||
w := doPart(h, partAt(s.ID, 0), "hello", 5)
|
||||
if w.Code != http.StatusConflict || decodeErr(t, w) != "upload_offset_mismatch" ||
|
||||
errMessage(t, w) != "the upload holds 5 bytes; send the part that starts there" {
|
||||
t.Fatalf("replayed part: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
|
||||
// Too early: refused before the world lock is asked for.
|
||||
w = commit(api, s.ID, `{}`)
|
||||
if w.Code != http.StatusConflict || decodeErr(t, w) != "upload_incomplete" ||
|
||||
files.calls != 0 || len(cl.acquired) != 0 {
|
||||
t.Fatalf("early commit: code = %d calls = %d acquired %v (%s)", w.Code, files.calls, cl.acquired, w.Body.String())
|
||||
}
|
||||
|
||||
if w := status(api, s.ID); w.Code != http.StatusOK || sessionAnswer(t, w) != (fileSessionView{
|
||||
ID: s.ID, Path: sessionPath, Size: 10, Received: 5, PartMaxBytes: fileedit.PartBytes}) {
|
||||
t.Fatalf("status: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if w := doPart(h, partAt(s.ID, 5), "world", 5); w.Code != http.StatusOK || sessionAnswer(t, w).Received != 10 {
|
||||
t.Fatalf("part 2: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if w := doPart(h, partAt(s.ID, 10), "!", 1); w.Code != http.StatusRequestEntityTooLarge || decodeErr(t, w) != "part_too_large" {
|
||||
t.Fatalf("a part past the size: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
|
||||
w = commit(api, s.ID, `{}`)
|
||||
if w.Code != http.StatusAccepted {
|
||||
t.Fatalf("commit: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
want := map[string]any{"op": map[string]any{
|
||||
"id": "op1", "op": "upload", "path": sessionPath, "state": "running",
|
||||
"started_at": "2026-09-28T10:00:00Z", "done": float64(0), "total": float64(0),
|
||||
}}
|
||||
if got := fileAnswer(t, w); !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("commit answer = %v, want %v", got, want)
|
||||
}
|
||||
digest := sha256.Sum256([]byte("helloworld"))
|
||||
sum := hex.EncodeToString(digest[:])
|
||||
src := files.gotSource
|
||||
if files.calls != 1 || files.gotOp != fileedit.OpUpload || files.gotServer != "survival" || files.gotPath != sessionPath ||
|
||||
src.URL != api.InternalBaseURL+"/api/v1/internal/file-uploads/"+s.ID ||
|
||||
src.Size != 10 || src.SHA256 != sum || len(src.Token) != 64 || files.gotOverwrite {
|
||||
t.Fatalf("executor saw calls=%d op %q server %q path %q source %+v overwrite %v",
|
||||
files.calls, files.gotOp, files.gotServer, files.gotPath, src, files.gotOverwrite)
|
||||
}
|
||||
if strings.Join(cl.acquired, ",") != "survival:"+maintenance.KindFileWrite || strings.Join(cl.released, ",") != "survival" {
|
||||
t.Fatalf("lock acquired %v, released %v", cl.acquired, cl.released)
|
||||
}
|
||||
onlyAudit(t, repo, "file.upload", "survival:"+sessionPath, `{"overwrite":false,"sha256":"`+sum+`","size_bytes":10}`)
|
||||
|
||||
fetched := fetchStaged(api, src)
|
||||
if fetched.Code != http.StatusOK || fetched.Body.String() != "helloworld" || fetched.Header().Get("Content-Length") != "10" {
|
||||
t.Fatalf("fetch: code = %d %q Content-Length %q", fetched.Code, fetched.Body.String(), fetched.Header().Get("Content-Length"))
|
||||
}
|
||||
if again := fetchStaged(api, src); again.Code != http.StatusNotFound {
|
||||
t.Fatalf("second fetch: code = %d", again.Code)
|
||||
}
|
||||
// Served whole, so the session is gone and its file with it.
|
||||
if w := status(api, s.ID); w.Code != http.StatusNotFound || decodeErr(t, w) != "upload_not_found" {
|
||||
t.Fatalf("status after the fetch: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
stageEmpty(t, api)
|
||||
})
|
||||
|
||||
t.Run("a Job that never fetched is committed again with a fresh token", func(t *testing.T) {
|
||||
api, repo, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, 3)
|
||||
doPart(api.ExternalHandler(), partAt(s.ID, 0), "abc", 3)
|
||||
if w := commit(api, s.ID, `{}`); w.Code != http.StatusAccepted {
|
||||
t.Fatalf("first commit: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
first := files.gotSource
|
||||
if w := commit(api, s.ID, `{"overwrite":true}`); w.Code != http.StatusAccepted || !files.gotOverwrite {
|
||||
t.Fatalf("second commit: code = %d overwrite %v (%s)", w.Code, files.gotOverwrite, w.Body.String())
|
||||
}
|
||||
second := files.gotSource
|
||||
if w := fetchStaged(api, first); w.Code != http.StatusNotFound {
|
||||
t.Fatalf("the first commit's token still opens it: code = %d", w.Code)
|
||||
}
|
||||
if w := fetchStaged(api, second); w.Code != http.StatusOK || w.Body.String() != "abc" {
|
||||
t.Fatalf("fetch: code = %d %q", w.Code, w.Body.String())
|
||||
}
|
||||
if len(repo.audits) != 2 || string(repo.audits[1].Payload) != `{"overwrite":true,"sha256":"`+second.SHA256+`","size_bytes":3}` {
|
||||
t.Fatalf("audits = %+v", repo.audits)
|
||||
}
|
||||
})
|
||||
|
||||
// The running Job's hold on the world refuses the commit before Seal, so the
|
||||
// token that Job carries still opens the file.
|
||||
t.Run("a commit while the world is held leaves the running Job's token alone", func(t *testing.T) {
|
||||
api, _, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, 3)
|
||||
doPart(api.ExternalHandler(), partAt(s.ID, 0), "abc", 3)
|
||||
commit(api, s.ID, `{}`)
|
||||
running := files.gotSource
|
||||
cl.maintErr["survival"] = &MaintenanceBusyError{Kind: maintenance.KindFileWrite}
|
||||
w := commit(api, s.ID, `{}`)
|
||||
if w.Code != http.StatusConflict || decodeErr(t, w) != "maintenance_in_progress" || files.calls != 1 {
|
||||
t.Fatalf("code = %d calls = %d (%s)", w.Code, files.calls, w.Body.String())
|
||||
}
|
||||
if w := fetchStaged(api, running); w.Code != http.StatusOK || w.Body.String() != "abc" {
|
||||
t.Fatalf("the running Job lost its file: code = %d %q", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a Job that could not start is not audited and lets go of the world", func(t *testing.T) {
|
||||
api, repo, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, 3)
|
||||
doPart(api.ExternalHandler(), partAt(s.ID, 0), "abc", 3)
|
||||
files.err = errors.New("the cluster said no")
|
||||
w := commit(api, s.ID, `{}`)
|
||||
if w.Code != http.StatusInternalServerError || len(repo.audits) != 0 || strings.Join(cl.released, ",") != "survival" {
|
||||
t.Fatalf("code = %d audits %+v released %v (%s)", w.Code, repo.audits, cl.released, w.Body.String())
|
||||
}
|
||||
if w := status(api, s.ID); w.Code != http.StatusOK || sessionAnswer(t, w).Received != 3 {
|
||||
t.Fatalf("the session went with the failed start: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a fetch cut short keeps the session for the next commit", func(t *testing.T) {
|
||||
api, _, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, 3)
|
||||
doPart(api.ExternalHandler(), partAt(s.ID, 0), "abc", 3)
|
||||
commit(api, s.ID, `{}`)
|
||||
src := files.gotSource
|
||||
r := httptest.NewRequest("GET", strings.TrimPrefix(src.URL, api.InternalBaseURL), nil)
|
||||
r.Header.Set("Authorization", "Bearer "+src.Token)
|
||||
cut := &brokenWriter{header: http.Header{}}
|
||||
api.InternalHandler().ServeHTTP(cut, r)
|
||||
if cut.code != http.StatusOK {
|
||||
t.Fatalf("cut fetch: code = %d", cut.code)
|
||||
}
|
||||
if w := status(api, s.ID); w.Code != http.StatusOK {
|
||||
t.Fatalf("status after a cut fetch: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
commit(api, s.ID, `{}`)
|
||||
if w := fetchStaged(api, files.gotSource); w.Code != http.StatusOK || w.Body.String() != "abc" {
|
||||
t.Fatalf("refetch: code = %d %q", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a part cut short leaves the session where it was", func(t *testing.T) {
|
||||
api, _, _, _ := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, 10)
|
||||
w := doPart(api.ExternalHandler(), partAt(s.ID, 0), "abc", 5)
|
||||
if w.Code != http.StatusBadRequest || decodeErr(t, w) != "upload_incomplete" {
|
||||
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if w := status(api, s.ID); sessionAnswer(t, w).Received != 0 {
|
||||
t.Fatalf("status = %s", w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("while a part arrives, the session takes nothing else", func(t *testing.T) {
|
||||
api, _, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
h := api.ExternalHandler()
|
||||
s := beginSession(t, api, 6)
|
||||
pr, pw := io.Pipe()
|
||||
arriving := make(chan int)
|
||||
go func() {
|
||||
r := httptest.NewRequest("PUT", partAt(s.ID, 0), pr)
|
||||
r.Header.Set("Content-Type", "application/octet-stream")
|
||||
r.ContentLength = 3
|
||||
w := httptest.NewRecorder()
|
||||
h.ServeHTTP(w, r)
|
||||
arriving <- w.Code
|
||||
}()
|
||||
// The write returns once the handler has read the byte, so the part is
|
||||
// being appended.
|
||||
if _, err := pw.Write([]byte("a")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for name, w := range map[string]*httptest.ResponseRecorder{
|
||||
"another part": doPart(h, partAt(s.ID, 0), "xyz", 3),
|
||||
"cancel": do(h, "DELETE", sessionsRoute+"/"+s.ID, "", nil),
|
||||
} {
|
||||
if w.Code != http.StatusConflict || decodeErr(t, w) != "upload_busy" {
|
||||
t.Errorf("%s: code = %d (%s)", name, w.Code, w.Body.String())
|
||||
}
|
||||
}
|
||||
pw.CloseWithError(errors.New("the client went away"))
|
||||
if code := <-arriving; code != http.StatusBadRequest {
|
||||
t.Fatalf("the broken part answered %d", code)
|
||||
}
|
||||
if w := status(api, s.ID); w.Code != http.StatusOK || sessionAnswer(t, w).Received != 0 || files.calls != 0 {
|
||||
t.Fatalf("status: code = %d calls = %d (%s)", w.Code, files.calls, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("parts refused before a byte is read", func(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
name, target string
|
||||
length int64
|
||||
code int
|
||||
errCode string
|
||||
}{
|
||||
{"no offset", sessionsRoute + "/%s", 3, http.StatusBadRequest, "bad_request"},
|
||||
{"an offset that is no number", sessionsRoute + "/%s?offset=abc", 3, http.StatusBadRequest, "bad_request"},
|
||||
{"no Content-Length", sessionsRoute + "/%s?offset=0", -1, http.StatusLengthRequired, "length_required"},
|
||||
{"a part over the cap", sessionsRoute + "/%s?offset=0", fileedit.PartBytes + 1, http.StatusRequestEntityTooLarge, "part_too_large"},
|
||||
// At the cap it is taken, and the three bytes behind it end short.
|
||||
{"a part at the cap", sessionsRoute + "/%s?offset=0", fileedit.PartBytes, http.StatusBadRequest, "upload_incomplete"},
|
||||
} {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
api, _, _, _ := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, fileedit.PartBytes+10)
|
||||
w := doPart(api.ExternalHandler(), fmt.Sprintf(c.target, s.ID), "abc", c.length)
|
||||
if w.Code != c.code || decodeErr(t, w) != c.errCode {
|
||||
t.Fatalf("code = %d (%s), want %d %s", w.Code, w.Body.String(), c.code, c.errCode)
|
||||
}
|
||||
if w := status(api, s.ID); sessionAnswer(t, w).Received != 0 {
|
||||
t.Fatalf("status = %s", w.Body.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
// The parts need only the account and the server: starting the server midway
|
||||
// costs the upload nothing but the commit.
|
||||
t.Run("parts, status and cancel go on while the server runs", func(t *testing.T) {
|
||||
api, _, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
h := api.ExternalHandler()
|
||||
s := beginSession(t, api, 3)
|
||||
cl.byName["survival"].Ready = true
|
||||
cl.byName["survival"].DesiredState = string(v1alpha1.DesiredRunning)
|
||||
if w := doPart(h, partAt(s.ID, 0), "abc", 3); w.Code != http.StatusOK {
|
||||
t.Fatalf("part: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if w := status(api, s.ID); w.Code != http.StatusOK {
|
||||
t.Fatalf("status: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if w := commit(api, s.ID, `{}`); w.Code != http.StatusConflict || decodeErr(t, w) != "not_stopped" || files.calls != 0 {
|
||||
t.Fatalf("commit: code = %d calls = %d (%s)", w.Code, files.calls, w.Body.String())
|
||||
}
|
||||
if w := do(h, "DELETE", sessionsRoute+"/"+s.ID, "", nil); w.Code != http.StatusNoContent {
|
||||
t.Fatalf("cancel: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if w := status(api, s.ID); w.Code != http.StatusNotFound || decodeErr(t, w) != "upload_not_found" {
|
||||
t.Fatalf("status after cancel: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
stageEmpty(t, api)
|
||||
})
|
||||
|
||||
t.Run("a session answers only the account and server it was begun for", func(t *testing.T) {
|
||||
api, repo, _, files := mkFiles(t)
|
||||
repo.byName["creative"] = &ServerRecord{Name: "creative", OwnerID: "owner1"}
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, 3)
|
||||
|
||||
elsewhere := strings.Replace(sessionsRoute, "survival", "creative", 1) + "/" + s.ID
|
||||
if w := do(api.ExternalHandler(), "GET", elsewhere, "", nil); w.Code != http.StatusNotFound || decodeErr(t, w) != "upload_not_found" {
|
||||
t.Fatalf("another server: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
|
||||
// Staff may reach the server, and still not someone else's session.
|
||||
api.External = staticExternal{p: fileAdmin}
|
||||
h := api.ExternalHandler()
|
||||
for name, w := range map[string]*httptest.ResponseRecorder{
|
||||
"status": status(api, s.ID),
|
||||
"part": doPart(h, partAt(s.ID, 0), "abc", 3),
|
||||
"cancel": do(h, "DELETE", sessionsRoute+"/"+s.ID, "", nil),
|
||||
"commit": commit(api, s.ID, `{}`),
|
||||
} {
|
||||
if w.Code != http.StatusNotFound || decodeErr(t, w) != "upload_not_found" {
|
||||
t.Errorf("admin %s: code = %d (%s)", name, w.Code, w.Body.String())
|
||||
}
|
||||
}
|
||||
if files.calls != 0 {
|
||||
t.Fatal("another account's commit reached the executor")
|
||||
}
|
||||
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
if w := status(api, s.ID); w.Code != http.StatusOK || sessionAnswer(t, w).Received != 0 {
|
||||
t.Fatalf("the owner's session was touched: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a stranger is refused on every session route", func(t *testing.T) {
|
||||
api, _, _, _ := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, 3)
|
||||
api.External = staticExternal{p: fileStranger}
|
||||
h := api.ExternalHandler()
|
||||
for name, w := range map[string]*httptest.ResponseRecorder{
|
||||
"begin": do(h, "POST", sessionsRoute+"?path=a.jar", `{"size":3}`, jsonHeader),
|
||||
"status": status(api, s.ID),
|
||||
"part": doPart(h, partAt(s.ID, 0), "abc", 3),
|
||||
"cancel": do(h, "DELETE", sessionsRoute+"/"+s.ID, "", nil),
|
||||
"commit": commit(api, s.ID, `{}`),
|
||||
} {
|
||||
if w.Code != http.StatusForbidden {
|
||||
t.Errorf("%s: code = %d (%s)", name, w.Code, w.Body.String())
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("begin refused", func(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
name, target, body string
|
||||
setup func(*API)
|
||||
code int
|
||||
errCode string
|
||||
}{
|
||||
{"no size", sessionsRoute + "?path=a.jar", `{}`, nil, http.StatusBadRequest, "bad_request"},
|
||||
{"a negative size", sessionsRoute + "?path=a.jar", `{"size":-1}`, nil, http.StatusBadRequest, "bad_request"},
|
||||
{"no path", sessionsRoute, `{"size":3}`, nil, http.StatusBadRequest, "bad_request"},
|
||||
{"a path leaving the world", sessionsRoute + "?path=../a.jar", `{"size":3}`, nil, http.StatusBadRequest, "bad_path"},
|
||||
{"an absolute path", sessionsRoute + "?path=/etc/a.jar", `{"size":3}`, nil, http.StatusBadRequest, "bad_path"},
|
||||
{"the world folder itself", sessionsRoute + "?path=plugins/..", `{"size":3}`, nil, http.StatusBadRequest, "bad_path"},
|
||||
{"a staging disk at its floor", sessionsRoute + "?path=a.jar", `{"size":3}`,
|
||||
func(a *API) { a.FileStage.MinFree = 1 }, http.StatusInsufficientStorage, "upload_staging_full"},
|
||||
{"no stage", sessionsRoute + "?path=a.jar", `{"size":3}`,
|
||||
func(a *API) { a.FileStage = nil }, http.StatusServiceUnavailable, "files_unavailable"},
|
||||
{"no internal URL", sessionsRoute + "?path=a.jar", `{"size":3}`,
|
||||
func(a *API) { a.InternalBaseURL = "" }, http.StatusServiceUnavailable, "files_unavailable"},
|
||||
{"a caller with no account", sessionsRoute + "?path=a.jar", `{"size":3}`,
|
||||
func(a *API) { a.External = staticExternal{p: &Principal{Role: "admin", ViaAdminAccess: true}} },
|
||||
http.StatusForbidden, "forbidden"},
|
||||
} {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
api, _, _, _ := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
if c.setup != nil {
|
||||
c.setup(api)
|
||||
}
|
||||
w := do(api.ExternalHandler(), "POST", c.target, c.body, jsonHeader)
|
||||
if w.Code != c.code || decodeErr(t, w) != c.errCode {
|
||||
t.Fatalf("code = %d (%s), want %d %s", w.Code, w.Body.String(), c.code, c.errCode)
|
||||
}
|
||||
if api.FileStage != nil {
|
||||
stageEmpty(t, api)
|
||||
}
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a fifth upload at once -> 429", func(t *testing.T) {
|
||||
api, _, _, _ := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
for i := 0; i < fileedit.MaxSessionsPerUser; i++ {
|
||||
beginSession(t, api, 1)
|
||||
}
|
||||
w := do(api.ExternalHandler(), "POST", sessionsRoute+"?path=a.jar", `{"size":1}`, jsonHeader)
|
||||
if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "too_many_uploads" {
|
||||
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a commit or cancel of no session -> 404", func(t *testing.T) {
|
||||
api, _, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
const id = "00112233445566778899aabbccddeeff"
|
||||
for name, w := range map[string]*httptest.ResponseRecorder{
|
||||
"commit": commit(api, id, `{}`),
|
||||
"cancel": do(api.ExternalHandler(), "DELETE", sessionsRoute+"/"+id, "", nil),
|
||||
"part": doPart(api.ExternalHandler(), partAt(id, 0), "abc", 3),
|
||||
} {
|
||||
if w.Code != http.StatusNotFound || decodeErr(t, w) != "upload_not_found" {
|
||||
t.Errorf("%s: code = %d (%s)", name, w.Code, w.Body.String())
|
||||
}
|
||||
}
|
||||
if files.calls != 0 || len(cl.acquired) != 0 {
|
||||
t.Fatalf("calls = %d acquired %v", files.calls, cl.acquired)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a commit needs a body", func(t *testing.T) {
|
||||
api, _, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
s := beginSession(t, api, 0)
|
||||
if w := commit(api, s.ID, ""); w.Code != http.StatusBadRequest || files.calls != 0 {
|
||||
t.Fatalf("code = %d calls = %d (%s)", w.Code, files.calls, w.Body.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// brokenWriter takes the headers and fails every write, as a connection that
|
||||
// dropped once the answer began does.
|
||||
type brokenWriter struct {
|
||||
header http.Header
|
||||
code int
|
||||
}
|
||||
|
||||
func (b *brokenWriter) Header() http.Header { return b.header }
|
||||
func (b *brokenWriter) WriteHeader(code int) { b.code = code }
|
||||
func (b *brokenWriter) Write(p []byte) (int, error) { return 0, errors.New("connection reset") }
|
||||
|
||||
// TestFileOpsWorldGates pins the stopped and world-volume gates on the routes
|
||||
// that begin or start a background op: each refuses before a byte is staged, a
|
||||
// lock is asked for, or a Job is created.
|
||||
func TestFileOpsWorldGates(t *testing.T) {
|
||||
routes := []struct{ name, target, body string }{
|
||||
{"begin", sessionsRoute + "?path=a.jar", `{"size":3}`},
|
||||
{"commit", sessionsRoute + "/00112233445566778899aabbccddeeff/commit", `{}`},
|
||||
{"unzip", unzipRoute + "?path=maps/a.zip", `{}`},
|
||||
}
|
||||
gates := []struct {
|
||||
name, code string
|
||||
set func(*fakeCluster)
|
||||
}{
|
||||
{"running", "not_stopped", func(c *fakeCluster) {
|
||||
c.byName["survival"].Ready = true
|
||||
c.byName["survival"].DesiredState = string(v1alpha1.DesiredRunning)
|
||||
}},
|
||||
{"starting", "not_stopped", func(c *fakeCluster) {
|
||||
c.byName["survival"].DesiredState = string(v1alpha1.DesiredRunning)
|
||||
}},
|
||||
{"no world volume", "no_world_volume", func(c *fakeCluster) { c.noWorld["survival"] = true }},
|
||||
}
|
||||
for _, rt := range routes {
|
||||
for _, g := range gates {
|
||||
t.Run(rt.name+" on a "+g.name+" server", func(t *testing.T) {
|
||||
api, _, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
g.set(cl)
|
||||
w := do(api.ExternalHandler(), "POST", rt.target, rt.body, jsonHeader)
|
||||
if w.Code != http.StatusConflict || decodeErr(t, w) != g.code {
|
||||
t.Fatalf("code = %d (%s), want 409 %s", w.Code, w.Body.String(), g.code)
|
||||
}
|
||||
if files.calls != 0 || len(cl.acquired) != 0 {
|
||||
t.Fatalf("calls = %d acquired %v", files.calls, cl.acquired)
|
||||
}
|
||||
stageEmpty(t, api)
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestFileUnzip(t *testing.T) {
|
||||
unzip := func(api *API, path, body string) *httptest.ResponseRecorder {
|
||||
return do(api.ExternalHandler(), "POST", unzipRoute+"?path="+path, body, jsonHeader)
|
||||
}
|
||||
|
||||
t.Run("starts the Job under the world lock and audits it", func(t *testing.T) {
|
||||
api, repo, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
files.op = fileedit.OpState{ID: "op2", Op: fileedit.OpUnzip, Path: "maps/a.zip",
|
||||
State: fileedit.OpRunning, Started: opStarted}
|
||||
w := unzip(api, "maps/a.zip", `{}`)
|
||||
if w.Code != http.StatusAccepted {
|
||||
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
want := map[string]any{"op": map[string]any{
|
||||
"id": "op2", "op": "unzip", "path": "maps/a.zip", "state": "running",
|
||||
"started_at": "2026-09-28T10:00:00Z", "done": float64(0), "total": float64(0),
|
||||
}}
|
||||
if got := fileAnswer(t, w); !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("answer = %v, want %v", got, want)
|
||||
}
|
||||
if files.calls != 1 || files.gotOp != fileedit.OpUnzip || files.gotServer != "survival" ||
|
||||
files.gotPath != "maps/a.zip" || files.gotOverwrite {
|
||||
t.Fatalf("executor saw calls=%d op %q server %q path %q overwrite %v",
|
||||
files.calls, files.gotOp, files.gotServer, files.gotPath, files.gotOverwrite)
|
||||
}
|
||||
if strings.Join(cl.acquired, ",") != "survival:"+maintenance.KindFileWrite || strings.Join(cl.released, ",") != "survival" {
|
||||
t.Fatalf("lock acquired %v, released %v", cl.acquired, cl.released)
|
||||
}
|
||||
onlyAudit(t, repo, "file.unzip", "survival:maps/a.zip", `{"overwrite":false}`)
|
||||
})
|
||||
|
||||
t.Run("overwrite reaches the executor, and the suffix is any case", func(t *testing.T) {
|
||||
api, repo, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
if w := unzip(api, "maps/A.ZIP", `{"overwrite":true}`); w.Code != http.StatusAccepted || !files.gotOverwrite || files.gotPath != "maps/A.ZIP" {
|
||||
t.Fatalf("code = %d overwrite %v path %q (%s)", w.Code, files.gotOverwrite, files.gotPath, w.Body.String())
|
||||
}
|
||||
onlyAudit(t, repo, "file.unzip", "survival:maps/A.ZIP", `{"overwrite":true}`)
|
||||
})
|
||||
|
||||
t.Run("refused before the lock", func(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
name, path, body string
|
||||
code int
|
||||
errCode string
|
||||
}{
|
||||
{"not a zip", "maps/a.tar.gz", `{}`, http.StatusBadRequest, "bad_path"},
|
||||
{"zip only inside the name", "maps/a.zip.bak", `{}`, http.StatusBadRequest, "bad_path"},
|
||||
{"no path", "", `{}`, http.StatusBadRequest, "bad_request"},
|
||||
{"no body", "maps/a.zip", "", http.StatusBadRequest, "bad_request"},
|
||||
} {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
api, repo, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
w := unzip(api, c.path, c.body)
|
||||
if w.Code != c.code || decodeErr(t, w) != c.errCode {
|
||||
t.Fatalf("code = %d (%s), want %d %s", w.Code, w.Body.String(), c.code, c.errCode)
|
||||
}
|
||||
if files.calls != 0 || len(cl.acquired) != 0 || len(repo.audits) != 0 {
|
||||
t.Fatalf("calls = %d acquired %v audits %+v", files.calls, cl.acquired, repo.audits)
|
||||
}
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a held world -> 409, no Job", func(t *testing.T) {
|
||||
api, _, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
cl.maintErr["survival"] = &MaintenanceBusyError{Kind: maintenance.KindBackup}
|
||||
if w := unzip(api, "maps/a.zip", `{}`); w.Code != http.StatusConflict || decodeErr(t, w) != "maintenance_in_progress" || files.calls != 0 {
|
||||
t.Fatalf("code = %d calls = %d (%s)", w.Code, files.calls, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a Job that could not start is not audited and lets go of the world", func(t *testing.T) {
|
||||
api, repo, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
files.err = errors.New("the cluster said no")
|
||||
w := unzip(api, "maps/a.zip", `{}`)
|
||||
if w.Code != http.StatusInternalServerError || len(repo.audits) != 0 || strings.Join(cl.released, ",") != "survival" {
|
||||
t.Fatalf("code = %d audits %+v released %v (%s)", w.Code, repo.audits, cl.released, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a stranger -> 403, an admin may", func(t *testing.T) {
|
||||
api, _, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileStranger}
|
||||
if w := unzip(api, "maps/a.zip", `{}`); w.Code != http.StatusForbidden || files.calls != 0 {
|
||||
t.Fatalf("stranger: code = %d calls = %d", w.Code, files.calls)
|
||||
}
|
||||
api.External = staticExternal{p: fileAdmin}
|
||||
if w := unzip(api, "maps/a.zip", `{}`); w.Code != http.StatusAccepted || files.calls != 1 {
|
||||
t.Fatalf("admin: code = %d calls = %d (%s)", w.Code, files.calls, w.Body.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestFileOps(t *testing.T) {
|
||||
ended := opStarted.Add(3 * time.Minute)
|
||||
list := func(api *API) *httptest.ResponseRecorder {
|
||||
return do(api.ExternalHandler(), "GET", opsRoute, "", nil)
|
||||
}
|
||||
|
||||
t.Run("each state and failure as the API shows it", func(t *testing.T) {
|
||||
api, _, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
base := func(id, op, state string) fileedit.OpState {
|
||||
s := fileedit.OpState{ID: id, Op: op, Path: "maps/a.zip", State: state, Started: opStarted}
|
||||
if state != fileedit.OpRunning {
|
||||
s.Finished = ended
|
||||
}
|
||||
return s
|
||||
}
|
||||
running := base("a", fileedit.OpUnzip, fileedit.OpRunning)
|
||||
running.Done, running.Total = 40, 100
|
||||
unzipped := base("b", fileedit.OpUnzip, fileedit.OpSucceeded)
|
||||
unzipped.Done, unzipped.Total = 100, 100
|
||||
unzipped.Result = &fileedit.Result{Files: 7, Bytes: 100}
|
||||
uploaded := base("c", fileedit.OpUpload, fileedit.OpSucceeded)
|
||||
uploaded.Result = &fileedit.Result{SHA256: testSum}
|
||||
exists := base("d", fileedit.OpUnzip, fileedit.OpFailed)
|
||||
exists.Result = &fileedit.Result{Code: fileedit.CodeExists, Error: "2 files are already there",
|
||||
Conflicts: []string{"maps/level.dat", "maps/r.0.0.mca"}, ConflictCount: 2}
|
||||
full := base("e", fileedit.OpUpload, fileedit.OpFailed)
|
||||
full.Result = &fileedit.Result{Code: fileedit.CodeNoSpace, Error: "no room", Need: 900, Avail: 100}
|
||||
changed := base("f", fileedit.OpUpload, fileedit.OpFailed)
|
||||
changed.Result = &fileedit.Result{Code: fileedit.CodeConflict, Error: "the bytes changed"}
|
||||
unsafe := base("g", fileedit.OpUnzip, fileedit.OpFailed)
|
||||
unsafe.Result = &fileedit.Result{Code: fileedit.CodeArchiveUnsafe, Error: "leaves the folder", Entry: "../x",
|
||||
Files: 3, Bytes: 9}
|
||||
deadline := base("h", fileedit.OpUnzip, fileedit.OpFailed)
|
||||
deadline.Reason = "DeadlineExceeded"
|
||||
killed := base("i", fileedit.OpUpload, fileedit.OpFailed)
|
||||
killed.Reason = "BackoffLimitExceeded"
|
||||
files.ops = []fileedit.OpState{running, unzipped, uploaded, exists, full, changed, unsafe, deadline, killed}
|
||||
|
||||
w := list(api)
|
||||
if w.Code != http.StatusOK || files.gotServer != "survival" {
|
||||
t.Fatalf("code = %d server %q (%s)", w.Code, files.gotServer, w.Body.String())
|
||||
}
|
||||
op := func(id, kind, state string, extra map[string]any) map[string]any {
|
||||
m := map[string]any{"id": id, "op": kind, "path": "maps/a.zip", "state": state,
|
||||
"started_at": "2026-09-28T10:00:00Z", "done": float64(0), "total": float64(0)}
|
||||
if state != "running" {
|
||||
m["finished_at"] = "2026-09-28T10:03:00Z"
|
||||
}
|
||||
for k, v := range extra {
|
||||
m[k] = v
|
||||
}
|
||||
return m
|
||||
}
|
||||
failure := func(code, msg string, extra map[string]any) map[string]any {
|
||||
m := map[string]any{"code": code, "message": msg}
|
||||
for k, v := range extra {
|
||||
m[k] = v
|
||||
}
|
||||
return map[string]any{"error": m}
|
||||
}
|
||||
want := map[string]any{"ops": []any{
|
||||
op("a", "unzip", "running", map[string]any{"done": float64(40), "total": float64(100)}),
|
||||
op("b", "unzip", "succeeded", map[string]any{"done": float64(100), "total": float64(100),
|
||||
"files": float64(7), "bytes": float64(100)}),
|
||||
op("c", "upload", "succeeded", nil),
|
||||
op("d", "unzip", "failed", failure("file_exists", "2 files are already there", map[string]any{
|
||||
"conflicts": []any{"maps/level.dat", "maps/r.0.0.mca"}, "conflict_count": float64(2)})),
|
||||
op("e", "upload", "failed", failure("volume_full", "no room", map[string]any{
|
||||
"need": float64(900), "avail": float64(100)})),
|
||||
op("f", "upload", "failed", failure("file_changed", "the bytes changed", nil)),
|
||||
op("g", "unzip", "failed", failure("archive_unsafe", "leaves the folder", map[string]any{"entry": "../x"})),
|
||||
op("h", "unzip", "failed", failure("job_failed",
|
||||
"the file operation ran out of time (DeadlineExceeded); run it again", nil)),
|
||||
op("i", "upload", "failed", failure("job_failed",
|
||||
"the file operation stopped before it could report how it went (BackoffLimitExceeded); run it again", nil)),
|
||||
}}
|
||||
if got := fileAnswer(t, w); !reflect.DeepEqual(got, want) {
|
||||
gotJSON, _ := json.MarshalIndent(got, "", " ")
|
||||
wantJSON, _ := json.MarshalIndent(want, "", " ")
|
||||
t.Fatalf("ops =\n%s\nwant\n%s", gotJSON, wantJSON)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("none is an empty list", func(t *testing.T) {
|
||||
api, _, _, _ := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
if w := list(api); w.Code != http.StatusOK || strings.TrimSpace(w.Body.String()) != `{"ops":[]}` {
|
||||
t.Fatalf("code = %d %s", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("answers while the server runs", func(t *testing.T) {
|
||||
api, _, cl, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
cl.byName["survival"].Ready = true
|
||||
cl.byName["survival"].DesiredState = string(v1alpha1.DesiredRunning)
|
||||
cl.noWorld["survival"] = true
|
||||
if w := list(api); w.Code != http.StatusOK || files.calls != 1 {
|
||||
t.Fatalf("code = %d calls = %d (%s)", w.Code, files.calls, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("who may look", func(t *testing.T) {
|
||||
api, repo, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileStranger}
|
||||
if w := list(api); w.Code != http.StatusForbidden || files.calls != 0 {
|
||||
t.Fatalf("stranger: code = %d calls = %d", w.Code, files.calls)
|
||||
}
|
||||
repo.byName["survival"].OwnerID = "someone-else"
|
||||
api.External = staticExternal{p: fileAdmin}
|
||||
if w := list(api); w.Code != http.StatusOK || files.calls != 1 {
|
||||
t.Fatalf("admin: code = %d calls = %d (%s)", w.Code, files.calls, w.Body.String())
|
||||
}
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
if w := do(api.ExternalHandler(), "GET", strings.Replace(opsRoute, "survival", "missing", 1), "", nil); w.Code != http.StatusNotFound {
|
||||
t.Fatalf("unknown server: code = %d", w.Code)
|
||||
}
|
||||
if w := do(api.ExternalHandler(), "GET", strings.Replace(opsRoute, "survival", "X", 1), "", nil); w.Code != http.StatusBadRequest || decodeErr(t, w) != "bad_name" {
|
||||
t.Fatalf("bad name: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("no executor -> 503, a failing one -> 500", func(t *testing.T) {
|
||||
api, _, _, files := mkFiles(t)
|
||||
api.External = staticExternal{p: fileOwner}
|
||||
files.err = errors.New("the cluster said no")
|
||||
if w := list(api); w.Code != http.StatusInternalServerError {
|
||||
t.Fatalf("failing: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
api.Files = nil
|
||||
if w := list(api); w.Code != http.StatusServiceUnavailable || decodeErr(t, w) != "files_unavailable" {
|
||||
t.Fatalf("nil: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -12,7 +12,6 @@ import (
|
||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||
"felis.lolicon.best/internal/fileedit"
|
||||
"felis.lolicon.best/internal/maintenance"
|
||||
"felis.lolicon.best/internal/naming"
|
||||
)
|
||||
|
||||
// FileEditor is the server-file-editor surface the API depends on: list a
|
||||
@@ -39,18 +38,23 @@ import (
|
||||
// 507 / 409.
|
||||
//
|
||||
// Read and Write return the file's SHA-256 (hex); Upload lands exactly the bytes
|
||||
// src describes or fails. Write's expect is the
|
||||
// src describes or fails. StartUpload and StartUnzip are the two that can run
|
||||
// longer than a request: they start their Job and return, and Ops reports it
|
||||
// (handlers_fileops.go). List also reports the room left on the volume. Write's expect is the
|
||||
// hash a client read the file at; when set, a file that changed since is refused
|
||||
// with ErrConflict instead of being overwritten. createOnly and a false overwrite
|
||||
// refuse an existing path with ErrExists.
|
||||
type FileEditor interface {
|
||||
List(ctx context.Context, server, path string) (entries []fileedit.Entry, truncated bool, err error)
|
||||
List(ctx context.Context, server, path string) (fileedit.Listing, error)
|
||||
Read(ctx context.Context, server, path string) (content []byte, sha256 string, err error)
|
||||
Write(ctx context.Context, server, path string, content []byte, expect string, createOnly bool) (sha256 string, err error)
|
||||
Mkdir(ctx context.Context, server, path string) error
|
||||
Delete(ctx context.Context, server, path string) error
|
||||
Rename(ctx context.Context, server, path, to string) error
|
||||
Upload(ctx context.Context, server, path string, src fileedit.UploadSource, overwrite bool) error
|
||||
StartUpload(ctx context.Context, server, path string, src fileedit.UploadSource, overwrite bool) (fileedit.OpState, error)
|
||||
StartUnzip(ctx context.Context, server, path string, overwrite bool) (fileedit.OpState, error)
|
||||
Ops(ctx context.Context, server string) ([]fileedit.OpState, error)
|
||||
}
|
||||
|
||||
// writeFileRequest is the PUT /servers/{name}/file body. Content is []byte, so
|
||||
@@ -96,16 +100,18 @@ func (a *API) handleListFiles(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
path := r.URL.Query().Get("path")
|
||||
|
||||
entries, truncated, err := a.Files.List(r.Context(), name, path)
|
||||
ls, err := a.Files.List(r.Context(), name, path)
|
||||
if err != nil {
|
||||
writeFileEditError(w, r, err)
|
||||
return
|
||||
}
|
||||
if entries == nil {
|
||||
entries = []fileedit.Entry{} // an empty directory is [], never null
|
||||
if ls.Entries == nil {
|
||||
ls.Entries = []fileedit.Entry{} // an empty directory is [], never null
|
||||
}
|
||||
// free_bytes lets the panel refuse an upload the volume cannot take before
|
||||
// sending any of it; the Job that lands it checks again.
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"path": path, "entries": entries, "truncated": truncated,
|
||||
"path": path, "entries": ls.Entries, "truncated": ls.Truncated, "free_bytes": ls.Free,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -351,7 +357,8 @@ func (a *API) handleUploadFile(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
if r.ContentLength > fileedit.MaxUploadBytes {
|
||||
writeError(w, r, newError(http.StatusRequestEntityTooLarge, "too_large",
|
||||
"the file is %d bytes; uploads are at most %d", r.ContentLength, fileedit.MaxUploadBytes))
|
||||
"the file is %d bytes; one request carries at most %d, so send it as an upload session (POST files/uploads)",
|
||||
r.ContentLength, fileedit.MaxUploadBytes))
|
||||
return
|
||||
}
|
||||
overwrite := r.URL.Query().Get("overwrite") == "true"
|
||||
@@ -397,7 +404,8 @@ func (a *API) handleUploadFile(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
// handleInternalFileUpload serves GET /api/v1/internal/file-uploads/{id} — the
|
||||
// staged bytes of one upload, to the one Job created to land them. It is Public on
|
||||
// staged bytes of one upload, to the one Job created to land them. A file sent
|
||||
// in parts (handlers_fileops.go) is served the same way. It is Public on
|
||||
// the internal face: the Job holds no service token (it holds no credential at
|
||||
// all), so the bearer token minted with the upload is the whole check, and it
|
||||
// opens that upload once. An unknown id, a wrong token and a spent one are the
|
||||
@@ -421,7 +429,11 @@ func (a *API) handleInternalFileUpload(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "application/octet-stream")
|
||||
w.Header().Set("Content-Length", strconv.FormatInt(size, 10))
|
||||
w.WriteHeader(http.StatusOK)
|
||||
_, _ = io.Copy(w, f)
|
||||
// A session sent whole is on the Job's side now; one cut short stays, so
|
||||
// committing it again does not mean sending it again.
|
||||
if n, err := io.Copy(w, f); err == nil && n == size {
|
||||
a.FileStage.Served(r.PathValue("id"))
|
||||
}
|
||||
}
|
||||
|
||||
// auditFile records a file change. The target is "<server>:<path>", as file.write
|
||||
@@ -463,20 +475,8 @@ var sha256Hex = regexp.MustCompile(`^[0-9a-f]{64}$`)
|
||||
// It returns the validated server name and false if it has already written a
|
||||
// response.
|
||||
func (a *API) authorizeFileOp(w http.ResponseWriter, r *http.Request) (string, bool) {
|
||||
name := r.PathValue("name")
|
||||
if err := naming.ValidateServerName(name); err != nil {
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_name", "invalid server name: %v", err))
|
||||
return "", false
|
||||
}
|
||||
|
||||
p := principalFromContext(r.Context())
|
||||
rec, err := a.Repo.ServerByName(r.Context(), name)
|
||||
if err != nil {
|
||||
a.writeLookupError(w, r, err)
|
||||
return "", false
|
||||
}
|
||||
if !a.isOwnerOrAdmin(p, rec) {
|
||||
writeError(w, r, errForbidden)
|
||||
name, ok := a.authorizeServerFiles(w, r) // ① ② ③
|
||||
if !ok {
|
||||
return "", false
|
||||
}
|
||||
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||
"felis.lolicon.best/internal/fileedit"
|
||||
@@ -42,14 +43,19 @@ type fakeFileEditor struct {
|
||||
|
||||
entries []fileedit.Entry
|
||||
truncated bool
|
||||
free int64
|
||||
content []byte
|
||||
sum string
|
||||
|
||||
// op is what StartUpload and StartUnzip answer (started), ops what Ops does.
|
||||
op fileedit.OpState
|
||||
ops []fileedit.OpState
|
||||
}
|
||||
|
||||
func (f *fakeFileEditor) List(_ context.Context, server, path string) ([]fileedit.Entry, bool, error) {
|
||||
func (f *fakeFileEditor) List(_ context.Context, server, path string) (fileedit.Listing, error) {
|
||||
f.calls++
|
||||
f.gotServer, f.gotPath = server, path
|
||||
return f.entries, f.truncated, f.err
|
||||
return fileedit.Listing{Entries: f.entries, Truncated: f.truncated, Free: f.free}, f.err
|
||||
}
|
||||
|
||||
func (f *fakeFileEditor) Read(_ context.Context, server, path string) ([]byte, string, error) {
|
||||
@@ -92,6 +98,33 @@ func (f *fakeFileEditor) Upload(_ context.Context, server, path string, src file
|
||||
return f.err
|
||||
}
|
||||
|
||||
func (f *fakeFileEditor) StartUpload(_ context.Context, server, path string, src fileedit.UploadSource, overwrite bool) (fileedit.OpState, error) {
|
||||
f.calls++
|
||||
f.gotOp, f.gotServer, f.gotPath, f.gotSource, f.gotOverwrite = fileedit.OpUpload, server, path, src, overwrite
|
||||
return f.started(fileedit.OpUpload, path), f.err
|
||||
}
|
||||
|
||||
func (f *fakeFileEditor) StartUnzip(_ context.Context, server, path string, overwrite bool) (fileedit.OpState, error) {
|
||||
f.calls++
|
||||
f.gotOp, f.gotServer, f.gotPath, f.gotOverwrite = fileedit.OpUnzip, server, path, overwrite
|
||||
return f.started(fileedit.OpUnzip, path), f.err
|
||||
}
|
||||
|
||||
// started is the op StartUpload and StartUnzip answer: f.op when a test set
|
||||
// one, else a running op as the real Editor answers it.
|
||||
func (f *fakeFileEditor) started(op, path string) fileedit.OpState {
|
||||
if f.op.ID != "" {
|
||||
return f.op
|
||||
}
|
||||
return fileedit.OpState{ID: "op" + strconv.Itoa(f.calls), Op: op, Path: path, State: fileedit.OpRunning, Started: time.Now()}
|
||||
}
|
||||
|
||||
func (f *fakeFileEditor) Ops(_ context.Context, server string) ([]fileedit.OpState, error) {
|
||||
f.calls++
|
||||
f.gotServer = server
|
||||
return f.ops, f.err
|
||||
}
|
||||
|
||||
// fileRouteHeader is the Content-Type a file route's body goes with: raw bytes
|
||||
// for an upload, JSON for any other body.
|
||||
func fileRouteHeader(name, body string) map[string]string {
|
||||
@@ -340,6 +373,7 @@ func TestFileEditorHandlers(t *testing.T) {
|
||||
api, _, _, files := mkFiles(t)
|
||||
files.entries = []fileedit.Entry{{Name: "paper.yml", Size: 12}, {Name: "sub", IsDir: true}}
|
||||
files.truncated = true
|
||||
files.free = 5 << 30
|
||||
api.External = staticExternal{p: owner}
|
||||
|
||||
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/files?path=config", "", nil)
|
||||
@@ -350,6 +384,7 @@ func TestFileEditorHandlers(t *testing.T) {
|
||||
Path string `json:"path"`
|
||||
Entries []fileedit.Entry `json:"entries"`
|
||||
Truncated bool `json:"truncated"`
|
||||
FreeBytes int64 `json:"free_bytes"`
|
||||
}
|
||||
if err := json.Unmarshal(w.Body.Bytes(), &resp); err != nil {
|
||||
t.Fatalf("body not JSON: %v (%s)", err, w.Body.String())
|
||||
@@ -357,7 +392,7 @@ func TestFileEditorHandlers(t *testing.T) {
|
||||
if files.gotPath != "config" {
|
||||
t.Fatalf("executor saw path %q, want the query value verbatim", files.gotPath)
|
||||
}
|
||||
if resp.Path != "config" || len(resp.Entries) != 2 || !resp.Truncated {
|
||||
if resp.Path != "config" || len(resp.Entries) != 2 || !resp.Truncated || resp.FreeBytes != 5<<30 {
|
||||
t.Fatalf("unexpected response %+v", resp)
|
||||
}
|
||||
})
|
||||
|
||||
@@ -14,7 +14,7 @@ import (
|
||||
// object an operator with kubectl could read. The route is the API-side outlet.
|
||||
type AsyncJob struct {
|
||||
Name string `json:"name"`
|
||||
Kind string `json:"kind"` // "backup" | "restore" | "export_world" | "export_backup"
|
||||
Kind string `json:"kind"` // "backup" | "restore" | "export_world" | "export_backup" | "export_files"
|
||||
State string `json:"state"` // "running" | "succeeded" | "failed"
|
||||
Message string `json:"message,omitempty"`
|
||||
StartedAt time.Time `json:"started_at,omitzero"`
|
||||
|
||||
@@ -230,11 +230,15 @@ func jobOutcome(j *batchv1.Job) (AsyncJob, bool) {
|
||||
case jobManagedByRestore:
|
||||
kind = "restore"
|
||||
case jobManagedByExport:
|
||||
// As maintenance.JobKind reads it: only a Job that says it reads a
|
||||
// backup is not a world export.
|
||||
kind = "export_world"
|
||||
if j.Labels[maintenance.LabelExportMode] == maintenance.ExportModeBackup {
|
||||
// As maintenance.JobKind reads it: a Job that names no mode this build
|
||||
// knows reads as a world export, the kind that holds the world.
|
||||
switch j.Labels[maintenance.LabelExportMode] {
|
||||
case maintenance.ExportModeBackup:
|
||||
kind = "export_backup"
|
||||
case maintenance.ExportModeFiles:
|
||||
kind = "export_files"
|
||||
default:
|
||||
kind = "export_world"
|
||||
}
|
||||
default:
|
||||
return AsyncJob{}, false
|
||||
|
||||
@@ -42,6 +42,9 @@ func TestOpenAPISchemasMatchWireStructs(t *testing.T) {
|
||||
"BackupView": BackupView{},
|
||||
"ExportTicket": exportTicketView{},
|
||||
"ExportStatus": exportStatusView{},
|
||||
"FileUploadSession": fileSessionView{},
|
||||
"FileOp": fileOpView{},
|
||||
"FileOpError": fileOpError{},
|
||||
"Schedule": Schedule{},
|
||||
"Build": build.Build{},
|
||||
"Image": build.Image{},
|
||||
|
||||
Reference in new issue
Block a user