TelemetryService class

Пакет: com.hypixel.hytale.server.core.telemetry

Файл: com/hypixel/hytale/server/core/telemetry/TelemetryService.java

Поля (61)

МодификаторыТипИмя
final HytaleLogger LOGGER
break
instance
instance
return
return
return
return
return
ServerAuthManager var1
int var1
Universe var1
TelemetryPackets.ServerStopPacket var10
PacketStatsRecorder.PacketStatsEntry var11
long var12
long var13
long var14
long var15
long var16
TelemetryPackets.P2PSessionBandwidthPacket var17
float var17
long var18
TelemetryPackets.ServerStartPacket var2
HttpRequest var2
int var2
var2
TelemetryService.NetworkStats var2
int var2
String var2
long var2
Optional var3
TelemetryPackets.ServerHeartbeatPacket var3
float var3
long var3
HttpRequest var3
HttpRequest var3
String var3
var3
int var4
Universe var4
long var4
long var4
ServerAuthManager var4
int var5
long var5
long var5
HttpResponse var5
Throwable var6
var6
boolean var7
var7
boolean var7
var7
var7
var7
long var8
boolean var8
var8
PacketHandler var8
ServerManager var9
PacketStatsRecorder var9

Методы (48)

МодификаторыВозвратСигнатура
record NetworkStatsrecord NetworkStats(float sentBytesPerSecond, float receivedBytesPerSecond, float totalSentKb, float totalReceivedKb)
for for(int var10 = 0; var10 < 512; var10++)
for for(int var4 = 0; var4 < 3; var4++)
private int getNextSequenceint getNextSequence()
private void handleAsyncErrorvoid handleAsyncError(@Nonnull Throwable var1, @Nonnull HttpRequest var2, @Nonnull String var3, int var4, @Nullable Consumer<HttpResponse<String>> var5)
if if(!var1)
if if(var2 == null)
if if(var1 != null)
if if(var4 > 0)
if if(var4 != null)
if if(var5 > this.peakPlayers)
if if(var8 >= 0L)
if if(var2 != null)
if if(var9 != null)
if if(var12 < 0L || var14 < 0L)
if if(var1 != null)
if if(var9 != null)
if if(var16 > 0L)
if if(var1 < this.minPerformance)
if if(var1 < 50.0F)
if if(this.heartbeatFuture != null)
if if(var3 != null)
if if(var6 != null)
if if(var4 != null)
if if(var6 instanceof HttpTimeoutException)
if if(var7 && var4 < 2)
if if(var3 != null)
if if(var4 < 2)
boolean isEnabledboolean isEnabled()
public void onAuthenticatedvoid onAuthenticated()
private void onServerStartCompletevoid onServerStartComplete(@Nullable HttpResponse<String> var1)
public void recordP2PAccessvoid recordP2PAccess(@Nonnull String var1)
void recordPlayerConnectvoid recordPlayerConnect(@Nonnull UUID var1)
private void scheduleRetryvoid scheduleRetry(@Nonnull HttpRequest var1, @Nonnull String var2, int var3, @Nullable Consumer<HttpResponse<String>> var4)
private void sendHeartbeatvoid sendHeartbeat()
private void sendP2PSessionBandwidthvoid sendP2PSessionBandwidth(int var1)
private void sendPacketAsyncvoid sendPacketAsync(@Nonnull Object var1, @Nonnull String var2)
private void sendPacketSyncvoid sendPacketSync(@Nonnull Object var1, @Nonnull String var2)
void sendServerStartvoid sendServerStart()
private void sendServerStartPacketvoid sendServerStartPacket(@Nonnull Object var1)
void sendServerStopvoid sendServerStop(@Nonnull String var1)
private void sendWithRetryAsyncvoid sendWithRetryAsync(@Nonnull HttpRequest var1, @Nonnull String var2, int var3, @Nullable Consumer<HttpResponse<String>> var4)
void setEnabledvoid setEnabled(boolean var1)
void shutdownvoid shutdown()
private void startHeartbeatvoid startHeartbeat(int var1)
private void stopHeartbeatvoid stopHeartbeat()
private void trackPerformancevoid trackPerformance(float var1)
private void updateNetworkStatsvoid updateNetworkStats()

Исходный код

