diff --git a/cmd/felis/api.go b/cmd/felis/api.go index 2d8d319..ce54cf9 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -465,6 +465,7 @@ func buildConfig(cfg *config.Config) build.Config { UserNamespaces: cfg.Registry.BuildUserNamespaces, UserNamespacesProbe: new(atomic.Bool), RuntimeClass: cfg.Registry.BuildRuntimeClass, + MaxConcurrent: cfg.Registry.MaxConcurrentBuilds, // Empty keeps Trivy's own default; an install with builds points this at // the internal DB mirror (see config.RegistryConfig.TrivyDBRepository). TrivyDBRepository: cfg.Registry.TrivyDBRepository, diff --git a/deploy/bootstrap.sh b/deploy/bootstrap.sh index 181d5e6..4195c7a 100644 --- a/deploy/bootstrap.sh +++ b/deploy/bootstrap.sh @@ -2352,7 +2352,7 @@ persisted_registry_block() { out="$(awk ' /^[[:space:]]*\[/ { sect = $0; next } sect ~ /^[[:space:]]*\[registry\][[:space:]]*$/ && - /^[[:space:]]*(kaniko_image|trivy_image|trivy_db_repository|trivy_java_db_repository|build_cpu_limit|build_mem_limit|build_disk_limit|build_user_namespaces|build_runtime_class|user_uploads_context)[[:space:]]*=/ { print } + /^[[:space:]]*(kaniko_image|trivy_image|trivy_db_repository|trivy_java_db_repository|build_cpu_limit|build_mem_limit|build_disk_limit|build_user_namespaces|build_runtime_class|max_concurrent_builds|user_uploads_context)[[:space:]]*=/ { print } sect ~ /^[[:space:]]*\[registry\.s3\][[:space:]]*$/ && /^[[:space:]]*[A-Za-z_]+[[:space:]]*=/ { if (!s3hdr) { printf "[registry.s3]\n"; s3hdr = 1 } print diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index e299046..e42d3d2 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -493,6 +493,7 @@ build_mem_limit = "4Gi" build_disk_limit = "12Gi" # §8f build_user_namespaces = "auto" # §8f: auto | on | off build_runtime_class = "" # §8f: e.g. "gvisor" +max_concurrent_builds = 2 # §8f: 1-6; later builds queue ``` Mirror the executor images into the registry once. On the node itself, push @@ -582,7 +583,8 @@ approved-but-hostile Dockerfile and the node is the pod around it: | Sandbox runtime | optional `build_runtime_class` (gVisor, Kata) | below | | Credentials | the registry credential lives only in the `push` container; the service token only in `context-fetch` | jobspec | | Resources | CPU, memory and ephemeral-storage limits per container; `activeDeadlineSeconds`; the context extraction stops at 4 GiB or 200 000 entries | jobspec, `felis fetch-context` | -| Namespace backstop | `felis-build-limits` LimitRange gives any container without limits 1 CPU / 1 GiB / 1 GiB disk | bundle | +| Namespace backstop | `felis-build-limits` LimitRange gives any container without limits 1 CPU / 1 GiB / 1 GiB disk; `felis-build-quota` allows 8 running pods and no PVCs | bundle | +| Concurrency | at most `[registry] max_concurrent_builds` (default 2, at most 6) builds run; later ones wait as `pending` (Queued) and start oldest first | `build.Builder` | **Reviewed bytes.** Every upload records the sha256 of the archive, and the review page shows it. The context download carries the same value in the diff --git a/internal/build/build.go b/internal/build/build.go index 08eaefe..0e52c80 100644 --- a/internal/build/build.go +++ b/internal/build/build.go @@ -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 diff --git a/internal/build/build_test.go b/internal/build/build_test.go index 4e8aa9b..bc75b8c 100644 --- a/internal/build/build_test.go +++ b/internal/build/build_test.go @@ -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 — diff --git a/internal/build/jobspec.go b/internal/build/jobspec.go index 7f1a0df..93154ce 100644 --- a/internal/build/jobspec.go +++ b/internal/build/jobspec.go @@ -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 == "" { diff --git a/internal/build/jobspec_test.go b/internal/build/jobspec_test.go index 0b5dfc8..5e0146c 100644 --- a/internal/build/jobspec_test.go +++ b/internal/build/jobspec_test.go @@ -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 { diff --git a/internal/build/k8sjobs.go b/internal/build/k8sjobs.go index bf69052..f88f792 100644 --- a/internal/build/k8sjobs.go +++ b/internal/build/k8sjobs.go @@ -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 diff --git a/internal/config/config.go b/internal/config/config.go index 838e7c4..5a2c3b5 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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. diff --git a/internal/platform/bundle.go b/internal/platform/bundle.go index f567dbc..b205409 100644 --- a/internal/platform/bundle.go +++ b/internal/platform/bundle.go @@ -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. diff --git a/panel/src/i18n/resources/en-US/admin.json b/panel/src/i18n/resources/en-US/admin.json index fa1b49b..f3e9338 100644 --- a/panel/src/i18n/resources/en-US/admin.json +++ b/panel/src/i18n/resources/en-US/admin.json @@ -31,6 +31,12 @@ "approve_needs_context": "Nothing to approve: the context has not been uploaded, or must be uploaded again", "base_image_label": "Base Image", "base_image_placeholder": "e.g. library/postgres:15", + "build_status_pending": "Queued", + "build_status_building": "Building", + "build_status_succeeded": "Succeeded", + "build_status_failed": "Failed", + "build_status_cancelled": "Cancelled", + "build_queued_hint": "The concurrent build limit ([registry] max_concurrent_builds) is reached; this build starts on its own when an earlier one finishes.", "view_logs_btn": "Logs", "cancel_build_btn": "Cancel", "no_builds_title": "No Builds Found", diff --git a/panel/src/i18n/resources/zh-CN/admin.json b/panel/src/i18n/resources/zh-CN/admin.json index ce2d160..77dbeff 100644 --- a/panel/src/i18n/resources/zh-CN/admin.json +++ b/panel/src/i18n/resources/zh-CN/admin.json @@ -31,6 +31,12 @@ "approve_needs_context": "没有可通过的上下文:提交者尚未上传或需要重新上传", "base_image_label": "基础镜像", "base_image_placeholder": "例如: library/postgres:15", + "build_status_pending": "排队中", + "build_status_building": "构建中", + "build_status_succeeded": "成功", + "build_status_failed": "失败", + "build_status_cancelled": "已取消", + "build_queued_hint": "同时运行的构建已达上限([registry] max_concurrent_builds),这个构建会在前面的构建结束后自动开始。", "view_logs_btn": "日志", "cancel_build_btn": "取消", "no_builds_title": "暂无构建任务", diff --git a/panel/src/pages/admin/ImageBuildPage.tsx b/panel/src/pages/admin/ImageBuildPage.tsx index caf5de7..e3c66b3 100644 --- a/panel/src/pages/admin/ImageBuildPage.tsx +++ b/panel/src/pages/admin/ImageBuildPage.tsx @@ -491,12 +491,13 @@ export function ImageBuildPage() {
- {b.status} + {t(`build_status_${b.status}`)}