Files
Felis/internal/fileedit/k8sjobs.go
T

345 lines
12 KiB
Go

package fileedit
import (
"bufio"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"sort"
"strings"
"time"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
)
// 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
// 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
// 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
// 4 MiB is generous headroom while still refusing to let a Pod that floods stderr
// pull felis-api's memory down with it.
const maxLogBytes = 4 << 20
// K8sRunner is the production Runner (spec §7, §16). It drives one file operation
// end to end using ONLY the three permissions felis-api already holds in the
// minecraft namespace, which is the entire point of the design:
//
// jobs:create → create the file Job
// pods:list → find its Pod and observe the phase (no pods:get, no pods:watch)
// pods/log:get → read the printed result back
//
// It deliberately takes a typed kubernetes.Interface rather than the
// controller-runtime client that internal/restore's K8sJobs uses: the log
// subresource (GetLogs(...).Stream) exists only on the typed CoreV1 client, and
// 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.
type K8sRunner struct {
cs kubernetes.Interface
}
// NewK8sRunner builds a Runner over cs. Every per-operation parameter — the
// namespace included — travels in the JobParams the Editor renders, so there is no
// Config to retain here.
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.
//
// 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) {
job, err := FilesJob(p)
if err != nil {
return nil, err
}
if _, err := k.cs.BatchV1().Jobs(p.Namespace).Create(ctx, job, metav1.CreateOptions{}); err != nil {
return nil, fmt.Errorf("fileedit: create file job: %w", err)
}
pod, err := k.awaitPod(ctx, p)
if err != nil {
return nil, err
}
log, err := k.podLog(ctx, p.Namespace, pod.Name)
if err != nil {
return nil, 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))
}
return payload, nil
}
// awaitPod polls until the operation's Pod reaches a terminal phase. It selects by
// the per-invocation LabelOpID, so it can never observe a different operation's Pod
// — the reason that label exists.
//
// Both Succeeded and Failed are terminal and BOTH return the Pod rather than an
// error, because a caller-fault result (a path that escapes the root, a file that
// 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.
func (k *K8sRunner) awaitPod(ctx context.Context, p JobParams) (*corev1.Pod, error) {
ticker := time.NewTicker(pollInterval)
defer ticker.Stop()
for {
pods, err := k.cs.CoreV1().Pods(p.Namespace).List(ctx, metav1.ListOptions{
LabelSelector: LabelOpID + "=" + p.OpID,
})
if err != nil {
return nil, fmt.Errorf("fileedit: find file job pod: %w", err)
}
for i := range pods.Items {
switch pods.Items[i].Status.Phase {
case corev1.PodSucceeded, corev1.PodFailed:
return &pods.Items[i], nil
}
}
select {
case <-ticker.C:
case <-ctx.Done():
// The Editor's Timeout (or the client disconnecting) fired. The Job is left
// alone deliberately: felis-api holds no jobs:delete, and the Job's own
// activeDeadlineSeconds plus ttlSecondsAfterFinished retire it without help.
return nil, fmt.Errorf("fileedit: timed out waiting for the file job to finish: %w", ctx.Err())
}
}
}
// podLog reads a finished Pod's log. Follow is off — the Pod has already
// terminated, so the log is complete and a follow would merely block until the
// stream closed.
func (k *K8sRunner) podLog(ctx context.Context, namespace, pod string) (string, error) {
stream, err := k.cs.CoreV1().Pods(namespace).GetLogs(pod, &corev1.PodLogOptions{
Container: containerName,
}).Stream(ctx)
if err != nil {
return "", fmt.Errorf("fileedit: read file job log: %w", err)
}
defer stream.Close()
b, err := io.ReadAll(io.LimitReader(stream, maxLogBytes))
if err != nil && !errors.Is(err, io.EOF) {
return "", fmt.Errorf("fileedit: read file job log: %w", err)
}
return string(b), nil
}
// extractResult finds the marked payload in a Pod log. It scans for the LAST line
// carrying ResultPrefix because pods/log returns stdout and stderr MERGED: a Go
// runtime warning or a libc message can appear anywhere in the stream, so the
// payload must be located by its marker rather than by position. Taking the last
// match rather than the first is the conservative choice — if a marker somehow
// appeared more than once, the final one is the operation's actual outcome.
func extractResult(log string) ([]byte, bool) {
var payload string
var found bool
sc := bufio.NewScanner(strings.NewReader(log))
sc.Buffer(make([]byte, 0, 64*1024), maxLogBytes)
for sc.Scan() {
if rest, ok := strings.CutPrefix(sc.Text(), ResultPrefix); ok {
payload, found = rest, true
}
}
if !found {
return nil, false
}
return []byte(payload), true
}
// tail returns the last few hundred bytes of a log for an error message, so a
// diagnostic is useful without embedding an entire log in an error string.
func tail(log string) string {
const n = 512
log = strings.TrimSpace(log)
if len(log) <= n {
return log
}
return "..." + log[len(log)-n:]
}
// Start creates the Job for p and returns without waiting on it. Ops reads it
// back.
func (k *K8sRunner) Start(ctx context.Context, p JobParams) error {
job, err := FilesJob(p)
if err != nil {
return err
}
if _, err := k.cs.BatchV1().Jobs(p.Namespace).Create(ctx, job, metav1.CreateOptions{}); err != nil {
return fmt.Errorf("fileedit: create file job: %w", err)
}
return nil
}
// maxOps bounds what Ops reports, and so how many Pod logs one call reads. Only
// one background op runs per server at a time (it holds the world volume), so
// this many are the one running and the latest that finished within AsyncTTL.
const maxOps = 10
// opLogLines is how much of an op's log Ops reads: the result line is the last
// thing the Job prints to stdout, the progress lines come before it, and a
// failed run ends with one line on stderr.
const opLogLines = 20
// Ops lists the server's background Jobs, newest first, with how far each has
// got and how it ended, read from the tail of its Pod's log. A Pod whose log
// cannot be read yet (still pulling its image) or any more (its node went away)
// reports no progress rather than failing the listing.
func (k *K8sRunner) Ops(ctx context.Context, namespace, server string) ([]OpState, error) {
sel := metav1.ListOptions{LabelSelector: LabelManagedBy + "=" + managedByValue + "," +
LabelServer + "=" + server + "," + LabelAsync + "=true"}
jobs, err := k.cs.BatchV1().Jobs(namespace).List(ctx, sel)
if err != nil {
return nil, fmt.Errorf("fileedit: list file jobs: %w", err)
}
pods, err := k.cs.CoreV1().Pods(namespace).List(ctx, sel)
if err != nil {
return nil, fmt.Errorf("fileedit: list file job pods: %w", err)
}
// backoffLimit is 0, so a Job has one Pod; the newest wins all the same.
podOf := map[string]*corev1.Pod{}
for i := range pods.Items {
pod := &pods.Items[i]
id := pod.Labels[LabelOpID]
if cur := podOf[id]; cur == nil || pod.CreationTimestamp.After(cur.CreationTimestamp.Time) {
podOf[id] = pod
}
}
items := jobs.Items
sort.SliceStable(items, func(i, j int) bool {
return items[i].CreationTimestamp.After(items[j].CreationTimestamp.Time)
})
if len(items) > maxOps {
items = items[:maxOps]
}
out := make([]OpState, 0, len(items))
for i := range items {
log := ""
pod := podOf[items[i].Labels[LabelOpID]]
if pod != nil && pod.Status.Phase != corev1.PodPending {
log, _ = k.logTail(ctx, namespace, pod.Name)
}
out = append(out, opState(&items[i], pod, log))
}
return out, nil
}
// logTail reads the last opLogLines lines of a Pod's log.
func (k *K8sRunner) logTail(ctx context.Context, namespace, pod string) (string, error) {
lines := int64(opLogLines)
stream, err := k.cs.CoreV1().Pods(namespace).GetLogs(pod, &corev1.PodLogOptions{
Container: containerName, TailLines: &lines,
}).Stream(ctx)
if err != nil {
return "", err
}
defer stream.Close()
b, err := io.ReadAll(io.LimitReader(stream, maxLogBytes))
return string(b), err
}
// opState reads one background Job, and the tail of its Pod's log, as an
// OpState. The printed result decides the outcome whatever the Job's condition
// says: a Job killed at its deadline just after printing did finish its work.
// A Job that ended without one failed, and the condition's reason says how
// (DeadlineExceeded, BackoffLimitExceeded), unless its Pod says the kernel
// killed it for memory (ReasonOOMKilled), which the condition does not tell; a
// Job that completed but whose log could not be read has an outcome no one can
// tell, ResultUnavailable.
func opState(job *batchv1.Job, pod *corev1.Pod, log string) OpState {
st := OpState{
ID: job.Labels[LabelOpID], Op: job.Labels[LabelMode], Path: job.Annotations[AnnotationPath],
State: OpRunning, Started: job.CreationTimestamp.Time,
}
if p, ok := lastProgress(log); ok {
st.Done, st.Total = p.Done, p.Total
}
ended, reason := false, ""
for _, c := range job.Status.Conditions {
if c.Status != corev1.ConditionTrue {
continue
}
switch c.Type {
case batchv1.JobComplete, batchv1.JobSuccessCriteriaMet:
ended, reason = true, "ResultUnavailable"
st.Finished = c.LastTransitionTime.Time
case batchv1.JobFailed, batchv1.JobFailureTarget:
ended, reason = true, c.Reason
st.Finished = c.LastTransitionTime.Time
}
}
if !ended {
return st
}
if payload, ok := extractResult(log); ok {
var res Result
if json.Unmarshal(payload, &res) == nil {
st.Result = &res
}
}
switch {
case st.Result != nil && st.Result.Code == "":
st.State = OpSucceeded
case st.Result != nil:
st.State = OpFailed
default:
st.State, st.Reason = OpFailed, reason
if killedForMemory(pod) {
st.Reason = ReasonOOMKilled
} else if st.Reason == "" {
st.Reason = "Failed"
}
}
return st
}
// ReasonOOMKilled is an op's Reason when the kernel killed its Job for going over
// the Job's memory limit: an archive of more entries than the Job can hold the
// list of.
const ReasonOOMKilled = "OOMKilled"
// killedForMemory reports whether pod's container was killed for going over its
// memory limit.
func killedForMemory(pod *corev1.Pod) bool {
if pod == nil {
return false
}
for _, cs := range pod.Status.ContainerStatuses {
if t := cs.State.Terminated; cs.Name == containerName && t != nil && t.Reason == ReasonOOMKilled {
return true
}
}
return false
}