Показать/скрыть
class="kw">package com.hypixel.hytale.server.core.telemetry;

class="kw">import com.hypixel.hytale.logger.HytaleLogger;
class="kw">import com.hypixel.hytale.protocol.io.PacketStatsRecorder;
class="kw">import com.hypixel.hytale.protocol.io.ServerListener;
class="kw">import com.hypixel.hytale.server.core.HytaleServer;
class="kw">import com.hypixel.hytale.server.core.Options;
class="kw">import com.hypixel.hytale.server.core.auth.AuthConfig;
class="kw">import com.hypixel.hytale.server.core.auth.ServerAuthManager;
class="kw">import com.hypixel.hytale.server.core.io.PacketHandler;
class="kw">import com.hypixel.hytale.server.core.io.ServerManager;
class="kw">import com.hypixel.hytale.server.core.io.transport.TransportType;
class="kw">import com.hypixel.hytale.server.core.universe.PlayerRef;
class="kw">import com.hypixel.hytale.server.core.universe.Universe;
class="kw">import java.io.IOException;
class="kw">import java.lang.management.GarbageCollectorMXBean;
class="kw">import java.lang.management.ManagementFactory;
class="kw">import java.net.URI;
class="kw">import java.net.http.HttpClient;
class="kw">import java.net.http.HttpRequest;
class="kw">import java.net.http.HttpResponse;
class="kw">import java.net.http.HttpTimeoutException;
class="kw">import java.net.http.HttpRequest.BodyPublishers;
class="kw">import java.net.http.HttpResponse.BodyHandlers;
class="kw">import java.time.Duration;
class="kw">import java.time.Instant;
class="kw">import java.util.Optional;
class="kw">import java.util.Set;
class="kw">import java.util.UUID;
class="kw">import java.util.concurrent.CompletionException;
class="kw">import java.util.concurrent.ConcurrentHashMap;
class="kw">import java.util.concurrent.ScheduledFuture;
class="kw">import java.util.concurrent.TimeUnit;
class="kw">import java.util.concurrent.atomic.AtomicBoolean;
class="kw">import java.util.concurrent.atomic.AtomicInteger;
class="kw">import java.util.concurrent.atomic.AtomicLong;
class="kw">import java.util.function.Consumer;
class="kw">import java.util.logging.Level;
class="kw">import javax.annotation.Nonnull;
class="kw">import javax.annotation.Nullable;

class="kw">public class="kw">final class TelemetryService {
   class="kw">private class="kw">static class="kw">final HytaleLogger LOGGER = HytaleLogger.forEnclosingClass();
   class="kw">private class="kw">static class="kw">final int DEFAULT_HEARTBEAT_INTERVAL_SECONDS = 300;
   class="kw">private class="kw">static class="kw">final String HEARTBEAT_INTERVAL_HEADER = "X-Telemetry-Heartbeat-Interval";
   class="kw">private class="kw">static class="kw">final int MAX_RETRY_ATTEMPTS = 3;
   class="kw">private class="kw">static class="kw">final int RETRY_DELAY_MS = 2500;
   class="kw">private class="kw">static class="kw">final Duration REQUEST_TIMEOUT = Duration.ofSeconds(10L);
   @Nullable
   class="kw">private class="kw">static TelemetryService instance;
   @Nonnull
   class="kw">private class="kw">final String sessionId;
   @Nonnull
   class="kw">private class="kw">final Instant bootTime;
   @Nonnull
   class="kw">private class="kw">final HttpClient httpClient;
   @Nonnull
   class="kw">private class="kw">final AtomicInteger sequenceNumber = new AtomicInteger(0);
   @Nonnull
   class="kw">private class="kw">final AtomicBoolean sessionStarted = new AtomicBoolean(false);
   @Nonnull
   class="kw">private class="kw">final AtomicBoolean pendingServerStart = new AtomicBoolean(false);
   @Nonnull
   class="kw">private class="kw">final AtomicBoolean enabled = new AtomicBoolean(true);
   @Nullable
   class="kw">private ScheduledFuture<?> heartbeatFuture;
   @Nonnull
   class="kw">private class="kw">final AtomicLong totalBytesSent = new AtomicLong();
   @Nonnull
   class="kw">private class="kw">final AtomicLong totalBytesReceived = new AtomicLong();
   @Nonnull
   class="kw">private class="kw">final AtomicLong lastBytesSent = new AtomicLong();
   @Nonnull
   class="kw">private class="kw">final AtomicLong lastBytesReceived = new AtomicLong();
   @Nonnull
   class="kw">private class="kw">final AtomicLong lastNetworkUpdateTime = new AtomicLong(System.nanoTime());
   class="kw">private class="kw">volatile float sentBytesPerSecond;
   class="kw">private class="kw">volatile float receivedBytesPerSecond;
   class="kw">private int peakPlayers;
   @Nonnull
   class="kw">private class="kw">final Set<UUID> uniquePlayers = ConcurrentHashMap.newKeySet();
   class="kw">private int totalConnections;
   @Nullable
   class="kw">private class="kw">volatile String p2pAccessLevel;
   class="kw">private float performanceSum;
   class="kw">private int performanceCount;
   class="kw">private float minPerformance = Float.MAX_VALUE;
   class="kw">private int performanceDropsBelow50;
   @Nonnull
   class="kw">private class="kw">volatile String shutdownReason = "normal";
   @Nonnull
   class="kw">private class="kw">final TelemetryStorage storage;

