> ## 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.

# Mesurer la latence

> Découvrez comment mesurer et analyser correctement la latence pour les flux gRPC en utilisant diverses méthodes de test.

<Warning>
  **AVERTISSEMENT CRITIQUE : Centre de données de production uniquement**

  Ces tests de latence sont conçus uniquement pour **les environnements de centre de données de production**. NE PAS effectuer ces tests sur des machines locales ou des connexions internet grand public. La bande passante locale ne peut pas gérer les abonnements lourds de Solana et produira des résultats sans signification qui ne reflètent pas les performances réelles.
</Warning>

<Warning>
  **EXIGENCE DE CO-LOCATION : Déployez près de votre point de terminaison LaserStream**

  Pour des mesures de latence significatives, vous **devez** co-localiser votre infrastructure de test dans la même région que le point de terminaison LaserStream choisi. La distance réseau dominera vos mesures - tester depuis un autre continent montrera la latence du réseau, pas la performance de LaserStream.
</Warning>

***

## Comprendre la latence dans les systèmes blockchain distribués

Lors de la manipulation de services de streaming blockchain, la mesure de la latence devient complexe car les systèmes distribués n'ont pas d'horloge universelle. Contrairement aux systèmes traditionnels où vous pouvez mesurer le temps d'aller-retour vers un seul serveur, les réseaux blockchain impliquent plusieurs validateurs, chacun recevant et traitant la même transaction à des moments différents.

**Le défi fondamental :** Les blockchains comme Solana n'ont pas de concept de temps absolu. Chaque nœud valideur recevra globalement la même transaction à des moments différents, et la confirmation dépend d'un pourcentage du cluster atteignant un consensus. Cela rend la mesure déterministe de la latence impossible au sens traditionnel.

## Niveaux d'engagement et priorités de latence

Solana offre trois niveaux d'engagement, chacun avec des caractéristiques de latence différentes :

* **Traité** : Le plus rapide, confirmation d'un seul validateur (\~400ms)
* **Confirmé** : Moyen, confirmation de supermajorité (\~2-3 secondes)
* **Finalisé** : Le plus lent, finalisation complète du réseau (\~15-30 secondes)

Pour les applications sensibles à la latence, l'**engagement traité** est généralement l'objectif. Tous les tests de ce guide utilisent le niveau d'engagement traité puisque la plupart des cas d'utilisation à haute fréquence privilégient la vitesse sur la finalité absolue.

## Trois approches pour mesurer la latence

### 1. Comparer les flux gRPC parallèles

**Méthode la plus fiable** - Compare deux flux indépendants vers la même source de données, mesurant lequel reçoit d'abord des événements identiques.

**Avantages :**

* Élimine les problèmes de synchronisation d'horloge
* Fournit une comparaison relative des performances
* Le plus précis pour comparer les services

### 2. Comparaison de l'horodatage local vs created\_at

**Fiabilité modérée** - Mesure la différence entre le moment où votre système reçoit un message et l'horodatage intégré dans le message par le service LaserStream.

**Limitations :**

* Représente seulement le moment où LaserStream a créé le message en interne
* Les retards en amont vers LaserStream ne seront pas capturés
* Moins précis que la Méthode 1 pour une latence de bout en bout réelle

### 3. Analyse de l'horodatage des blocs (non recommandé)

**Non recommandé** - Compare le temps de réception local contre l'horodatage du bloc de Solana.

**Limitations importantes :**

* Les horodatages des blocs n'ont qu'une granularité au niveau de la seconde
* Solana produit des blocs toutes les 400ms
* Fournit peu d'informations utiles

***

## Exigences de configuration

### Co-location régionale

Pour des mesures de latence significatives, déployez votre infrastructure de test dans le même centre de données ou région que votre point de terminaison LaserStream.

**Régions LaserStream disponibles :**

* **ewr** : New York, US (côte Est) - `https://laserstream-mainnet-ewr.helius-rpc.com`
* **pitt** : Pittsburgh, US (central) - `https://laserstream-mainnet-pitt.helius-rpc.com`
* **slc** : Salt Lake City, US (côte Ouest) - `https://laserstream-mainnet-slc.helius-rpc.com`
* **ams** : Amsterdam, Europe - `https://laserstream-mainnet-ams.helius-rpc.com`
* **fra** : Francfort, Europe - `https://laserstream-mainnet-fra.helius-rpc.com`
* **tyo** : Tokyo, Asie - `https://laserstream-mainnet-tyo.helius-rpc.com`
* **sgp** : Singapour, Asie - `https://laserstream-mainnet-sgp.helius-rpc.com`

