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 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 send(String payload, String trustedKey); } private record Session(long epoch, long observation) { } private static final class Pending { final DecisionRequest request; final CompletableFuture answer = new CompletableFuture<>(); String model, payload; long inputUpperBound; BigDecimal reservation = BigDecimal.ZERO; CompletableFuture wire; boolean cancelled; final java.util.List webSources = new java.util.ArrayList<>(); Pending(DecisionRequest request) { this.request = request; } } private final Config config; private final Supplier keySource; private final StrategicGate strategicGate; private final Transport transport; private final LongSupplier clock; private final ScheduledExecutorService scheduler; private final BudgetLedger ledger; private final ArrayDeque queue = new ArrayDeque<>(); private final Map active = new HashMap<>(); private final Map sessions = new HashMap<>(); private final Map lastCalls = new HashMap<>(); private final ArrayDeque dispatched = new ArrayDeque<>(); private final ArrayDeque webDispatched = new ArrayDeque<>(); private final Map 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 trustedKey, StrategicGate strategicGate) { this(config, trustedKey, strategicGate, new HttpTransport(config.timeout()), System::currentTimeMillis); } public DecisionGateway(Config config, Supplier 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 trustedKey, StrategicGate strategicGate, Transport transport, LongSupplier clock) { this(config, trustedKey, strategicGate, transport, clock, BudgetLedger.memory()); } public DecisionGateway(Config config, Supplier 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 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 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 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 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())); } } }