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 SETTINGS = new WeakHashMap<>(); private static final Map> LAST = new WeakHashMap<>(); private static final Map>> REQUESTS = new WeakHashMap<>(); private static final Set STOPPING = Collections.synchronizedSet(Collections.newSetFromMap(new WeakHashMap<>())); private static final Set 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 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 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 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 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 result = new CompletableFuture<>(); REQUESTS.computeIfAbsent(server, key -> new HashMap<>()).put(bodyId, result); CompletableFuture 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 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 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 reserved) { List audience = audience(entity).stream().filter(player -> reserved.contains(player.getUUID())).toList(); Component message = Component.literal(" " + text); audience.forEach(player -> player.sendSystemMessage(message)); return !audience.isEmpty(); } // TRANSPORTE: sin redirecciones, credenciales o rutas proporcionadas por Luna. private static CompletableFuture 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 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 { private final CompletableFuture 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 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 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()); } } }