perf(panel): speed up file browsing and preserve loading forms

This commit is contained in:
Lemon-miaow committed 2026-10-04 20:18:11 +08:00
1 parent 4d8e9af240
commit 2c98e8b2ab
27 files changed
+1387 -145

No files matched your search

+266
View File
@@ -0,0 +1,266 @@
package fileedit
import (
"bytes"
"context"
"crypto/sha256"
"crypto/subtle"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"strings"
"sync"
"time"
)
// Browse runs inside the read-only Job and reuses Execute for every command.
// Refuse write commands even if the other end violates the protocol.
func Browse(ctx context.Context, worldRoot, url, token string) error {
if token == "" {
return fmt.Errorf("fileedit: browser token is required")
}
client := &http.Client{Timeout: browserIdle + 30*time.Second,
CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }}
answer := BrowseAnswer{}
for {
body, err := json.Marshal(answer)
if err != nil {
return err
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(body))
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("Content-Type", "application/json")
resp, err := client.Do(req)
if err != nil {
return err
}
if resp.StatusCode == http.StatusNoContent {
resp.Body.Close()
return nil
}
if resp.StatusCode != http.StatusOK {
resp.Body.Close()
return fmt.Errorf("fileedit: browser exchange returned %s", resp.Status)
}
var command BrowseCommand
err = json.NewDecoder(io.LimitReader(resp.Body, 64<<10)).Decode(&command)
resp.Body.Close()
if err != nil {
return err
}
if command.Op != OpList && command.Op != OpRead {
return fmt.Errorf("fileedit: browser refused operation %q", command.Op)
}
res, err := Execute(worldRoot, Request{Op: command.Op, Path: command.Path})
answer = BrowseAnswer{ID: command.ID}
if err != nil {
answer.Error = err.Error()
continue
}
answer.Result, err = json.Marshal(res)
if err != nil {
return err
}
}
}
const (
BrowserTokenEnv = "FELIS_FILE_BROWSER_TOKEN"
BrowserRoute = "/api/v1/internal/file-browser/"
browserIdle = 45 * time.Second
browserLifetime = 4 * time.Minute
maxBrowsers = 4
)
var errBrowserFull = errors.New("fileedit: browser capacity reached")
// Browser reuses a short-lived, read-only Job per world. Jobs pull commands from
// the existing internal API face: no inbound Pod port or new RBAC is needed.
// Every browser request is still authorized by the normal file API handlers.
type Browser struct {
BaseURL string
mu sync.Mutex
sessions map[string]*browseSession
}
type BrowseCommand struct {
ID string `json:"id"`
Op string `json:"op"`
Path string `json:"path"`
}
type BrowseAnswer struct {
ID string `json:"id"`
Result json.RawMessage `json:"result,omitempty"`
Error string `json:"error,omitempty"`
}
type browseSession struct {
id, key string
tokenHash [sha256.Size]byte
commands chan BrowseCommand
answers chan BrowseAnswer
serial chan struct{}
done chan struct{}
once sync.Once
timer *time.Timer
mu sync.Mutex
connected bool
pending string
}
func (b *Browser) close(s *browseSession) {
b.mu.Lock()
if b.sessions[s.key] == s {
delete(b.sessions, s.key)
}
b.mu.Unlock()
s.once.Do(func() {
if s.timer != nil {
s.timer.Stop()
}
close(s.done)
})
}
func (b *Browser) session(ctx context.Context, p JobParams, start func(context.Context, JobParams) error) (*browseSession, error) {
key := strings.Join([]string{p.Namespace, p.Server, p.WorldPVC, p.NodeName, p.Image}, "\x00")
b.mu.Lock()
defer b.mu.Unlock()
if s := b.sessions[key]; s != nil {
return s, nil
}
if len(b.sessions) >= maxBrowsers {
return nil, errBrowserFull
}
token, err := randomHex(32)
if err != nil {
return nil, err
}
s := &browseSession{id: p.OpID, key: key, tokenHash: sha256.Sum256([]byte(token)),
commands: make(chan BrowseCommand, 1), answers: make(chan BrowseAnswer, 1), serial: make(chan struct{}, 1), done: make(chan struct{})}
if b.sessions == nil {
b.sessions = make(map[string]*browseSession)
}
b.sessions[key] = s
p.BrowserURL, p.BrowserToken = strings.TrimRight(b.BaseURL, "/")+BrowserRoute+s.id, token
p.Deadline = browserLifetime
if err := start(ctx, p); err != nil {
delete(b.sessions, key)
return nil, err
}
if err := ctx.Err(); err != nil {
delete(b.sessions, key)
return nil, err
}
s.timer = time.AfterFunc(browserLifetime, func() { b.close(s) })
return s, nil
}
func (b *Browser) Run(ctx context.Context, p JobParams, start func(context.Context, JobParams) error) ([]byte, error) {
s, err := b.session(ctx, p, start)
if err != nil {
return nil, err
}
select {
case s.serial <- struct{}{}:
case <-ctx.Done():
return nil, ctx.Err()
case <-s.done:
return nil, fmt.Errorf("fileedit: browser closed")
}
defer func() { <-s.serial }()
s.mu.Lock()
s.pending = p.OpID
s.mu.Unlock()
select {
case s.commands <- BrowseCommand{ID: p.OpID, Op: p.Op, Path: p.Path}:
case <-ctx.Done():
b.close(s)
return nil, ctx.Err()
case <-s.done:
return nil, fmt.Errorf("fileedit: browser closed")
}
select {
case answer := <-s.answers:
if answer.Error != "" {
return nil, fmt.Errorf("fileedit: browser read: %s", answer.Error)
}
return answer.Result, nil
case <-ctx.Done():
b.close(s)
return nil, ctx.Err()
case <-s.done:
return nil, fmt.Errorf("fileedit: browser closed before returning a result")
}
}
// ServeHTTP accepts one worker's result and long-polls its next command. The
// random token opens only this world/session, and lives only in the Job's env.
func (b *Browser) ServeHTTP(w http.ResponseWriter, r *http.Request) {
id := strings.TrimPrefix(r.URL.Path, BrowserRoute)
b.mu.Lock()
var s *browseSession
for _, candidate := range b.sessions {
if candidate.id == id {
s = candidate
break
}
}
b.mu.Unlock()
token, bearer := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer ")
sum := sha256.Sum256([]byte(token))
if s == nil || !bearer || subtle.ConstantTimeCompare(sum[:], s.tokenHash[:]) != 1 {
http.Error(w, "unknown browser", http.StatusNotFound)
return
}
var answer BrowseAnswer
decoder := json.NewDecoder(http.MaxBytesReader(w, r.Body, maxLogBytes))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&answer); err != nil {
http.Error(w, "invalid result", http.StatusBadRequest)
return
}
s.mu.Lock()
valid := !s.connected && answer.ID == "" || s.connected && answer.ID != "" && answer.ID == s.pending
if valid {
s.connected = true
if answer.ID != "" {
s.pending = ""
}
}
s.mu.Unlock()
if !valid {
http.Error(w, "unexpected result", http.StatusConflict)
return
}
if answer.ID != "" {
select {
case s.answers <- answer:
case <-s.done:
w.WriteHeader(http.StatusGone)
return
}
}
timer := time.NewTimer(browserIdle)
defer timer.Stop()
select {
case command := <-s.commands:
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(command); err != nil {
b.close(s)
}
case <-timer.C:
b.close(s)
w.WriteHeader(http.StatusNoContent)
case <-s.done:
w.WriteHeader(http.StatusNoContent)
case <-r.Context().Done():
b.close(s)
}
}
+199
View File
@@ -0,0 +1,199 @@
package fileedit
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
)
func TestBrowserReusesWorkerAndReadsCurrentBytes(t *testing.T) {
root := t.TempDir()
if err := os.WriteFile(filepath.Join(root, "config.yml"), []byte("enabled: true\n"), 0600); err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
b := &Browser{}
srv := httptest.NewServer(b)
b.BaseURL = srv.URL
defer srv.Close()
defer cancel()
var workers atomic.Int32
start := func(_ context.Context, p JobParams) error {
workers.Add(1)
job, err := FilesJob(p)
if err != nil {
return err
}
pod := job.Spec.Template.Spec
if pod.AutomountServiceAccountToken == nil || *pod.AutomountServiceAccountToken || len(pod.Volumes) != 1 || !pod.Containers[0].VolumeMounts[0].ReadOnly {
t.Error("browser weakened the file Job's isolation")
}
go func() {
if err := Browse(ctx, root, p.BrowserURL, p.BrowserToken); err != nil && ctx.Err() == nil {
t.Errorf("worker: %v", err)
}
}()
return nil
}
read := func(op, path string) Result {
p := testParams(op)
p.Path = path
p.OpID, _ = newOpID()
payload, err := b.Run(ctx, p, start)
if err != nil {
t.Fatal(err)
}
var result Result
if err := json.Unmarshal(payload, &result); err != nil {
t.Fatal(err)
}
return result
}
if got := read(OpRead, "config.yml"); string(got.Content) != "enabled: true\n" {
t.Fatalf("first read: %+v", got)
}
if err := os.WriteFile(filepath.Join(root, "config.yml"), []byte("enabled: false\n"), 0600); err != nil {
t.Fatal(err)
}
if got := read(OpRead, "config.yml"); string(got.Content) != "enabled: false\n" {
t.Fatalf("stale read: %+v", got)
}
if got := read(OpList, ""); len(got.Entries) != 1 {
t.Fatalf("listing: %+v", got)
}
if got := read(OpRead, "../outside"); got.Code != CodeBadPath {
t.Fatalf("containment: %+v", got)
}
if workers.Load() != 1 {
t.Fatalf("started %d workers for repeated reads", workers.Load())
}
// Simultaneous readers still receive their own result, never the other path's.
var wg sync.WaitGroup
for _, path := range []string{"config.yml", "missing.yml"} {
wg.Add(1)
go func() {
defer wg.Done()
p := testParams(OpRead)
p.Path = path
p.OpID, _ = newOpID()
payload, err := b.Run(ctx, p, start)
if err != nil {
t.Error(err)
return
}
var got Result
json.Unmarshal(payload, &got)
if path == "missing.yml" && got.Code != CodeNotFound || path == "config.yml" && string(got.Content) != "enabled: false\n" {
t.Errorf("%s: %+v", path, got)
}
}()
}
wg.Wait()
}
func TestBrowserCancellationRevokesToken(t *testing.T) {
b := &Browser{BaseURL: "http://internal"}
var started JobParams
ctx, cancel := context.WithCancel(context.Background())
p := testParams(OpRead)
_, err := b.Run(ctx, p, func(_ context.Context, p JobParams) error { started = p; cancel(); return nil })
if err == nil {
t.Fatal("cancelled read succeeded")
}
r := httptest.NewRequest(http.MethodPost, started.BrowserURL, strings.NewReader(`{}`))
r.Header.Set("Authorization", "Bearer "+started.BrowserToken)
w := httptest.NewRecorder()
b.ServeHTTP(w, r)
if w.Code != http.StatusNotFound {
t.Fatalf("cancelled token returned %d", w.Code)
}
}
func TestBrowserRejectsWrongTokenAndStaleResult(t *testing.T) {
b := &Browser{BaseURL: "http://internal"}
p := testParams(OpRead)
var started JobParams
s, err := b.session(context.Background(), p, func(_ context.Context, p JobParams) error { started = p; return nil })
if err != nil {
t.Fatal(err)
}
defer b.close(s)
for _, header := range []string{"", "Bearer wrong", started.BrowserToken} {
r := httptest.NewRequest(http.MethodPost, started.BrowserURL, strings.NewReader(`{}`))
r.Header.Set("Authorization", header)
w := httptest.NewRecorder()
b.ServeHTTP(w, r)
if w.Code != http.StatusNotFound {
t.Fatalf("wrong token returned %d", w.Code)
}
}
r := httptest.NewRequest(http.MethodPost, started.BrowserURL, strings.NewReader(`{"id":"another-read","result":{}}`))
r.Header.Set("Authorization", "Bearer "+started.BrowserToken)
w := httptest.NewRecorder()
b.ServeHTTP(w, r)
if w.Code != http.StatusConflict {
t.Fatalf("stale result returned %d", w.Code)
}
}
func TestBrowserCapacityReleasedAfterClosingSession(t *testing.T) {
b := &Browser{BaseURL: "http://internal"}
defer func() {
for _, s := range b.sessions {
b.close(s)
}
}()
started := 0
start := func(context.Context, JobParams) error { started++; return nil }
for i := 0; i < maxBrowsers; i++ {
p := testParams(OpList)
p.Server = string(rune('a' + i))
p.OpID = p.Server
if _, err := b.session(context.Background(), p, start); err != nil {
t.Fatal(err)
}
}
p := testParams(OpRead)
if _, err := b.session(context.Background(), p, start); err != errBrowserFull || started != maxBrowsers {
t.Fatalf("capacity: err=%v, started=%d", err, started)
}
for _, s := range b.sessions {
b.close(s)
break
}
if _, err := b.session(context.Background(), p, start); err != nil || started != maxBrowsers+1 {
t.Fatalf("reopen: err=%v, started=%d", err, started)
}
}
func TestBrowseRefusesMutation(t *testing.T) {
root := t.TempDir()
path := filepath.Join(root, "keep.txt")
os.WriteFile(path, []byte("keep"), 0600)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(BrowseCommand{ID: "x", Op: OpDelete, Path: "keep.txt"})
}))
defer srv.Close()
if err := Browse(context.Background(), root, srv.URL, "token"); err == nil {
t.Fatal("worker accepted a write command")
}
got, err := os.ReadFile(path)
if err != nil || string(got) != "keep" {
t.Fatalf("file changed: %q, %v", got, err)
}
p := testParams(OpWrite)
p.BrowserURL = srv.URL
p.BrowserToken = "token"
if _, err := FilesJob(p); err == nil {
t.Fatal("writable browser Job was allowed")
}
}
+6 -6
View File
@@ -32,14 +32,14 @@
// upload fetched by the Job from felis-api's internal face (see Stage), because a
// 64 MiB jar fits in neither a Job spec nor an environment.
//
// The price is latency: every operation is a Pod schedule + image pull, so a
// listing takes seconds rather than milliseconds. That is inherent to RWO plus a
// stopped server, not a property of this transport: the file manager works on a
// stopped server, one operation per Job.
// Interactive reads reuse a bounded, short-lived, read-only Job through Browser.
// Its command channel uses the existing internal API face. At capacity, reads
// use one-shot Jobs and stream results before the Pod's terminal phase. Changes
// still wait for termination to preserve the world lock's ordering.
//
// The Editor depends on the Runner interface, so the orchestration and the error
// mapping are unit-tested against an in-memory fake; the client-go implementation
// (k8sjobs.go) compiles here but is exercised only against a live cluster.
// mapping are unit-tested against an in-memory fake; the client-go transport
// (k8sjobs.go) is tested against the Kubernetes HTTP API shape.
package fileedit
import (
+16 -5
View File
@@ -69,8 +69,10 @@ type JobParams struct {
// Async marks a Job felis-api does not wait on: it carries LabelAsync and
// AnnotationPath, which Ops reads it back by. Only an upload or an unzip
// runs so.
Async bool
WorldPVC string
Async bool
WorldPVC string
BrowserURL string
BrowserToken string
Namespace string
ServiceAccount string
@@ -128,7 +130,7 @@ func filesLabels(p JobParams) map[string]string {
// - mounts EXACTLY ONE volume — the world PVC — and NO Secret, NO ConfigMap, and
// NO backup PVC. It is therefore strictly blinder than the backup Pod, which
// mounts the config Secret to self-record its row: a file-editor Pod has nothing
// to record, so it is handed no database URL and no credential of any kind (the
// to record, so it is handed no database URL or platform credential (the
// four-power red line, spec §22);
// - mounts that one volume READ-ONLY for list and read (see mutates), so the
// two operations that only look physically cannot change anything — the
@@ -147,8 +149,12 @@ func filesLabels(p JobParams) map[string]string {
//
// The container runs `/usr/local/bin/felis files` (cmd/felis), which performs the
// operation under os.Root containment and prints the marked JSON Result line that
// felis-api reads back through pods/log.
// felis-api reads back through pods/log. A browser instead pulls read commands
// from the internal API and posts each result there using a scoped token.
func FilesJob(p JobParams) (*batchv1.Job, error) {
if p.BrowserURL != "" && (mutates(p.Op) || p.BrowserToken == "" || p.Async) {
return nil, fmt.Errorf("fileedit: a browser must be read-only and have a token")
}
if p.Image == "" {
return nil, fmt.Errorf("fileedit: image is empty")
}
@@ -246,7 +252,7 @@ func FilesJob(p JobParams) (*batchv1.Job, error) {
// so the spec is the sole channel into the Pod; base64 keeps arbitrary bytes —
// CRLF line endings, a UTF-8 BOM, a binary blob — intact through a field that
// must be a valid string. Content is set ONLY for a write and the token ONLY
// for an upload, so no other Job spec carries either.
// for an upload or read-only browser, so no other Job spec carries either.
switch p.Op {
case OpWrite:
parts := splitContent(p.Content)
@@ -257,6 +263,11 @@ func FilesJob(p JobParams) (*batchv1.Job, error) {
case OpUpload:
container.Env = []corev1.EnvVar{{Name: UploadTokenEnv, Value: p.UploadToken}}
}
if p.BrowserURL != "" {
container.Args = []string{"--browse-url", p.BrowserURL, "--worlds-root", p.WorldsRoot}
container.Env = []corev1.EnvVar{{Name: BrowserTokenEnv, Value: p.BrowserToken}}
container.Resources.Requests = corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("10m"), corev1.ResourceMemory: resource.MustParse("32Mi")}
}
job := &batchv1.Job{
ObjectMeta: metav1.ObjectMeta{
+48 -28
View File
@@ -18,12 +18,10 @@ import (
)
// pollInterval is how often the runner re-Lists Pods while waiting for the file
// Job to finish. It is a POLL rather than a Watch because felis-api holds
// container to start. It is a POLL rather than a Watch because felis-api holds
// pods:list and NOT pods:watch (internal/platform.APIMinecraftRole) — establishing
// a watch would need a permission this design exists to avoid. Half a second is
// well inside the human-perceptible floor for an operation already dominated by
// Pod scheduling, while keeping the request count on a slow image pull modest.
const pollInterval = 500 * time.Millisecond
// a watch would need a permission this design exists to avoid.
const pollInterval = 150 * time.Millisecond
// maxLogBytes bounds what the runner will buffer from a Pod's log. The payload is
// at most a base64-encoded MaxReadBytes (≈4/3 of 1 MiB) plus the JSON envelope, so
@@ -45,12 +43,11 @@ const maxLogBytes = 4 << 20
// the Job create is available on both — so one client covers all three calls
// instead of the binding carrying two.
//
// INTEGRATION-ONLY: like K8sLogStreamer and K8sCluster this needs a live cluster;
// it compiles here but is exercised only against one, never by the hermetic test
// suite. The Oracle verifies the layer above it (Editor orchestration and error
// mapping) against a fake Runner, and the Job shape via the pure jobspec.
// The transport is tested against the Kubernetes HTTP API shape; the Editor
// tests separately verify orchestration and error mapping against a fake Runner.
type K8sRunner struct {
cs kubernetes.Interface
cs kubernetes.Interface
Browser *Browser
}
// NewK8sRunner builds a Runner over cs. Every per-operation parameter — the
@@ -60,14 +57,23 @@ func NewK8sRunner(cs kubernetes.Interface) *K8sRunner {
return &K8sRunner{cs: cs}
}
// Run creates the file Job, waits for its Pod to reach a terminal phase, and
// returns the JSON payload from the ResultPrefix line of that Pod's log.
// Run streams read results as soon as they are printed. Mutations still wait
// for termination so returning success does not leave the next write racing
// the cluster-side world lock.
//
// The Job name carries a fresh random OpID (FilesJobName), so a create collision is
// not an expected condition the way it is for restore — an AlreadyExists here means
// a 64-bit collision inside one TTL window and is reported rather than absorbed,
// because absorbing it would mean returning ANOTHER operation's output.
func (k *K8sRunner) Run(ctx context.Context, p JobParams) ([]byte, error) {
if k.Browser != nil && !mutates(p.Op) {
payload, err := k.Browser.Run(ctx, p, k.Start)
// At capacity, retain the bounded one-shot path instead of starting
// more idle workers. All reads still use the same containment checks.
if !errors.Is(err, errBrowserFull) {
return payload, err
}
}
job, err := FilesJob(p)
if err != nil {
return nil, err
@@ -81,24 +87,30 @@ func (k *K8sRunner) Run(ctx context.Context, p JobParams) ([]byte, error) {
return nil, err
}
log, err := k.podLog(ctx, p.Namespace, pod.Name)
stream, err := k.cs.CoreV1().Pods(p.Namespace).GetLogs(pod.Name, &corev1.PodLogOptions{
Container: containerName, Follow: true,
}).Stream(ctx)
if err != nil {
return nil, err
return nil, fmt.Errorf("fileedit: read file job log: %w", err)
}
payload, ok := extractResult(log)
if !ok {
// No marked line: the entrypoint died before printing (an unmountable volume,
// an OOM kill, a deadline). The log tail travels in the error for the operator's
// benefit — this error reaches felis-api's logs, while the caller gets the
// generic 500 writeError produces, so no node detail leaks to the browser.
return nil, fmt.Errorf("fileedit: file job %s produced no result (phase %s): %s",
job.Name, pod.Status.Phase, tail(log))
defer stream.Close()
sc := bufio.NewScanner(io.LimitReader(stream, maxLogBytes))
sc.Buffer(make([]byte, 0, 64*1024), maxLogBytes)
var last string
for sc.Scan() {
if payload, ok := strings.CutPrefix(sc.Text(), ResultPrefix); ok {
return []byte(payload), nil
}
last = sc.Text()
}
return payload, nil
if err := sc.Err(); err != nil {
return nil, fmt.Errorf("fileedit: read file job result: %w", err)
}
return nil, fmt.Errorf("fileedit: file job %s produced no result (phase %s): %s",
job.Name, pod.Status.Phase, tail(last))
}
// awaitPod polls until the operation's Pod reaches a terminal phase. It selects by
// awaitPod polls until the file container has started (or failed). It selects by
// the per-invocation LabelOpID, so it can never observe a different operation's Pod
// — the reason that label exists.
//
@@ -107,7 +119,7 @@ func (k *K8sRunner) Run(ctx context.Context, p JobParams) ([]byte, error) {
// is too large) is printed and then exited on cleanly, and even a genuinely failed
// Pod may have printed a diagnosable result first. Deciding what the outcome MEANS
// is the caller's job (Run reads the printed result); this function only decides
// when there is nothing left to wait for.
// when its log can be opened without a ContainerCreating refusal.
func (k *K8sRunner) awaitPod(ctx context.Context, p JobParams) (*corev1.Pod, error) {
ticker := time.NewTicker(pollInterval)
defer ticker.Stop()
@@ -120,9 +132,17 @@ func (k *K8sRunner) awaitPod(ctx context.Context, p JobParams) (*corev1.Pod, err
return nil, fmt.Errorf("fileedit: find file job pod: %w", err)
}
for i := range pods.Items {
switch pods.Items[i].Status.Phase {
pod := &pods.Items[i]
if !mutates(p.Op) {
for _, c := range pod.Status.ContainerStatuses {
if c.Name == containerName && (c.State.Running != nil || c.State.Terminated != nil) {
return pod, nil
}
}
}
switch pod.Status.Phase {
case corev1.PodSucceeded, corev1.PodFailed:
return &pods.Items[i], nil
return pod, nil
}
}
+67
View File
@@ -3,6 +3,10 @@ package fileedit
import (
"context"
"fmt"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
@@ -10,10 +14,73 @@ import (
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/kubernetes/fake"
"k8s.io/client-go/rest"
k8stesting "k8s.io/client-go/testing"
)
func TestRunReadsResultBeforePodTermination(t *testing.T) {
for _, op := range []string{OpList, OpRead} {
t.Run(op, func(t *testing.T) {
closed := make(chan struct{})
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/jobs"):
w.Header().Set("Content-Type", "application/json")
io.WriteString(w, `{"apiVersion":"batch/v1","kind":"Job","metadata":{"name":"files-test"}}`)
case strings.HasSuffix(r.URL.Path, "/pods"):
if !strings.Contains(r.URL.Query().Get("labelSelector"), LabelOpID+"=") {
t.Error("pod lookup did not select this operation")
}
w.Header().Set("Content-Type", "application/json")
io.WriteString(w, `{"apiVersion":"v1","kind":"PodList","items":[{"metadata":{"name":"file-pod"},"status":{"phase":"Running","containerStatuses":[{"name":"files","state":{"running":{}}}]}}]}`)
case strings.HasSuffix(r.URL.Path, "/log"):
if r.URL.Query().Get("follow") != "true" {
t.Error("result log was not streamed")
}
io.WriteString(w, "diagnostic line\n"+ResultPrefix+`{"entries":[]}`+"\n")
w.(http.Flusher).Flush()
// The log stays open and the Pod stays Running. Run must return
// on the result line and close the stream, not await either EOF.
<-r.Context().Done()
close(closed)
default:
t.Errorf("unexpected cluster call: %s %s", r.Method, r.URL)
w.WriteHeader(http.StatusNotFound)
}
}))
defer srv.Close()
cs, err := kubernetes.NewForConfig(&rest.Config{Host: srv.URL})
if err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
got, err := NewK8sRunner(cs).Run(ctx, testParams(op))
if err != nil || string(got) != `{"entries":[]}` {
t.Fatalf("Run = %s, %v", got, err)
}
select {
case <-closed:
case <-ctx.Done():
t.Fatal("result stream was not closed")
}
})
}
}
func TestAwaitPodKeepsWritesWaitingForTermination(t *testing.T) {
p := testParams(OpWrite)
pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "file-pod", Namespace: p.Namespace, Labels: filesLabels(p)},
Status: corev1.PodStatus{Phase: corev1.PodRunning, ContainerStatuses: []corev1.ContainerStatus{{Name: containerName, State: corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}}}}}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond)
defer cancel()
if got, err := NewK8sRunner(fake.NewSimpleClientset(pod)).awaitPod(ctx, p); err == nil || got != nil {
t.Fatalf("a running writer was treated as finished: %v, %v", got, err)
}
}
var (
opCreated = time.Date(2026, 9, 28, 10, 0, 0, 0, time.UTC)
opEnded = opCreated.Add(3 * time.Minute)