Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
Expand Down
14 changes: 11 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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-<short>` | one commit; immutable |
Expand All @@ -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
```
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -971,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
```
Expand Down
5 changes: 5 additions & 0 deletions cmd/easyp-svc/start.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions deploy/.env.dev.example
Original file line number Diff line number Diff line change
Expand Up @@ -85,15 +85,15 @@ 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-<short>` to pin one commit. Bump this and
# run `docker compose pull && docker compose up -d` to upgrade.
#
# Not `latest`: it is published only on a release tag, so it silently means
# "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
Expand Down
6 changes: 3 additions & 3 deletions deploy/charts/easyp-service/Chart.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
2 changes: 1 addition & 1 deletion deploy/charts/easyp-service/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
```
Expand Down
4 changes: 2 additions & 2 deletions deploy/docker-compose.dev.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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-<short>` 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}"
Expand Down Expand Up @@ -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-<short>` 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}"
Expand Down
64 changes: 49 additions & 15 deletions internal/core/pool.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 одновременных генераций, держит очередь
Expand Down Expand Up @@ -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,
Expand All @@ -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()
Expand Down Expand Up @@ -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() {
Expand Down
71 changes: 71 additions & 0 deletions internal/core/pool_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
43 changes: 43 additions & 0 deletions internal/license/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
}

Expand Down Expand Up @@ -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()
Expand All @@ -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,
)
}
Loading
Loading