> ## Documentation Index
> Fetch the complete documentation index at: https://www.helius.dev/docs/llms.txt
> Use this file to discover all available pages before exploring further.

# Latenzmessung

> Erfahren Sie, wie Sie die Latenz für gRPC-Streams korrekt messen und analysieren, indem Sie verschiedene Testmethoden verwenden.

<Warning>
  **WICHTIGER HAFTUNGSAUSSCHLUSS: Nur für produktive Rechenzentren**

  Diese Latenztests sind nur für **produktive Rechenzentrumsumgebungen** konzipiert. Führen Sie diese Tests NICHT auf lokalen Maschinen oder mit Internetverbindungen für Verbraucher durch. Lokale Bandbreiten können schwere Solana-Abonnements nicht bewältigen und liefern bedeutungslose Ergebnisse, die keine realen Leistungsdaten widerspiegeln.
</Warning>

<Warning>
  **ANFORDERUNG DER KO-LOKATION: In der Nähe Ihres LaserStream-Endpunkts bereitstellen**

  Für sinnvolle Latenzmessungen **müssen** Sie Ihre Testinfrastruktur in derselben Region wie Ihr gewählter LaserStream-Endpunkt bereitstellen. Die Netzwerkentfernung wird Ihre Messungen dominieren - Tests von einem anderen Kontinent aus zeigen die Netzwerklatenz und nicht die Leistung von LaserStream.
</Warning>

***

## Verständnis der Latenz in verteilten Blockchain-Systemen

Beim Arbeiten mit Blockchain-Streaming-Diensten wird die Latenzmessung komplex, da verteilte Systeme keine universelle Uhr haben. Im Gegensatz zu traditionellen Systemen, bei denen Sie die Round-Trip-Zeit zu einem einzelnen Server messen können, beinhalten Blockchain-Netzwerke mehrere Validatoren, die dieselbe Transaktion zu unterschiedlichen Zeiten empfangen und verarbeiten.

**Die grundlegende Herausforderung:** Blockchains wie Solana haben kein Konzept der absoluten Zeit. Jeder Validator-Knoten erhält weltweit dieselbe Transaktion zu unterschiedlichen Zeiten, und die Bestätigung hängt davon ab, dass ein Prozentsatz des Clusters Konsens erreicht. Dies macht eine deterministische Latenzmessung im traditionellen Sinn unmöglich.

## Verpflichtungsstufen und Latenzprioritäten

Solana bietet drei Verpflichtungsstufen, jeweils mit unterschiedlichen Latenzeigenschaften:

* **Verarbeitet**: Schnellste, Bestätigung durch einen einzelnen Validator (\~400ms)
* **Bestätigt**: Mittel, Bestätigung durch eine Supermehrheit (\~2-3 Sekunden)
* **Finalisiert**: Langsamste, komplette Netzwerk-Finalisierung (\~15-30 Sekunden)

Für latenzempfindliche Anwendungen ist **verarbeitete Verpflichtung** typischerweise das Ziel. Alle Tests in diesem Leitfaden verwenden die verarbeitete Verpflichtungsstufe, da die meisten Hochfrequenzanwendungen Geschwindigkeit vor absoluter Endgültigkeit priorisieren.

## Drei Ansätze zur Latenzmessung

### 1. Vergleich paralleler gRPC-Streams

**Zuverlässigste Methode** - Vergleicht zwei unabhängige Streams zur gleichen Datenquelle und misst, welcher identische Ereignisse zuerst empfängt.

**Vorteile:**

* Beseitigt Probleme mit der Uhrensynchronisation
* Bietet einen relativen Leistungsvergleich
* Am genauesten für den Vergleich von Diensten

### 2. Vergleich lokaler Zeitstempel vs. erstellt\_am

**Mäßige Zuverlässigkeit** - Misst den Unterschied zwischen dem Empfang einer Nachricht durch Ihr System und dem Zeitstempel, der in der Nachricht durch den LaserStream-Dienst eingebettet ist.

**Einschränkungen:**

* Repräsentiert nur, wann LaserStream die Nachricht intern erstellt hat
* Verzögerungen bis zu LaserStream werden nicht erfasst
* Weniger genau als Methode 1 für echte Ende-zu-Ende-Latenz

### 3. Block-Zeitstempel-Analyse (nicht empfohlen)

**Nicht empfohlen** - Vergleicht lokale Empfangszeit mit Solanas Block-Zeitstempel.

**Bedeutende Einschränkungen:**

* Block-Zeitstempel haben nur Granularität auf Sekundenebene
* Solana erzeugt alle 400ms Blöcke
* Bietet minimale nützliche Informationen

