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

# Medición de la latencia

> Aprende a medir y analizar correctamente la latencia de los streams de gRPC mediante varios métodos de prueba.

<Warning>
  **AVISO CRÍTICO: Solo para centros de datos de producción**

  Estas pruebas de latencia están diseñadas **exclusivamente para entornos de centros de datos de producción**. NO ejecutes estas pruebas en máquinas locales ni con conexiones de Internet residenciales. El ancho de banda local no puede manejar suscripciones intensivas de Solana y producirá resultados sin valor que no reflejan el rendimiento real.
</Warning>

<Warning>
  **REQUISITO DE COUBICACIÓN: Despliega cerca de tu endpoint de LaserStream**

  Para obtener mediciones de latencia útiles, **debes** coubicar tu infraestructura de pruebas en la misma región que el endpoint de LaserStream que elegiste. La distancia de red dominará las mediciones: si haces pruebas desde otro continente, medirás la latencia de la red, no el rendimiento de LaserStream.
</Warning>

***

## Comprender la latencia en sistemas blockchain distribuidos

Cuando trabajas con servicios de streaming de blockchain, medir la latencia se vuelve complejo porque los sistemas distribuidos no tienen un reloj universal. A diferencia de los sistemas tradicionales, donde puedes medir el tiempo de ida y vuelta a un solo servidor, las redes blockchain incluyen varios validadores, y cada uno recibe y procesa la misma transacción en momentos diferentes.

**El desafío fundamental:** Las blockchains como Solana no tienen un concepto de tiempo absoluto. Cada nodo validador del mundo recibe la misma transacción en un momento diferente, y la confirmación depende de que un porcentaje del clúster alcance un consenso. Esto hace imposible medir la latencia de forma determinista en el sentido tradicional.

## Niveles de compromiso y prioridades de latencia

Solana ofrece tres niveles de compromiso, cada uno con distintas características de latencia:

* **Procesado**: El más rápido, confirmación de un solo validador (\~400 ms)
* **Confirmado**: Intermedio, confirmación por supermayoría (\~2-3 segundos)
* **Finalizado**: El más lento, finalización completa de la red (\~15-30 segundos)

Para aplicaciones sensibles a la latencia, el **nivel de compromiso procesado** suele ser el objetivo. Todas las pruebas de esta guía usan el nivel de compromiso procesado, ya que la mayoría de los casos de uso de alta frecuencia priorizan la velocidad sobre la finalidad absoluta.

## Tres enfoques para medir la latencia

### 1. Comparar streams de gRPC en paralelo

**El método más confiable**: compara dos streams independientes de la misma fuente de datos y mide cuál recibe primero los eventos idénticos.

**Ventajas:**

* Elimina los problemas de sincronización de relojes
* Permite comparar el rendimiento relativo
* Es el más preciso para comparar servicios

### 2. Comparar la marca de tiempo local con created\_at

**Confiabilidad moderada**: mide la diferencia entre el momento en que tu sistema recibe un mensaje y la marca de tiempo que el servicio LaserStream incorporó en el mensaje.

**Limitaciones:**

* Solo representa el momento en que LaserStream creó el mensaje internamente
* No captura los retrasos anteriores a LaserStream
* Es menos preciso que el método 1 para medir la latencia real de extremo a extremo

### 3. Análisis de la marca de tiempo del bloque (no recomendado)

**No recomendado**: compara la hora de recepción local con la marca de tiempo del bloque de Solana.

**Limitaciones importantes:**

* Las marcas de tiempo de los bloques solo tienen una granularidad de segundos
* Solana produce bloques cada 400 ms
* Proporciona muy poca información útil

***

## Requisitos de configuración

### Coubicación regional

Para obtener mediciones de latencia útiles, despliega tu infraestructura de pruebas en el mismo centro de datos o región que tu endpoint de LaserStream.

**Regiones disponibles de LaserStream:**