Pour les tests sur devnet, utilisez : `https://laserstream-devnet-ewr.helius-rpc.com`

Consultez la [documentation gRPC de LaserStream](/docs/fr/laserstream/grpc) pour des instructions complètes de configuration et des directives de sélection de point de terminaison.

### Configuration de l'environnement Rust

Tous les scripts de mesure utilisent Rust avec Cargo. Configuration de base :

```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"
```

Créez un fichier `.env` avec vos identifiants :

```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
```

Obtenez votre clé API Helius depuis le [Tableau de bord Helius](https://dashboard.helius.dev/). LaserStream devnet est disponible sur tous les plans. L'accès au réseau principal nécessite un plan Business ou Professional.

***

## Méthode 1 : Comparaison des flux parallèles

Ce script établit deux connexions indépendantes à différents points de terminaison gRPC et mesure lequel reçoit d'abord les mêmes messages `BlockMeta`. Cette approche élimine les problèmes de synchronisation d'horloge en utilisant le timing relatif.

```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
}
```

**Ce que cela mesure :** La différence de performance relative entre deux services de streaming. Le delta montre quel service fournit en premier les mêmes informations de slot.

**Principales métriques :**

* **Delta positif** : Premier service (YS) plus lent que le deuxième service (LS) - LaserStream est plus rapide
* **Delta négatif** : Premier service (YS) plus rapide que le deuxième service (LS) - LaserStream est plus lent
* **Moyenne/Médiane** : Différence de performance moyenne
* **P95** : Différence de latence au 95e percentile

**Exécution du test :**

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

**Exemple de sortie :**

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

La sortie montre les différences de latence en temps réel et les statistiques périodiques. Un delta moyen positif indique que le deuxième service (LaserStream) fournit constamment les données plus rapidement.

***

## Méthode 2 : Analyse des horodatages créés

Cette approche compare l'horodatage `created_at` intégré dans les messages par rapport à votre heure système locale lors de leur réception.

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

**Limitation importante :** Cette méthode ne mesure que depuis le moment où LaserStream a créé le message jusqu'à ce que vous l'ayez reçu. Elle ne tient pas compte des retards en amont entre l'événement blockchain et le traitement par LaserStream.

**Exécution du test :**

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

**Exemple de sortie :**

```
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
---
```

Cette méthode fournit des informations sur la latence du réseau et de traitement entre LaserStream et votre application, mais devrait être utilisée en conjonction avec la Méthode 1 pour une analyse complète.

***

## Meilleures pratiques pour le test de latence

### Principes clés

* **Co-location** : Déployez les tests dans la même région que votre point de terminaison LaserStream pour minimiser la latence réseau
* **Méthodes multiples** : Utilisez la comparaison de flux parallèles (Méthode 1) comme métrique principale, complétée par l'analyse des horodatages
* **Surveillance à long terme** : Effectuez des tests sur de longues périodes pour capturer différentes conditions de réseau et congestions blockchain
* **Analyse statistique** : Concentrez-vous sur les percentiles (P95, P99) plutôt que juste les moyennes pour comprendre la latence de la traîne

### Interprétation des résultats

1. **Établir une base de référence** : Réalisez des tests pendant au moins 1 heure pour établir les performances de base dans des conditions normales
2. **Identifier les motifs** : Recherchez des motifs dans les pics de latence - correspondent-ils à une activité blockchain élevée ou à une congestion réseau ?
3. **Comparer les percentiles** : La latence P95 est souvent plus importante que la latence moyenne pour l'expérience utilisateur
4. **Surveiller la cohérence** : Des performances cohérentes sont souvent plus précieuses qu'une latence minimale absolue

Rappelez-vous que la latence blockchain est intrinsèquement variable en raison des exigences de consensus du réseau. Concentrez-vous sur les différences de performance relative et la cohérence plutôt que sur les chiffres absolus.
