From a76161862f8943489644a203742e438f6c461fc2 Mon Sep 17 00:00:00 2001 From: Edgar Sipki Date: Thu, 1 Oct 2026 17:19:59 +0300 Subject: [PATCH 1/2] pool, license: close the shutdown race, say when the pool ignores a tier change MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WorkerPool.Get checked an atomic flag and then sent on the job channel, while Shutdown set the flag and then closed it; a Get between the two sent on a closed channel and panicked, and a second Shutdown closed it twice. Normal shutdown never got there only because GracefulStop drains handlers before the deferred Shutdown — an ordering the pool did not own, and one serve.GRPC skipped when Serve itself failed. The check and the send now share a read lock, Shutdown is idempotent, and a failed Serve drains its handlers too. The pool reads its licence ceilings once at startup. A licence lapsing past grace moved audit and the plugin cap at once and the pool only on the next restart, silently. The licence manager now logs a Warn when that gap opens, counts MaxGenerations as a change, and the README and AGENTS.md say so. Co-Authored-By: Claude Opus 5.5 (1M context) --- AGENTS.md | 3 ++ README.md | 8 ++++ cmd/easyp-svc/start.go | 5 +++ internal/core/pool.go | 64 +++++++++++++++++++++------- internal/core/pool_test.go | 71 ++++++++++++++++++++++++++++++++ internal/license/manager.go | 43 +++++++++++++++++++ internal/license/manager_test.go | 54 ++++++++++++++++++++++++ internal/serve/grpc.go | 5 +++ 8 files changed, 238 insertions(+), 15 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 8c7375a..bf5dcbf 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -59,6 +59,9 @@ The rest of the difference is ceilings, and Enterprise simply has none: A ceiling caps the resolved configuration at startup rather than rejecting the request, so an over-ambitious Community config starts and logs what it got. +The pool ceilings are only read then: a licence lapsing mid-run moves audit and +the plugin cap at once, the pool on the next restart, and `license.Manager` logs +a Warn when that gap opens. A token names a tier and nothing else. Which features that tier unlocks is decided in `core.EnterpriseLicenseClaims`, in the release, so extending the diff --git a/README.md b/README.md index 7f672e3..71c12ac 100644 --- a/README.md +++ b/README.md @@ -500,6 +500,14 @@ configuration above it is lowered at startup, and the service logs which setting was lowered and to what — it is a ceiling, not a substitution, so asking for less than the tier permits gives you less. +The two worker-pool ceilings are read **once, at startup**. A licence that lapses +past its grace period while the service runs turns audit off and the plugin cap +on at the next licence refresh, but the pool keeps the capacity it started with +until the process restarts — and the restart then lowers it. The service logs +`licence tier changed while running; the worker pool keeps the ceilings it +started with and applies the new ones on restart` at the moment the tier changes, +so the drop does not arrive unannounced with an unrelated deploy. + Enterprise needs two things — a token and the public key it is verified against: ```bash diff --git a/cmd/easyp-svc/start.go b/cmd/easyp-svc/start.go index f213557..111ba63 100644 --- a/cmd/easyp-svc/start.go +++ b/cmd/easyp-svc/start.go @@ -581,6 +581,11 @@ func checkServiceTier(configured string, actualTier func() string, log *slog.Log // // setting is the dotted name of the field being capped, so the log line names // what an operator would have to change rather than what the code calls it. +// +// The cap is taken once, here: the pool is sized from it and never resized. A +// tier change while running — a licence lapsing past its grace period — moves +// audit and the plugin cap at once but reaches the pool only on restart, which +// license.Manager announces at Warn when it happens. func cappedByLicence(setting string, configured, licenseLimit int, log *slog.Logger) int { if licenseLimit <= 0 || licenseLimit >= configured { return configured diff --git a/internal/core/pool.go b/internal/core/pool.go index c7f58dd..696f1d5 100644 --- a/internal/core/pool.go +++ b/internal/core/pool.go @@ -65,7 +65,6 @@ type WorkerPool struct { logger *slog.Logger metrics Metrics wg sync.WaitGroup - closed atomic.Bool activeWorkers prometheus.Gauge rejectedTotal prometheus.Counter jobsTotal prometheus.Counter @@ -76,6 +75,14 @@ type WorkerPool struct { // этого не годятся: воркер освобождается, отдав plugin в канал, а Generate // вызывается уже горутиной вызывающего. gen *genLimiter + + // mu and closed make "not shut down yet" and "send the job" one step. + // An atomic flag checked before the send left a gap in which Shutdown + // could close the channel, and the send then panicked. Get holds the read + // lock across a send that never blocks, so Shutdown's write lock waits at + // most for a channel operation. + mu sync.RWMutex + closed bool } // genLimiter пропускает не более cap одновременных генераций, держит очередь @@ -325,10 +332,6 @@ func (p *WorkerPool) Get(ctx context.Context, pluginGroup, pluginName, pluginVer )) defer span.End() - if p.closed.Load() { - return nil, ErrShuttingDown - } - resultCh := make(chan jobResult, 1) jobItem := job{ ctx: ctx, @@ -338,16 +341,13 @@ func (p *WorkerPool) Get(ctx context.Context, pluginGroup, pluginName, pluginVer result: resultCh, } - select { - case p.jobs <- jobItem: - // Job принят в очередь - p.jobsTotal.Inc() - default: - p.logger.Warn("job queue full, rejecting request") - p.rejectedTotal.Inc() - span.AddEvent("pool.rejected") + err := p.enqueue(jobItem) + if err != nil { + if errors.Is(err, ErrServerOverloaded) { + span.AddEvent("pool.rejected") + } - return nil, ErrServerOverloaded + return nil, err } queueStart := time.Now() @@ -482,11 +482,45 @@ func isTransient(err error) bool { return false } +// enqueue hands a job to the workers without waiting for room, refusing it +// once Shutdown has started. +func (p *WorkerPool) enqueue(work job) error { //nolint:funcorder // helper for Get above + p.mu.RLock() + defer p.mu.RUnlock() + + if p.closed { + return ErrShuttingDown + } + + select { + case p.jobs <- work: + p.jobsTotal.Inc() + + return nil + default: + p.logger.Warn("job queue full, rejecting request") + p.rejectedTotal.Inc() + + return ErrServerOverloaded + } +} + // Shutdown закрывает канал заданий и ожидает завершения воркеров. // Возвращает количество потерянных заданий. +// +// A second call returns 0 and does nothing: closing the channel twice would +// panic, and Shutdown is reached from a defer. func (p *WorkerPool) Shutdown(timeout time.Duration) int { - p.closed.Store(true) + p.mu.Lock() + if p.closed { + p.mu.Unlock() + + return 0 + } + + p.closed = true close(p.jobs) + p.mu.Unlock() done := make(chan struct{}) go func() { diff --git a/internal/core/pool_test.go b/internal/core/pool_test.go index beee6df..ddd7487 100644 --- a/internal/core/pool_test.go +++ b/internal/core/pool_test.go @@ -272,3 +272,74 @@ func TestMaxRetriesIsHonouredIncludingZero(t *testing.T) { }) } } + +// TestGetRacingShutdownNeverPanics pins the shutdown handshake between Get and +// Shutdown. Get used to check an atomic flag and then send on the job channel, +// and Shutdown to set the flag and then close the channel; a Get that passed the +// check before the flag was set sent on a closed channel and panicked. Nothing +// in the pool prevented it — only the order of the deferred calls in +// cmd/easyp-svc, which a failing listener skips. +// +// Each round races a burst of callers against one Shutdown. Panics are +// recovered and counted so the failure is an assertion, not a crashed binary. +func TestGetRacingShutdownNeverPanics(t *testing.T) { + t.Parallel() + + const ( + rounds = 200 + callers = 16 + ) + + var panics atomic.Int64 + + for range rounds { + pool := newTestPool(t, &countingPlugin{}, WorkerPoolConfig{ + Workers: 2, + QueueSize: callers, + }) + + start := make(chan struct{}) + + var wg sync.WaitGroup + + for range callers { + wg.Go(func() { + defer func() { + if recover() != nil { + panics.Add(1) + } + }() + + <-start + + for range 4 { + _, err := pool.Get(t.Context(), "test", "plugin", "v1.0.0") + if err != nil && !errors.Is(err, ErrShuttingDown) && !errors.Is(err, ErrServerOverloaded) { + t.Errorf("unexpected error from Get during shutdown: %v", err) + } + } + }) + } + + close(start) + pool.Shutdown(5 * time.Second) + wg.Wait() + } + + require.Zero(t, panics.Load(), "Get sent on the job channel after Shutdown closed it") +} + +// TestShutdownIsIdempotent covers a second Shutdown, which used to close the +// job channel twice and panic. Shutdown sits in a defer; a caller reaching it +// along two paths must not turn a clean stop into a crash. +func TestShutdownIsIdempotent(t *testing.T) { + t.Parallel() + + pool := newTestPool(t, &countingPlugin{}, WorkerPoolConfig{Workers: 1}) + + require.Zero(t, pool.Shutdown(time.Second)) + require.NotPanics(t, func() { require.Zero(t, pool.Shutdown(time.Second)) }) + + _, err := pool.Get(t.Context(), "test", "plugin", "v1.0.0") + require.ErrorIs(t, err, ErrShuttingDown) +} diff --git a/internal/license/manager.go b/internal/license/manager.go index bffd187..5443b2c 100644 --- a/internal/license/manager.go +++ b/internal/license/manager.go @@ -30,6 +30,12 @@ type Manager struct { logger *slog.Logger metrics *Metrics guard *safe.Guard + + // startup is what the first fetch returned, nil until it has. The worker + // pool reads its ceilings once, right after that fetch, and never again; a + // later tier change moves every other limit at once but leaves the pool + // where it started until the process restarts. Kept so refresh can say so. + startup *core.LicenseClaims } // NewManager creates a Manager backed by the given LicenseClient. @@ -64,6 +70,9 @@ func NewManager( // Perform initial fetch so claims are populated before the first request. lm.refresh(ctx) + startup := lm.Claims() + lm.startup = &startup + return lm, nil } @@ -119,6 +128,7 @@ func (lm *Manager) refresh(ctx context.Context) { changed := lm.claims.Tier != claims.Tier || lm.claims.MaxWorkers != claims.MaxWorkers || lm.claims.MaxPlugins != claims.MaxPlugins || + lm.claims.MaxGenerations != claims.MaxGenerations || len(lm.claims.Features) != len(claims.Features) lm.claims = claims lm.mu.Unlock() @@ -134,8 +144,41 @@ func (lm *Manager) refresh(ctx context.Context) { "tier", claims.Tier, "max_workers", claims.MaxWorkers, "max_plugins", claims.MaxPlugins, + "max_generations", claims.MaxGenerations, "features_count", len(claims.Features), ) + if changed { + lm.warnPoolKeepsStartupCeilings(ctx, claims) + } + lm.metrics.observe(claims) } + +// warnPoolKeepsStartupCeilings reports a tier change the worker pool will not +// follow until restart. +// +// The case that matters is a licence lapsing past its grace period: audit and +// the plugin cap switch to community on this refresh, while the pool keeps +// running with the capacity it was started with. Without this line the drop +// arrives silently with the next deploy, weeks after the licence that caused it. +func (lm *Manager) warnPoolKeepsStartupCeilings(ctx context.Context, claims core.LicenseClaims) { + if lm.startup == nil { + return + } + + if claims.MaxWorkers == lm.startup.MaxWorkers && claims.MaxGenerations == lm.startup.MaxGenerations { + return + } + + lm.logger.WarnContext(ctx, + "licence tier changed while running; the worker pool keeps the ceilings it started with "+ + "and applies the new ones on restart", + "tier", claims.Tier, + "startup_tier", lm.startup.Tier, + "max_workers", claims.MaxWorkers, + "startup_max_workers", lm.startup.MaxWorkers, + "max_generations", claims.MaxGenerations, + "startup_max_generations", lm.startup.MaxGenerations, + ) +} diff --git a/internal/license/manager_test.go b/internal/license/manager_test.go index 72306e3..a080359 100644 --- a/internal/license/manager_test.go +++ b/internal/license/manager_test.go @@ -75,3 +75,57 @@ func TestRefreshKeepsPreviousClaimsOnError(t *testing.T) { require.Equal(t, core.LicenseTierEnterprise, manager.Claims().Tier, "a failed refresh must not downgrade a licence that was valid a moment ago") } + +// TestCeilingChangeWarnsThatPoolNeedsRestart pins the one signal an operator +// gets that the worker pool did not follow a tier change. The pool reads its +// ceilings once at startup; a licence lapsing mid-run otherwise drops capacity +// silently on the next restart. +func TestCeilingChangeWarnsThatPoolNeedsRestart(t *testing.T) { + t.Parallel() + + const warning = "applies the new ones on restart" + + var buf strings.Builder + + logger := slog.New(slog.NewTextHandler(&buf, &slog.HandlerOptions{Level: slog.LevelInfo})) + client := &stubClient{claims: core.EnterpriseLicenseClaims(time.Now().Add(30*24*time.Hour), false)} + + manager, err := NewManager(t.Context(), client, Config{CacheTTL: time.Hour}, + logger, prometheus.NewRegistry(), "test") + require.NoError(t, err) + + require.NotContains(t, buf.String(), warning, + "the first fetch is what the pool is sized from, so there is nothing to warn about") + + client.claims = core.CommunityLicenseClaims() + manager.refresh(t.Context()) + manager.refresh(t.Context()) + + require.Equal(t, 1, strings.Count(buf.String(), warning), + "one warning per transition, not one per refresh") +} + +// TestGenerationsCeilingAloneCountsAsChange covers the field the change check +// used to skip: a licence differing only in MaxGenerations refreshed at Debug, +// so a change to the ceiling that bounds throughput went unannounced. +func TestGenerationsCeilingAloneCountsAsChange(t *testing.T) { + t.Parallel() + + var buf strings.Builder + + logger := slog.New(slog.NewTextHandler(&buf, &slog.HandlerOptions{Level: slog.LevelInfo})) + claims := core.CommunityLicenseClaims() + client := &stubClient{claims: claims} + + manager, err := NewManager(t.Context(), client, Config{CacheTTL: time.Hour}, + logger, prometheus.NewRegistry(), "test") + require.NoError(t, err) + + before := strings.Count(buf.String(), "license refreshed") + + claims.MaxGenerations++ + client.claims = claims + manager.refresh(t.Context()) + + require.Equal(t, before+1, strings.Count(buf.String(), "license refreshed")) +} diff --git a/internal/serve/grpc.go b/internal/serve/grpc.go index ea28864..dd14352 100644 --- a/internal/serve/grpc.go +++ b/internal/serve/grpc.go @@ -27,6 +27,11 @@ func GRPC(log *slog.Logger, host string, port uint16, srv *grpc.Server) func(con select { case err = <-errc: + // Serve returning on its own stops accepting, not serving: the + // connections it already holds keep running handlers. Draining them + // here keeps them from outliving the worker pool and the database, + // which the caller closes as soon as this returns. + srv.GracefulStop() case <-ctx.Done(): srv.GracefulStop() } From 31aeec291309d3d07fce529eda4ad14577853f72 Mon Sep 17 00:00:00 2001 From: Edgar Sipki Date: Thu, 1 Oct 2026 17:19:59 +0300 Subject: [PATCH 2/2] release: v1.0.5 WorkerPool shutdown no longer panics on a racing Get or a second Shutdown; the licence manager warns when a tier change will reach the pool only on restart. Co-Authored-By: Claude Opus 5.5 (1M context) --- AGENTS.md | 2 +- README.md | 6 +++--- deploy/.env.dev.example | 4 ++-- deploy/charts/easyp-service/Chart.yaml | 6 +++--- deploy/charts/easyp-service/README.md | 2 +- deploy/docker-compose.dev.yml | 4 ++-- 6 files changed, 12 insertions(+), 12 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index bf5dcbf..e4a8260 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -33,7 +33,7 @@ tidy. `api/` and `sdk/` carry their own `go.mod` and their own Apache-2.0 before the split a client importing the SDK had its licence scanner report Elastic-2.0 on their own build. -They are tagged separately — `v1.0.4`, `api/v1.0.4`, `sdk/v1.0.4` — and released +They are tagged separately — `v1.0.5`, `api/v1.0.5`, `sdk/v1.0.5` — and released in lockstep, so that one version number means one thing. A change to the wire contract therefore touches a module whose version is a promise to people outside this repository. diff --git a/README.md b/README.md index 71c12ac..2ad9f39 100644 --- a/README.md +++ b/README.md @@ -25,7 +25,7 @@ Every release publishes an image, a Helm chart and the binaries. Pick one. | Tag | Meaning | |-----|---------| -| `v1.0.4` | a release; immutable | +| `v1.0.5` | a release; immutable | | `latest` | the newest release — only moves on a release tag | | `edge` | the tip of `master`; moves on every push | | `sha-` | one commit; immutable | @@ -38,7 +38,7 @@ kubectl create secret generic easyp-env \ --from-literal=DB_POSTGRES_DSN='postgres://user:pass@host:5432/easyp?sslmode=require' helm install easyp oci://ghcr.io/easyp-tech/charts/easyp-service \ - --version 1.0.4 \ + --version 1.0.5 \ --set secrets.existingSecret=easyp-env \ --set tls.enabled=false ``` @@ -979,7 +979,7 @@ A release tag produces, in this order: To verify an image before running it: ```bash -cosign verify ghcr.io/easyp-tech/service:v1.0.4 \ +cosign verify ghcr.io/easyp-tech/service:v1.0.5 \ --certificate-identity-regexp '^https://github.com/easyp-tech/service/\.github/workflows/release\.yml@refs/tags/v' \ --certificate-oidc-issuer https://token.actions.githubusercontent.com ``` diff --git a/deploy/.env.dev.example b/deploy/.env.dev.example index eab27af..3f1a2a1 100644 --- a/deploy/.env.dev.example +++ b/deploy/.env.dev.example @@ -85,7 +85,7 @@ REGISTRY_S3_SECRET_ACCESS_KEY= # EASYP_GID=1000 # --- Service image --- -# Which published image the two tier containers run. A release tag (v1.0.4), +# Which published image the two tier containers run. A release tag (v1.0.5), # `edge` for the tip of master, or `sha-` to pin one commit. Bump this and # run `docker compose pull && docker compose up -d` to upgrade. # @@ -93,7 +93,7 @@ REGISTRY_S3_SECRET_ACCESS_KEY= # "the last release", and compose will happily keep a cached layer from months # ago. A config file in this repository and the binary that reads it have to # agree about which defaults exist. -# EASYP_SERVICE_VERSION=v1.0.4 +# EASYP_SERVICE_VERSION=v1.0.5 # --- Observability (docker-compose.observability.yml) --- # The telemetry backends write to an object store. These may point at the same diff --git a/deploy/charts/easyp-service/Chart.yaml b/deploy/charts/easyp-service/Chart.yaml index 631bc42..e83a6b8 100644 --- a/deploy/charts/easyp-service/Chart.yaml +++ b/deploy/charts/easyp-service/Chart.yaml @@ -51,14 +51,14 @@ type: application # YAML, an alert incapable of producing anything. tests/render.sh now runs # `promtool test rules` against real series rather than only `check rules`, # which reads expressions without knowing whether one can ever fire. -version: 1.0.4 +version: 1.0.5 # appVersion tracks the service release this chart was validated against and is # the default image tag — so it has to be a tag that exists. It was "0.9.0", # and the registry publishes "v0.9.0": a default install pulled a tag that has # never existed and stopped at ImagePullBackOff; `helm template` renders a -# wrong tag as happily as a right one. v1.0.4 is the release this chart ships +# wrong tag as happily as a right one. v1.0.5 is the release this chart ships # with — do not publish the chart before that tag exists. -appVersion: "v1.0.4" +appVersion: "v1.0.5" home: https://github.com/easyp-tech/service sources: diff --git a/deploy/charts/easyp-service/README.md b/deploy/charts/easyp-service/README.md index 39cb6aa..df704f5 100644 --- a/deploy/charts/easyp-service/README.md +++ b/deploy/charts/easyp-service/README.md @@ -34,7 +34,7 @@ checkout is not required: ```bash helm install easyp oci://ghcr.io/easyp-tech/charts/easyp-service \ - --version 1.0.4 \ + --version 1.0.5 \ --set secrets.existingSecret=easyp-env \ --set tls.enabled=false ``` diff --git a/deploy/docker-compose.dev.yml b/deploy/docker-compose.dev.yml index ac82820..66894b5 100644 --- a/deploy/docker-compose.dev.yml +++ b/deploy/docker-compose.dev.yml @@ -135,7 +135,7 @@ services: # loader learned to supply defaults, with an error naming a key whose default # is documented. Set EASYP_SERVICE_VERSION to a release tag, to `edge` for # the tip of master, or to `sha-` to pin the exact commit under test. - image: ghcr.io/easyp-tech/service:${EASYP_SERVICE_VERSION:-v1.0.4} + image: ghcr.io/easyp-tech/service:${EASYP_SERVICE_VERSION:-v1.0.5} container_name: easyp-api-community restart: always user: "${EASYP_UID:-65532}:${EASYP_GID:-65532}" @@ -222,7 +222,7 @@ services: # loader learned to supply defaults, with an error naming a key whose default # is documented. Set EASYP_SERVICE_VERSION to a release tag, to `edge` for # the tip of master, or to `sha-` to pin the exact commit under test. - image: ghcr.io/easyp-tech/service:${EASYP_SERVICE_VERSION:-v1.0.4} + image: ghcr.io/easyp-tech/service:${EASYP_SERVICE_VERSION:-v1.0.5} container_name: easyp-api-enterprise restart: always user: "${EASYP_UID:-65532}:${EASYP_GID:-65532}"