   TelemetryService() {
      this.sessionId = UUID.randomUUID().toString();
      this.bootTime = Instant.now();
      this.httpClient = HttpClient.newBuilder().connectTimeout(REQUEST_TIMEOUT).build();
      this.storage = new TelemetryStorage(this.sessionId);
      instance = this;
      LOGGER.at(Level.INFO).log("Telemetry service initialized with session ID: %s", this.sessionId);
   }

   @Nullable
   class="kw">public class="kw">static TelemetryService get() {
      class="kw">return instance;
   }

   boolean isEnabled() {
      class="kw">return this.enabled.get();
   }

   void setEnabled(boolean var1) {
      this.enabled.set(var1);
      if (!var1) {
         this.stopHeartbeat();
      }

      LOGGER.at(Level.INFO).log("Telemetry %s", var1 ? "enabled" : "disabled");
   }

   class="kw">public void onAuthenticated() {
      if (this.pendingServerStart.get()) {
         LOGGER.at(Level.FINE).log("Server authenticated, sending pending server_start telemetry");
         this.sendServerStart();
      }
   }

   void sendServerStart() {
      if (this.enabled.get() && !this.sessionStarted.get()) {
         ServerAuthManager var1 = ServerAuthManager.getInstance();
         if (!var1.hasSessionToken()) {
            this.pendingServerStart.set(true);
            LOGGER.at(Level.FINE).log("Server not yet authenticated, telemetry will start after authentication");
         } else if (this.sessionStarted.compareAndSet(false, true)) {
            this.pendingServerStart.set(false);
            TelemetryPackets.ServerStartPacket var2 = TelemetryDataCollector.collectServerStart(this.sessionId, getUtcTimestamp(), this.getNextSequence());
            this.sendServerStartPacket(var2);
         }
      }
   }

   class="kw">private void sendServerStartPacket(@Nonnull Object var1) {
      HttpRequest var2 = this.buildRequest(var1, "server_start");
      if (var2 == null) {
         this.startHeartbeat(300);
      } else {
         this.sendWithRetryAsync(var2, "server_start", 0, this::onServerStartComplete);
      }
   }

   class="kw">private void onServerStartComplete(@Nullable HttpResponse<String> var1) {
      int var2 = 300;
      if (var1 != null) {
         Optional var3 = var1.headers().firstValue("X-Telemetry-Heartbeat-Interval");
         if (var3.isPresent()) {
            try {
               int var4 = Integer.parseInt(var3.get());
               if (var4 > 0) {
                  var2 = var4;
                  LOGGER.at(Level.FINE).log("Heartbeat interval set to %d seconds from server response", var2);
               }
            } catch (NumberFormatException var5) {
               LOGGER.at(Level.WARNING).log("Invalid heartbeat interval header value: %s", var3.get());
            }
         }
      }

      this.startHeartbeat(var2);
   }