***

## Anforderungen an das Setup

### Regionale Ko-Lokation

Für sinnvolle Latenzmessungen richten Sie Ihre Testinfrastruktur im gleichen Rechenzentrum oder in der gleichen Region wie Ihr LaserStream-Endpunkt ein.

**Verfügbare LaserStream-Regionen:**

* **ewr**: New York, US (Ostküste)
* **pitt**: Pittsburgh, US (Zentral)
* **slc**: Salt Lake City, US (Westküste)
* **ams**: Amsterdam, Europa
* **fra**: Frankfurt, Europa
* **tyo**: Tokio, Asien
* **sgp**: Singapur, Asien

Für Devnet-Tests verwenden: `https://laserstream-devnet-ewr.helius-rpc.com`

Siehe die [LaserStream gRPC-Dokumentation](/docs/de/laserstream/grpc) für vollständige Anweisungen zur Einrichtung und Auswahl von Endpunkten.

### Einrichtung der Rust-Umgebung

Alle Messskripte verwenden Rust mit Cargo. Grundlegende Einrichtung:

```bash theme={"system"}
# Install Rust
curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh

# Setup project
cargo new latency-testing
cd latency-testing

# Add dependencies to Cargo.toml
[dependencies]
# Async runtime
tokio = { version = "1", features = ["full"] }
# Futures utilities for StreamExt, SinkExt
futures = "0.3"
# Environment variable loading
dotenvy = "0.15"
# Logging
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["fmt", "env-filter"] }
# Yellowstone gRPC client and protocol
yellowstone-grpc-client = "13.3.0"
yellowstone-grpc-proto = "12.6.0"
# For timestamp analysis
prost-types = "0.14"
```

Erstellen Sie eine `.env` Datei mit Ihren Anmeldeinformationen:

```bash theme={"system"}
YS_GRPC_URL=your-comparison-endpoint
YS_API_KEY=your-comparison-api-key
LS_GRPC_URL=your-laserstream-endpoint  
LS_API_KEY=your-helius-api-key
```

