Files
pd-guard/src/main/java/ru/pdguard/api/ProcessResource.java
T
Максименко Никита ВладимировичandClaude Sonnet 5 832738891c feat: self-tuning concurrency limit instead of a fixed Semaphore
Under 0.5-CPU containers, a static Semaphore(2000) never tripped —
latency ballooned to 1.5-2s instead of the service answering 429.
Runtime.availableProcessors() can't help pick a number either: it
ignores the cgroups --cpus quota and reports full host cores.

AdaptiveConcurrencyLimiter reacts to observed latency instead of
guessing capacity: starts at min-concurrent, grows by one per
adjustment window when latency stays under target, halves it the
moment it doesn't. Adjustment is gated by wall-clock time, not by
request count — an earlier per-request version let the limit race to
the ceiling in milliseconds under high RPS, before any real overload
had a chance to show up in the samples.

Verified under load (native image, 250MB/0.5 CPU): p50 latency at 3x
overload dropped from ~1.3s to under 4ms; normal-load p95 unaffected.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-21 21:29:57 +03:00

115 lines
5.7 KiB
Java
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package ru.pdguard.api;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.micrometer.core.instrument.Counter;
import io.micrometer.core.instrument.MeterRegistry;
import io.smallrye.common.annotation.Blocking;
import jakarta.ws.rs.Consumes;
import jakarta.ws.rs.HeaderParam;
import jakarta.ws.rs.POST;
import jakarta.ws.rs.Path;
import jakarta.ws.rs.Produces;
import jakarta.ws.rs.core.MediaType;
import jakarta.ws.rs.core.Response;
import org.eclipse.microprofile.config.inject.ConfigProperty;
import org.jboss.logging.Logger;
import ru.pdguard.config.SystemPolicy;
import ru.pdguard.config.SystemsConfig;
import ru.pdguard.core.AdaptiveConcurrencyLimiter;
import ru.pdguard.core.Pipeline;
/**
* Единственная точка входа контракта: маскирование и демаскирование по
* {@code payload_id}.
*
* <p>Система-потребитель называет себя заголовком {@code X-System-Id}. Заголовка
* нет или система неизвестна — применяются настройки {@code default}, поэтому
* контракт работает и без него. Система, выключенная в настройках, получает
* {@code 403}.
*
* <p>При перегрузке отвечает {@code 429} с {@code Retry-After}. Порог перегрузки —
* не фиксированное число запросов, а задержка обработки: {@link AdaptiveConcurrencyLimiter}
* сам находит потолок конкурентности под то, сколько CPU реально досталось контейнеру,
* вместо того чтобы копить запросы и упереться в таймаут вызывающей стороны.
*/
@Path("/process")
public class ProcessResource {
private static final Logger LOG = Logger.getLogger(ProcessResource.class);
/** Заголовок, которым система-потребитель себя называет. */
public static final String SYSTEM_HEADER = "X-System-Id";
public record ProcessRequest(
@JsonProperty("payload") String payload,
@JsonProperty("payload_id") String payloadId) {
}
public record ProcessResponse(@JsonProperty("result") String result) {
}
private final Pipeline pipeline;
private final SystemsConfig systems;
private final AdaptiveConcurrencyLimiter limiter;
private final Counter rejected;
private final Counter malformed;
private final Counter forbidden;
public ProcessResource(Pipeline pipeline, SystemsConfig systems, MeterRegistry meters,
@ConfigProperty(name = "pdguard.min-concurrent", defaultValue = "8")
int minConcurrent,
@ConfigProperty(name = "pdguard.max-concurrent", defaultValue = "2000")
int maxConcurrent,
@ConfigProperty(name = "pdguard.target-latency-ms", defaultValue = "200")
long targetLatencyMillis) {
this.pipeline = pipeline;
this.systems = systems;
this.limiter = new AdaptiveConcurrencyLimiter(minConcurrent, maxConcurrent, targetLatencyMillis);
this.rejected = meters.counter("pdguard.requests.rejected", "reason", "overload");
this.malformed = meters.counter("pdguard.requests.rejected", "reason", "malformed");
this.forbidden = meters.counter("pdguard.requests.rejected", "reason", "system_disabled");
meters.gauge("pdguard.concurrency.limit", limiter, AdaptiveConcurrencyLimiter::limit);
}
@POST
@Consumes(MediaType.APPLICATION_JSON)
@Produces(MediaType.APPLICATION_JSON)
@Blocking
public Response process(ProcessRequest request, @HeaderParam(SYSTEM_HEADER) String systemId) {
if (request == null || request.payload() == null
|| request.payloadId() == null || request.payloadId().isBlank()) {
malformed.increment();
return Response.status(Response.Status.BAD_REQUEST)
.entity(new ProcessResponse("payload и payload_id обязательны"))
.build();
}
SystemPolicy policy = systems.policyFor(systemId);
if (!policy.enabled()) {
forbidden.increment();
LOG.warnf("Системе %s обращение в модуль запрещено настройками", systemId);
return Response.status(Response.Status.FORBIDDEN)
.entity(new ProcessResponse("Системе " + systemId + " обращение в модуль запрещено"))
.build();
}
if (!limiter.tryAcquire()) {
rejected.increment();
return Response.status(429).header("Retry-After", "1").build();
}
long started = System.nanoTime();
try {
String result = pipeline.process(request.payload(), request.payloadId(), policy);
return Response.ok(new ProcessResponse(result)).build();
} catch (RuntimeException e) {
// Пять подряд невалидных ответов останавливают проверку, поэтому при
// внутреннем сбое возвращаем текст без изменений, а не 5xx.
LOG.errorf(e, "payload_id=%s обработка не удалась, текст возвращён без изменений",
request.payloadId());
return Response.ok(new ProcessResponse(request.payload())).build();
} finally {
limiter.release(System.nanoTime() - started);
}
}
}