Volver al índice
src/main/java/ar/com/companeros/decision/DecisionGateway.java
package ar.com.companeros.decision;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.ArrayDeque;
import java.util.HashMap;
import java.util.Iterator;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
/** Proveedor asíncrono de decisiones: su salida nunca modifica directamente Minecraft. */
public final class DecisionGateway implements AutoCloseable {
public static final String LUNA = "gpt-6-luna";
public static final String SOL = "gpt-6.1-sol";
private static final URI ENDPOINT = URI.create("https://api.openai.com/v1/responses");
// Configuración confiable del operador. Desactivado hasta verificar la base física y credenciales.
public record Config(boolean enabled, boolean foundationVerified, boolean solEnabled,
BigDecimal maxSessionUsd, int maxRequestsHour, int maxQueue, int maxConcurrent,
int maxOutputTokens, int minAgentIntervalSeconds, Duration timeout, boolean webEnabled,
int maxWebRequestsHour, int minWebIntervalSeconds, java.util.Set<String> webDomains) {
public Config(boolean enabled, boolean foundationVerified, boolean solEnabled, BigDecimal maxSessionUsd,
int maxRequestsHour, int maxQueue, int maxConcurrent, int maxOutputTokens, int minAgentIntervalSeconds,
Duration timeout, boolean webEnabled) {
this(enabled, foundationVerified, solEnabled, maxSessionUsd, maxRequestsHour, maxQueue, maxConcurrent,
maxOutputTokens, minAgentIntervalSeconds, timeout, webEnabled, 4, 900,
ar.com.companeros.runtime.ResearchConfig.defaults().allowedDomains());
}
public Config(boolean enabled, boolean foundationVerified, boolean solEnabled,
BigDecimal maxSessionUsd, int maxRequestsHour, int maxQueue, int maxConcurrent,
int maxOutputTokens, int minAgentIntervalSeconds, Duration timeout) {
this(enabled, foundationVerified, solEnabled, maxSessionUsd, maxRequestsHour, maxQueue,
maxConcurrent, maxOutputTokens, minAgentIntervalSeconds, timeout, false);
}
public Config {
new ar.com.companeros.runtime.ResearchConfig(webEnabled, maxWebRequestsHour, minWebIntervalSeconds, webDomains);
webDomains = java.util.Set.copyOf(webDomains);
java.util.Objects.requireNonNull(maxSessionUsd); java.util.Objects.requireNonNull(timeout);
if (maxSessionUsd.signum() < 0 || maxSessionUsd.precision() > 20 || Math.abs((long) maxSessionUsd.scale()) > 12
|| maxRequestsHour < 1 || maxRequestsHour > 120
|| maxQueue < 1 || maxQueue > 64 || maxConcurrent < 1 || maxConcurrent > 2
|| maxOutputTokens < 256 || maxOutputTokens > 2000 || minAgentIntervalSeconds < 60
|| timeout.isNegative() || timeout.isZero() || timeout.compareTo(Duration.ofSeconds(20)) > 0)
throw new IllegalArgumentException("Configuración de decisiones fuera de límite");
}
public static Config disabled() {
return new Config(false, false, false, BigDecimal.ZERO, 30, 16, 2, 2000, 60, Duration.ofSeconds(20));
}
}
public record BudgetUsage(long dispatchedCalls, int callsLastHour, int inFlight, int queued,
long inputTokens, long outputTokens, long reasoningTokens, long ambiguousCalls,
BigDecimal chargedUsd, BigDecimal reservedUsd, BigDecimal remainingUsd) { }
public record WireReply(int statusCode, String body) { }
@FunctionalInterface public interface Transport {
CompletableFuture<WireReply> send(String payload, String trustedKey);
}
private record Session(long epoch, long observation) { }
private static final class Pending {
final DecisionRequest request;
final CompletableFuture<DecisionResponse> answer = new CompletableFuture<>();
String model, payload;
long inputUpperBound;
BigDecimal reservation = BigDecimal.ZERO;
CompletableFuture<WireReply> wire;
boolean cancelled;
final java.util.List<String> webSources = new java.util.ArrayList<>();
Pending(DecisionRequest request) { this.request = request; }
}
private final Config config;
private final Supplier<String> keySource;
private final StrategicGate strategicGate;
private final Transport transport;
private final LongSupplier clock;
private final ScheduledExecutorService scheduler;
private final BudgetLedger ledger;
private final ArrayDeque<Pending> queue = new ArrayDeque<>();
private final Map<UUID, Pending> active = new HashMap<>();
private final Map<UUID, Session> sessions = new HashMap<>();
private final Map<UUID, Long> lastCalls = new HashMap<>();
private final ArrayDeque<Long> dispatched = new ArrayDeque<>();
private final ArrayDeque<Long> webDispatched = new ArrayDeque<>();
private final Map<UUID, Long> lastWebCalls = new HashMap<>();
private boolean closed;
private String fatalReason;
private long accountRetryAt;
private static final long ACCOUNT_RETRY_MILLIS = 5 * 60_000L;
private long dispatchedCalls, inputTokens, outputTokens, reasoningTokens, ambiguousCalls;
private BigDecimal charged = BigDecimal.ZERO, reserved = BigDecimal.ZERO;
public DecisionGateway(Config config, Supplier<String> trustedKey, StrategicGate strategicGate) {
this(config, trustedKey, strategicGate, new HttpTransport(config.timeout()), System::currentTimeMillis);
}
public DecisionGateway(Config config, Supplier<String> trustedKey, StrategicGate strategicGate, BudgetLedger ledger) {
this(config, trustedKey, strategicGate, new HttpTransport(config.timeout()), System::currentTimeMillis, ledger);
}
/** Inyección de transporte/reloj para pruebas sin red ni claves; no se ofrece como herramienta a agentes. */
public DecisionGateway(Config config, Supplier<String> trustedKey, StrategicGate strategicGate,
Transport transport, LongSupplier clock) {
this(config, trustedKey, strategicGate, transport, clock, BudgetLedger.memory());
}
public DecisionGateway(Config config, Supplier<String> trustedKey, StrategicGate strategicGate,
Transport transport, LongSupplier clock, BudgetLedger ledger) {
this.config = java.util.Objects.requireNonNull(config);
this.keySource = java.util.Objects.requireNonNull(trustedKey);
this.strategicGate = java.util.Objects.requireNonNull(strategicGate);
this.transport = java.util.Objects.requireNonNull(transport); this.clock = java.util.Objects.requireNonNull(clock);
this.ledger = java.util.Objects.requireNonNull(ledger);
try { charged = ledger.load(); }
catch (java.io.IOException | RuntimeException unavailable) {
charged = config.maxSessionUsd(); fatalReason = "BUDGET_STORAGE_ERROR";
}
this.scheduler = Executors.newSingleThreadScheduledExecutor(runnable -> {
Thread thread = new Thread(runnable, "companeros-decision"); thread.setDaemon(true); return thread;
});
scheduler.scheduleWithFixedDelay(this::pumpSafely, 100, 100, TimeUnit.MILLISECONDS);
}
// API del runtime. Las respuestas se revalidan nuevamente en el hilo lógico antes de ejecutar.
public synchronized CompletionStage<DecisionResponse> decide(DecisionRequest request) {
if (closed) return completed(request, DecisionResponse.Status.CLOSED, "PROVIDER_CLOSED");
if (disabled()) return completed(request, DecisionResponse.Status.DISABLED, disabledReason());
if (request.kind() == DecisionRequest.Kind.WEB_RESEARCH && !config.webEnabled())
return completed(request, DecisionResponse.Status.DISABLED, "WEB_DISABLED");
if (sessions.size() >= 256 && !sessions.containsKey(request.actorUUID()))
return completed(request, DecisionResponse.Status.THROTTLED, "SESSION_LIMIT");
sessions.putIfAbsent(request.actorUUID(), new Session(request.sessionEpoch(), request.observationSeq()));
if (!current(request)) return completed(request, DecisionResponse.Status.STALE, "SESSION_CHANGED");
if (active.containsKey(request.actorUUID()) || queue.stream().anyMatch(p -> p.request.actorUUID().equals(request.actorUUID())))
return completed(request, DecisionResponse.Status.THROTTLED, "DECISION_ALREADY_PENDING");
if (queue.size() >= config.maxQueue()) return completed(request, DecisionResponse.Status.THROTTLED, "QUEUE_FULL");
Pending pending = new Pending(request); queue.addLast(pending); pump();
return pending.answer;
}
public synchronized void updateSession(UUID actor, long epoch, long observation) {
if (epoch < 0 || observation < 0) throw new IllegalArgumentException("Revisión de sesión inválida");
if (sessions.size() >= 256 && !sessions.containsKey(actor))
throw new IllegalStateException("Registro de sesiones fuera de límite");
sessions.put(actor, new Session(epoch, observation));
Iterator<Pending> iterator = queue.iterator();
while (iterator.hasNext()) {
Pending pending = iterator.next();
if (pending.request.actorUUID().equals(actor) && !current(pending.request)) {
iterator.remove(); pending.answer.complete(response(pending, DecisionResponse.Status.STALE, null,
DecisionResponse.Usage.zero(), "SESSION_CHANGED"));
}
}
Pending pending = active.get(actor);
if (pending != null && !current(pending.request)) {
pending.cancelled = true;
// Mantener el transporte hasta respuesta/timeout conserva cupo y contabilización real.
}
}
public synchronized void cancel(UUID actor, long epoch) {
updateSession(actor, epoch + 1, 0);
}
public synchronized boolean closed() { return closed; }
public synchronized boolean disabled() {
// Un saldo recargado debe poder recuperarse sin cambiar archivos ni reiniciar.
if (("ACCOUNT_NO_CREDIT".equals(fatalReason) || "MODEL_OR_ACCOUNT_ACCESS".equals(fatalReason))
&& clock.getAsLong() >= accountRetryAt) {
fatalReason = null; accountRetryAt = 0;
}
return !config.enabled() || !config.foundationVerified() || config.maxSessionUsd().signum() == 0
|| fatalReason != null || trustedKey().isBlank();
}
public synchronized BudgetUsage budgetUsage() {
prune(clock.getAsLong());
return new BudgetUsage(dispatchedCalls, dispatched.size(), active.size(), queue.size(), inputTokens,
outputTokens, reasoningTokens, ambiguousCalls, charged, reserved,
config.maxSessionUsd().subtract(charged).subtract(reserved).max(BigDecimal.ZERO));
}
// Cola, cupos y reserva conservadora antes de cualquier solicitud de red.
private void pumpSafely() {
synchronized (this) { if (!closed) pump(); }
}
private void pump() {
if (closed) return;
long now = clock.getAsLong(); prune(now);
if (disabled()) {
while (!queue.isEmpty()) {
Pending pending = queue.removeFirst(); pending.answer.complete(response(pending,
DecisionResponse.Status.DISABLED, null, DecisionResponse.Usage.zero(), disabledReason()));
}
return;
}
while (active.size() < config.maxConcurrent() && !queue.isEmpty()) {
Pending eligible = null;
for (Pending candidate : queue) {
Long last = lastCalls.get(candidate.request.actorUUID());
if (last == null || now - last >= config.minAgentIntervalSeconds() * 1000L) { eligible = candidate; break; }
}
if (eligible == null) return;
queue.remove(eligible);
if (eligible.request.kind() == DecisionRequest.Kind.WEB_RESEARCH) {
while (!webDispatched.isEmpty() && now - webDispatched.peekFirst() >= 3_600_000) webDispatched.removeFirst();
Long previous = lastWebCalls.get(eligible.request.actorUUID());
if (webDispatched.size() >= config.maxWebRequestsHour() || previous != null && now - previous < config.minWebIntervalSeconds() * 1000L) {
eligible.answer.complete(response(eligible, DecisionResponse.Status.THROTTLED, null,
DecisionResponse.Usage.zero(), "WEB_RESEARCH_COOLDOWN")); continue;
}
}
if (!current(eligible.request)) {
eligible.answer.complete(response(eligible, DecisionResponse.Status.STALE, null,
DecisionResponse.Usage.zero(), "SESSION_CHANGED")); continue;
}
if (dispatched.size() >= config.maxRequestsHour()) {
eligible.answer.complete(response(eligible, DecisionResponse.Status.THROTTLED, null,
DecisionResponse.Usage.zero(), "HOURLY_LIMIT")); continue;
}
eligible.model = eligible.request.kind() == DecisionRequest.Kind.STRATEGIC ? SOL : LUNA;
try {
eligible.payload = DecisionCodec.payload(eligible.request, eligible.model, config.maxOutputTokens(), config.webDomains());
// Un token de texto no supera la cantidad total de bytes UTF-8. Se suma margen de protocolo.
eligible.inputUpperBound = eligible.payload.getBytes(StandardCharsets.UTF_8).length + 1024L;
// La búsqueda incorpora tokens del proveedor: reserva amplia y conservadora.
if (eligible.request.kind() == DecisionRequest.Kind.WEB_RESEARCH) eligible.inputUpperBound += 1_000_000;
eligible.reservation = estimate(eligible.model, eligible.inputUpperBound, 0,
eligible.inputUpperBound, config.maxOutputTokens(), 0, true).usd();
if (eligible.request.kind() == DecisionRequest.Kind.WEB_RESEARCH)
eligible.reservation = eligible.reservation.max(new BigDecimal("0.30")).add(new BigDecimal("0.01"));
} catch (RuntimeException invalid) {
eligible.answer.complete(response(eligible, DecisionResponse.Status.INVALID_RESPONSE, null,
DecisionResponse.Usage.zero(), "INVALID_REQUEST")); continue;
}
if (charged.add(reserved).add(eligible.reservation).compareTo(config.maxSessionUsd()) > 0) {
eligible.answer.complete(response(eligible, DecisionResponse.Status.BUDGET_EXHAUSTED, null,
DecisionResponse.Usage.zero(), "SESSION_BUDGET")); continue;
}
if (eligible.request.kind() == DecisionRequest.Kind.STRATEGIC
&& (!config.solEnabled() || !strategicGate.consume(eligible.request.strategicPermit(), eligible.request.actorUUID()))) {
eligible.answer.complete(response(eligible, DecisionResponse.Status.STRATEGIC_DENIED, null,
DecisionResponse.Usage.zero(), "STRATEGIC_PERMISSION")); continue;
}
Pending pending = eligible;
reserved = reserved.add(pending.reservation);
if (!persistBudget()) {
reserved = reserved.subtract(pending.reservation);
pending.answer.complete(response(pending, DecisionResponse.Status.DISABLED, null,
DecisionResponse.Usage.zero(), "BUDGET_STORAGE_ERROR"));
return;
}
active.put(pending.request.actorUUID(), pending);
lastCalls.put(pending.request.actorUUID(), now); dispatched.addLast(now); dispatchedCalls++;
if (pending.request.kind() == DecisionRequest.Kind.WEB_RESEARCH) {
lastWebCalls.put(pending.request.actorUUID(), now); webDispatched.addLast(now);
}
try {
pending.wire = transport.send(pending.payload, trustedKey());
pending.wire.orTimeout(config.timeout().toMillis(), TimeUnit.MILLISECONDS)
.whenComplete((reply, error) -> settle(pending, reply, error));
} catch (RuntimeException error) { settle(pending, null, error); }
}
}
private synchronized void settle(Pending pending, WireReply reply, Throwable error) {
if (active.get(pending.request.actorUUID()) != pending) return;
active.remove(pending.request.actorUUID()); reserved = reserved.subtract(pending.reservation);
DecisionResponse.Status status = DecisionResponse.Status.NETWORK_ERROR;
ToolDecision plan = null; String reason = "NETWORK_FAILURE_OR_TIMEOUT";
DecisionResponse.Usage usage = estimate(pending.model, pending.inputUpperBound, 0,
pending.inputUpperBound, config.maxOutputTokens(), 0, true);
if (error == null && reply != null) {
if (reply.statusCode() != 200) {
status = DecisionResponse.Status.HTTP_ERROR; reason = "HTTP_" + reply.statusCode();
// Rechazos explícitos no son solicitudes ejecutadas; fallos 5xx siguen siendo ambiguos.
if (reply.statusCode() >= 400 && reply.statusCode() < 500) usage = DecisionResponse.Usage.zero();
if (reply.statusCode() == 401 || reply.statusCode() == 403 || reply.statusCode() == 404) {
fatalReason = "MODEL_OR_ACCOUNT_ACCESS";
accountRetryAt = clock.getAsLong() + ACCOUNT_RETRY_MILLIS;
}
if ((reply.statusCode() == 429 || reply.statusCode() == 402) && reply.body().length() <= 262_144) {
try {
JsonObject rejection = DecisionCodec.strictObject(reply.body(), 262_144).getAsJsonObject("error");
String code = rejection.get("code").getAsString();
if (java.util.Set.of("credit_balance_exhausted", "insufficient_quota", "billing_hard_limit_reached").contains(code)) {
fatalReason = "ACCOUNT_NO_CREDIT";
accountRetryAt = clock.getAsLong() + ACCOUNT_RETRY_MILLIS;
}
} catch (RuntimeException malformedError) { /* Un límite transitorio no cambia las credenciales. */ }
}
} else {
try {
if (reply.body().length() > 262_144) throw new IllegalArgumentException("Respuesta fuera de límite");
JsonObject root = DecisionCodec.strictObject(reply.body(), 262_144);
if (root.has("usage") && root.get("usage").isJsonObject()) usage = parseUsage(root.getAsJsonObject("usage"), pending.model);
String returnedModel = root.has("model") ? root.get("model").getAsString() : "";
if (!returnedModel.equals(pending.model) && !returnedModel.startsWith(pending.model + "-"))
throw new IllegalArgumentException("Modelo no solicitado");
if (!root.has("status") || !root.get("status").getAsString().equals("completed")) {
status = DecisionResponse.Status.INCOMPLETE; reason = "RESPONSE_NOT_COMPLETED";
} else {
StringBuilder text = new StringBuilder(); boolean refused = false;
for (JsonElement item : root.getAsJsonArray("output")) {
JsonObject output = item.getAsJsonObject();
String type = output.get("type").getAsString();
if (type.equals("reasoning")) continue;
if (type.equals("web_search_call") && pending.request.kind() == DecisionRequest.Kind.WEB_RESEARCH) {
JsonObject action = output.getAsJsonObject("action");
if (action != null && action.has("sources")) for (JsonElement source : action.getAsJsonArray("sources")) {
String url = source.getAsJsonObject().get("url").getAsString();
URI uri = URI.create(url);
String host = uri.getHost() == null ? "" : uri.getHost().toLowerCase(java.util.Locale.ROOT);
if (url.length() <= 512 && ("https".equals(uri.getScheme()) || "http".equals(uri.getScheme()))
&& uri.getUserInfo() == null && config.webDomains().stream().anyMatch(domain -> host.equals(domain) || host.endsWith("." + domain))
&& pending.webSources.size() < 4 && !pending.webSources.contains(url)) pending.webSources.add(url);
}
continue;
}
if (!type.equals("message")) throw new IllegalArgumentException("Salida no permitida");
for (JsonElement content : output.getAsJsonArray("content")) {
JsonObject value = content.getAsJsonObject();
String contentType = value.get("type").getAsString();
if (contentType.equals("refusal")) refused = true;
else if (contentType.equals("output_text")) text.append(value.get("text").getAsString());
else throw new IllegalArgumentException("Contenido no permitido");
}
}
if (refused) { status = DecisionResponse.Status.REFUSED; reason = "MODEL_REFUSAL"; }
else {
plan = DecisionCodec.validate(text.toString(), pending.request);
if (pending.request.kind() == DecisionRequest.Kind.WEB_RESEARCH && pending.webSources.isEmpty())
throw new IllegalArgumentException("Investigación sin fuentes");
status = DecisionResponse.Status.ACCEPTED; reason = "VALIDATED_PROPOSAL";
}
}
} catch (RuntimeException invalid) { status = DecisionResponse.Status.INVALID_RESPONSE; reason = "INVALID_PROVIDER_RESPONSE"; }
}
}
if (pending.request.kind() == DecisionRequest.Kind.WEB_RESEARCH && (reply == null || reply.statusCode() == 200 || reply.statusCode() >= 500))
usage = new DecisionResponse.Usage(usage.inputTokens(), usage.outputTokens(), usage.reasoningTokens(),
usage.cachedTokens(), usage.cacheWriteTokens(), usage.estimatedBillable(), usage.usd().add(new BigDecimal("0.01")));
charged = charged.add(usage.usd()); inputTokens += usage.inputTokens(); outputTokens += usage.outputTokens();
reasoningTokens += usage.reasoningTokens(); if (usage.estimatedBillable()) ambiguousCalls++;
persistBudget();
if (closed || pending.cancelled || !current(pending.request)) {
status = closed ? DecisionResponse.Status.CANCELLED : DecisionResponse.Status.STALE;
plan = null; reason = "SESSION_CHANGED_OR_CANCELLED";
}
pending.answer.complete(response(pending, status, plan, usage, reason));
if (!closed) pump();
}
// output_tokens incluye razonamiento. Cache writes y lecturas tienen tarifas separadas.
private static DecisionResponse.Usage parseUsage(JsonObject usage, String model) {
if (!usage.has("input_tokens") || !usage.has("output_tokens")
|| usage.get("input_tokens").isJsonNull() || usage.get("output_tokens").isJsonNull())
throw new IllegalArgumentException("Uso de tokens incompleto");
long input = number(usage, "input_tokens"), output = number(usage, "output_tokens");
JsonObject details = usage.has("input_tokens_details") && usage.get("input_tokens_details").isJsonObject()
? usage.getAsJsonObject("input_tokens_details") : new JsonObject();
JsonObject outputDetails = usage.has("output_tokens_details") && usage.get("output_tokens_details").isJsonObject()
? usage.getAsJsonObject("output_tokens_details") : new JsonObject();
long cached = number(details, "cached_tokens"), written = number(details, "cache_write_tokens");
long reasoning = number(outputDetails, "reasoning_tokens");
if (cached + written > input || reasoning > output) throw new IllegalArgumentException("Uso de tokens inválido");
return estimate(model, input, cached, written, output, reasoning, false);
}
private static long number(JsonObject object, String key) {
if (!object.has(key) || object.get(key).isJsonNull()) return 0;
long value = object.get(key).getAsBigDecimal().longValueExact();
if (value < 0 || value > 10_000_000) throw new IllegalArgumentException("Uso fuera de límite");
return value;
}
public static DecisionResponse.Usage estimate(String model, long input, long cached, long written,
long output, long reasoning, boolean ambiguous) {
if (!model.equals(LUNA) && !model.equals(SOL)) throw new IllegalArgumentException("Modelo no permitido");
if (input < 0 || cached < 0 || written < 0 || output < 0 || reasoning < 0
|| cached + written > input || reasoning > output) throw new IllegalArgumentException("Uso inválido");
// Standard, contexto corto, valores documentados el 4-10-2026; Fast/servicios externos no habilitados.
BigDecimal inputRate = new BigDecimal(model.equals(LUNA) ? "0.10" : "2.00");
BigDecimal cachedRate = new BigDecimal(model.equals(LUNA) ? "0.01" : "0.10");
BigDecimal writtenRate = new BigDecimal(model.equals(LUNA) ? "0.125" : "2.50");
BigDecimal outputRate = new BigDecimal(model.equals(LUNA) ? "0.50" : "10.00");
BigDecimal price = inputRate.multiply(BigDecimal.valueOf(input - cached - written))
.add(cachedRate.multiply(BigDecimal.valueOf(cached)))
.add(writtenRate.multiply(BigDecimal.valueOf(written)))
.add(outputRate.multiply(BigDecimal.valueOf(output)))
.divide(BigDecimal.valueOf(1_000_000), 9, RoundingMode.CEILING);
return new DecisionResponse.Usage(input, output, reasoning, cached, written, ambiguous, price);
}
private boolean current(DecisionRequest request) {
Session session = sessions.get(request.actorUUID());
return session != null && session.epoch() == request.sessionEpoch() && session.observation() == request.observationSeq();
}
// También guarda reservas: apagar o cambiar de mundo nunca restablece el presupuesto.
private boolean persistBudget() {
try { ledger.save(charged.add(reserved)); return true; }
catch (java.io.IOException | RuntimeException failure) {
fatalReason = "BUDGET_STORAGE_ERROR"; return false;
}
}
private void prune(long now) { while (!dispatched.isEmpty() && now - dispatched.peekFirst() >= 3_600_000) dispatched.removeFirst(); }
private String trustedKey() {
if (!config.enabled() || !config.foundationVerified()) return "";
try { String key = keySource.get(); return key == null ? "" : key.strip(); }
catch (RuntimeException unavailable) { return ""; }
}
private String disabledReason() {
if (!config.enabled()) return "OPERATOR_DISABLED";
if (!config.foundationVerified()) return "FOUNDATION_NOT_VERIFIED";
if (fatalReason != null) return fatalReason;
if (config.maxSessionUsd().signum() == 0) return "NO_BUDGET";
return "NO_CREDENTIAL";
}
private static CompletionStage<DecisionResponse> completed(DecisionRequest request, DecisionResponse.Status status, String reason) {
return CompletableFuture.completedFuture(new DecisionResponse(request.requestId(), request.actorUUID(),
request.sessionEpoch(), request.observationSeq(), status, "", null, DecisionResponse.Usage.zero(), reason));
}
private static DecisionResponse response(Pending pending, DecisionResponse.Status status,
ToolDecision plan, DecisionResponse.Usage usage, String reason) {
return new DecisionResponse(pending.request.requestId(), pending.request.actorUUID(),
pending.request.sessionEpoch(), pending.request.observationSeq(), status,
pending.model == null ? "" : pending.model, plan, usage, reason, pending.webSources);
}
@Override public synchronized void close() {
if (closed) return; closed = true;
while (!queue.isEmpty()) {
Pending pending = queue.removeFirst(); pending.answer.complete(response(pending,
DecisionResponse.Status.CANCELLED, null, DecisionResponse.Usage.zero(), "PROVIDER_CLOSED"));
}
for (Pending pending : new java.util.ArrayList<>(active.values())) {
pending.cancelled = true;
settle(pending, null, new IllegalStateException("Closed"));
if (pending.wire != null) pending.wire.cancel(true);
}
scheduler.shutdownNow();
ledger.close();
}
// HTTPS solamente al endpoint fijo; sin SDK, red del agente o claves en la observación.
private static final class HttpTransport implements Transport {
private final Duration timeout;
private HttpClient client;
HttpTransport(Duration timeout) { this.timeout = timeout; }
@Override public synchronized CompletableFuture<WireReply> send(String payload, String key) {
if (client == null) client = HttpClient.newBuilder().connectTimeout(timeout)
.followRedirects(HttpClient.Redirect.NEVER).build();
HttpRequest request = HttpRequest.newBuilder(ENDPOINT).timeout(timeout)
.header("Content-Type", "application/json").header("Authorization", "Bearer " + key)
.POST(HttpRequest.BodyPublishers.ofString(payload, StandardCharsets.UTF_8)).build();
return client.sendAsync(request, HttpResponse.BodyHandlers.ofString(StandardCharsets.UTF_8))
.thenApply(response -> new WireReply(response.statusCode(), response.body()));
}
}
}