feat(api): expose async backup/restore job status (fixes #7)
Backup and restore only enqueue a cluster Job; a later failure left its
only trace in that Job object, invisible without kubectl. Add
GET /api/v1/servers/{name}/jobs (owner-or-admin) projecting the newest
20 managed Jobs (felis-backup / felis-restore) as
running|succeeded|failed with message and timestamps. Nil reader -> 503
jobs_unavailable, mirroring the backup/restore feature gates. RBAC gains
jobs:list; OpenAPI parity updated.
This commit is contained in:
8 files changed
+365
-2
No files matched your search
@@ -63,6 +63,11 @@ type API struct {
|
||||
// authorization boundary is exercised before the backup-Job executor is wired.
|
||||
Backuper Backuper
|
||||
|
||||
// JobStatus reads the latest backup/restore Job outcomes for GET
|
||||
// /servers/{name}/jobs — the status outlet for async failures (the enqueue
|
||||
// endpoints only answer 202). Optional: nil → that route reports 503.
|
||||
JobStatus JobStatusReader
|
||||
|
||||
// Files is the server file editor (list / read / write a file in a stopped
|
||||
// server's world volume — the "one wrong line in server.properties" repair).
|
||||
// Like Restorer and Backuper it is optional: when nil the file routes report
|
||||
@@ -407,6 +412,7 @@ func (a *API) externalAPIRoutes() []apiRoute {
|
||||
// owned), and restore is gated by owner-or-admin PLUS a former-owner match, so
|
||||
// neither sits behind adminOnly.
|
||||
{Method: "GET", Pattern: "/api/v1/backups", h: a.handleListBackups},
|
||||
{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},
|
||||
// Server file editor: list / read / write a file in a STOPPED server's world
|
||||
|
||||
@@ -0,0 +1,65 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/naming"
|
||||
)
|
||||
|
||||
// AsyncJob is the observable outcome of one asynchronous world operation. The API
|
||||
// only ENQUEUES backup/restore Jobs — the work, and any failure, happens in the
|
||||
// cluster — so without this projection a failed job left its only trace in a Job
|
||||
// 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"
|
||||
State string `json:"state"` // "running" | "succeeded" | "failed"
|
||||
Message string `json:"message,omitempty"`
|
||||
StartedAt time.Time `json:"started_at,omitempty"`
|
||||
FinishedAt time.Time `json:"finished_at,omitempty"`
|
||||
}
|
||||
|
||||
// JobStatusReader reads the newest backup/restore Jobs for a server, newest
|
||||
// first. Optional like Restorer/Backuper: when nil the route answers 503.
|
||||
type JobStatusReader interface {
|
||||
LatestJobs(ctx context.Context, serverName string) ([]AsyncJob, error)
|
||||
}
|
||||
|
||||
// handleServerJobs serves GET /api/v1/servers/{name}/jobs — the latest async world
|
||||
// operations for one server, so a 202 that later failed is visible without kubectl.
|
||||
// Authorization mirrors the backup/restore gates' front half (owner-or-admin); the
|
||||
// message is free-form Job text and can name paths the owner already sees through
|
||||
// the file editor.
|
||||
func (a *API) handleServerJobs(w http.ResponseWriter, r *http.Request) {
|
||||
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
|
||||
}
|
||||
rec, err := a.Repo.ServerByName(r.Context(), name)
|
||||
if err != nil {
|
||||
a.writeLookupError(w, r, err)
|
||||
return
|
||||
}
|
||||
if !a.isOwnerOrAdmin(p, rec) {
|
||||
writeError(w, r, errForbidden)
|
||||
return
|
||||
}
|
||||
if a.JobStatus == nil {
|
||||
writeError(w, r, newError(http.StatusServiceUnavailable, "jobs_unavailable",
|
||||
"job status is not configured"))
|
||||
return
|
||||
}
|
||||
jobs, err := a.JobStatus.LatestJobs(r.Context(), name)
|
||||
if err != nil {
|
||||
writeError(w, r, err)
|
||||
return
|
||||
}
|
||||
if jobs == nil {
|
||||
jobs = []AsyncJob{}
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{"server": name, "jobs": jobs})
|
||||
}
|
||||
@@ -0,0 +1,147 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
batchv1 "k8s.io/api/batch/v1"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
)
|
||||
|
||||
type fakeJobStatus struct {
|
||||
jobs []AsyncJob
|
||||
err error
|
||||
got string
|
||||
}
|
||||
|
||||
func (f *fakeJobStatus) LatestJobs(_ context.Context, server string) ([]AsyncJob, error) {
|
||||
f.got = server
|
||||
return f.jobs, f.err
|
||||
}
|
||||
|
||||
// TestServerJobsHandler covers GET /api/v1/servers/{name}/jobs: the owner sees the
|
||||
// recorded outcomes, strangers are refused, a nil reader is an honest 503, and a
|
||||
// reader error surfaces as 500.
|
||||
func TestServerJobsHandler(t *testing.T) {
|
||||
owner := &Principal{UserID: "owner1", Email: "[email protected]", Role: "user"}
|
||||
stranger := &Principal{UserID: "other", Email: "[email protected]", Role: "user"}
|
||||
|
||||
mk := func(reader JobStatusReader) (*API, *fakeJobStatus) {
|
||||
repo := newFakeRepo()
|
||||
repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner1"}
|
||||
a := newTestAPI(repo, newFakeCluster())
|
||||
a.External = staticExternal{p: owner}
|
||||
var fjs *fakeJobStatus
|
||||
if reader != nil {
|
||||
a.JobStatus = reader
|
||||
}
|
||||
if fj, ok := reader.(*fakeJobStatus); ok {
|
||||
fjs = fj
|
||||
}
|
||||
return a, fjs
|
||||
}
|
||||
|
||||
t.Run("owner sees failed job", func(t *testing.T) {
|
||||
fjs := &fakeJobStatus{jobs: []AsyncJob{{Name: "restore-survival", Kind: "restore", State: "failed", Message: "backoff limit exceeded"}}}
|
||||
a, _ := mk(fjs)
|
||||
w := do(a.ExternalHandler(), "GET", "/api/v1/servers/survival/jobs", "", nil)
|
||||
if w.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d, want 200 (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
var resp struct {
|
||||
Server string `json:"server"`
|
||||
Jobs []AsyncJob `json:"jobs"`
|
||||
}
|
||||
if err := json.Unmarshal(w.Body.Bytes(), &resp); err != nil {
|
||||
t.Fatalf("bad JSON: %v", err)
|
||||
}
|
||||
if resp.Server != "survival" || len(resp.Jobs) != 1 || resp.Jobs[0].State != "failed" {
|
||||
t.Fatalf("unexpected payload %+v", resp)
|
||||
}
|
||||
if fjs.got != "survival" {
|
||||
t.Fatalf("reader asked for %q", fjs.got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("stranger -> 403", func(t *testing.T) {
|
||||
fjs := &fakeJobStatus{}
|
||||
a, _ := mk(fjs)
|
||||
a.External = staticExternal{p: stranger}
|
||||
if w := do(a.ExternalHandler(), "GET", "/api/v1/servers/survival/jobs", "", nil); w.Code != http.StatusForbidden {
|
||||
t.Fatalf("code = %d, want 403", w.Code)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("nil reader -> 503", func(t *testing.T) {
|
||||
a, _ := mk(nil)
|
||||
if w := do(a.ExternalHandler(), "GET", "/api/v1/servers/survival/jobs", "", nil); w.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("code = %d, want 503", w.Code)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("reader error -> 500", func(t *testing.T) {
|
||||
fjs := &fakeJobStatus{err: errors.New("apiserver down")}
|
||||
a, _ := mk(fjs)
|
||||
if w := do(a.ExternalHandler(), "GET", "/api/v1/servers/survival/jobs", "", nil); w.Code != http.StatusInternalServerError {
|
||||
t.Fatalf("code = %d, want 500", w.Code)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("unknown server -> 404", func(t *testing.T) {
|
||||
a, _ := mk(&fakeJobStatus{})
|
||||
if w := do(a.ExternalHandler(), "GET", "/api/v1/servers/ghost/jobs", "", nil); w.Code != http.StatusNotFound {
|
||||
t.Fatalf("code = %d, want 404", w.Code)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// TestJobToAsyncJob pins the Job→AsyncJob projection: kinds come from managed-by,
|
||||
// terminal conditions decide state, and unknown owners are dropped.
|
||||
func TestJobToAsyncJob(t *testing.T) {
|
||||
start := metav1.NewTime(time.Date(2026, 9, 22, 10, 0, 0, 0, time.UTC))
|
||||
done := metav1.NewTime(time.Date(2026, 9, 22, 10, 5, 0, 0, time.UTC))
|
||||
|
||||
failed := &batchv1.Job{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: "restore-survival", Labels: map[string]string{jobManagedByLabel: jobManagedByRestore}},
|
||||
Status: batchv1.JobStatus{
|
||||
StartTime: &start,
|
||||
Conditions: []batchv1.JobCondition{{
|
||||
Type: batchv1.JobFailed, Status: corev1.ConditionTrue,
|
||||
Reason: "BackoffLimitExceeded", Message: "Job has reached the specified backoff limit",
|
||||
LastTransitionTime: done,
|
||||
}},
|
||||
},
|
||||
}
|
||||
aj, ok := jobToAsyncJob(failed)
|
||||
if !ok || aj.Kind != "restore" || aj.State != "failed" || aj.Name != "restore-survival" {
|
||||
t.Fatalf("failed job projection = %+v ok=%v", aj, ok)
|
||||
}
|
||||
if !aj.FinishedAt.Equal(done.Time) || !aj.StartedAt.Equal(start.Time) {
|
||||
t.Fatalf("timestamps = %+v", aj)
|
||||
}
|
||||
|
||||
complete := &batchv1.Job{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: "backup-survival-1", Labels: map[string]string{jobManagedByLabel: jobManagedByBackup}},
|
||||
Status: batchv1.JobStatus{Conditions: []batchv1.JobCondition{{
|
||||
Type: batchv1.JobComplete, Status: corev1.ConditionTrue, LastTransitionTime: done,
|
||||
}}},
|
||||
}
|
||||
if aj, ok := jobToAsyncJob(complete); !ok || aj.Kind != "backup" || aj.State != "succeeded" {
|
||||
t.Fatalf("complete job projection = %+v ok=%v", aj, ok)
|
||||
}
|
||||
|
||||
running := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{jobManagedByLabel: jobManagedByBackup}}}
|
||||
if aj, ok := jobToAsyncJob(running); !ok || aj.State != "running" {
|
||||
t.Fatalf("running job projection = %+v ok=%v", aj, ok)
|
||||
}
|
||||
|
||||
foreign := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{jobManagedByLabel: "someone-else"}}}
|
||||
if _, ok := jobToAsyncJob(foreign); ok {
|
||||
t.Fatal("foreign job must be dropped")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,92 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sort"
|
||||
|
||||
batchv1 "k8s.io/api/batch/v1"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
)
|
||||
|
||||
// Labels the backup/restore executors apply (internal/backupjob and
|
||||
// internal/restore keep their own copies). Literals on purpose: api defines the
|
||||
// read seam and must not import the executors — their dependency direction is
|
||||
// "executors implement api's interfaces", never the reverse.
|
||||
const (
|
||||
jobServerLabel = "felis.lolicon.best/server"
|
||||
jobManagedByLabel = "app.kubernetes.io/managed-by"
|
||||
|
||||
jobManagedByBackup = "felis-backup"
|
||||
jobManagedByRestore = "felis-restore"
|
||||
)
|
||||
|
||||
// K8sJobStatus reads the async Jobs the executors created, by the server label
|
||||
// both apply.
|
||||
type K8sJobStatus struct {
|
||||
c client.Client
|
||||
namespace string
|
||||
}
|
||||
|
||||
// NewK8sJobStatus builds the reader over the cluster client and the namespace the
|
||||
// server workloads (and their Jobs) live in.
|
||||
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.
|
||||
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),
|
||||
client.MatchingLabels{jobServerLabel: serverName}); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make([]AsyncJob, 0, len(list.Items))
|
||||
for i := range list.Items {
|
||||
if aj, ok := jobToAsyncJob(&list.Items[i]); ok {
|
||||
out = append(out, aj)
|
||||
}
|
||||
}
|
||||
sort.Slice(out, func(i, j int) bool { return out[i].StartedAt.After(out[j].StartedAt) })
|
||||
if len(out) > 20 {
|
||||
out = out[:20]
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// jobToAsyncJob projects one Job onto its kind/state/message. Complete condition →
|
||||
// succeeded, Failed → failed with its reason (Job conditions carry the generic
|
||||
// "backoff limit exceeded" text; the pod log holds the underlying error), anything
|
||||
// else is still running.
|
||||
func jobToAsyncJob(j *batchv1.Job) (AsyncJob, bool) {
|
||||
kind := ""
|
||||
switch j.Labels[jobManagedByLabel] {
|
||||
case jobManagedByBackup:
|
||||
kind = "backup"
|
||||
case jobManagedByRestore:
|
||||
kind = "restore"
|
||||
default:
|
||||
return AsyncJob{}, false
|
||||
}
|
||||
aj := AsyncJob{Name: j.Name, Kind: kind, State: "running", StartedAt: j.CreationTimestamp.Time}
|
||||
if j.Status.StartTime != nil {
|
||||
aj.StartedAt = j.Status.StartTime.Time
|
||||
}
|
||||
for _, c := range j.Status.Conditions {
|
||||
switch {
|
||||
case c.Type == batchv1.JobComplete && c.Status == corev1.ConditionTrue:
|
||||
aj.State = "succeeded"
|
||||
aj.FinishedAt = c.LastTransitionTime.Time
|
||||
case c.Type == batchv1.JobFailed && c.Status == corev1.ConditionTrue:
|
||||
aj.State = "failed"
|
||||
aj.FinishedAt = c.LastTransitionTime.Time
|
||||
aj.Message = c.Message
|
||||
if aj.Message == "" {
|
||||
aj.Message = c.Reason
|
||||
}
|
||||
}
|
||||
}
|
||||
return aj, true
|
||||
}
|
||||
Reference in new issue
Block a user