   class="kw">private void sendHeartbeat() {
      if (this.enabled.get() && this.sessionStarted.get()) {
         this.updateNetworkStats();
         int var1 = (int)Duration.between(this.bootTime, Instant.now()).toSeconds();
         TelemetryService.NetworkStats var2 = new TelemetryService.NetworkStats(
            this.sentBytesPerSecond, this.receivedBytesPerSecond, (float)this.totalBytesSent.get() / 1024.0F, (float)this.totalBytesReceived.get() / 1024.0F
         );
         TelemetryPackets.ServerHeartbeatPacket var3 = TelemetryDataCollector.collectHeartbeat(
            this.sessionId, getUtcTimestamp(), this.getNextSequence(), var1, var2
         );
         this.trackPerformance(var3.performance().tickPerformancePercent());
         Universe var4 = Universe.get();
         if (var4 != null) {
            int var5 = var4.getPlayerCount();
            if (var5 > this.peakPlayers) {
               this.peakPlayers = var5;
            }
         }

         this.sendPacketAsync(var3, "heartbeat");
      }
   }

   void sendServerStop(@Nonnull String var1) {
      if (this.enabled.get() && this.sessionStarted.get()) {
         this.shutdownReason = var1;
         this.stopHeartbeat();
         int var2 = (int)Duration.between(this.bootTime, Instant.now()).toSeconds();
         float var3 = this.performanceCount > 0 ? this.performanceSum / this.performanceCount : 0.0F;
         long var4 = 0L;

         for (GarbageCollectorMXBean var7 : ManagementFactory.getGarbageCollectorMXBeans()) {
            long var8 = var7.getCollectionTime();
            if (var8 >= 0L) {
               var4 += var8;
            }
         }

         TelemetryPackets.ServerStopPacket var10 = new TelemetryPackets.ServerStopPacket(
            this.sessionId,
            getUtcTimestamp(),
            this.getNextSequence(),
            new TelemetryPackets.ServerSessionSummary(var2, this.shutdownReason, this.peakPlayers, this.uniquePlayers.size()),
            new TelemetryPackets.ServerPerformanceSummary(
               var3, this.minPerformance == Float.MAX_VALUE ? 0.0F : this.minPerformance, this.performanceDropsBelow50, var4
            ),
            new TelemetryPackets.ServerNetworkSummary(
               (float)this.totalBytesSent.get() / 1048576.0F, (float)this.totalBytesReceived.get() / 1048576.0F, this.totalConnections
            )
         );
         this.sendPacketSync(var10, "server_stop");
         this.sendP2PSessionBandwidth(var2);
         this.sessionStarted.set(false);
         LOGGER.at(Level.FINE).log("Telemetry session ended. Uptime: %ds", var2);
      }
   }

   void recordPlayerConnect(@Nonnull UUID var1) {
      this.uniquePlayers.add(var1);
      this.totalConnections++;
   }

   class="kw">public void recordP2PAccess(@Nonnull String var1) {
      this.p2pAccessLevel = var1;
   }

   class="kw">private void sendP2PSessionBandwidth(int var1) {
      String var2 = this.p2pAccessLevel;
      if (var2 != null) {
         long var3 = 0L;
         long var5 = 0L;
         boolean var7 = true;
         boolean var8 = false;
         ServerManager var9 = ServerManager.get();
         if (var9 != null) {
            for (ServerListener var11 : var9.getListeners()) {
               var8 = true;
               long var12 = var11.wireBytesSent();
               long var14 = var11.wireBytesReceived();
               if (var12 < 0L || var14 < 0L) {
                  var7 = false;
                  break;
               }

               var3 += var12;
               var5 += var14;
            }
         }

         var7 &= var8;
         TelemetryPackets.P2PSessionBandwidthPacket var17 = new TelemetryPackets.P2PSessionBandwidthPacket(
            this.sessionId,
            getUtcTimestamp(),
            this.getNextSequence(),
            var2,
            ((TransportType)Options.getOptionSet().valueOf(Options.TRANSPORT)).toString(),
            var1,
            this.peakPlayers,
            var7 ? var3 : -1L,
            var7 ? var5 : -1L,
            this.totalBytesSent.get(),
            this.totalBytesReceived.get()
         );
         this.sendPacketSync(var17, "p2p_session_bandwidth");
      }
   }

