Volver al índice

src/main/java/ar/com/companeros/horror/HorrorSpeech.java

package ar.com.companeros.horror;

import com.google.gson.*;
import java.io.*;
import java.net.URI;
import java.net.http.*;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.nio.file.*;
import java.time.Duration;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.Flow;
import net.minecraft.network.chat.Component;
import net.minecraft.server.MinecraftServer;
import net.minecraft.server.level.ServerLevel;
import net.minecraft.server.level.ServerPlayer;
import net.minecraftforge.common.MinecraftForge;
import net.minecraftforge.event.server.ServerStoppingEvent;
import net.minecraftforge.fml.loading.FMLPaths;

/** Texto decidido por Luna -> servicio de voz permitido -> observadores humanos. */
public final class HorrorSpeech {
    public static final int MAX_AUDIO_BYTES = 192 * 1024;
    private static final Gson JSON = new GsonBuilder().setPrettyPrinting().create();
    private static final Map<MinecraftServer, Settings> SETTINGS = new WeakHashMap<>();
    private static final Map<MinecraftServer, Map<UUID, Long>> LAST = new WeakHashMap<>();
    private static final Map<MinecraftServer, Map<UUID, CompletableFuture<Boolean>>> REQUESTS = new WeakHashMap<>();
    private static final Set<MinecraftServer> STOPPING = Collections.synchronizedSet(Collections.newSetFromMap(new WeakHashMap<>()));
    private static final Set<UUID> PENDING = new HashSet<>();
    private static final HttpClient CLIENT = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(5))
        .followRedirects(HttpClient.Redirect.NEVER).build();
    private static boolean installed;
    private HorrorSpeech() { }

    // CONFIGURACIÓN: la URL sólo procede del operador, nunca de una frase o herramienta.
    public record Settings(boolean enabled, String endpoint, boolean textFallback,
                           int minIntervalSeconds, int timeoutSeconds) {
        public Settings {
            Objects.requireNonNull(endpoint);
            if (minIntervalSeconds < 10 || minIntervalSeconds > 600 || timeoutSeconds < 1 || timeoutSeconds > 30)
                throw new IllegalArgumentException("Límites de voz inválidos");
            if (enabled) permittedEndpoint(endpoint);
        }
        public static Settings defaults() { return new Settings(false, "", true, 15, 12); }
    }
    public static URI permittedEndpoint(String endpoint) {
        URI uri = URI.create(endpoint);
        String host = uri.getHost();
        boolean loopback = "127.0.0.1".equals(host) || "localhost".equalsIgnoreCase(host) || "[::1]".equals(host);
        if (host == null || uri.getUserInfo() != null || uri.getFragment() != null
                || !("https".equalsIgnoreCase(uri.getScheme()) || (loopback && "http".equalsIgnoreCase(uri.getScheme()))))
            throw new IllegalArgumentException("La voz requiere HTTPS o un servicio local");
        return uri;
    }
    public static Settings load(Path root) throws IOException {
        Files.createDirectories(root);
        Path file = root.resolve("voice.json");
        if (!Files.exists(file)) {
            JsonObject defaults = new JsonObject(); defaults.addProperty("schemaVersion", 1);
            Settings s = Settings.defaults(); defaults.addProperty("enabled", s.enabled()); defaults.addProperty("endpoint", s.endpoint());
            defaults.addProperty("textFallback", s.textFallback()); defaults.addProperty("minIntervalSeconds", s.minIntervalSeconds());
            defaults.addProperty("timeoutSeconds", s.timeoutSeconds());
            Files.writeString(file, JSON.toJson(defaults), StandardCharsets.UTF_8, StandardOpenOption.CREATE_NEW);
        }
        if (Files.isSymbolicLink(file) || Files.size(file) > 4096) throw new IOException("voice.json fuera de límite");
        try {
            JsonObject data = JsonParser.parseString(Files.readString(file, StandardCharsets.UTF_8)).getAsJsonObject();
            if (!Set.of("schemaVersion", "enabled", "endpoint", "textFallback", "minIntervalSeconds", "timeoutSeconds").equals(data.keySet())
                    || data.get("schemaVersion").getAsInt() != 1) throw new IllegalArgumentException("Formato desconocido");
            return new Settings(data.get("enabled").getAsBoolean(), data.get("endpoint").getAsString(),
                data.get("textFallback").getAsBoolean(), data.get("minIntervalSeconds").getAsInt(), data.get("timeoutSeconds").getAsInt());
        } catch (RuntimeException invalid) { throw new IOException("No se pudo leer voice.json", invalid); }
    }
    public static void install() {
        if (installed) return; installed = true;
        MinecraftForge.EVENT_BUS.addListener((ServerStoppingEvent event) -> {
            STOPPING.add(event.getServer());
            var requests = REQUESTS.remove(event.getServer());
            if (requests != null) for (var request : List.copyOf(requests.entrySet())) {
                PENDING.remove(request.getKey()); request.getValue().cancel(true);
            }
            SETTINGS.remove(event.getServer()); LAST.remove(event.getServer());
            HorrorModule.findParasite(event.getServer()).ifPresent(entity -> PENDING.remove(entity.getUUID()));
        });
    }

    // PETICIÓN: un cuerpo canónico, una frase acotada y una petición pendiente como máximo.
    public static CompletableFuture<Boolean> speak(ParasiteEntity entity, String text) {
        if (!(entity.level() instanceof ServerLevel level) || !validText(text)) return CompletableFuture.completedFuture(false);
        MinecraftServer server = level.getServer();
        if (STOPPING.contains(server)) return CompletableFuture.completedFuture(false);
        if (!entity.isAlive() || HorrorModule.findParasite(server).filter(body -> body == entity).isEmpty())
            return CompletableFuture.completedFuture(false);
        List<ServerPlayer> audience = audience(entity);
        if (audience.isEmpty()) return CompletableFuture.completedFuture(false);
        Settings settings;
        try { settings = SETTINGS.computeIfAbsent(server, key -> {
            try { return load(FMLPaths.CONFIGDIR.get().resolve("raps")); }
            catch (IOException error) { throw new IllegalStateException(error); }
        }); } catch (RuntimeException invalid) { return CompletableFuture.completedFuture(false); }
        long now = level.getGameTime(); Map<UUID, Long> last = LAST.computeIfAbsent(server, key -> new HashMap<>());
        if (PENDING.contains(entity.getUUID()) || now - last.getOrDefault(entity.getUUID(), -12000L) < settings.minIntervalSeconds() * 20L)
            return CompletableFuture.completedFuture(false);
        Set<UUID> reserved = new HashSet<>();
        for (ServerPlayer player : audience) if (HorrorPacing.get(server).reserveSound(player.getUUID(), now)) reserved.add(player.getUUID());
        if (reserved.isEmpty()) return CompletableFuture.completedFuture(false);
        last.put(entity.getUUID(), now);
        if (!settings.enabled()) return CompletableFuture.completedFuture(settings.textFallback() && deliverText(entity, text, reserved));
        PENDING.add(entity.getUUID()); UUID bodyId = entity.getUUID();
        CompletableFuture<Boolean> result = new CompletableFuture<>();
        REQUESTS.computeIfAbsent(server, key -> new HashMap<>()).put(bodyId, result);
        CompletableFuture<byte[]> request = request(settings, text);
        result.whenComplete((success, error) -> { if (result.isCancelled()) request.cancel(true); });
        request.whenComplete((audio, error) -> {
          if (STOPPING.contains(server)) { result.cancel(true); return; }
          server.execute(() -> {
            PENDING.remove(bodyId);
            var tracked = REQUESTS.get(server);
            if (tracked != null) tracked.remove(bodyId, result);
            if (result.isCancelled()) return;
            if (!entity.isAlive() || HorrorModule.findParasite(server).filter(body -> body == entity).isEmpty()) { result.complete(false); return; }
            List<ServerPlayer> current = audience(entity).stream().filter(player -> reserved.contains(player.getUUID())).toList();
            if (error == null && validAudio(audio) && !current.isEmpty()) {
                for (ServerPlayer player : current) HorrorNetwork.sendSpeech(player, entity, audio);
                result.complete(true);
            } else result.complete(settings.textFallback() && deliverText(entity, text, reserved));
          });
        });
        return result;
    }
    public static boolean validText(String text) {
        return text != null && !text.isBlank() && text.length() <= 320
            && text.codePoints().noneMatch(c -> Character.isISOControl(c) && c != '\n');
    }
    public static boolean validAudio(byte[] data) {
        return HorrorNetwork.validSpeechBytes(data);
    }
    private static List<ServerPlayer> audience(ParasiteEntity entity) {
        ServerLevel level = (ServerLevel)entity.level();
        return level.players().stream().filter(player -> HorrorModule.mayTarget(level.getServer(), player))
            .filter(player -> HorrorPacing.get(level.getServer()).available(player.getUUID(), level.getGameTime()))
            .filter(player -> player.distanceToSqr(entity) <= 24 * 24).toList();
    }
    private static boolean deliverText(ParasiteEntity entity, String text, Set<UUID> reserved) {
        List<ServerPlayer> audience = audience(entity).stream().filter(player -> reserved.contains(player.getUUID())).toList();
        Component message = Component.literal("<Huésped> " + text);
        audience.forEach(player -> player.sendSystemMessage(message));
        return !audience.isEmpty();
    }

    // TRANSPORTE: sin redirecciones, credenciales o rutas proporcionadas por Luna.
    private static CompletableFuture<byte[]> request(Settings settings, String text) {
        JsonObject payload = new JsonObject(); payload.addProperty("text", text);
        HttpRequest request = HttpRequest.newBuilder(permittedEndpoint(settings.endpoint()))
            .timeout(Duration.ofSeconds(settings.timeoutSeconds())).header("Content-Type", "application/json")
            .POST(HttpRequest.BodyPublishers.ofString(JSON.toJson(payload), StandardCharsets.UTF_8)).build();
        var transport = CLIENT.sendAsync(request, info -> new LimitedAudio(settings.timeoutSeconds()));
        CompletableFuture<byte[]> result = new CompletableFuture<>();
        transport.whenComplete((response, error) -> {
            if (error != null) result.completeExceptionally(error);
            else if (response.statusCode() != 200 || !response.headers().firstValue("Content-Type").orElse("").startsWith("audio/ogg")
                    || !validAudio(response.body())) result.completeExceptionally(new IOException("Respuesta de voz inválida"));
            else result.complete(response.body());
        });
        result.orTimeout(settings.timeoutSeconds(), TimeUnit.SECONDS).whenComplete((audio, error) -> {
            if (error != null || result.isCancelled()) transport.cancel(true);
        });
        return result;
    }
    private static final class LimitedAudio implements HttpResponse.BodySubscriber<byte[]> {
        private final CompletableFuture<byte[]> body = new CompletableFuture<>();
        private final ByteArrayOutputStream bytes = new ByteArrayOutputStream();
        private volatile Flow.Subscription subscription;
        private LimitedAudio(int timeoutSeconds) {
            CompletableFuture.delayedExecutor(timeoutSeconds, TimeUnit.SECONDS).execute(() -> {
                if (body.completeExceptionally(new TimeoutException("Tiempo de voz agotado")) && subscription != null) subscription.cancel();
            });
        }
        @Override public CompletionStage<byte[]> getBody() { return body; }
        @Override public void onSubscribe(Flow.Subscription value) {
            subscription = value;
            if (body.isDone()) value.cancel(); else value.request(1);
        }
        @Override public void onNext(List<ByteBuffer> buffers) {
            if (body.isDone()) { subscription.cancel(); return; }
            for (ByteBuffer buffer : buffers) {
                if (bytes.size() + buffer.remaining() > MAX_AUDIO_BYTES) {
                    subscription.cancel(); body.completeExceptionally(new IOException("Audio fuera de límite")); return;
                }
                byte[] chunk = new byte[buffer.remaining()]; buffer.get(chunk); bytes.writeBytes(chunk);
            }
            if (body.isDone()) subscription.cancel(); else subscription.request(1);
        }
        @Override public void onError(Throwable error) { body.completeExceptionally(error); }
        @Override public void onComplete() { body.complete(bytes.toByteArray()); }
    }
}