fix(build): 构建并发上限与排队,构建命名空间加 ResourceQuota,SyncAll 逐个容错
This commit is contained in:
13 files changed
+273
-17
No files matched your search
+117
-12
@@ -37,6 +37,7 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
@@ -267,6 +268,11 @@ type Config struct {
|
||||
// when set. The class must exist on the cluster.
|
||||
RuntimeClass string
|
||||
|
||||
// MaxConcurrent caps how many build Jobs run at once. A build submitted past
|
||||
// the cap stays pending, and SyncAll starts queued builds oldest first as
|
||||
// running ones finish. Zero applies the default; MaxConcurrentLimit bounds it.
|
||||
MaxConcurrent int
|
||||
|
||||
// ContextOrigin is the scheme://host[:port] of the platform's internal API
|
||||
// face, the only host an http(s) ContextRef may name: the fetch step presents
|
||||
// the service token to it. Empty refuses every http(s) context.
|
||||
@@ -302,8 +308,14 @@ const (
|
||||
defaultCPULimit = "2"
|
||||
defaultMemLimit = "4Gi"
|
||||
defaultDiskLimit = "12Gi"
|
||||
defaultMaxConcurrent = 2
|
||||
)
|
||||
|
||||
// MaxConcurrentLimit is the highest MaxConcurrent the platform accepts. The
|
||||
// build namespace's pod quota (BuildResourceQuota) leaves room for exactly this
|
||||
// many build pods plus the user-namespace probe.
|
||||
const MaxConcurrentLimit = 6
|
||||
|
||||
// withDefaults returns a copy of c with zero fields filled, so a partially
|
||||
// configured Config (or the zero value, in tests) is always usable.
|
||||
func (c Config) withDefaults() Config {
|
||||
@@ -337,16 +349,25 @@ func (c Config) withDefaults() Config {
|
||||
if c.UserNamespaces == "" {
|
||||
c.UserNamespaces = UserNamespacesAuto
|
||||
}
|
||||
if c.MaxConcurrent <= 0 {
|
||||
c.MaxConcurrent = defaultMaxConcurrent
|
||||
}
|
||||
if c.MaxConcurrent > MaxConcurrentLimit {
|
||||
c.MaxConcurrent = MaxConcurrentLimit
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// Builder orchestrates the build subsystem. It holds no mutable state; the
|
||||
// clock and id generator are injectable for hermetic tests.
|
||||
// Builder orchestrates the build subsystem. The clock and id generator are
|
||||
// injectable for hermetic tests. Its only state is the lock that serializes
|
||||
// starting Jobs, so two submissions cannot both take the last free slot.
|
||||
type Builder struct {
|
||||
Store Store
|
||||
Jobs Jobs
|
||||
Config Config
|
||||
|
||||
startMu sync.Mutex
|
||||
|
||||
// Now is the clock, injectable for tests. Defaults to time.Now.
|
||||
Now func() time.Time
|
||||
// IDGen mints build ids. Defaults to a time-based generator.
|
||||
@@ -368,10 +389,11 @@ func (b *Builder) newID() string {
|
||||
}
|
||||
|
||||
// Submit validates req, records a pending build, and starts the Kaniko+Trivy
|
||||
// Job (spec §16). The build is returned in the building state once the Job is
|
||||
// created; if Job creation fails the build is marked failed so it never lingers
|
||||
// pending. The caller (felis-api) drives the build to a terminal state by
|
||||
// polling Sync / SyncAll.
|
||||
// Job (spec §16) when fewer than MaxConcurrent builds are running. Past the cap
|
||||
// the build is returned pending and waits in the queue SyncAll drains. A started
|
||||
// build is returned building; if Job creation fails it is marked failed so it
|
||||
// never lingers pending. The caller (felis-api) drives the build to a terminal
|
||||
// state by polling Sync / SyncAll.
|
||||
func (b *Builder) Submit(ctx context.Context, req Request) (*Build, error) {
|
||||
cfg := b.Config.withDefaults()
|
||||
if err := Validate(req, cfg); err != nil {
|
||||
@@ -394,6 +416,34 @@ func (b *Builder) Submit(ctx context.Context, req Request) (*Build, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
b.startMu.Lock()
|
||||
defer b.startMu.Unlock()
|
||||
running, err := b.running(ctx)
|
||||
if err != nil || running >= cfg.MaxConcurrent {
|
||||
// Queued. When the count could not be read, SyncAll retries the start.
|
||||
return bld, nil
|
||||
}
|
||||
return b.start(ctx, bld, cfg)
|
||||
}
|
||||
|
||||
// running counts builds whose Job has been started and not yet reconciled to a
|
||||
// terminal state.
|
||||
func (b *Builder) running(ctx context.Context) (int, error) {
|
||||
builds, err := b.Store.ListUnfinishedBuilds(ctx)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
n := 0
|
||||
for i := range builds {
|
||||
if builds[i].Status == StatusBuilding {
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// start creates bld's Job and records it. The caller holds startMu.
|
||||
func (b *Builder) start(ctx context.Context, bld *Build, cfg Config) (*Build, error) {
|
||||
jobName, err := b.Jobs.CreateBuildJob(ctx, b.jobParams(bld, cfg))
|
||||
if err != nil {
|
||||
// The pending row exists; mark it failed so it is not reconciled forever.
|
||||
@@ -466,8 +516,12 @@ func (b *Builder) Sync(ctx context.Context, id string) (*Build, error) {
|
||||
return bld, nil
|
||||
}
|
||||
if bld.JobName == "" {
|
||||
// Created but the Job name was never recorded; treat as failed rather
|
||||
// than reconcile forever against a phantom Job.
|
||||
if bld.Status == StatusPending {
|
||||
// Queued behind MaxConcurrent; SyncAll starts it.
|
||||
return bld, nil
|
||||
}
|
||||
// Building with no Job name recorded; treat as failed rather than
|
||||
// reconcile forever against a phantom Job.
|
||||
return b.finish(ctx, bld, StatusFailed, "no build job recorded")
|
||||
}
|
||||
|
||||
@@ -499,25 +553,76 @@ func (b *Builder) Sync(ctx context.Context, id string) (*Build, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// SyncAll reconciles every unfinished build and returns the count advanced to a
|
||||
// SyncAll reconciles every running build, then starts queued builds oldest
|
||||
// first while fewer than MaxConcurrent run. It returns the count advanced to a
|
||||
// terminal state. felis-api calls this periodically (spec §16: the scan gate is
|
||||
// observed, not pushed by the build Pod).
|
||||
// observed, not pushed by the build Pod). One build's error does not stop the
|
||||
// rest; every error comes back joined.
|
||||
func (b *Builder) SyncAll(ctx context.Context) (int, error) {
|
||||
builds, err := b.Store.ListUnfinishedBuilds(ctx)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
advanced := 0
|
||||
var errs []error
|
||||
for i := range builds {
|
||||
if builds[i].Status == StatusPending && builds[i].JobName == "" {
|
||||
continue // queued; startQueued below
|
||||
}
|
||||
bld, err := b.Sync(ctx, builds[i].ID)
|
||||
if err != nil {
|
||||
return advanced, err
|
||||
errs = append(errs, fmt.Errorf("build %s: %w", builds[i].ID, err))
|
||||
continue
|
||||
}
|
||||
if bld.Status.terminal() {
|
||||
advanced++
|
||||
}
|
||||
}
|
||||
return advanced, nil
|
||||
started, err := b.startQueued(ctx)
|
||||
if err != nil {
|
||||
errs = append(errs, err)
|
||||
}
|
||||
advanced += started
|
||||
return advanced, errors.Join(errs...)
|
||||
}
|
||||
|
||||
// startQueued starts pending builds, oldest first, while fewer than
|
||||
// MaxConcurrent run. It returns how many it failed outright (a Job that could
|
||||
// not be created), which count as advanced to a terminal state.
|
||||
func (b *Builder) startQueued(ctx context.Context) (int, error) {
|
||||
cfg := b.Config.withDefaults()
|
||||
b.startMu.Lock()
|
||||
defer b.startMu.Unlock()
|
||||
builds, err := b.Store.ListUnfinishedBuilds(ctx)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
running := 0
|
||||
for i := range builds {
|
||||
if builds[i].Status == StatusBuilding {
|
||||
running++
|
||||
}
|
||||
}
|
||||
failed := 0
|
||||
var errs []error
|
||||
for i := range builds {
|
||||
if running >= cfg.MaxConcurrent {
|
||||
break
|
||||
}
|
||||
bld := &builds[i]
|
||||
if bld.Status != StatusPending || bld.JobName != "" {
|
||||
continue
|
||||
}
|
||||
if _, err := b.start(ctx, bld, cfg); err != nil {
|
||||
errs = append(errs, fmt.Errorf("build %s: %w", bld.ID, err))
|
||||
if bld.Status == StatusFailed {
|
||||
failed++
|
||||
}
|
||||
continue
|
||||
}
|
||||
running++
|
||||
}
|
||||
return failed, errors.Join(errs...)
|
||||
}
|
||||
|
||||
// Cancel stops an in-flight build: delete its Job and mark it cancelled. A
|
||||
|
||||
@@ -3,6 +3,7 @@ package build
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sort"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -77,6 +78,13 @@ func (f *fakeStore) ListUnfinishedBuilds(_ context.Context) ([]Build, error) {
|
||||
out = append(out, *b)
|
||||
}
|
||||
}
|
||||
// Oldest first, like the SQL; the frozen test clock ties, so the id breaks it.
|
||||
sort.Slice(out, func(i, j int) bool {
|
||||
if !out[i].CreatedAt.Equal(out[j].CreatedAt) {
|
||||
return out[i].CreatedAt.Before(out[j].CreatedAt)
|
||||
}
|
||||
return out[i].ID < out[j].ID
|
||||
})
|
||||
return out, nil
|
||||
}
|
||||
|
||||
@@ -114,6 +122,9 @@ type fakeJobs struct {
|
||||
phase JobPhase
|
||||
phaseErr error
|
||||
createErr error
|
||||
// Per-job overrides of phase / phaseErr, keyed by Job name.
|
||||
phases map[string]JobPhase
|
||||
phaseErrs map[string]error
|
||||
|
||||
created []JobParams
|
||||
cancelled []string
|
||||
@@ -127,7 +138,13 @@ func (f *fakeJobs) CreateBuildJob(_ context.Context, p JobParams) (string, error
|
||||
return BuildJobName(p.BuildID), nil
|
||||
}
|
||||
|
||||
func (f *fakeJobs) JobPhase(_ context.Context, _ string) (JobPhase, error) {
|
||||
func (f *fakeJobs) JobPhase(_ context.Context, name string) (JobPhase, error) {
|
||||
if err, ok := f.phaseErrs[name]; ok {
|
||||
return JobUnknown, err
|
||||
}
|
||||
if p, ok := f.phases[name]; ok {
|
||||
return p, nil
|
||||
}
|
||||
return f.phase, f.phaseErr
|
||||
}
|
||||
|
||||
@@ -445,6 +462,76 @@ func TestSyncAllAdvancesUnfinishedBuilds(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// Builds past MaxConcurrent wait pending and start oldest first as running ones
|
||||
// finish (build-supply-chain-12).
|
||||
func TestSubmitQueuesPastMaxConcurrent(t *testing.T) {
|
||||
b, st, jb := newBuilder()
|
||||
b.Config.MaxConcurrent = 1
|
||||
ctx := context.Background()
|
||||
first, err := b.Submit(ctx, goodRequest())
|
||||
if err != nil || first.Status != StatusBuilding {
|
||||
t.Fatalf("first = %+v, %v; want building", first, err)
|
||||
}
|
||||
var queued []*Build
|
||||
for range 2 {
|
||||
q, err := b.Submit(ctx, goodRequest())
|
||||
if err != nil || q.Status != StatusPending || q.JobName != "" {
|
||||
t.Fatalf("queued = %+v, %v; want pending with no Job", q, err)
|
||||
}
|
||||
queued = append(queued, q)
|
||||
}
|
||||
if len(jb.created) != 1 {
|
||||
t.Fatalf("jobs created = %d, want 1", len(jb.created))
|
||||
}
|
||||
|
||||
// A queued build is not a phantom: Sync leaves it alone.
|
||||
if got, err := b.Sync(ctx, queued[0].ID); err != nil || got.Status != StatusPending {
|
||||
t.Fatalf("Sync(queued) = %+v, %v; want still pending", got, err)
|
||||
}
|
||||
// Nothing frees up while the first runs.
|
||||
if _, err := b.SyncAll(ctx); err != nil || len(jb.created) != 1 {
|
||||
t.Fatalf("SyncAll while full: err %v, jobs %d", err, len(jb.created))
|
||||
}
|
||||
|
||||
jb.phases = map[string]JobPhase{first.JobName: JobSucceeded}
|
||||
n, err := b.SyncAll(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("SyncAll: %v", err)
|
||||
}
|
||||
if n != 1 || len(jb.created) != 2 || jb.created[1].BuildID != queued[0].ID {
|
||||
t.Fatalf("advanced %d, jobs %v; want the first finished and the oldest queued started", n, jb.created)
|
||||
}
|
||||
if st.builds[queued[0].ID].Status != StatusBuilding || st.builds[queued[1].ID].Status != StatusPending {
|
||||
t.Fatalf("statuses = %s, %s; want building, pending",
|
||||
st.builds[queued[0].ID].Status, st.builds[queued[1].ID].Status)
|
||||
}
|
||||
}
|
||||
|
||||
// One build whose Job cannot be read does not stall the others.
|
||||
func TestSyncAllContinuesPastOneBuildsError(t *testing.T) {
|
||||
b, st, jb := newBuilder()
|
||||
b.Config.MaxConcurrent = 3
|
||||
ctx := context.Background()
|
||||
var ids []string
|
||||
for range 3 {
|
||||
bld, err := b.Submit(ctx, goodRequest())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ids = append(ids, bld.ID)
|
||||
}
|
||||
jb.phase = JobSucceeded
|
||||
jb.phaseErrs = map[string]error{BuildJobName(ids[0]): errors.New("apiserver hiccup")}
|
||||
n, err := b.SyncAll(ctx)
|
||||
if err == nil || !strings.Contains(err.Error(), ids[0]) {
|
||||
t.Fatalf("err = %v, want it to name %s", err, ids[0])
|
||||
}
|
||||
if n != 2 || st.builds[ids[1]].Status != StatusSucceeded || st.builds[ids[2]].Status != StatusSucceeded {
|
||||
t.Fatalf("advanced %d; statuses %s %s; want the other two succeeded", n,
|
||||
st.builds[ids[1]].Status, st.builds[ids[2]].Status)
|
||||
}
|
||||
}
|
||||
|
||||
// TestImageBuildFailuresMetricCountsFailedBuilds asserts felis_image_build_failures_total
|
||||
// (spec §23) advances exactly once per failed build from BOTH terminal-failure
|
||||
// producers, and never on a successful build. There are two distinct Inc sites —
|
||||
|
||||
@@ -656,6 +656,27 @@ func BuildLimitRange(namespace string) *corev1.LimitRange {
|
||||
}
|
||||
}
|
||||
|
||||
// BuildResourceQuota caps how many pods can run in the build namespace at once:
|
||||
// MaxConcurrentLimit builds, the user-namespace probe, and one to spare. The
|
||||
// Builder's queue keeps builds under it; the quota holds even when something
|
||||
// else creates pods there. Finished pods do not count.
|
||||
func BuildResourceQuota(namespace string) *corev1.ResourceQuota {
|
||||
return &corev1.ResourceQuota{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: "felis-build-quota",
|
||||
Namespace: namespace,
|
||||
Labels: map[string]string{
|
||||
LabelManagedBy: managedByValue,
|
||||
LabelComponent: componentValue,
|
||||
},
|
||||
},
|
||||
Spec: corev1.ResourceQuotaSpec{Hard: corev1.ResourceList{
|
||||
corev1.ResourcePods: *resource.NewQuantity(MaxConcurrentLimit+2, resource.DecimalSI),
|
||||
corev1.ResourcePersistentVolumeClaims: *resource.NewQuantity(0, resource.DecimalSI),
|
||||
}},
|
||||
}
|
||||
}
|
||||
|
||||
// resourceLimits parses the CPU/memory limits into a ResourceList.
|
||||
func resourceLimits(cpu, mem string) (corev1.ResourceList, error) {
|
||||
if cpu == "" {
|
||||
|
||||
@@ -641,6 +641,19 @@ func TestBuildLimitRangeCoversEphemeralStorage(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// The pod quota must fit the most builds the config allows plus the probe pod,
|
||||
// or the queue admits a build whose pod the API server then refuses.
|
||||
func TestBuildResourceQuotaFitsMaxConcurrent(t *testing.T) {
|
||||
rq := BuildResourceQuota("felis-build")
|
||||
pods := rq.Spec.Hard[corev1.ResourcePods]
|
||||
if rq.Namespace != "felis-build" || pods.Value() < MaxConcurrentLimit+1 {
|
||||
t.Fatalf("quota = %#v; want pods ≥ %d", rq.Spec.Hard, MaxConcurrentLimit+1)
|
||||
}
|
||||
if c := (Config{MaxConcurrent: 100}).withDefaults().MaxConcurrent; c != MaxConcurrentLimit {
|
||||
t.Fatalf("MaxConcurrent 100 resolved to %d, want the %d cap", c, MaxConcurrentLimit)
|
||||
}
|
||||
}
|
||||
|
||||
func initNames(cs []corev1.Container) []string {
|
||||
names := make([]string, 0, len(cs))
|
||||
for _, c := range cs {
|
||||
|
||||
@@ -36,6 +36,11 @@ func (k *K8sJobs) CreateBuildJob(ctx context.Context, p JobParams) (string, erro
|
||||
return "", err
|
||||
}
|
||||
if err := k.c.Create(ctx, job); err != nil {
|
||||
// The name is the build's; an existing Job is this build's own, created
|
||||
// by a start that stopped before recording it.
|
||||
if apierrors.IsAlreadyExists(err) {
|
||||
return job.Name, nil
|
||||
}
|
||||
return "", err
|
||||
}
|
||||
return job.Name, nil
|
||||
|
||||
@@ -171,6 +171,9 @@ type RegistryConfig struct {
|
||||
// BuildRuntimeClass runs build pods under a sandbox RuntimeClass such as
|
||||
// gVisor or Kata. Empty runs them under the node's default runtime.
|
||||
BuildRuntimeClass string `toml:"build_runtime_class"`
|
||||
// MaxConcurrentBuilds caps how many builds run at once; later ones queue.
|
||||
// Zero keeps 2; at most 6 (the build namespace's pod quota).
|
||||
MaxConcurrentBuilds int `toml:"max_concurrent_builds"`
|
||||
// TrivyDBRepository points Trivy at an OCI repository holding the
|
||||
// vulnerability DB (--db-repository). Trivy's default fetches from
|
||||
// mirror.gcr.io/ghcr.io, which the build egress lock denies — so on a
|
||||
@@ -463,6 +466,9 @@ func (c *Config) Validate() error {
|
||||
default:
|
||||
return fmt.Errorf("config: [registry] build_user_namespaces %q must be auto, on or off", c.Registry.BuildUserNamespaces)
|
||||
}
|
||||
if n := c.Registry.MaxConcurrentBuilds; n < 0 || n > 6 {
|
||||
return fmt.Errorf("config: [registry] max_concurrent_builds %d must be 1-6 (0 keeps 2)", n)
|
||||
}
|
||||
// [smtp] is optional as a whole, but once a host is named the block must be
|
||||
// deliverable: a From address (relays reject MAIL FROM:<>) and a sane port.
|
||||
// Fail at load, not at the first OTP a player is waiting on.
|
||||
|
||||
@@ -27,7 +27,7 @@ type Object interface {
|
||||
// namespaced Roles + RoleBindings — felis-api and felis-operator always, plus the
|
||||
// destructive felis-reaper identity only when retention is enabled, gated with its
|
||||
// CronJob), the two weak Job SAs (build/restore, which have NO Role anywhere), the
|
||||
// build-namespace egress NetworkPolicy and LimitRange, and the minecraft-namespace
|
||||
// build-namespace egress NetworkPolicy, LimitRange and ResourceQuota, and the minecraft-namespace
|
||||
// ingress NetworkPolicies.
|
||||
//
|
||||
// Scope: this is the authorization + network fence (spec §21, §22) plus the
|
||||
@@ -96,6 +96,9 @@ func Objects(p Params) []Object {
|
||||
buildLR := build.BuildLimitRange(p.BuildNamespace)
|
||||
buildLR.TypeMeta = metav1.TypeMeta{APIVersion: "v1", Kind: "LimitRange"}
|
||||
objs = append(objs, buildLR)
|
||||
buildRQ := build.BuildResourceQuota(p.BuildNamespace)
|
||||
buildRQ.TypeMeta = metav1.TypeMeta{APIVersion: "v1", Kind: "ResourceQuota"}
|
||||
objs = append(objs, buildRQ)
|
||||
|
||||
// Minecraft-namespace ingress fence (default-deny + RCON + game), the server
|
||||
// egress fence, and the registry's ingress fence.
|
||||
|
||||
Reference in new issue
Block a user