   class="kw">private void updateNetworkStats() {
      Universe var1 = Universe.get();
      if (var1 != null) {
         long var2 = 0L;
         long var4 = 0L;

         for (PlayerRef var7 : var1.getPlayers()) {
            PacketHandler var8 = var7.getPacketHandler();
            PacketStatsRecorder var9 = var8.getPacketStatsRecorder();
            if (var9 != null) {
               for (int var10 = 0; var10 < 512; var10++) {
                  PacketStatsRecorder.PacketStatsEntry var11 = var9.getEntry(var10);
                  if (var11.hasData()) {
                     var2 += var11.getSentCompressedTotal();
                     var4 += var11.getReceivedCompressedTotal();
                  }
               }
            }
         }

         long var15 = System.nanoTime();
         long var16 = var15 - this.lastNetworkUpdateTime.get();
         if (var16 > 0L) {
            float var17 = (float)var16 / 1.0E9F;
            long var18 = var2 - this.lastBytesSent.get();
            long var13 = var4 - this.lastBytesReceived.get();
            this.sentBytesPerSecond = (float)var18 / var17;
            this.receivedBytesPerSecond = (float)var13 / var17;
            this.lastNetworkUpdateTime.set(var15);
         }

         this.lastBytesSent.set(var2);
         this.lastBytesReceived.set(var4);
         this.totalBytesSent.set(var2);
         this.totalBytesReceived.set(var4);
      }
   }

   class="kw">private void trackPerformance(float var1) {
      this.performanceSum += var1;
      this.performanceCount++;
      if (var1 < this.minPerformance) {
         this.minPerformance = var1;
      }

      if (var1 < 50.0F) {
         this.performanceDropsBelow50++;
      }
   }

   class="kw">private void startHeartbeat(int var1) {
      this.stopHeartbeat();
      LOGGER.at(Level.FINE).log("Starting telemetry heartbeat with interval: %d seconds", var1);
      this.heartbeatFuture = HytaleServer.SCHEDULED_EXECUTOR.scheduleWithFixedDelay(this::sendHeartbeat, var1, var1, TimeUnit.SECONDS);
   }

   class="kw">private void stopHeartbeat() {
      if (this.heartbeatFuture != null) {
         this.heartbeatFuture.cancel(false);
         this.heartbeatFuture = null;
      }
   }

   class="kw">private int getNextSequence() {
      class="kw">return this.sequenceNumber.getAndIncrement();
   }

   @Nonnull
   class="kw">private class="kw">static String getUtcTimestamp() {
      class="kw">return Instant.now().toString();
   }

   class="kw">private void sendPacketAsync(@Nonnull Object var1, @Nonnull String var2) {
      HttpRequest var3 = this.buildRequest(var1, var2);
      if (var3 != null) {
         this.sendWithRetryAsync(var3, var2, 0, null);
      }
   }

   class="kw">private void sendWithRetryAsync(@Nonnull HttpRequest var1, @Nonnull String var2, int var3, @Nullable Consumer<HttpResponse<String>> var4) {
      this.httpClient.sendAsync(var1, BodyHandlers.ofString()).whenComplete((var5, var6) -> {
         if (var6 != null) {
            this.handleAsyncError(var6, var1, var2, var3, var4);
         } else if (var5.statusCode() >= 200 && var5.statusCode() < 300) {
            if (var4 != null) {
               var4.accept((HttpResponse<String>)var5);
            }
         } else {
            LOGGER.at(Level.WARNING).log("Telemetry request failed: %s (status %d: %s)", var2, var5.statusCode(), var5.body());
            if (var5.statusCode() >= 500 && var3 < 2) {
               this.scheduleRetry(var1, var2, var3, var4);
            } else if (var4 != null) {
               var4.accept(null);
            }
         }
      });
   }