* **ewr**: Nueva York, EE. UU. (costa este) - `https://laserstream-mainnet-ewr.helius-rpc.com`
* **pitt**: Pittsburgh, EE. UU. (centro) - `https://laserstream-mainnet-pitt.helius-rpc.com`
* **slc**: Salt Lake City, EE. UU. (costa oeste) - `https://laserstream-mainnet-slc.helius-rpc.com`
* **ams**: Ámsterdam, Europa - `https://laserstream-mainnet-ams.helius-rpc.com`
* **fra**: Fráncfort, Europa - `https://laserstream-mainnet-fra.helius-rpc.com`
* **tyo**: Tokio, Asia - `https://laserstream-mainnet-tyo.helius-rpc.com`
* **sgp**: Singapur, Asia - `https://laserstream-mainnet-sgp.helius-rpc.com`

Para hacer pruebas en devnet, usa: `https://laserstream-devnet-ewr.helius-rpc.com`

Consulta la [documentación de gRPC de LaserStream](/docs/es/laserstream/grpc) para obtener instrucciones completas de configuración y pautas para seleccionar un endpoint.

### Configuración del entorno de Rust

Todos los scripts de medición usan Rust con Cargo. Configuración básica:

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

Crea un archivo `.env` con tus credenciales:

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

Obtén tu clave de API de Helius en el [panel de Helius](https://dashboard.helius.dev/). LaserStream devnet está disponible en todos los planes. El acceso a mainnet requiere un plan Business o Professional.

***

## Método 1: Comparación de streams en paralelo

Este script establece dos conexiones independientes con distintos endpoints de gRPC y mide cuál recibe primero los mismos mensajes `BlockMeta`. Este enfoque elimina los problemas de sincronización de relojes mediante tiempos relativos.

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

**Qué mide:** La diferencia de rendimiento relativo entre dos servicios de streaming. La diferencia muestra qué servicio entrega primero la misma información del slot.

**Métricas clave:**

* **Diferencia positiva**: El primer servicio (YS) es más lento que el segundo (LS); LaserStream es más rápido
* **Diferencia negativa**: El primer servicio (YS) es más rápido que el segundo (LS); LaserStream es más lento
* **Media/mediana**: Diferencia de rendimiento promedio
* **P95**: Diferencia de latencia del percentil 95

**Ejecutar la prueba:**

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

**Ejemplo de salida:**

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

La salida muestra diferencias de latencia en tiempo real y estadísticas periódicas. Una diferencia media positiva indica que el segundo servicio (LaserStream) entrega los datos de forma sistemáticamente más rápida.

***

## Método 2: Análisis de la marca de tiempo de creación

Este enfoque compara la marca de tiempo `created_at` incorporada en los mensajes con la hora de tu sistema local al recibirlos.

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

**Limitación importante:** Este método solo mide desde que LaserStream creó el mensaje hasta que lo recibiste. No considera ningún retraso anterior entre el evento de blockchain y el procesamiento de LaserStream.

**Ejecutar la prueba:**

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

**Ejemplo de salida:**

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

Este método proporciona información sobre la latencia de red y procesamiento entre LaserStream y tu aplicación, pero debes usarlo junto con el método 1 para realizar un análisis completo.

***

## Prácticas recomendadas para las pruebas de latencia

### Principios clave

* **Coubicación**: Despliega las pruebas en la misma región que tu endpoint de LaserStream para minimizar la latencia de red
* **Varios métodos**: Usa la comparación de streams en paralelo (método 1) como métrica principal y compleméntala con el análisis de marcas de tiempo
* **Monitoreo a largo plazo**: Ejecuta las pruebas durante periodos prolongados para capturar distintas condiciones de red y niveles de congestión de la blockchain
* **Análisis estadístico**: Céntrate en los percentiles (P95, P99), no solo en los promedios, para comprender la latencia de cola

### Interpretar los resultados

1. **Establece una referencia**: Ejecuta las pruebas durante al menos 1 hora para establecer el rendimiento de referencia en condiciones normales
2. **Identifica patrones**: Busca patrones en los picos de latencia. ¿Se correlacionan con una alta actividad de la blockchain o con la congestión de la red?
3. **Compara percentiles**: La latencia P95 suele ser más importante que la latencia promedio para la experiencia del usuario
4. **Monitorea la consistencia**: Un rendimiento consistente suele ser más valioso que la latencia mínima absoluta

Recuerda que la latencia de la blockchain es variable por naturaleza debido a los requisitos de consenso de la red. Céntrate en las diferencias de rendimiento relativo y en la consistencia, no en las cifras absolutas.
