feat(backups): 主人和管理员可下载单个备份、导出停服世界,字节经一次性票据从导出 Job 流式转给浏览器,备份按 sha256 核对

This commit is contained in:
Lemon-miaow committed 2026-09-28 00:58:42 +08:00
1 parent 34b81ee8fe
commit 5dffadb40d
35 files changed
+4018 -70

No files matched your search

+24
View File
@@ -98,6 +98,12 @@ type API struct {
FileStage *fileedit.Stage
InternalBaseURL string
// Exporter starts the Job behind a world or backup download (exports.go),
// which PUTs the archive to InternalBaseURL. Optional like Restorer: the
// export routes report 503 unless both are set, after the owner-or-admin,
// backup and stopped gates.
Exporter Exporter
// Schedules stores the servers' scheduled tasks (schedules.go), which
// RunSchedules fires. Optional: when nil the schedule routes report 503 and
// RunSchedules does nothing.
@@ -222,6 +228,9 @@ type API struct {
streamCapOnce sync.Once
streamCap *streamLimiter
exportsOnce sync.Once
exports *exportRegistry
authDoorOnce sync.Once
authDoorBuckets *bucketSet
mailOnce sync.Once
@@ -459,6 +468,13 @@ func (a *API) internalAPIRoutes() []apiRoute {
// because that Job holds no service token; the one-time bearer token minted
// with the upload is the check (handlers_files.go).
{Method: "GET", Pattern: "/api/v1/internal/file-uploads/{id}", Public: true, h: a.handleInternalFileUpload},
// An export Job's archive, held open until the owner's browser downloads
// it. Public for the same reason as file uploads: the Job holds no service
// token, and the one-time bearer token minted with the export is the check
// (exports.go). The body is read at the browser's pace, so the handler
// lifts the minimum-rate body deadline and applies its own stall bound.
{Method: "PUT", Pattern: "/api/v1/internal/exports/{id}", Public: true, h: a.handleInternalExportUpload},
}
}
@@ -558,6 +574,14 @@ func (a *API) externalAPIRoutes() []apiRoute {
{Method: "GET", Pattern: "/api/v1/servers/{name}/jobs", h: a.handleServerJobs},
{Method: "POST", Pattern: "/api/v1/servers/{name}/restore-backup", h: a.handleRestoreBackup},
{Method: "POST", Pattern: "/api/v1/servers/{name}/backup", h: a.handleBackupNow},
// World export (exports.go): download a backup, or a stopped server's world
// as it is now, straight to the browser. App-tier with the restore gate
// inside each start route; the ticket routes answer only the user who
// started the export.
{Method: "POST", Pattern: "/api/v1/servers/{name}/backups/{id}/export", h: a.handleExportBackup},
{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,
// exactly like the backup pair above and for the same reason — every route
+4 -3
View File
@@ -255,8 +255,9 @@ type fakeSession struct {
// fakeBackup mirrors a world_backups row: the client-facing view plus the
// server-side backup_ref the list queries never expose.
type fakeBackup struct {
view BackupView
ref string
view BackupView
ref string
sha256 string
}
// fakeLinkCode mirrors an account_link_codes row.
@@ -1051,7 +1052,7 @@ func (f *fakeRepo) BackupByID(_ context.Context, id string) (*BackupRecord, erro
return &BackupRecord{
ID: b.view.ID, ServerName: b.view.ServerName,
FormerOwner: b.view.FormerOwner, BackupRef: b.ref,
SizeBytes: b.view.SizeBytes, Corrupt: b.view.Corrupt,
SizeBytes: b.view.SizeBytes, Corrupt: b.view.Corrupt, SHA256: b.sha256,
}, nil
}
}
+677
View File
@@ -0,0 +1,677 @@
package api
import (
"context"
"crypto/rand"
"crypto/sha256"
"crypto/subtle"
"encoding/hex"
"errors"
"fmt"
"hash"
"io"
"log"
"mime"
"net/http"
"strconv"
"strings"
"sync"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"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:
//
// 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.
// 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.
// 3. The panel polls GET /exports/{ticket} until it reads ready, then points a
// hidden <a download> at GET /exports/{ticket}/download, which claims the
// waiting PUT and copies its body into the response 32 KiB at a time. The
// 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.
//
// 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.
// Exporter starts the Job that archives a world or a backup and hands it to the
// internal upload route (internal/worldexport). Optional: when nil the export
// routes answer 503.
type Exporter interface {
Start(ctx context.Context, r worldexport.Request) (job string, err error)
}
// 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.
const (
exportMaxActive = 2 // admitted and not yet over, install-wide
exportMaxPerUser = 1
exportPerHour = 6 // started by one user in any hour
exportCopyBuffer = 32 << 10
)
// Timings. Vars only so a test can shrink them.
var (
// exportClaimTTL is how long the Job's upload waits for the browser.
exportClaimTTL = 90 * time.Second
// exportPendingTTL is how long an export may take to reach ready: the Job
// is scheduled, pulls its image and connects well inside it.
exportPendingTTL = 10 * time.Minute
// exportStall is the longest either end may go without moving a byte.
exportStall = 2 * time.Minute
// exportKeepSpent is how long a finished ticket still answers 410
// export_expired before it reads as unknown.
exportKeepSpent = 10 * time.Minute
)
// States of an export.
const (
exportPending = "pending" // the Job has not connected yet
exportReady = "ready" // its upload is waiting for the browser
exportStreaming = "streaming" // the browser is downloading
exportFailed = "failed" // the Job died before it connected
exportSpent = "spent" // downloaded, or given up
)
// exportTicketView answers both export routes.
type exportTicketView struct {
Ticket string `json:"ticket"`
State string `json:"state"`
Filename string `json:"filename"`
}
// exportStatusView answers GET /exports/{ticket}.
type exportStatusView struct {
State string `json:"state"`
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")
}
func errNoExport() error {
return newError(http.StatusNotFound, "not_found", "no such export")
}
type exportEntry struct {
ticket string
id string
tokenHash [32]byte
userID string
server string
mode string
sha256 string // what a backup must hash to; empty for a world
filename 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.
type exportUpload struct {
body io.Reader
size int64 // -1 when the Job streams it chunked
claimed chan struct{}
done chan error // how the download ended; buffered
}
func (e *exportEntry) active() bool {
return e.state == exportPending || e.state == exportReady || e.state == exportStreaming
}
type exportRegistry struct {
mu sync.Mutex
byTicket map[string]*exportEntry
byID map[string]*exportEntry
starts map[string][]time.Time // per user, oldest first
}
func (a *API) exportTickets() *exportRegistry {
a.exportsOnce.Do(func() {
a.exports = &exportRegistry{
byTicket: map[string]*exportEntry{},
byID: map[string]*exportEntry{},
starts: map[string][]time.Time{},
}
})
return a.exports
}
// sweepLocked expires what has waited too long and forgets what ended long ago.
func (g *exportRegistry) sweepLocked(now time.Time) {
for t, e := range g.byTicket {
switch {
case e.state == exportPending && now.Sub(e.at) >= exportPendingTTL:
e.state, e.at = exportSpent, now
case !e.active() && now.Sub(e.at) >= exportKeepSpent:
delete(g.byTicket, t)
delete(g.byID, e.id)
}
}
for u, ts := range g.starts {
for len(ts) > 0 && now.Sub(ts[0]) >= time.Hour {
ts = ts[1:]
}
if len(ts) == 0 {
delete(g.starts, u)
} else {
g.starts[u] = ts
}
}
}
func randomHex(n int) string {
b := make([]byte, n)
_, _ = rand.Read(b)
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.
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() {
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))
}
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)
return token, nil
}
// drop forgets an export that never got its Job, and gives its start back.
func (g *exportRegistry) drop(e *exportEntry) {
g.mu.Lock()
defer g.mu.Unlock()
delete(g.byTicket, e.ticket)
delete(g.byID, e.id)
ts := g.starts[e.userID]
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:]...)
break
}
}
}
func (g *exportRegistry) started(e *exportEntry, job string) {
g.mu.Lock()
defer g.mu.Unlock()
e.job = job
}
// lookup returns a copy of userID's export by its ticket. Another user's reads
// as unknown, and a spent one as export_expired.
func (g *exportRegistry) lookup(ticket, userID string, now time.Time) (exportEntry, error) {
g.mu.Lock()
defer g.mu.Unlock()
g.sweepLocked(now)
e, ok := g.byTicket[ticket]
switch {
case !ok || e.userID != userID:
return exportEntry{}, errNoExport()
case e.state == exportSpent || e.state == exportStreaming:
return exportEntry{}, errExportExpired()
}
return *e, nil
}
// fail records that the export's Job died before it connected, and reports
// false when the export has moved on since.
func (g *exportRegistry) fail(ticket, message string, now time.Time) bool {
g.mu.Lock()
defer g.mu.Unlock()
e, ok := g.byTicket[ticket]
if !ok || e.state != exportPending {
return false
}
e.state, e.message, e.at = exportFailed, message, now
return true
}
// arrive parks the Job's upload on the export the id and token open. Unknown
// ids, wrong tokens and exports past pending all read the same.
func (g *exportRegistry) arrive(id, token string, up *exportUpload, now time.Time) bool {
g.mu.Lock()
defer g.mu.Unlock()
g.sweepLocked(now)
e, ok := g.byID[id]
if !ok {
return false
}
sum := sha256.Sum256([]byte(token))
if subtle.ConstantTimeCompare(sum[:], e.tokenHash[:]) != 1 || e.state != exportPending {
return false
}
e.state, e.at, e.upload = exportReady, now, up
return true
}
// giveUp spends an export whose upload nobody claimed, and reports false when a
// claim got there first.
func (g *exportRegistry) giveUp(up *exportUpload, now time.Time) bool {
g.mu.Lock()
defer g.mu.Unlock()
for _, e := range g.byID {
if e.upload == up && e.state == exportReady {
e.state, e.at = exportSpent, now
return true
}
}
return false
}
// claim hands userID's ready export to its download: from here on the ticket is
// spent whatever happens to the download.
func (g *exportRegistry) claim(ticket, userID string, now time.Time) (exportEntry, error) {
g.mu.Lock()
defer g.mu.Unlock()
g.sweepLocked(now)
e, ok := g.byTicket[ticket]
switch {
case !ok || e.userID != userID:
return exportEntry{}, errNoExport()
case e.state == exportPending:
return exportEntry{}, newError(http.StatusConflict, "export_not_ready",
"the export is still being prepared; wait until its status reads ready")
case e.state != exportReady:
return exportEntry{}, errExportExpired()
}
e.state, e.at = exportStreaming, now
close(e.upload.claimed)
return *e, nil
}
func (g *exportRegistry) finish(ticket string, now time.Time) {
g.mu.Lock()
defer g.mu.Unlock()
if e, ok := g.byTicket[ticket]; ok {
e.state, e.at = exportSpent, now
}
}
// handleExportBackup starts the download of one backup (POST
// /api/v1/servers/{name}/backups/{id}/export). The gate is restore's and
// delete's together:
//
// ① name validation; ② an unknown server is 404; ③ owner-or-admin, else 403
// ④ the backup must exist, and a non-admin must be its former owner: any
// other id reads as unknown (404 no_backup), as it does in their list
// ⑤ cross-server guard: the backup must be this server's (403)
// ⑥ an archive that failed its read-back is 409 backup_corrupt
//
// It needs no stopped server and takes no world lock: the Job reads the backup
// store only.
func (a *API) handleExportBackup(w http.ResponseWriter, r *http.Request) {
p := principalFromContext(r.Context())
name, rec, ok := a.exportGate(w, r)
if !ok {
return
}
backup, err := a.Repo.BackupByID(r.Context(), r.PathValue("id"))
if err == nil && !p.IsAdmin() && backup.FormerOwner != p.UserID {
err = ErrNotFound
}
if errors.Is(err, ErrNotFound) {
writeError(w, r, newError(http.StatusNotFound, "no_backup", "no matching backup exists"))
return
}
if err != nil {
writeError(w, r, err)
return
}
if backup.ServerName != name {
writeError(w, r, errForbidden)
return
}
if backup.Corrupt {
writeError(w, r, newError(http.StatusConflict, "backup_corrupt",
"this backup did not read back intact and cannot be exported; pick another"))
return
}
if a.Exporter == nil || a.InternalBaseURL == "" {
writeError(w, r, errExportUnavailable())
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}
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) {
return
}
a.auditEntry(r, AuditEntry{Actor: auditActor(p), ActorUserID: p.UserID, Action: "backup.export", ServerName: rec.Name,
Payload: auditPayload(map[string]any{"backup_id": backup.ID, "size_bytes": backup.SizeBytes})})
writeJSON(w, http.StatusAccepted, exportTicketView{Ticket: e.ticket, State: exportPending, Filename: e.filename})
}
// handleExportWorld starts the download of a stopped server's world as it is
// now (POST /api/v1/servers/{name}/world/export). The gate is a backup's: owner
// or admin, the server fully stopped, a world volume to read, and the
// world-volume lock, which the export Job then holds (as KindExport) until the
// download ends, so the server cannot start under it and tear the archive.
func (a *API) handleExportWorld(w http.ResponseWriter, r *http.Request) {
p := principalFromContext(r.Context())
name, rec, ok := a.exportGate(w, r)
if !ok {
return
}
const notStopped = "stop the server before exporting its world"
info, err := a.Cluster.GetServer(r.Context(), name)
if err != nil {
a.writeLookupError(w, r, err)
return
}
if info.Ready || info.DesiredState != string(v1alpha1.DesiredStopped) {
writeError(w, r, newError(http.StatusConflict, "not_stopped", notStopped))
return
}
if exists, err := a.Cluster.WorldVolumeExists(r.Context(), name); err != nil {
writeError(w, r, err)
return
} else if !exists {
writeError(w, r, errNoWorldVolume())
return
}
if a.Exporter == nil || a.InternalBaseURL == "" {
writeError(w, r, errExportUnavailable())
return
}
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"))}
token, err := reg.admit(e, a.now())
if err != nil {
writeError(w, r, err)
return
}
release, ok := a.acquireWorld(w, r, name, maintenance.KindExport, notStopped)
if !ok {
reg.drop(e)
return
}
defer release()
if !a.startExport(w, r, e, token, "") {
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})
}
func errExportUnavailable() error {
return newError(http.StatusServiceUnavailable, "export_unavailable", "world export is not configured")
}
// exportGate is the front half both export routes share: a valid name, a known
// server, and a caller who owns it or is an admin. The ticket is bound to the
// caller's account, so a principal without one cannot start an export.
func (a *API) exportGate(w http.ResponseWriter, r *http.Request) (string, *ServerRecord, bool) {
p := principalFromContext(r.Context())
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 "", nil, false
}
rec, err := a.Repo.ServerByName(r.Context(), name)
if err != nil {
a.writeLookupError(w, r, err)
return "", nil, false
}
if p.UserID == "" || !a.isOwnerOrAdmin(p, rec) {
writeError(w, r, errForbidden)
return "", nil, false
}
return name, rec, true
}
// 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,
})
if err != nil {
a.exportTickets().drop(e)
writeError(w, r, err)
return false
}
a.exportTickets().started(e, job)
return true
}
// handleExportStatus answers GET /api/v1/exports/{ticket}: pending while the
// Job gets going, ready once its upload waits for the browser, failed (with the
// Job's error) when it died first. A downloaded or expired ticket is 410
// export_expired; an unknown one, or another user's, is 404.
func (a *API) handleExportStatus(w http.ResponseWriter, r *http.Request) {
p := principalFromContext(r.Context())
reg := a.exportTickets()
e, err := reg.lookup(r.PathValue("ticket"), p.UserID, a.now())
if err != nil {
writeError(w, r, err)
return
}
// A Job list that fails leaves the export pending: the next poll asks again,
// and exportPendingTTL ends it either way.
if e.state == exportPending && e.job != "" && a.JobStatus != nil {
jobs, _ := a.JobStatus.LatestJobs(r.Context(), e.server)
for _, j := range jobs {
if j.Name == e.job && j.State == "failed" && reg.fail(e.ticket, j.Message, a.now()) {
e.state, e.message = exportFailed, j.Message
}
}
}
writeJSON(w, http.StatusOK, exportStatusView{State: e.state, Message: e.message})
}
// handleExportDownload streams a ready export to its browser (GET
// /api/v1/exports/{ticket}/download). The claim spends the ticket before the
// first byte moves, so a second tab, a retry or a HEAD probe never gets a
// second copy; a HEAD is refused outright, since claiming on it would spend the
// ticket on a response with no body.
func (a *API) handleExportDownload(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
w.Header().Set("Allow", http.MethodGet)
writeError(w, r, newError(http.StatusMethodNotAllowed, "method_not_allowed",
"%s is not allowed here; use GET", r.Method))
return
}
p := principalFromContext(r.Context())
reg := a.exportTickets()
e, err := reg.claim(r.PathValue("ticket"), p.UserID, a.now())
if err != nil {
writeError(w, r, err)
return
}
defer reg.finish(e.ticket, a.now())
h := w.Header()
h.Set("Content-Type", "application/gzip")
if cd := mime.FormatMediaType("attachment", map[string]string{"filename": e.filename}); cd != "" {
h.Set("Content-Disposition", cd)
} else {
h.Set("Content-Disposition", "attachment")
}
h.Set("X-Content-Type-Options", "nosniff")
h.Set("Cache-Control", "no-store")
if e.upload.size >= 0 {
h.Set("Content-Length", strconv.FormatInt(e.upload.size, 10))
}
w.WriteHeader(http.StatusOK)
err = copyExport(w, e.upload.body, e.sha256)
e.upload.done <- err
if err != nil {
log.Printf("api: export %s of %s ended early: %v", e.id, e.server, err)
// Abort the response rather than end it: a download cut short must
// read as failed in the browser, never as a complete file.
panic(http.ErrAbortHandler)
}
}
// 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)
}
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
}
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
}
// handleInternalExportUpload takes an export Job's archive (PUT
// /api/v1/internal/exports/{id}). The one-time bearer token minted with the
// 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.
func (a *API) handleInternalExportUpload(w http.ResponseWriter, r *http.Request) {
token, ok := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer ")
if !ok {
writeError(w, r, errNoExport())
return
}
rc := takeBodyDeadline(w, r)
up := &exportUpload{
body: &stallBody{r: r.Body, rc: rc}, size: r.ContentLength,
claimed: make(chan struct{}), done: make(chan error, 1),
}
reg := a.exportTickets()
if !reg.arrive(r.PathValue("id"), token, up, a.now()) {
writeError(w, r, errNoExport())
return
}
timer := time.NewTimer(exportClaimTTL)
defer timer.Stop()
select {
case <-up.claimed:
case <-timer.C:
if reg.giveUp(up, a.now()) {
writeError(w, r, errExportExpired())
return
}
case <-r.Context().Done():
if reg.giveUp(up, a.now()) {
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:
writeError(w, r, newError(http.StatusGone, "export_expired", "the download ended before the archive did: %v", err))
}
}
// stallBody restarts the connection's read deadline on every read, so reading
// fails only once no byte has come for exportStall, and clears it when the body
// ends, as deadlineBody does.
type stallBody struct {
r io.Reader
rc *http.ResponseController
done bool
}
func (b *stallBody) Read(p []byte) (int, error) {
if b.done {
return b.r.Read(p)
}
_ = b.rc.SetReadDeadline(time.Now().Add(exportStall))
n, err := b.r.Read(p)
if err != nil {
b.done = true
_ = b.rc.SetReadDeadline(time.Time{})
}
return n, err
}
+921
View File
@@ -0,0 +1,921 @@
package api
import (
"bytes"
"context"
"crypto/rand"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"io"
"net/http"
"net/http/httptest"
"regexp"
"strconv"
"strings"
"sync/atomic"
"testing"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/maintenance"
"felis.lolicon.best/internal/worldexport"
batchv1 "k8s.io/api/batch/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
)
type fakeExporter struct {
reqs []worldexport.Request
err error
}
func (f *fakeExporter) Start(_ context.Context, r worldexport.Request) (string, error) {
f.reqs = append(f.reqs, r)
if f.err != nil {
return "", f.err
}
return worldexport.JobName(r.Server, r.ID), nil
}
// asUser authenticates each request as the principal its X-Test-User header
// names, so one handler serves several users.
type asUser map[string]*Principal
func (u asUser) Authenticate(r *http.Request) (*Principal, error) {
if p, ok := u[r.Header.Get("X-Test-User")]; ok {
return p, nil
}
return nil, errUnauthorized
}
var exportUsers = asUser{
"owner1": {UserID: "owner1", Email: "[email protected]", Role: "user"},
"owner3": {UserID: "owner3", Email: "[email protected]", Role: "user"},
"newowner": {UserID: "newowner", Email: "[email protected]", Role: "user"},
"stranger": {UserID: "stranger", Email: "[email protected]", Role: "user"},
"admin1": {UserID: "admin1", Email: "[email protected]", Role: "admin", ViaAdminAccess: true},
"nouser": {Email: "[email protected]", Role: "admin", ViaAdminAccess: true},
}
func as(user string) map[string]string { return map[string]string{"X-Test-User": user} }
const (
exportBase = "http://felis-api-internal.felis.svc:8081"
worldPath = "/api/v1/servers/survival/world/export"
backupPath = "/api/v1/servers/survival/backups/bk1/export"
exportStamp = "20231114-221320" // newTestAPI's clock, UTC
)
var (
hex64 = regexp.MustCompile(`^[0-9a-f]{64}$`)
hex16 = regexp.MustCompile(`^[0-9a-f]{16}$`)
)
// exportFixture: owner1 owns survival (stopped) and creative, owner3 owns gamma
// (stopped); bk1 is survival's backup, bk2 creative's, both formerly owner1's.
func exportFixture() (*API, *fakeRepo, *fakeCluster, *fakeExporter) {
repo := newFakeRepo()
for name, owner := range map[string]string{"survival": "owner1", "creative": "owner1", "gamma": "owner3"} {
repo.byName[name] = &ServerRecord{Name: name, OwnerID: owner}
}
repo.backups = []fakeBackup{
{view: BackupView{ID: "bk1", ServerName: "survival", FormerOwner: "owner1", Status: "present", SizeBytes: 1024},
ref: "/backups/survival-bk1.tar.gz"},
{view: BackupView{ID: "bk2", ServerName: "creative", FormerOwner: "owner1", Status: "present", SizeBytes: 2048},
ref: "/backups/creative-bk2.tar.gz"},
}
cl := newFakeCluster()
for _, name := range []string{"survival", "gamma"} {
cl.byName[name] = &ServerInfo{Name: name, Phase: "Stopped", DesiredState: string(v1alpha1.DesiredStopped)}
}
ex := &fakeExporter{}
a := newTestAPI(repo, cl)
a.External = exportUsers
a.Exporter, a.InternalBaseURL = ex, exportBase
return a, repo, cl, ex
}
func beginExport(t *testing.T, h http.Handler, path, user string) exportTicketView {
t.Helper()
w := do(h, "POST", path, "", as(user))
if w.Code != http.StatusAccepted {
t.Fatalf("POST %s as %s = %d, want 202 (%s)", path, user, w.Code, w.Body.String())
}
var v exportTicketView
if err := json.Unmarshal(w.Body.Bytes(), &v); err != nil {
t.Fatalf("ticket not JSON: %v (%s)", err, w.Body.String())
}
return v
}
func exportState(t *testing.T, h http.Handler, ticket, user string) (int, exportStatusView) {
t.Helper()
w := do(h, "GET", "/api/v1/exports/"+ticket, "", as(user))
var v exportStatusView
if w.Code == http.StatusOK {
if err := json.Unmarshal(w.Body.Bytes(), &v); err != nil {
t.Fatalf("status not JSON: %v (%s)", err, w.Body.String())
}
} else {
v.State = decodeErr(t, w)
}
return w.Code, v
}
func waitExportReady(t *testing.T, h http.Handler, ticket, user string) {
t.Helper()
for deadline := time.Now().Add(5 * time.Second); ; {
if code, v := exportState(t, h, ticket, user); code == http.StatusOK && v.State == exportReady {
return
}
if time.Now().After(deadline) {
t.Fatal("the export never read ready")
}
time.Sleep(2 * time.Millisecond)
}
}
// uploadExport serves the Job's PUT on the internal face in the background.
func uploadExport(h http.Handler, id, token string, body io.Reader) <-chan *httptest.ResponseRecorder {
out := make(chan *httptest.ResponseRecorder, 1)
go func() {
r := httptest.NewRequest("PUT", "/api/v1/internal/exports/"+id, body)
r.Header.Set("Authorization", "Bearer "+token)
w := httptest.NewRecorder()
h.ServeHTTP(w, r)
out <- w
}()
return out
}
func awaitUpload(t *testing.T, up <-chan *httptest.ResponseRecorder) *httptest.ResponseRecorder {
t.Helper()
select {
case w := <-up:
return w
case <-time.After(5 * time.Second):
t.Fatal("the upload never answered")
return nil
}
}
// doSoon is do for a request that must be refused outright. A download that
// claimed the export by mistake would block on the parked body, and an upload
// let in by mistake would wait for a browser, so either fails the test after a
// bound instead of hanging it.
func doSoon(t *testing.T, h http.Handler, method, path, body string, headers map[string]string) *httptest.ResponseRecorder {
t.Helper()
out := make(chan *httptest.ResponseRecorder, 1)
go func() { out <- do(h, method, path, body, headers) }()
select {
case w := <-out:
return w
case <-time.After(2 * time.Second):
t.Fatalf("%s %s never answered: it claimed the export", method, path)
return nil
}
}
func randomBytes(n int) []byte {
b := make([]byte, n)
_, _ = rand.Read(b)
return b
}
func TestExportBackupGate(t *testing.T) {
t.Run("former owner starts a backup export", func(t *testing.T) {
a, repo, cl, ex := exportFixture()
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)
}
if len(ex.reqs) != 1 {
t.Fatalf("exporter started %d Jobs, want 1", len(ex.reqs))
}
r := ex.reqs[0]
if r.Server != "survival" || r.Mode != worldexport.ModeBackup || r.BackupRef != "/backups/survival-bk1.tar.gz" ||
!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)
}
if len(cl.acquired) != 0 {
t.Fatalf("a backup export took the world lock: %v", cl.acquired)
}
if len(repo.audits) != 1 || repo.audits[0].Action != "backup.export" || repo.audits[0].ServerName != "survival" ||
repo.audits[0].ActorUserID != "owner1" || string(repo.audits[0].Payload) != `{"backup_id":"bk1","size_bytes":1024}` {
t.Fatalf("audit = %+v", repo.audits)
}
})
for _, tc := range []struct {
name, user, path string
edit func(*API, *fakeRepo)
code int
errCode string
}{
{name: "admin exports another's world", user: "admin1", path: backupPath,
edit: func(_ *API, r *fakeRepo) { r.backups[0].view.FormerOwner = "someone" }, code: http.StatusAccepted},
{name: "stranger", user: "stranger", path: backupPath, code: http.StatusForbidden, errCode: "forbidden"},
{name: "principal without an account", user: "nouser", path: backupPath, code: http.StatusForbidden, errCode: "forbidden"},
{name: "current owner who is not the former owner", user: "newowner", path: backupPath,
edit: func(_ *API, r *fakeRepo) { r.byName["survival"].OwnerID = "newowner" },
code: http.StatusNotFound, errCode: "no_backup"},
{name: "unknown backup", user: "owner1", path: "/api/v1/servers/survival/backups/nope/export",
code: http.StatusNotFound, errCode: "no_backup"},
{name: "another server's backup", user: "owner1", path: "/api/v1/servers/survival/backups/bk2/export",
code: http.StatusForbidden, errCode: "forbidden"},
{name: "corrupt backup", user: "owner1", path: backupPath,
edit: func(_ *API, r *fakeRepo) { r.backups[0].view.Corrupt = true }, code: http.StatusConflict, errCode: "backup_corrupt"},
{name: "bad name", user: "owner1", path: "/api/v1/servers/Bad_Name/backups/bk1/export",
code: http.StatusBadRequest, errCode: "bad_name"},
{name: "unknown server", user: "owner1", path: "/api/v1/servers/nosuch/backups/bk1/export",
code: http.StatusNotFound, errCode: "not_found"},
{name: "no exporter", user: "owner1", path: backupPath,
edit: func(a *API, _ *fakeRepo) { a.Exporter = nil }, code: http.StatusServiceUnavailable, errCode: "export_unavailable"},
{name: "no internal URL", user: "owner1", path: backupPath,
edit: func(a *API, _ *fakeRepo) { a.InternalBaseURL = "" }, code: http.StatusServiceUnavailable, errCode: "export_unavailable"},
} {
t.Run(tc.name, func(t *testing.T) {
a, repo, _, ex := exportFixture()
if tc.edit != nil {
tc.edit(a, repo)
}
w := do(a.ExternalHandler(), "POST", tc.path, "", as(tc.user))
if w.Code != tc.code {
t.Fatalf("code = %d, want %d (%s)", w.Code, tc.code, w.Body.String())
}
if tc.code == http.StatusAccepted {
if len(ex.reqs) != 1 {
t.Fatalf("exporter started %d Jobs, want 1", len(ex.reqs))
}
return
}
if got := decodeErr(t, w); got != tc.errCode {
t.Errorf("error code = %q, want %q", got, tc.errCode)
}
if len(ex.reqs) != 0 || len(repo.audits) != 0 {
t.Errorf("a refused export started %d Jobs and wrote %d audits", len(ex.reqs), len(repo.audits))
}
})
}
}
func TestExportWorldGate(t *testing.T) {
t.Run("owner exports a stopped world under the world lock", func(t *testing.T) {
a, repo, cl, ex := exportFixture()
v := beginExport(t, a.ExternalHandler(), worldPath, "owner1")
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" {
t.Fatalf("export requests = %+v", ex.reqs)
}
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 != "world.export" || repo.audits[0].ServerName != "survival" {
t.Fatalf("audit = %+v", repo.audits)
}
})
for _, tc := range []struct {
name, user string
edit func(*API, *fakeCluster)
code int
errCode string
}{
{name: "admin", user: "admin1", code: http.StatusAccepted},
{name: "stranger", user: "stranger", code: http.StatusForbidden, errCode: "forbidden"},
{name: "running", user: "owner1", edit: func(_ *API, c *fakeCluster) { c.byName["survival"].Ready = true },
code: http.StatusConflict, errCode: "not_stopped"},
{name: "coming up", user: "owner1",
edit: func(_ *API, c *fakeCluster) { c.byName["survival"].DesiredState = string(v1alpha1.DesiredRunning) },
code: http.StatusConflict, errCode: "not_stopped"},
{name: "no world volume", user: "owner1", edit: func(_ *API, c *fakeCluster) { c.noWorld["survival"] = true },
code: http.StatusConflict, errCode: "no_world_volume"},
{name: "world busy", user: "owner1",
edit: func(_ *API, c *fakeCluster) {
c.maintErr["survival"] = &MaintenanceBusyError{Kind: maintenance.KindRestore}
},
code: http.StatusConflict, errCode: "maintenance_in_progress"},
{name: "no exporter", user: "owner1", edit: func(a *API, _ *fakeCluster) { a.Exporter = nil },
code: http.StatusServiceUnavailable, errCode: "export_unavailable"},
{name: "no internal URL", user: "owner1", edit: func(a *API, _ *fakeCluster) { a.InternalBaseURL = "" },
code: http.StatusServiceUnavailable, errCode: "export_unavailable"},
} {
t.Run(tc.name, func(t *testing.T) {
a, repo, cl, ex := exportFixture()
if tc.edit != nil {
tc.edit(a, cl)
}
w := do(a.ExternalHandler(), "POST", worldPath, "", as(tc.user))
if w.Code != tc.code {
t.Fatalf("code = %d, want %d (%s)", w.Code, tc.code, w.Body.String())
}
if tc.code == http.StatusAccepted {
if len(ex.reqs) != 1 {
t.Fatalf("exporter started %d Jobs, want 1", len(ex.reqs))
}
return
}
if got := decodeErr(t, w); got != tc.errCode {
t.Errorf("error code = %q, want %q", got, tc.errCode)
}
if len(ex.reqs) != 0 || len(repo.audits) != 0 || len(cl.released) != 0 {
t.Errorf("a refused export started %d Jobs, wrote %d audits, released %v", len(ex.reqs), len(repo.audits), cl.released)
}
})
}
// What anyone refused by a running world export hears.
t.Run("a world export holding the lock is named", func(t *testing.T) {
a, _, cl, _ := exportFixture()
cl.maintErr["survival"] = &MaintenanceBusyError{Kind: maintenance.KindExport}
w := do(a.ExternalHandler(), "POST", worldPath, "", as("admin1"))
var raw map[string]map[string]string
_ = json.Unmarshal(w.Body.Bytes(), &raw)
if want := "a world export is running on this server's world; retry once it finishes"; w.Code != http.StatusConflict ||
raw["error"]["code"] != "maintenance_in_progress" || raw["error"]["message"] != want {
t.Fatalf("busy = %d %s, want 409 %q", w.Code, w.Body.String(), want)
}
})
// A refusal after admission gives the export back: the owner is not left
// blocked by an export that never got a Job.
t.Run("lock conflict and a failed Job are refunded", func(t *testing.T) {
a, _, cl, ex := exportFixture()
h := a.ExternalHandler()
cl.maintErr["survival"] = &MaintenanceBusyError{Kind: maintenance.KindBackup}
if w := do(h, "POST", worldPath, "", as("owner1")); w.Code != http.StatusConflict {
t.Fatalf("busy world: code = %d, want 409", w.Code)
}
delete(cl.maintErr, "survival")
ex.err = errors.New("jobs is forbidden")
if w := do(h, "POST", worldPath, "", as("owner1")); w.Code != http.StatusInternalServerError {
t.Fatalf("failed Job: code = %d, want 500 (%s)", w.Code, w.Body.String())
}
if strings.Join(cl.released, ",") != "survival" {
t.Fatalf("a failed Job left the lock held: released %v", cl.released)
}
ex.err = nil
beginExport(t, h, worldPath, "owner1")
if n := len(a.exportTickets().starts["owner1"]); n != 1 {
t.Fatalf("hourly starts = %d, want only the export that got a Job", n)
}
})
}
func TestExportLimits(t *testing.T) {
t.Run("one per user", func(t *testing.T) {
a, _, _, ex := exportFixture()
h := a.ExternalHandler()
beginExport(t, h, worldPath, "owner1")
w := do(h, "POST", backupPath, "", as("owner1"))
if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "export_busy" || w.Header().Get("Retry-After") != "90" {
t.Fatalf("second export: %d %s Retry-After %q", w.Code, w.Body.String(), w.Header().Get("Retry-After"))
}
if len(ex.reqs) != 1 {
t.Fatalf("exporter started %d Jobs, want 1", len(ex.reqs))
}
})
t.Run("two across the install", func(t *testing.T) {
a, _, cl, ex := exportFixture()
h := a.ExternalHandler()
beginExport(t, h, worldPath, "owner1")
beginExport(t, h, "/api/v1/servers/creative/backups/bk2/export", "admin1")
w := do(h, "POST", "/api/v1/servers/gamma/world/export", "", as("owner3"))
if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "export_busy" || w.Header().Get("Retry-After") != "60" {
t.Fatalf("third export: %d %s Retry-After %q", w.Code, w.Body.String(), w.Header().Get("Retry-After"))
}
if len(ex.reqs) != 2 || strings.Join(cl.acquired, ",") != "survival:export" {
t.Fatalf("a refused export started a Job or took gamma's lock: %d Jobs, acquired %v", len(ex.reqs), cl.acquired)
}
})
t.Run("six per user per hour", func(t *testing.T) {
defer func(old time.Duration) { exportPendingTTL = old }(exportPendingTTL)
exportPendingTTL = time.Minute
a, _, _, _ := exportFixture()
var clock atomic.Int64
clock.Store(1_700_000_000)
a.Now = func() time.Time { return time.Unix(clock.Load(), 0) }
h := a.ExternalHandler()
for i := range exportPerHour {
clock.Store(1_700_000_000 + int64(i)*120) // each start outlives the last one's pending TTL
beginExport(t, h, worldPath, "owner1")
}
clock.Store(1_700_000_000 + 12*60)
w := do(h, "POST", worldPath, "", as("owner1"))
if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "export_busy" || w.Header().Get("Retry-After") != "2880" {
t.Fatalf("seventh export: %d %s Retry-After %q", w.Code, w.Body.String(), w.Header().Get("Retry-After"))
}
clock.Store(1_700_000_000 + 3600)
beginExport(t, h, worldPath, "owner1")
})
}
// TestExportRendezvous walks one world export end to end through the handlers:
// the Job's PUT parks until the owner's download claims it, the archive streams
// through untouched, the ticket works once and for its user only.
func TestExportRendezvous(t *testing.T) {
a, repo, _, ex := exportFixture()
var clock atomic.Int64
clock.Store(1_700_000_000)
a.Now = func() time.Time { return time.Unix(clock.Load(), 0) }
ext, in := a.ExternalHandler(), a.InternalHandler()
v := beginExport(t, ext, worldPath, "owner1")
job := ex.reqs[0]
download := "/api/v1/exports/" + v.Ticket + "/download"
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusOK || s != (exportStatusView{State: "pending"}) {
t.Fatalf("fresh status = %d %+v", code, s)
}
for _, user := range []string{"stranger", "admin1"} {
if code, s := exportState(t, ext, v.Ticket, user); code != http.StatusNotFound || s.State != "not_found" {
t.Fatalf("%s reads the owner's ticket: %d %+v", user, code, s)
}
}
if w := do(ext, "GET", download, "", as("owner1")); w.Code != http.StatusConflict || decodeErr(t, w) != "export_not_ready" {
t.Fatalf("download while pending = %d %s", w.Code, w.Body.String())
}
for name, hdr := range map[string]map[string]string{
"no token": nil,
"wrong token": {"Authorization": "Bearer " + strings.Repeat("0", 64)},
"the ticket": {"Authorization": "Bearer " + v.Ticket},
} {
if w := doSoon(t, in, "PUT", "/api/v1/internal/exports/"+job.ID, "x", hdr); w.Code != http.StatusNotFound || decodeErr(t, w) != "not_found" {
t.Fatalf("PUT with %s = %d %s", name, w.Code, w.Body.String())
}
}
if w := doSoon(t, in, "PUT", "/api/v1/internal/exports/ffffffffffffffff", "x", map[string]string{"Authorization": "Bearer " + job.Token}); w.Code != http.StatusNotFound {
t.Fatalf("PUT to another id = %d", w.Code)
}
archive := randomBytes(3*exportCopyBuffer + 4321)
pr, pw := io.Pipe()
up := uploadExport(in, job.ID, job.Token, pr)
waitExportReady(t, ext, v.Ticket, "owner1")
if w := doSoon(t, in, "PUT", "/api/v1/internal/exports/"+job.ID, "x", map[string]string{"Authorization": "Bearer " + job.Token}); w.Code != http.StatusNotFound {
t.Fatalf("a second PUT with the spent token = %d, want 404", w.Code)
}
// Neither a HEAD nor another user's GET claims it.
if w := doSoon(t, ext, "HEAD", download, "", as("owner1")); w.Code != http.StatusMethodNotAllowed || w.Header().Get("Allow") != "GET" {
t.Fatalf("HEAD = %d, Allow %q", w.Code, w.Header().Get("Allow"))
}
if w := doSoon(t, ext, "GET", download, "", as("stranger")); w.Code != http.StatusNotFound {
t.Fatalf("stranger's download = %d", w.Code)
}
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusOK || s.State != exportReady {
t.Fatalf("after HEAD and a stranger: %d %+v, want still ready", code, s)
}
got := make(chan *httptest.ResponseRecorder, 1)
go func() { got <- do(ext, "GET", download, "", as("owner1")) }()
_, _ = pw.Write(archive[:10_000]) // returns once the download has read it
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusGone || s.State != "export_expired" {
t.Fatalf("status mid-download = %d %+v, want 410 export_expired", code, s)
}
if w := doSoon(t, ext, "GET", download, "", as("owner1")); w.Code != http.StatusGone || decodeErr(t, w) != "export_expired" {
t.Fatalf("a second download mid-stream = %d %s", w.Code, w.Body.String())
}
go func() {
for rest := archive[10_000:]; len(rest) > 0; {
n := min(len(rest), 10_000)
_, _ = pw.Write(rest[:n])
rest = rest[n:]
}
pw.Close()
}()
w := <-got
if w.Code != http.StatusOK || !bytes.Equal(w.Body.Bytes(), archive) {
t.Fatalf("download = %d, %d bytes; want 200 and the %d archive bytes", w.Code, w.Body.Len(), len(archive))
}
h := w.Header()
if h.Get("Content-Type") != "application/gzip" || h.Get("X-Content-Type-Options") != "nosniff" || h.Get("Cache-Control") != "no-store" ||
h.Get("Content-Disposition") != "attachment; filename=survival-world-"+exportStamp+".tar.gz" || h.Get("Content-Length") != "" {
t.Fatalf("download headers = %v", h)
}
if w := awaitUpload(t, up); w.Code != http.StatusNoContent {
t.Fatalf("upload answered %d %s, want 204", w.Code, w.Body.String())
}
if w := do(ext, "GET", download, "", as("owner1")); w.Code != http.StatusGone || decodeErr(t, w) != "export_expired" {
t.Fatalf("second download = %d %s", w.Code, w.Body.String())
}
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusGone || s.State != "export_expired" {
t.Fatalf("status after download = %d %+v", code, s)
}
clock.Add(int64(exportKeepSpent / time.Second))
if code, _ := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusNotFound {
t.Fatalf("status long after = %d, want 404", code)
}
if len(repo.audits) != 1 {
t.Fatalf("audits = %+v, want the start only", repo.audits)
}
beginExport(t, ext, worldPath, "owner1") // the finished export no longer counts
}
func TestExportExpiry(t *testing.T) {
t.Run("a Job that never connects", func(t *testing.T) {
a, _, _, ex := exportFixture()
var clock atomic.Int64
clock.Store(1_700_000_000)
a.Now = func() time.Time { return time.Unix(clock.Load(), 0) }
ext := a.ExternalHandler()
v := beginExport(t, ext, worldPath, "owner1")
clock.Add(int64(exportPendingTTL/time.Second) - 1)
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusOK || s.State != "pending" {
t.Fatalf("just inside the pending TTL: %d %+v", code, s)
}
clock.Add(1)
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusGone || s.State != "export_expired" {
t.Fatalf("past the pending TTL: %d %+v", code, s)
}
job := ex.reqs[0]
if w := doSoon(t, a.InternalHandler(), "PUT", "/api/v1/internal/exports/"+job.ID, "x", map[string]string{"Authorization": "Bearer " + job.Token}); w.Code != http.StatusNotFound {
t.Fatalf("a late Job's PUT = %d, want 404", w.Code)
}
})
t.Run("a browser that never comes", func(t *testing.T) {
defer func(old time.Duration) { exportClaimTTL = old }(exportClaimTTL)
exportClaimTTL = 30 * time.Millisecond
a, _, _, ex := exportFixture()
ext := a.ExternalHandler()
v := beginExport(t, ext, backupPath, "owner1")
job := ex.reqs[0]
w := awaitUpload(t, uploadExport(a.InternalHandler(), job.ID, job.Token, strings.NewReader("archive")))
if w.Code != http.StatusGone || decodeErr(t, w) != "export_expired" {
t.Fatalf("unclaimed upload = %d %s, want 410 export_expired", w.Code, w.Body.String())
}
if w := do(ext, "GET", "/api/v1/exports/"+v.Ticket+"/download", "", as("owner1")); w.Code != http.StatusGone {
t.Fatalf("download after the claim TTL = %d, want 410", w.Code)
}
})
}
// TestExportStatusReportsFailedJob: a Job that dies before it connects turns
// the export failed with the Job's own error, and frees the owner's slot.
func TestExportStatusReportsFailedJob(t *testing.T) {
a, _, _, ex := exportFixture()
js := &fakeJobStatus{}
a.JobStatus = js
ext := a.ExternalHandler()
v := beginExport(t, ext, worldPath, "owner1")
job := ex.reqs[0]
name := worldexport.JobName("survival", job.ID)
js.jobs = []AsyncJob{
{Name: "export-survival-0000000000000000", Kind: "export_world", State: "failed", Message: "someone else's"},
{Name: name, Kind: "export_world", State: "running"},
}
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusOK || s != (exportStatusView{State: "pending"}) {
t.Fatalf("running Job: %d %+v", code, s)
}
js.jobs[1] = AsyncJob{Name: name, Kind: "export_world", State: "failed", Message: "felis export: open /world: permission denied"}
want := exportStatusView{State: "failed", Message: "felis export: open /world: permission denied"}
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusOK || s != want || js.got != "survival" {
t.Fatalf("failed Job: %d %+v (asked about %q)", code, s, js.got)
}
js.jobs = nil
if code, s := exportState(t, ext, v.Ticket, "owner1"); code != http.StatusOK || s != want {
t.Fatalf("failed export once the Job is gone: %d %+v", code, s)
}
if w := doSoon(t, a.InternalHandler(), "PUT", "/api/v1/internal/exports/"+job.ID, "x", map[string]string{"Authorization": "Bearer " + job.Token}); w.Code != http.StatusNotFound {
t.Fatalf("PUT on a failed export = %d, want 404", w.Code)
}
beginExport(t, ext, backupPath, "owner1")
}
// TestExportRegistryRaces: two transitions the handlers reach only in a race.
// A Job reported failed just as its upload arrived must not fail the ready
// export, and the claim timer firing as the browser claims must not spend the
// export under its download.
func TestExportRegistryRaces(t *testing.T) {
now := time.Unix(1_700_000_000, 0)
reg := (&API{}).exportTickets()
e := &exportEntry{userID: "owner1", server: "survival", mode: worldexport.ModeWorld, filename: "survival.tar.gz"}
token, err := reg.admit(e, now)
if err != nil {
t.Fatal(err)
}
up := &exportUpload{body: strings.NewReader("x"), size: 1, claimed: make(chan struct{}), done: make(chan error, 1)}
if !reg.arrive(e.id, token, up, now) {
t.Fatal("the upload did not arrive")
}
if reg.fail(e.ticket, "felis export: killed", now) {
t.Fatal("a late Job failure failed a ready export")
}
if got, err := reg.lookup(e.ticket, "owner1", now); err != nil || got.state != exportReady || got.message != "" {
t.Fatalf("after a late Job failure: state %q, message %q, err %v; want ready", got.state, got.message, err)
}
if _, err := reg.claim(e.ticket, "owner1", now); err != nil {
t.Fatalf("claim: %v", err)
}
if reg.giveUp(up, now) {
t.Fatal("the claim timer spent an export being downloaded")
}
if e.state != exportStreaming {
t.Fatalf("state after a late giveUp = %q, want streaming", e.state)
}
}
// exportServers runs both faces for real, so deadlines apply and an aborted
// download reaches the client as an error.
func exportServers(t *testing.T, a *API) (ext, in *httptest.Server) {
ext, in = httptest.NewServer(a.ExternalHandler()), httptest.NewServer(a.InternalHandler())
t.Cleanup(func() { ext.Close(); in.Close() })
return ext, in
}
func realUpload(t *testing.T, in *httptest.Server, job worldexport.Request, body io.Reader, size int64) <-chan *http.Response {
req, err := http.NewRequest("PUT", in.URL+"/api/v1/internal/exports/"+job.ID, body)
if err != nil {
t.Fatal(err)
}
req.ContentLength = size
req.Header.Set("Authorization", "Bearer "+job.Token)
out := make(chan *http.Response, 1)
go func() {
resp, err := http.DefaultClient.Do(req)
if err != nil {
out <- &http.Response{StatusCode: -1, Status: err.Error(), Body: io.NopCloser(strings.NewReader(""))}
return
}
out <- resp
}()
return out
}
// realDownload fetches the owner's download: the response (nil when the request
// failed outright), the bytes read and the error that ended the read.
func realDownload(ext *httptest.Server, ticket string) (*http.Response, []byte, error) {
req, _ := http.NewRequest("GET", ext.URL+"/api/v1/exports/"+ticket+"/download", nil)
req.Header.Set("X-Test-User", "owner1")
resp, err := http.DefaultClient.Do(req)
if err != nil {
return nil, nil, err
}
defer resp.Body.Close()
got, err := io.ReadAll(resp.Body)
return resp, got, err
}
func awaitResponse(t *testing.T, ch <-chan *http.Response) *http.Response {
t.Helper()
select {
case resp := <-ch:
t.Cleanup(func() { resp.Body.Close() })
return resp
case <-time.After(5 * time.Second):
t.Fatal("the upload never answered")
return nil
}
}
// 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) {
archive := randomBytes(3*exportCopyBuffer + 4321)
sum := sha256.Sum256(archive)
for _, tc := range []struct {
name string
digest string
}{
{"recorded digest matches", hex.EncodeToString(sum[:])},
{"no digest recorded", ""},
{"digest mismatch", strings.Repeat("ab", 32)},
} {
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)
}
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"))
}
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
}
if err != nil || !bytes.Equal(got, archive) {
t.Fatalf("read %d of %d bytes, err %v", len(got), len(archive), err)
}
if upResp.StatusCode != http.StatusNoContent {
t.Fatalf("upload answered %d, want 204", upResp.StatusCode)
}
})
}
// 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())
}
})
}
// zeros reads as an endless run of zero bytes.
type zeros struct{}
func (zeros) Read(p []byte) (int, error) {
clear(p)
return len(p), nil
}
func mustRead(t *testing.T, r io.Reader) []byte {
t.Helper()
b, err := io.ReadAll(r)
if err != nil {
t.Fatal(err)
}
return b
}
// TestExportUploadPace: the Job's upload waits for the browser longer than the
// body-deadline grace, and then moves at the browser's pace, without being cut
// off; a body that stops moving for exportStall is.
func TestExportUploadPace(t *testing.T) {
defer func(g time.Duration, r float64, s time.Duration) { bodyGrace, bodyMinRate, exportStall = g, r, s }(bodyGrace, bodyMinRate, exportStall)
bodyGrace, bodyMinRate, exportStall = 100*time.Millisecond, 1<<30, 300*time.Millisecond
t.Run("a slow claim is not a slow body", func(t *testing.T) {
a, _, _, ex := exportFixture()
ext, in := exportServers(t, a)
v := beginExport(t, a.ExternalHandler(), worldPath, "owner1")
archive := randomBytes(64 << 10)
pr, pw := io.Pipe()
go func() {
_, _ = pw.Write(archive[:1000])
time.Sleep(3 * bodyGrace)
_, _ = pw.Write(archive[1000:])
pw.Close()
}()
up := realUpload(t, in, ex.reqs[0], pr, -1)
waitExportReady(t, a.ExternalHandler(), v.Ticket, "owner1")
time.Sleep(3 * bodyGrace)
_, got, err := realDownload(ext, v.Ticket)
if err != nil || !bytes.Equal(got, archive) {
t.Fatalf("read %d of %d bytes, err %v", len(got), len(archive), err)
}
if resp := awaitResponse(t, up); resp.StatusCode != http.StatusNoContent {
t.Fatalf("upload answered %d, want 204", resp.StatusCode)
}
})
t.Run("a stalled body is cut off", func(t *testing.T) {
a, _, _, ex := exportFixture()
ext, in := exportServers(t, a)
v := beginExport(t, a.ExternalHandler(), worldPath, "owner1")
pr, pw := io.Pipe()
defer pw.Close()
go func() { _, _ = pw.Write(make([]byte, 1000)) }() // then nothing, ever
up := realUpload(t, in, ex.reqs[0], pr, -1)
waitExportReady(t, a.ExternalHandler(), v.Ticket, "owner1")
done := make(chan error, 1)
go func() { _, _, err := realDownload(ext, v.Ticket); done <- err }()
select {
case err := <-done:
if err == nil {
t.Fatal("a stalled archive downloaded as complete")
}
case <-time.After(5 * time.Second):
t.Fatal("a stalled upload was never cut off")
}
pw.CloseWithError(errors.New("test over"))
awaitResponse(t, up)
})
t.Run("a browser that stops reading is cut off", func(t *testing.T) {
a, _, _, ex := exportFixture()
ext, in := exportServers(t, a)
v := beginExport(t, a.ExternalHandler(), worldPath, "owner1")
realUpload(t, in, ex.reqs[0], io.LimitReader(zeros{}, 1<<30), -1)
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)
}
t.Cleanup(func() { resp.Body.Close() }) // before the servers close, so a held handler ends
// Read nothing: the socket buffers fill and the next write waits.
reg := a.exportTickets()
for deadline := time.Now().Add(5 * time.Second); ; {
reg.mu.Lock()
state := reg.byTicket[v.Ticket].state
reg.mu.Unlock()
if state == exportSpent {
break
}
if time.Now().After(deadline) {
t.Fatalf("a browser that stopped reading held the download open: state %q", state)
}
time.Sleep(10 * time.Millisecond)
}
if _, err := io.ReadAll(resp.Body); err == nil {
t.Fatal("the cut-off download read as complete")
}
})
}
// TestK8sExportJobs: export Jobs show in the jobs list by what they archive,
// and never count as world work the backup scheduler waits for.
func TestK8sExportJobs(t *testing.T) {
scheme := runtime.NewScheme()
if err := clientgoscheme.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
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",
TargetURL: exportBase + "/x", Token: "t", Namespace: "minecraft", Image: "felis:1",
})
if err != nil {
t.Fatal(err)
}
return j
}
world, backupJob := job("1111111111111111", worldexport.ModeWorld), job("2222222222222222", worldexport.ModeBackup)
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).
WithStatusSubresource(&batchv1.Job{}).Build()
k := NewK8sJobStatus(c, "minecraft")
ctx := context.Background()
if n, err := k.RunningWorldJobs(ctx); err != nil || n != 1 {
t.Fatalf("RunningWorldJobs = %d, %v; want the restore only", n, err)
}
jobs, err := k.LatestJobs(ctx, "survival")
if err != nil {
t.Fatal(err)
}
kinds := map[string]string{}
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"}
if len(kinds) != len(want) {
t.Fatalf("jobs = %v, want %v", kinds, want)
}
for name, k := range want {
if kinds[name] != k {
t.Errorf("%s = %q, want %q", name, kinds[name], k)
}
}
// The kinds agree with maintenance.JobKind: a world export holds 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 _, holds := maintenance.JobKind(backupJob); holds {
t.Error("a backup export holds the world")
}
}
+1 -1
View File
@@ -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"
Kind string `json:"kind"` // "backup" | "restore" | "export_world" | "export_backup"
State string `json:"state"` // "running" | "succeeded" | "failed"
Message string `json:"message,omitempty"`
StartedAt time.Time `json:"started_at,omitzero"`
+16 -5
View File
@@ -25,6 +25,7 @@ const (
jobManagedByBackup = "felis-backup"
jobManagedByRestore = "felis-restore"
jobManagedByExport = "felis-export"
// jobBackupReasonLabel marks the backup Job of a scheduled restore point
// (backupjob.LabelReason).
@@ -45,9 +46,9 @@ func NewK8sJobStatus(c client.Client, namespace string) *K8sJobStatus {
return &K8sJobStatus{c: c, namespace: namespace}
}
// LatestJobs lists this server's backup/restore Jobs newest-first, capped so a
// long history cannot balloon the response. Jobs the label selector catches but
// another component created (unknown managed-by) are dropped.
// LatestJobs lists this server's backup, restore and export Jobs newest-first,
// capped so a long history cannot balloon the response. Jobs the label selector
// catches but another component created (unknown managed-by) are dropped.
func (k *K8sJobStatus) LatestJobs(ctx context.Context, serverName string) ([]AsyncJob, error) {
var list batchv1.JobList
if err := k.c.List(ctx, &list, client.InNamespace(k.namespace),
@@ -147,7 +148,10 @@ func jobToAsyncJob(j *batchv1.Job) (AsyncJob, bool) {
}
// RunningWorldJobs counts the backup and restore Jobs of every server that
// have yet to finish (BackupScheduler waits for them).
// have yet to finish (BackupScheduler waits for them). Export Jobs are left
// out: one moves at its browser's pace, a download left running for an hour
// must not hold every world's restore point back, and the export limits
// (exportMaxActive) already bound the load they add.
func (k *K8sJobStatus) RunningWorldJobs(ctx context.Context) (int, error) {
var list batchv1.JobList
if err := k.c.List(ctx, &list, client.InNamespace(k.namespace), client.HasLabels{jobManagedByLabel}); err != nil {
@@ -156,7 +160,7 @@ func (k *K8sJobStatus) RunningWorldJobs(ctx context.Context) (int, error) {
n := 0
for i := range list.Items {
j := &list.Items[i]
if _, ok := jobOutcome(j); ok && !maintenance.JobFinished(j) {
if _, ok := jobOutcome(j); ok && j.Labels[jobManagedByLabel] != jobManagedByExport && !maintenance.JobFinished(j) {
n++
}
}
@@ -225,6 +229,13 @@ func jobOutcome(j *batchv1.Job) (AsyncJob, bool) {
kind = "backup"
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 {
kind = "export_backup"
}
default:
return AsyncJob{}, false
}
+5 -3
View File
@@ -11,9 +11,9 @@ import (
)
// maintenanceError maps the world-volume lock's refusals onto their 409s:
// maintenance_in_progress while a restore, backup or file write holds the volume,
// not_stopped (with the caller's wording) while the server is not fully down.
// Anything else passes through unchanged.
// maintenance_in_progress while a restore, backup, file write or world export
// holds the volume, not_stopped (with the caller's wording) while the server is
// not fully down. Anything else passes through unchanged.
func maintenanceError(err error, notStopped string) error {
var busy *MaintenanceBusyError
switch {
@@ -37,6 +37,8 @@ func maintenanceLabel(kind string) string {
return "a backup"
case maintenance.KindFileWrite:
return "a file write"
case maintenance.KindExport:
return "a world export"
case maintenance.KindReap:
return "the idle-world reaper"
}
+12
View File
@@ -210,6 +210,18 @@ func (b *deadlineBody) Read(p []byte) (int, error) {
return n, err
}
// takeBodyDeadline lifts withBodyDeadline from a request whose body is read at
// someone else's pace (an export's upload waits for the browser, then moves at
// its speed), and returns the controller the handler bounds it with instead.
func takeBodyDeadline(w http.ResponseWriter, r *http.Request) *http.ResponseController {
if b, ok := r.Body.(*deadlineBody); ok {
r.Body = b.ReadCloser
_ = b.rc.SetReadDeadline(time.Time{})
return b.rc
}
return http.NewResponseController(w)
}
// requireInternal enforces service-token auth for the internal face and stashes
// the caller the token belongs to, which callersOnly checks against the route.
// It never applies Zero Trust (spec §14 red line).
+2
View File
@@ -40,6 +40,8 @@ func TestOpenAPISchemasMatchWireStructs(t *testing.T) {
"MyServerView": MyServerView{},
"AllowlistEntry": AllowlistEntry{},
"BackupView": BackupView{},
"ExportTicket": exportTicketView{},
"ExportStatus": exportStatusView{},
"Schedule": Schedule{},
"Build": build.Build{},
"Image": build.Image{},
+2 -2
View File
@@ -958,11 +958,11 @@ func (p *PGRepo) LatestBackup(ctx context.Context, serverName string) (*BackupRe
// BackupByID returns a single present backup by its id, or ErrNotFound.
func (p *PGRepo) BackupByID(ctx context.Context, id string) (*BackupRecord, error) {
const q = `SELECT id, server_name, COALESCE(former_owner, ''), backup_ref, COALESCE(size_bytes, 0),
corrupt_at IS NOT NULL
corrupt_at IS NOT NULL, COALESCE(sha256, '')
FROM world_backups WHERE id = $1 AND status = 'present'`
var b BackupRecord
switch err := p.db.QueryRowContext(ctx, q, id).Scan(
&b.ID, &b.ServerName, &b.FormerOwner, &b.BackupRef, &b.SizeBytes, &b.Corrupt); {
&b.ID, &b.ServerName, &b.FormerOwner, &b.BackupRef, &b.SizeBytes, &b.Corrupt, &b.SHA256); {
case errors.Is(err, sql.ErrNoRows):
return nil, ErrNotFound
case err != nil:
+4
View File
@@ -143,6 +143,10 @@ type BackupRecord struct {
SizeBytes int64
// Corrupt reports that the archive failed a read-back (see BackupView).
Corrupt bool
// SHA256 is the digest recorded when the archive was written, which an
// export checks the bytes against as they stream. Empty for an archive from
// before digests were kept.
SHA256 string
}
// StaffUser is the login-side projection of a users row (spec §B passwordless
+8
View File
@@ -426,6 +426,14 @@ func writeTarGz(ctx context.Context, w io.Writer, srcDir string) (tarStats, erro
return st, nil
}
// WriteTarGz archives srcDir into w laid out exactly as Archive lays out a
// backup, so an exported world restores like any other archive, and returns the
// entries it left out. The world export Job streams it straight into its upload.
func WriteTarGz(ctx context.Context, w io.Writer, srcDir string) ([]string, error) {
st, err := writeTarGz(ctx, w, srcDir)
return st.skipped, err
}
// dirMeta is a directory's recorded permission bits and modification time,
// applied once nothing more is written into it.
type dirMeta struct {
+27 -8
View File
@@ -1,6 +1,6 @@
// Package maintenance is the per-server mutual exclusion between a game server
// and the Jobs that write or snapshot its world volume (restore, backup, file
// write). The world PVC is ReadWriteOnce, and RWO is exclusive per NODE: on a
// write, world export). The world PVC is ReadWriteOnce, and RWO is exclusive per NODE: on a
// single-node cluster the game pod and a restore pod mount it side by side, so
// the access mode alone guards nothing. A server woken mid-restore boots on a
// half-extracted world and the restore then prunes what it wrote; a server woken
@@ -29,7 +29,10 @@
//
// File reads and listings are not holders. They mount the volume read-only for a
// second or two, and a server starting beside one cannot hurt either side, so
// nobody waits for them.
// nobody waits for them. A world export mounts it read-only too, but it holds:
// it archives the world for as long as the download takes, and a server started
// beside it would hand the owner a torn archive. Exporting a backup reads only
// the backup store and holds nothing.
package maintenance
import (
@@ -50,14 +53,18 @@ const (
Grace = 2 * time.Minute
// LabelServer / LabelManagedBy are the labels every maintenance executor puts
// on its Job (internal/restore, internal/backupjob, internal/fileedit keep
// their own copies; maintenance_test pins them against these).
// on its Job (internal/restore, internal/backupjob, internal/fileedit and
// internal/worldexport keep their own copies; maintenance_test pins them
// against these).
LabelServer = "felis.lolicon.best/server"
LabelManagedBy = "app.kubernetes.io/managed-by"
// LabelFilesMode is the file operation (list, read, write, mkdir, delete,
// rename, upload) a files Job performs. Every one but list and read holds the
// volume.
LabelFilesMode = "felis.lolicon.best/files-mode"
// LabelExportMode is what an export Job archives: ExportModeWorld (the live
// world, which holds the volume) or ExportModeBackup (a stored archive).
LabelExportMode = "felis.lolicon.best/export-mode"
// LabelThenRestore marks a backup Job that is the safety snapshot in front of
// a restore. Its value is the chain's state: ThenRestorePending until felis-api
@@ -85,6 +92,7 @@ const (
KindRestore = "restore"
KindBackup = "backup"
KindFileWrite = "file-write"
KindExport = "export"
// KindReap is the reaper archiving an idle world and reclaiming its volume.
// It runs no Job: the reaper holds the Annotation itself and rewrites it
// well inside Grace for as long as it works on the world.
@@ -98,11 +106,18 @@ const (
FilesModeRead = "read"
)
// The LabelExportMode values.
const (
ExportModeWorld = "world"
ExportModeBackup = "backup"
)
// JobKind names the holder a Job represents, or reports false for a Job that
// holds nothing (a file read, a build, anything else in the namespace). A files
// Job holds unless it is a list or a read, so an operation this build does not
// know — and a Job without LabelFilesMode, which can only be an old one still
// inside its TTL — counts as a change: over-counting is the safe side.
// holds nothing (a file read, a backup export, a build, anything else in the
// namespace). A files Job holds unless it is a list or a read, so an operation
// this build does not know — and a Job without LabelFilesMode, which can only be
// an old one still inside its TTL — counts as a change: over-counting is the
// safe side. An export Job holds unless it names a backup, for the same reason.
func JobKind(j *batchv1.Job) (string, bool) {
switch j.Labels[LabelManagedBy] {
case "felis-restore":
@@ -113,6 +128,10 @@ func JobKind(j *batchv1.Job) (string, bool) {
if mode := j.Labels[LabelFilesMode]; mode != FilesModeList && mode != FilesModeRead {
return KindFileWrite, true
}
case "felis-export":
if j.Labels[LabelExportMode] != ExportModeBackup {
return KindExport, true
}
}
return "", false
}
+31
View File
@@ -7,6 +7,7 @@ import (
"felis.lolicon.best/internal/backupjob"
"felis.lolicon.best/internal/fileedit"
"felis.lolicon.best/internal/restore"
"felis.lolicon.best/internal/worldexport"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
)
@@ -66,6 +67,20 @@ func filesJob(t *testing.T, server, op string) batchv1.Job {
return *j
}
func exportJob(t *testing.T, server, mode string) batchv1.Job {
t.Helper()
j, err := worldexport.ExportJob(worldexport.JobParams{
Server: server, ID: "0011223344556677", Mode: mode, WorldPVC: "world-" + server + "-0",
BackupPVC: "felis-backups", BackupRef: "/backups/a.tar.gz", TargetURL: "http://api/x", Token: "t",
Namespace: "minecraft", ServiceAccount: "felis-restore", Image: "felis:1",
BackupRoot: "/backups", WorldsRoot: "/world",
})
if err != nil {
t.Fatalf("ExportJob: %v", err)
}
return *j
}
func finished(j batchv1.Job, cond batchv1.JobConditionType) batchv1.Job {
j.Status.Conditions = append(j.Status.Conditions, batchv1.JobCondition{Type: cond, Status: corev1.ConditionTrue})
return j
@@ -89,6 +104,8 @@ func TestJobKindMatchesTheExecutors(t *testing.T) {
{"file upload", filesJob(t, "survival", fileedit.OpUpload), KindFileWrite, true},
{"file read", filesJob(t, "survival", fileedit.OpRead), "", false},
{"file list", filesJob(t, "survival", fileedit.OpList), "", false},
{"world export", exportJob(t, "survival", worldexport.ModeWorld), KindExport, true},
{"backup export", exportJob(t, "survival", worldexport.ModeBackup), "", false},
} {
kind, ok := JobKind(&tc.job)
if kind != tc.kind || ok != tc.ok {
@@ -208,3 +225,17 @@ func TestHolderFromLock(t *testing.T) {
}
}
}
// An export Job that does not say it reads a backup holds, like a files Job
// with an unknown mode: a world export a newer build labels differently must
// not let the server start under it.
func TestExportJobWithoutModeHolds(t *testing.T) {
j := exportJob(t, "survival", worldexport.ModeBackup)
if j.Labels[LabelExportMode] != ExportModeBackup {
t.Fatalf("export mode label = %q", j.Labels[LabelExportMode])
}
delete(j.Labels, LabelExportMode)
if kind, ok := JobKind(&j); !ok || kind != KindExport {
t.Fatalf("an export Job without a mode = %q, %v; want export", kind, ok)
}
}
+2 -2
View File
@@ -344,8 +344,8 @@ func TestBackupReadBack(t *testing.T) {
if b, err := repo.BackupByID(ctx, newer); err != nil || !b.Corrupt {
t.Fatalf("BackupByID(corrupt) = (%+v, %v); want Corrupt", b, err)
}
if b, err := repo.BackupByID(ctx, older); err != nil || b.Corrupt {
t.Fatalf("BackupByID(intact) = (%+v, %v); want not Corrupt", b, err)
if b, err := repo.BackupByID(ctx, older); err != nil || b.Corrupt || b.SHA256 != "aa" {
t.Fatalf("BackupByID(intact) = (%+v, %v); want not Corrupt, sha256 aa", b, err)
}
views, _, err := repo.AllBackups(ctx, api.BackupListOpts{Server: name, Limit: api.MaxBackupListLimit})
+237
View File
@@ -0,0 +1,237 @@
package worldexport
import (
"fmt"
"time"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
// Label keys applied to export objects, the same ones the other executors use
// (internal/maintenance and internal/api keep their own copies;
// maintenance_test pins them against these).
const (
LabelManagedBy = "app.kubernetes.io/managed-by"
LabelComponent = "app.kubernetes.io/component"
LabelServer = "felis.lolicon.best/server"
// LabelMode is what the Job archives (ModeWorld or ModeBackup). A world
// export holds the world volume; a backup export does not.
LabelMode = "felis.lolicon.best/export-mode"
managedByValue = "felis-export"
componentValue = "world-export"
worldVolume = "world"
backupVolume = "backup"
containerName = "export"
felisBinaryPath = "/usr/local/bin/felis"
)
// The two things an export can archive.
const (
ModeWorld = "world"
ModeBackup = "backup"
)
// TokenEnv carries the one-time upload token into the Pod. It is the only
// secret the Pod holds, so it rides the environment and not argv, where a
// process listing on the node would show it.
const TokenEnv = "FELIS_EXPORT_TOKEN"
// JobParams are the rendered inputs to an export Job. ExportJob is a pure
// function of them, so the Job shape is unit-tested without a cluster.
type JobParams struct {
Server string
// ID names this export: it is the tail of the Job name and of the internal
// upload path, so two exports of one server never collide.
ID string
Mode string
// WorldPVC is mounted for ModeWorld, BackupPVC and BackupRef for ModeBackup.
WorldPVC string
BackupPVC string
BackupRef string
// TargetURL is where the Pod PUTs the archive (felis-api's internal face),
// and Token the one-time bearer token that opens it.
TargetURL string
Token string
Namespace string
ServiceAccount string
Image string
BackupRoot string
WorldsRoot string
Deadline time.Duration
CPULimit string
MemLimit string
RunAsUser int64
RunAsGroup int64
FSGroup int64
TTLAfterFinished time.Duration
}
// JobName is the Job an export runs as.
func JobName(server, id string) string { return "export-" + server + "-" + id }
func exportLabels(p JobParams) map[string]string {
return map[string]string{
LabelManagedBy: managedByValue,
LabelComponent: componentValue,
LabelServer: p.Server,
LabelMode: p.Mode,
}
}
// ExportJob renders the export Job. Every isolation guarantee lives here and is
// asserted by jobspec_test.go:
//
// - the weak felis-restore SA with its token auto-mount disabled, so the Pod
// cannot reach the K8s API;
// - EXACTLY one volume, read-only in the claim and in the mount: the world PVC
// for a world export, the backup PVC for a backup export. No Secret or
// ConfigMap: the one credential in the Pod is the upload token, which opens
// this one export and nothing else;
// - root with DAC_OVERRIDE and nothing more: the world is the game uid's
// mode-0600 files, and Pod Security baseline, which the minecraft namespace
// enforces, admits DAC_OVERRIDE but not the narrower DAC_READ_SEARCH. No
// privilege or escalation, a read-only root filesystem;
// - backoffLimit 0 (the token is spent by the first attempt, a retry could
// only fail) and activeDeadlineSeconds, so a download left hanging cannot
// hold the world forever; ttlSecondsAfterFinished GCs the finished Job.
//
// The container runs `/usr/local/bin/felis export` (cmd/felis). The backup ref is
// an absolute path, so the backup PVC is mounted at BackupRoot, the path the
// archives were written under, as the restore Job mounts it.
func ExportJob(p JobParams) (*batchv1.Job, error) {
if p.Image == "" {
return nil, fmt.Errorf("worldexport: image is empty")
}
if p.ID == "" || p.TargetURL == "" || p.Token == "" {
return nil, fmt.Errorf("worldexport: an export needs an id, a target URL and a token")
}
var (
args = []string{"--mode", p.Mode, "--server", p.Server, "--target-url", p.TargetURL}
volume corev1.Volume
mount corev1.VolumeMount
)
switch p.Mode {
case ModeWorld:
if p.WorldPVC == "" {
return nil, fmt.Errorf("worldexport: world PVC name is required")
}
args = append(args, "--worlds-root", p.WorldsRoot)
volume = readOnlyClaim(worldVolume, p.WorldPVC)
mount = corev1.VolumeMount{Name: worldVolume, MountPath: p.WorldsRoot, ReadOnly: true}
case ModeBackup:
if p.BackupPVC == "" || p.BackupRef == "" {
return nil, fmt.Errorf("worldexport: backup PVC name and archive ref are required")
}
args = append(args, "--ref", p.BackupRef, "--backup-root", p.BackupRoot)
volume = readOnlyClaim(backupVolume, p.BackupPVC)
mount = corev1.VolumeMount{Name: backupVolume, MountPath: p.BackupRoot, ReadOnly: true}
default:
return nil, fmt.Errorf("worldexport: unknown mode %q", p.Mode)
}
limits, err := resourceLimits(p.CPULimit, p.MemLimit)
if err != nil {
return nil, err
}
deadline := int64(p.Deadline / time.Second)
if deadline <= 0 {
deadline = int64(defaultDeadline / time.Second)
}
ttl := int32(p.TTLAfterFinished / time.Second)
if ttl <= 0 {
ttl = int32(defaultTTL / time.Second)
}
container := corev1.Container{
Name: containerName,
Image: p.Image,
Command: []string{felisBinaryPath, "export"},
Args: args,
Env: []corev1.EnvVar{{Name: TokenEnv, Value: p.Token}},
VolumeMounts: []corev1.VolumeMount{mount},
Resources: corev1.ResourceRequirements{Limits: limits, Requests: limits},
SecurityContext: &corev1.SecurityContext{
Privileged: boolPtr(false),
AllowPrivilegeEscalation: boolPtr(false),
ReadOnlyRootFilesystem: boolPtr(true),
Capabilities: &corev1.Capabilities{
Drop: []corev1.Capability{"ALL"},
Add: []corev1.Capability{"DAC_OVERRIDE"},
},
},
// The exit error reaches the export's status and GET
// /servers/{name}/jobs through the terminated state.
TerminationMessagePolicy: corev1.TerminationMessageFallbackToLogsOnError,
}
sc := &corev1.PodSecurityContext{
RunAsNonRoot: boolPtr(false),
RunAsUser: int64Ptr(p.RunAsUser),
RunAsGroup: int64Ptr(p.RunAsGroup),
}
if p.FSGroup > 0 {
sc.FSGroup = int64Ptr(p.FSGroup)
}
return &batchv1.Job{
ObjectMeta: metav1.ObjectMeta{
Name: JobName(p.Server, p.ID),
Namespace: p.Namespace,
Labels: exportLabels(p),
},
Spec: batchv1.JobSpec{
BackoffLimit: int32Ptr(0),
ActiveDeadlineSeconds: int64Ptr(deadline),
TTLSecondsAfterFinished: int32Ptr(ttl),
Template: corev1.PodTemplateSpec{
ObjectMeta: metav1.ObjectMeta{Labels: exportLabels(p)},
Spec: corev1.PodSpec{
RestartPolicy: corev1.RestartPolicyNever,
ServiceAccountName: p.ServiceAccount,
AutomountServiceAccountToken: boolPtr(false),
SecurityContext: sc,
Containers: []corev1.Container{container},
Volumes: []corev1.Volume{volume},
},
},
},
}, nil
}
func readOnlyClaim(name, claim string) corev1.Volume {
return corev1.Volume{
Name: name,
VolumeSource: corev1.VolumeSource{
PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{ClaimName: claim, ReadOnly: true},
},
}
}
// resourceLimits parses the CPU/memory limits into a ResourceList.
func resourceLimits(cpu, mem string) (corev1.ResourceList, error) {
if cpu == "" {
cpu = defaultCPULimit
}
if mem == "" {
mem = defaultMemLimit
}
cpuQty, err := resource.ParseQuantity(cpu)
if err != nil {
return nil, fmt.Errorf("worldexport: invalid cpu limit %q: %w", cpu, err)
}
memQty, err := resource.ParseQuantity(mem)
if err != nil {
return nil, fmt.Errorf("worldexport: invalid memory limit %q: %w", mem, err)
}
return corev1.ResourceList{corev1.ResourceCPU: cpuQty, corev1.ResourceMemory: memQty}, nil
}
func boolPtr(b bool) *bool { return &b }
func int32Ptr(i int32) *int32 { return &i }
func int64Ptr(i int64) *int64 { return &i }
+178
View File
@@ -0,0 +1,178 @@
package worldexport
import (
"context"
"slices"
"strings"
"testing"
"time"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes/fake"
)
const secretToken = "5ecret5ecret5ecret5ecret5ecret5ecret5ecret5ecret5ecret5ecret5ecr"
func params(mode string) JobParams {
return JobParams{
Server: "survival", ID: "0011223344556677", Mode: mode,
WorldPVC: "world-survival-0", BackupPVC: "felis-backups",
BackupRef: "/backups/survival-1.tar.gz",
TargetURL: "http://felis-api-internal.felis.svc.cluster.local:8081/api/v1/internal/exports/0011223344556677",
Token: secretToken,
Namespace: "minecraft", ServiceAccount: "felis-restore", Image: "felis:1",
BackupRoot: "/backups", WorldsRoot: "/world", Deadline: time.Hour,
CPULimit: "1", MemLimit: "256Mi", TTLAfterFinished: 5 * time.Minute,
}
}
// The export Pod reads a world or an archive and hands it to felis-api. It gets
// exactly the one volume it reads, read-only at both ends, no Secret, no API
// token, and no capability beyond DAC_OVERRIDE: a compromised export can read
// that one volume and talk to that one upload URL, nothing more.
func TestExportJobIsolation(t *testing.T) {
for _, tc := range []struct {
mode, volume, claim, mountPath string
args []string
}{
{ModeWorld, worldVolume, "world-survival-0", "/world", []string{"--worlds-root", "/world"}},
{ModeBackup, backupVolume, "felis-backups", "/backups", []string{"--ref", "/backups/survival-1.tar.gz", "--backup-root", "/backups"}},
} {
job, err := ExportJob(params(tc.mode))
if err != nil {
t.Fatalf("%s: ExportJob: %v", tc.mode, err)
}
if job.Name != "export-survival-0011223344556677" || job.Namespace != "minecraft" {
t.Errorf("%s: job = %s/%s", tc.mode, job.Namespace, job.Name)
}
for _, labels := range []map[string]string{job.Labels, job.Spec.Template.Labels} {
if labels[LabelManagedBy] != "felis-export" || labels[LabelServer] != "survival" || labels[LabelMode] != tc.mode {
t.Errorf("%s: labels = %v", tc.mode, labels)
}
}
spec := job.Spec.Template.Spec
if spec.ServiceAccountName != "felis-restore" {
t.Errorf("%s: service account = %q", tc.mode, spec.ServiceAccountName)
}
if spec.AutomountServiceAccountToken == nil || *spec.AutomountServiceAccountToken {
t.Errorf("%s: the SA token is mounted", tc.mode)
}
if spec.RestartPolicy != corev1.RestartPolicyNever {
t.Errorf("%s: restart policy = %q", tc.mode, spec.RestartPolicy)
}
if len(spec.Volumes) != 1 {
t.Fatalf("%s: volumes = %+v, want exactly one", tc.mode, spec.Volumes)
}
v := spec.Volumes[0]
if v.Name != tc.volume || v.PersistentVolumeClaim == nil || v.PersistentVolumeClaim.ClaimName != tc.claim || !v.PersistentVolumeClaim.ReadOnly {
t.Errorf("%s: volume = %+v, want %s read-only", tc.mode, v, tc.claim)
}
if len(spec.Containers) != 1 || len(spec.InitContainers) != 0 {
t.Fatalf("%s: containers = %d, init = %d", tc.mode, len(spec.Containers), len(spec.InitContainers))
}
c := spec.Containers[0]
if len(c.VolumeMounts) != 1 || c.VolumeMounts[0] != (corev1.VolumeMount{Name: tc.volume, MountPath: tc.mountPath, ReadOnly: true}) {
t.Errorf("%s: mounts = %+v", tc.mode, c.VolumeMounts)
}
if len(c.EnvFrom) != 0 || len(c.Env) != 1 || c.Env[0] != (corev1.EnvVar{Name: TokenEnv, Value: secretToken}) {
t.Errorf("%s: env = %+v, envFrom = %+v; want only the token", tc.mode, c.Env, c.EnvFrom)
}
if strings.Contains(strings.Join(c.Args, " "), secretToken) {
t.Errorf("%s: the token rides argv: %v", tc.mode, c.Args)
}
wantArgs := append([]string{"--mode", tc.mode, "--server", "survival", "--target-url", params(tc.mode).TargetURL}, tc.args...)
if !slices.Equal(c.Command, []string{felisBinaryPath, "export"}) || !slices.Equal(c.Args, wantArgs) {
t.Errorf("%s: command = %v %v", tc.mode, c.Command, c.Args)
}
sc := c.SecurityContext
if sc == nil || sc.Privileged == nil || *sc.Privileged || sc.AllowPrivilegeEscalation == nil || *sc.AllowPrivilegeEscalation ||
sc.ReadOnlyRootFilesystem == nil || !*sc.ReadOnlyRootFilesystem {
t.Errorf("%s: container security context = %+v", tc.mode, sc)
}
if sc.Capabilities == nil || !slices.Equal(sc.Capabilities.Drop, []corev1.Capability{"ALL"}) ||
!slices.Equal(sc.Capabilities.Add, []corev1.Capability{"DAC_OVERRIDE"}) {
t.Errorf("%s: capabilities = %+v", tc.mode, sc.Capabilities)
}
if c.TerminationMessagePolicy != corev1.TerminationMessageFallbackToLogsOnError {
t.Errorf("%s: termination message policy = %q", tc.mode, c.TerminationMessagePolicy)
}
if psc := spec.SecurityContext; psc == nil || psc.RunAsUser == nil || *psc.RunAsUser != 0 || psc.FSGroup != nil {
t.Errorf("%s: pod security context = %+v", tc.mode, psc)
}
if b := job.Spec.BackoffLimit; b == nil || *b != 0 {
t.Errorf("%s: backoffLimit = %v, want 0", tc.mode, b)
}
if d := job.Spec.ActiveDeadlineSeconds; d == nil || *d != 3600 {
t.Errorf("%s: activeDeadlineSeconds = %v, want 3600", tc.mode, d)
}
if ttl := job.Spec.TTLSecondsAfterFinished; ttl == nil || *ttl != 300 {
t.Errorf("%s: ttlSecondsAfterFinished = %v, want 300", tc.mode, ttl)
}
}
}
func TestExportJobRefusesIncompleteParams(t *testing.T) {
for name, edit := range map[string]func(*JobParams){
"no image": func(p *JobParams) { p.Image = "" },
"no token": func(p *JobParams) { p.Token = "" },
"no target": func(p *JobParams) { p.TargetURL = "" },
"no id": func(p *JobParams) { p.ID = "" },
"unknown mode": func(p *JobParams) { p.Mode = "both" },
"world without claim": func(p *JobParams) { p.WorldPVC = "" },
} {
p := params(ModeWorld)
edit(&p)
if _, err := ExportJob(p); err == nil {
t.Errorf("%s: ExportJob accepted %+v", name, p)
}
}
p := params(ModeBackup)
p.BackupRef = ""
if _, err := ExportJob(p); err == nil {
t.Error("a backup export without a ref was accepted")
}
}
// TestExportJobDefaults: a caller that leaves the deadline and the TTL unset
// still gets a Job that ends and is collected.
func TestExportJobDefaults(t *testing.T) {
p := params(ModeWorld)
p.Deadline, p.TTLAfterFinished = 0, 0
job, err := ExportJob(p)
if err != nil {
t.Fatal(err)
}
if d := job.Spec.ActiveDeadlineSeconds; d == nil || *d != 7200 {
t.Errorf("activeDeadlineSeconds = %v, want 7200", d)
}
if ttl := job.Spec.TTLSecondsAfterFinished; ttl == nil || *ttl != 600 {
t.Errorf("ttlSecondsAfterFinished = %v, want 600", ttl)
}
}
func TestStartCreatesTheJob(t *testing.T) {
cs := fake.NewSimpleClientset()
e := New(cs, Config{Image: "felis:1", BackupPVC: "felis-backups"})
name, err := e.Start(context.Background(), Request{
Server: "survival", Mode: ModeWorld, ID: "0011223344556677",
TargetURL: "http://api:8081/api/v1/internal/exports/0011223344556677", Token: secretToken,
})
if err != nil || name != "export-survival-0011223344556677" {
t.Fatalf("Start = %q, %v", name, err)
}
job, err := cs.BatchV1().Jobs("minecraft").Get(context.Background(), name, metav1.GetOptions{})
if err != nil {
t.Fatalf("the Job is not in the minecraft namespace: %v", err)
}
if claim := job.Spec.Template.Spec.Volumes[0].PersistentVolumeClaim.ClaimName; claim != "world-survival-0" {
t.Errorf("world claim = %q, want world-survival-0", claim)
}
if d := *job.Spec.ActiveDeadlineSeconds; d != int64(defaultDeadline/time.Second) {
t.Errorf("deadline = %d, want the %s default", d, defaultDeadline)
}
if _, err := e.Start(context.Background(), Request{Server: "survival", Mode: ModeWorld, ID: "0011223344556677",
TargetURL: "http://api:8081/x", Token: secretToken}); err == nil {
t.Error("a second Job of the same name was reported as created")
}
}
+125
View File
@@ -0,0 +1,125 @@
// Package worldexport is the executor behind the world export routes (POST
// /servers/{name}/world/export and POST /servers/{name}/backups/{id}/export): a
// one-shot Job in the minecraft namespace that reads a stopped server's world,
// or one of its stored archives, and PUTs the tar.gz to felis-api's internal
// face, which streams it on to the owner's browser as it arrives
// (internal/api/exports.go). Nothing is staged on the way: felis-api never
// mounts a world or the backup store, and the archive never lands on a disk it
// owns.
//
// Trust model as in internal/restore: the Pod runs under the weak felis-restore
// SA with no API token, mounts one volume read-only, and holds no database URL
// or Secret. The only credential it gets is the one-time token of this one
// upload. felis-api makes every authorization decision before the Job exists.
package worldexport
import (
"context"
"fmt"
"time"
"felis.lolicon.best/internal/naming"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
)
// Request is one export as felis-api admitted it.
type Request struct {
Server string
Mode string // ModeWorld or ModeBackup
// BackupRef is the archive a ModeBackup export reads.
BackupRef string
// ID names the export (JobName) and TargetURL/Token are where and how the
// Pod hands the archive over.
ID string
TargetURL string
Token string
}
// Config parameterises the executor. Image and BackupPVC have no safe default:
// cmd/felis leaves the API's Exporter nil when either is missing, and the
// routes answer 503.
type Config struct {
Namespace string
ServiceAccount string
Image string
BackupPVC string
// BackupRoot MUST be the path the archives were written under
// (cfg.Archive.LocalPath): the stored refs are absolute paths.
BackupRoot string
WorldsRoot string
// Deadline caps the Pod's wall-clock. The archive moves at the browser's
// pace, so it is longer than a backup's, and it is also the longest a world
// export can keep the server from starting.
Deadline time.Duration
CPULimit string
MemLimit string
RunAsUser int64
RunAsGroup int64
FSGroup int64
TTLAfterFinished time.Duration
}
const (
defaultNamespace = "minecraft"
defaultServiceAccount = "felis-restore"
defaultBackupRoot = "/backups"
defaultWorldsRoot = "/world"
defaultDeadline = 2 * time.Hour
defaultCPULimit = "1"
defaultMemLimit = "256Mi"
defaultTTL = 10 * time.Minute
)
func (c Config) withDefaults() Config {
if c.Namespace == "" {
c.Namespace = defaultNamespace
}
if c.ServiceAccount == "" {
c.ServiceAccount = defaultServiceAccount
}
if c.BackupRoot == "" {
c.BackupRoot = defaultBackupRoot
}
if c.WorldsRoot == "" {
c.WorldsRoot = defaultWorldsRoot
}
if c.Deadline <= 0 {
c.Deadline = defaultDeadline
}
return c
}
// Exporter is the production internal/api.Exporter. It creates the Job with
// jobs:create, which felis-api already holds in the minecraft namespace.
type Exporter struct {
cs kubernetes.Interface
cfg Config
}
// New builds an Exporter over the typed clientset.
func New(cs kubernetes.Interface, cfg Config) *Exporter {
return &Exporter{cs: cs, cfg: cfg.withDefaults()}
}
// Start creates the export Job for r and returns its name.
func (e *Exporter) Start(ctx context.Context, r Request) (string, error) {
c := e.cfg
job, err := ExportJob(JobParams{
Server: r.Server, ID: r.ID, Mode: r.Mode,
WorldPVC: naming.WorldPVCName(r.Server), BackupPVC: c.BackupPVC, BackupRef: r.BackupRef,
TargetURL: r.TargetURL, Token: r.Token,
Namespace: c.Namespace, ServiceAccount: c.ServiceAccount, Image: c.Image,
BackupRoot: c.BackupRoot, WorldsRoot: c.WorldsRoot, Deadline: c.Deadline,
CPULimit: c.CPULimit, MemLimit: c.MemLimit,
RunAsUser: c.RunAsUser, RunAsGroup: c.RunAsGroup, FSGroup: c.FSGroup,
TTLAfterFinished: c.TTLAfterFinished,
})
if err != nil {
return "", err
}
if _, err := e.cs.BatchV1().Jobs(c.Namespace).Create(ctx, job, metav1.CreateOptions{}); err != nil {
return "", fmt.Errorf("worldexport: create export job: %w", err)
}
return job.Name, nil
}