   class="kw">private void handleAsyncError(
      @Nonnull Throwable var1, @Nonnull HttpRequest var2, @Nonnull String var3, int var4, @Nullable Consumer<HttpResponse<String>> var5
   ) {
      Throwable var6 = var1;
      if (var1 class="kw">instanceof CompletionException && var1.getCause() != null) {
         var6 = var1.getCause();
      }

      boolean var7;
      if (var6 class="kw">instanceof HttpTimeoutException) {
         LOGGER.at(Level.WARNING).log("Telemetry %s timed out (attempt %d/%d)", var3, var4 + 1, 3);
         var7 = true;
      } else if (var6 class="kw">instanceof IOException) {
         LOGGER.at(Level.WARNING).log("Telemetry %s failed (attempt %d/%d): %s: %s", var3, var4 + 1, 3, var6.getClass().getSimpleName(), var6.getMessage());
         var7 = true;
      } else {
         LOGGER.at(Level.WARNING).log("Telemetry %s failed with unexpected error: %s", var3, var6.getClass().getSimpleName() + ": " + var6.getMessage());
         var7 = false;
      }

      if (var7 && var4 < 2) {
         this.scheduleRetry(var2, var3, var4, var5);
      } else if (var5 != null) {
         var5.accept(null);
      }
   }

   class="kw">private void scheduleRetry(@Nonnull HttpRequest var1, @Nonnull String var2, int var3, @Nullable Consumer<HttpResponse<String>> var4) {
      long var5 = 2500L * (1L << var3);
      HytaleServer.SCHEDULED_EXECUTOR.schedule(() -> this.sendWithRetryAsync(var1, var2, var3 + 1, var4), var5, TimeUnit.MILLISECONDS);
   }

   class="kw">private void sendPacketSync(@Nonnull Object var1, @Nonnull String var2) {
      HttpRequest var3 = this.buildRequest(var1, var2);
      if (var3 != null) {
         for (int var4 = 0; var4 < 3; var4++) {
            try {
               HttpResponse var5 = this.httpClient.send(var3, BodyHandlers.ofString());
               if (var5.statusCode() >= 200 && var5.statusCode() < 300) {
                  class="kw">return;
               }

               LOGGER.at(Level.WARNING).log("Telemetry request failed: %s (status %d: %s)", var2, var5.statusCode(), var5.body());
               if (var5.statusCode() < 500) {
                  class="kw">return;
               }
            } catch (HttpTimeoutException var7) {
               LOGGER.at(Level.WARNING).log("Telemetry request timed out (attempt %d/%d)", var4 + 1, 3);
            } catch (IOException var8) {
               LOGGER.at(Level.WARNING).log("Telemetry request failed (attempt %d/%d): %s", var4 + 1, 3, var8.getMessage());
            } catch (InterruptedException var9) {
               Thread.currentThread().interrupt();
               class="kw">return;
            } catch (Exception var10) {
               ((HytaleLogger.Api)LOGGER.at(Level.WARNING).withCause(var10)).log("Unexpected error sending telemetry");
               class="kw">return;
            }

            if (var4 < 2) {
               try {
                  Thread.sleep(Duration.ofMillis(2500L * (1L << var4)));
               } catch (InterruptedException var6) {
                  Thread.currentThread().interrupt();
                  class="kw">return;
               }
            }
         }

         LOGGER.at(Level.WARNING).log("Telemetry packet %s failed after %d attempts", var2, 3);
      }
   }

   void shutdown() {
      this.stopHeartbeat();
      this.storage.closeAndCompress();
      instance = null;
   }

   @Nullable
   class="kw">private HttpRequest buildRequest(@Nonnull Object var1, @Nonnull String var2) {
      String var3;
      try {
         var3 = TelemetryJsonSerializer.serialize(var1);
      } catch (Exception var5) {
         ((HytaleLogger.Api)LOGGER.at(Level.WARNING).withCause(var5)).log("Failed to serialize telemetry packet: %s", var2);
         class="kw">return null;
      }

      this.storage.writeEntry(var3);
      ServerAuthManager var4 = ServerAuthManager.getInstance();
      class="kw">return !var4.hasSessionToken()
         ? null
         : HttpRequest.newBuilder()
            .uri(URI.create("https://telemetry.hytale.com/telemetry/server"))
            .header("Content-Type", "application/json")
            .header("Accept", "application/json")
            .header("User-Agent", AuthConfig.USER_AGENT)
            .header("Authorization", "Bearer " + var4.getSessionToken())
            .timeout(REQUEST_TIMEOUT)
            .POST(BodyPublishers.ofString(var3))
            .build();
   }

   record NetworkStats(float sentBytesPerSecond, float receivedBytesPerSecond, float totalSentKb, float totalReceivedKb) {
      NetworkStats {
      }
   }
}