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()));
        }
    }
}