Holen Sie sich Ihren Helius-API-Schlüssel vom [Helius-Dashboard](https://dashboard.helius.dev/). LaserStream-Devnet ist in allen Tarifen verfügbar. Mainnet-Zugang erfordert einen Business- oder Professional-Plan.

***

## Methode 1: Vergleich paralleler Streams

Dieses Skript stellt zwei unabhängige Verbindungen zu verschiedenen gRPC-Endpunkten her und misst, welches die gleichen `BlockMeta` Nachrichten zuerst erhält. Dieser Ansatz beseitigt Probleme mit der Uhrensynchronisation durch Verwendung relativer Zeitmessung.

```rust [expandable] theme={"system"}
use std::collections::HashMap;
use std::time::SystemTime;

use dotenvy::dotenv;
use futures::{StreamExt, SinkExt};
use tokio::sync::mpsc;
use tracing::{debug, error, info};

use yellowstone_grpc_client::{ClientTlsConfig, GeyserGrpcClient};
use yellowstone_grpc_proto::prelude::{
    subscribe_update::UpdateOneof, CommitmentLevel, SubscribeRequest, 
    SubscribeRequestFilterBlocksMeta, SubscribeUpdate,
};

#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
enum Source {
    Yellowstone,
    Laserstream,
}

#[derive(Default, Debug)]
struct SlotTimings {
    ys_recv_ms: Option<i128>,
    ls_recv_ms: Option<i128>,
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let _ = dotenv();
    tracing_subscriber::fmt().with_env_filter("info").init();

    let ys_url = std::env::var("YS_GRPC_URL").expect("YS_GRPC_URL env variable not set");
    let ls_url = std::env::var("LS_GRPC_URL").expect("LS_GRPC_URL env variable not set");
    let ys_api_key = std::env::var("YS_API_KEY").ok();
    let ls_api_key = std::env::var("LS_API_KEY").ok();

    let commitment = CommitmentLevel::Processed;
    info!(?commitment, "Starting latency comparison");

    // Establish both clients
    let mut ys_client = GeyserGrpcClient::build_from_shared(ys_url.clone())
        .expect("invalid YS url")
        .x_token(ys_api_key.clone())?
        .tls_config(ClientTlsConfig::new().with_native_roots())?
        .max_decoding_message_size(10 * 1024 * 1024)
        .connect()
        .await?;

    let mut ls_client = GeyserGrpcClient::build_from_shared(ls_url.clone())
        .expect("invalid LS url")
        .x_token(ls_api_key.clone())?
        .tls_config(ClientTlsConfig::new().with_native_roots())?
        .max_decoding_message_size(10 * 1024 * 1024)
        .connect()
        .await?;

    let (mut ys_tx, mut ys_rx) = ys_client.subscribe().await?;
    let (mut ls_tx, mut ls_rx) = ls_client.subscribe().await?;
    
    let subscribe_request = SubscribeRequest {
        blocks_meta: {
            let mut m = HashMap::<String, SubscribeRequestFilterBlocksMeta>::new();
            m.insert("all".to_string(), SubscribeRequestFilterBlocksMeta::default());
            m
        },
        commitment: Some(commitment as i32),
        ..Default::default()
    };

    ys_tx.send(subscribe_request.clone()).await?;
    ls_tx.send(subscribe_request).await?;

    let (agg_tx, mut agg_rx) = mpsc::unbounded_channel::<(Source, u64, i128)>();

    // Spawn task for Yellowstone stream
    {
        let agg_tx = agg_tx.clone();
        tokio::spawn(async move {
            while let Some(update_res) = ys_rx.next().await {
                match update_res {
                    Ok(update) => handle_update(Source::Yellowstone, update, &agg_tx).await,
                    Err(e) => {
                        error!(target: "ys", "stream error: {:?}", e);
                        break;
                    }
                }
            }
        });
    }

    // Spawn task for Laserstream stream
    {
        let agg_tx = agg_tx.clone();
        tokio::spawn(async move {
            while let Some(update_res) = ls_rx.next().await {
                match update_res {
                    Ok(update) => handle_update(Source::Laserstream, update, &agg_tx).await,
                    Err(e) => {
                        error!(target: "ls", "stream error: {:?}", e);
                        break;
                    }
                }
            }
        });
    }

    // Aggregator – collect latencies per slot and print once we have both sources
    let mut timings: HashMap<u64, SlotTimings> = HashMap::new();
    let mut deltas: Vec<i128> = Vec::new();
    let mut count = 0;

    println!("slot,ys_recv_ms,ls_recv_ms,delta_ms");

    while let Some((source, slot, latency_ms)) = agg_rx.recv().await {
        let entry = timings.entry(slot).or_default();
        match source {
            Source::Yellowstone => entry.ys_recv_ms = Some(latency_ms),
            Source::Laserstream => entry.ls_recv_ms = Some(latency_ms),
        }

        if let (Some(ys), Some(ls)) = (entry.ys_recv_ms, entry.ls_recv_ms) {
            let delta = ys - ls; // positive => YS arrived later
            println!("{slot},{ys},{ls},{delta}");
            
            deltas.push(delta);
            count += 1;
            
            if count % 100 == 0 {
                print_statistics(&deltas, count);
            }
            
            timings.remove(&slot);
        }
    }

    Ok(())
}

async fn handle_update(
    source: Source,
    update: SubscribeUpdate,
    agg_tx: &mpsc::UnboundedSender<(Source, u64, i128)>,
) {
    if let Some(UpdateOneof::BlockMeta(block_meta)) = update.update_oneof {
        let slot = block_meta.slot;
        let recv_ms = system_time_to_millis(SystemTime::now());
        debug!(?source, slot, recv_ms, "BlockMeta received");
        let _ = agg_tx.send((source, slot, recv_ms));
    }
}

fn print_statistics(deltas: &[i128], count: usize) {
    if deltas.is_empty() {
        return;
    }
    
    let mut sorted_deltas = deltas.to_vec();
    sorted_deltas.sort();
    
    let median = if sorted_deltas.len() % 2 == 0 {
        let mid = sorted_deltas.len() / 2;
        (sorted_deltas[mid - 1] + sorted_deltas[mid]) / 2
    } else {
        sorted_deltas[sorted_deltas.len() / 2]
    };
    
    let min = *sorted_deltas.first().unwrap();
    let max = *sorted_deltas.last().unwrap();
    let sum: i128 = sorted_deltas.iter().sum();
    let mean = sum / sorted_deltas.len() as i128;
    
    let p25_idx = (sorted_deltas.len() as f64 * 0.25) as usize;
    let p75_idx = (sorted_deltas.len() as f64 * 0.75) as usize;
    let p95_idx = (sorted_deltas.len() as f64 * 0.95) as usize;
    
    let p25 = sorted_deltas[p25_idx.min(sorted_deltas.len() - 1)];
    let p75 = sorted_deltas[p75_idx.min(sorted_deltas.len() - 1)];
    let p95 = sorted_deltas[p95_idx.min(sorted_deltas.len() - 1)];
    
    eprintln!("--- Statistics after {} slots ---", count);
    eprintln!("Delta (YS - LS) in milliseconds:");
    eprintln!("  Min: {}, Max: {}", min, max);
    eprintln!("  Mean: {}, Median: {}", mean, median);
    eprintln!("  P25: {}, P75: {}, P95: {}", p25, p75, p95);
    eprintln!("  Positive deltas (YS slower): {}/{} ({:.1}%)", 
              sorted_deltas.iter().filter(|&&x| x > 0).count(),
              sorted_deltas.len(),
              sorted_deltas.iter().filter(|&&x| x > 0).count() as f64 / sorted_deltas.len() as f64 * 100.0);
    eprintln!("---");
}

fn system_time_to_millis(st: SystemTime) -> i128 {
    st.duration_since(SystemTime::UNIX_EPOCH)
        .unwrap()
        .as_millis() as i128
}
```

**Was dies misst:** Den relativen Leistungsunterschied zwischen zwei Streaming-Diensten. Das Delta zeigt, welcher Dienst die gleichen Slot-Informationen zuerst liefert.

**Wichtige Kennzahlen:**

* **Positives Delta**: Erster Dienst (YS) langsamer als zweiter Dienst (LS) - LaserStream ist schneller
* **Negatives Delta**: Erster Dienst (YS) schneller als zweiter Dienst (LS) - LaserStream ist langsamer
* **Mittelwert/Median**: Durchschnittliche Leistungsunterschied
* **P95**: 95. Perzentil Latenzunterschied

**Durchführung des Tests:**

```bash theme={"system"}
cargo run --bin latency-comparison
```

**Beispielausgabe:**

```
slot,ys_recv_ms,ls_recv_ms,delta_ms
352416939,1752168399141,1752168399140,1
352416940,1752168399526,1752168399512,14
352416941,1752168399890,1752168399877,13
```

Die Ausgabe zeigt Echtzeitlatenzunterschiede und periodische Statistiken. Ein positiver mittlerer Delta-Wert zeigt an, dass der zweite Dienst (LaserStream) Daten konsistent schneller liefert.

***

## Methode 2: Analyse des Erstellungs-Zeitstempels

Dieser Ansatz vergleicht den `created_at`-Zeitstempel, der in Nachrichten eingebettet ist, mit der lokalen Systemzeit beim Empfang.

```rust [expandable] theme={"system"}
use std::time::{Duration, SystemTime};
use dotenvy::dotenv;
use tracing::{debug, error, info};
use yellowstone_grpc_proto::prost_types::Timestamp;
use futures::StreamExt;
use futures::SinkExt;
use std::collections::HashMap;

use yellowstone_grpc_client::{ClientTlsConfig, GeyserGrpcClient};
use yellowstone_grpc_proto::prelude::{
    subscribe_update::UpdateOneof, CommitmentLevel, SubscribeRequest,
    SubscribeRequestFilterTransactions, SubscribeUpdate,
};

const ACCOUNTS_INCLUDE: &[&str] = &["BB5dnY55FXS1e1NXqZDwCzgdYJdMCj3B92PU6Q5Fb6DT"];
const COMMITMENT_LEVEL: &str = "processed";

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let _ = dotenv();
    tracing_subscriber::fmt().with_env_filter("info").init();

    let grpc_url = std::env::var("YS_GRPC_URL").expect("GRPC_URL env variable not set");
    let api_key = std::env::var("YS_API_KEY").ok();

    info!("Connecting to {} …", grpc_url);
    debug!(accounts_include = ?ACCOUNTS_INCLUDE, "Subscribing with accountsInclude filter");

    let mut client = GeyserGrpcClient::build_from_shared(grpc_url.clone())
        .expect("invalid URL")
        .x_token(api_key.clone())?
        .tls_config(ClientTlsConfig::new().with_native_roots())?
        .connect()
        .await?;

    let (mut subscribe_tx, mut subscribe_rx) = client.subscribe().await?;

    let mut tx_filter_map = HashMap::new();
    tx_filter_map.insert(
        "latency".to_string(),
        SubscribeRequestFilterTransactions {
            account_include: ACCOUNTS_INCLUDE.iter().map(|s| s.to_string()).collect(),
            vote: Some(false),
            failed: Some(false),
            ..Default::default()
        },
    );

    subscribe_tx
        .send(SubscribeRequest {
            transactions: tx_filter_map,
            commitment: Some(CommitmentLevel::Processed as i32),
            ..Default::default()
        })
        .await?;

    let mut latencies: Vec<i64> = Vec::new();
    let mut count = 0;

    while let Some(update_res) = subscribe_rx.next().await {
        match update_res {
            Ok(update) => {
                if let Some(UpdateOneof::Transaction(tx)) = update.update_oneof {
                    let recv_time = SystemTime::now();
                    
                    if let Some(created_at) = update.created_at {
                        let created_time = SystemTime::UNIX_EPOCH + Duration::new(
                            created_at.seconds as u64,
                            created_at.nanos as u32,
                        );
                        
                        if let Ok(latency) = recv_time.duration_since(created_time) {
                            let latency_ms = latency.as_millis() as i64;
                            latencies.push(latency_ms);
                            count += 1;
                            
                            println!("Transaction latency: {}ms", latency_ms);
                            
                            if count % 100 == 0 {
                                print_statistics(&latencies, count);
                            }
                        }
                    }
                }
            }
            Err(e) => {
                error!("Stream error: {:?}", e);
                break;
            }
        }
    }

    Ok(())
}

fn print_statistics(latencies: &[i64], count: usize) {
    if latencies.is_empty() {
        return;
    }
    
    let mut sorted = latencies.to_vec();
    sorted.sort();
    
    let median = if sorted.len() % 2 == 0 {
        let mid = sorted.len() / 2;
        (sorted[mid - 1] + sorted[mid]) / 2
    } else {
        sorted[sorted.len() / 2]
    };
    
    let min = *sorted.first().unwrap();
    let max = *sorted.last().unwrap();
    let sum: i64 = sorted.iter().sum();
    let mean = sum / sorted.len() as i64;
    
    let p95_idx = (sorted.len() as f64 * 0.95) as usize;
    let p95 = sorted[p95_idx.min(sorted.len() - 1)];
    
    eprintln!("--- Statistics after {} transactions ---", count);
    eprintln!("Latency (created_at to receive) in milliseconds:");
    eprintln!("  Min: {}, Max: {}", min, max);
    eprintln!("  Mean: {}, Median: {}", mean, median);
    eprintln!("  P95: {}", p95);
    eprintln!("---");
}
```

**Wichtige Einschränkung:** Diese Methode misst nur von dem Zeitpunkt, an dem LaserStream die Nachricht erstellte, bis zu dem Zeitpunkt, an dem Sie sie erhielten. Verzögerungen zwischen dem Blockchain-Ereignis und der Verarbeitung durch LaserStream werden nicht berücksichtigt.

**Durchführung des Tests:**

```bash theme={"system"}
cargo run --bin timestamp-analysis
```

**Beispielausgabe:**

```
Transaction latency: 45ms
Transaction latency: 52ms
Transaction latency: 38ms
--- Statistics after 100 transactions ---
Latency (created_at to receive) in milliseconds:
  Min: 28, Max: 89
  Mean: 47, Median: 45
  P95: 72
---
```

Diese Methode bietet Einblicke in die Netzwerk- und Verarbeitungsverzögerung zwischen LaserStream und Ihrer Anwendung, sollte jedoch in Verbindung mit Methode 1 für eine umfassende Analyse verwendet werden.

***

## Beste Praktiken für Latenztests

### Schlüsselprinzipien

* **Ko-Lokation**: Tests in derselben Region wie Ihr LaserStream-Endpunkt bereitstellen, um Netzwerklatenz zu minimieren
* **Mehrere Methoden**: Verwenden Sie den Vergleich paralleler Streams (Methode 1) als Ihre Hauptmetrik, ergänzt durch Zeitstempelanalyse
* **Langzeitüberwachung**: Tests über längere Zeiträume durchführen, um verschiedene Netzwerkbedingungen und Blockchain-Staus zu erfassen
* **Statistische Analyse**: Fokus auf Perzentilen (P95, P99) statt nur Durchschnittswerten, um Tail-Latenz zu verstehen

### Interpretieren der Ergebnisse

1. **Basislinie festlegen**: Tests mindestens 1 Stunde lang durchführen, um die Basisleistung unter normalen Bedingungen festzulegen
2. **Muster identifizieren**: Nach Mustern bei Latenzspitzen suchen - korrelieren sie mit hoher Blockchain-Aktivität oder Netzwerkstau?
3. **Perzentile vergleichen**: P95-Latenz ist oft wichtiger als durchschnittliche Latenz für das Benutzererlebnis
4. **Konsistenz überwachen**: Konsistente Leistung ist oft wertvoller als absolute minimale Latenz

Denken Sie daran, dass die Blockchain-Latenz aufgrund von Netzwerk-Konsenserfordernissen inhärent variabel ist. Konzentrieren Sie sich auf relative Leistungsunterschiede und Konsistenz statt auf absolute Zahlen.
