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

# JetStream-клиент

> Декодированные шреды OrbitFlare как gRPC-стримы с фильтрами транзакций и аккаунтов.

## Установка

```bash theme={null}
cargo add orbitflare-sdk --features jetstream
```

## Сборка клиента

```rust theme={null}
use orbitflare_sdk::{JetstreamClientBuilder, RetryPolicy, Result};
use std::time::Duration;

let client = JetstreamClientBuilder::new()
    .url("http://ny.jetstream.orbitflare.com")
    .fallback_url("http://fra.jetstream.orbitflare.com")
    .retry(RetryPolicy {
        initial_delay: Duration::from_millis(100),
        max_delay: Duration::from_secs(30),
        multiplier: 2.0,
        max_attempts: 0,
    })
    .timeout_secs(30)
    .keepalive_secs(60)
    .ping_interval_secs(10)
    .max_missed_pongs(3)
    .channel_capacity(4096)
    .build()?;
```

Минимально:

```rust theme={null}
let client = JetstreamClientBuilder::new()
    .url("http://ny.jetstream.orbitflare.com")
    .build()?;
```

URL по умолчанию берётся из `ORBITFLARE_JETSTREAM_URL`. Все методы билдера те же, что у [gRPC-клиента](/ru/sdk/rust-grpc) — те же дефолты и поведение.

## Написание YAML-конфига

JetStream поддерживает фильтры транзакций и аккаунтов. Слоты, блоки и commitment — специфичны для Yellowstone, здесь их нет.

```yaml theme={null}
# jetstream.yml
transactions:
  raydium:
    account_include:
      - "675kPX9MHTjS2zt1qfr1NYHuzeLXfQM9H24wFSUt1Mp8"
  pumpfun:
    account_include:
      - "6EF8rrecthR5Dkzon8Nwu78hRvfCKubJ14M5uBEwF6P"

accounts:
  my_wallet:
    account:
      - "YOUR_WALLET_ADDRESS"
    owner:
      - "TokenkegQfeZyiNwAJbNbGKPFXCWuBvf9Ss623VQ5DA"
```

### Справочник фильтров YAML

**`transactions`** — именованные фильтры. `account_include` совпадает с транзакциями, где задействованы эти адреса. `account_exclude` исключает. `account_required` — все перечисленные адреса должны быть в транзакции.

**`accounts`** — следить за конкретными адресами через `account` или за всеми аккаунтами программы через `owner`.

Поддерживается подстановка `${ENV_VAR}`.

## Подписка и чтение событий

### Из YAML

```rust theme={null}
let mut stream = client.subscribe_yaml("jetstream.yml")?;
```

### Программно

```rust theme={null}
use std::collections::HashMap;
use orbitflare_sdk::proto::jetstream::*;

let mut filters = HashMap::new();
filters.insert("target".into(), SubscribeRequestFilterTransactions {
    account_include: vec![some_address.to_string()],
    account_exclude: vec![],
    account_required: vec![],
});

let request = SubscribeRequest {
    transactions: filters,
    accounts: HashMap::new(),
    ping: Some(SubscribeRequestPing { id: 1 }),
};

let mut stream = client.subscribe(request);
```

### С типизированным билдером

`SubscribeRequestBuilder` и `TransactionFilter` собирают тот же запрос без ручного написания proto. Каждый `TransactionFilter` предоставляет `account_include`, `account_exclude` и `account_required` и принимает любой итератор строкоподобных значений.

```rust theme={null}
use orbitflare_sdk::jetstream::{SubscribeRequestBuilder, TransactionFilter};

let request = SubscribeRequestBuilder::new()
    .transactions(
        "pumpfun",
        TransactionFilter::new()
            .account_include(["6EF8rrecthR5Dkzon8Nwu78hRvfCKubJ14M5uBEwF6P"])
            .account_required(["So11111111111111111111111111111111111111112"]),
    )
    .build();

let mut stream = client.subscribe(request);
```

### Чтение потока

```rust theme={null}
use orbitflare_sdk::proto::jetstream::subscribe_update::UpdateOneof;

while let Some(update) = stream.next().await {
    let update = update?;
    match update.update_oneof {
        Some(UpdateOneof::Transaction(tx)) => {
            // tx.slot - the slot number
            // tx.transaction - transaction info with signature, account_keys,
            //   instructions, address_table_lookups
        }
        Some(UpdateOneof::Account(acct)) => {
            // acct.slot - the slot
            // acct.account - account info (pubkey, lamports, owner, data)
            // acct.is_startup - true during initial snapshot
        }
        _ => {}
    }
}
```

Закрытие, несколько потоков, переподключение и ping/pong работают так же, как у [gRPC-клиента](/ru/sdk/rust-grpc).

## Полный пример

Поток следит за свопами Raydium AMM и выводит подпись и число инструкций каждой транзакции.

```rust theme={null}
use orbitflare_sdk::{JetstreamClientBuilder, Result};
use orbitflare_sdk::proto::jetstream::subscribe_update::UpdateOneof;

#[tokio::main]
async fn main() -> Result<()> {
    let client = JetstreamClientBuilder::new()
        .url("http://ny.jetstream.orbitflare.com")
        .build()?;

    let mut stream = client.subscribe_yaml("jetstream.yml")?;
    let mut count: u64 = 0;

    println!("streaming raydium txs...");

    while let Some(update) = stream.next().await {
        let update = update?;

        if let Some(UpdateOneof::Transaction(tx)) = update.update_oneof {
            count += 1;
            if let Some(info) = &tx.transaction {
                let sig = bs58::encode(&info.signature).into_string();
                let num_ix = info.instructions.len();
                let num_accounts = info.account_keys.len();

                println!(
                    "#{count} slot={} sig={}... ix={num_ix} accounts={num_accounts}",
                    tx.slot,
                    &sig[..16],
                );
            }
        }
    }

    Ok(())
}
```

С этим `jetstream.yml`:

```yaml theme={null}
transactions:
  raydium:
    account_include:
      - "675kPX9MHTjS2zt1qfr1NYHuzeLXfQM9H24wFSUt1Mp8"
```

## JetStream v2

JetStream v2 (`orbitflare_sdk::jetstream::v2`) работает на **тех же эндпоинтах и с той же аутентификацией**, что и v1, и полностью аддитивен: v1 продолжает работать без изменений. Что добавляет v2:

* **Фильтры, управляемые во время работы** — добавляйте и удаляйте фильтры на активном потоке без переподключения.
* **Порядковые номера для каждого сообщения** — каждый ответ несёт монотонный `sequence`, так что вы можете обнаруживать потерянные сообщения.
* **Опциональное обогащение** — запрашивайте плательщика комиссии, идентификаторы программ, цену за compute-unit, лимит вычислений, разрешённые адреса из address-таблиц и другое, для каждой транзакции.
* **События жизненного цикла слота** — отдельный серверный поток событий слота alive/complete/dead.

Клиент v2 находится в `jetstream::v2`. Билдер идентичен v1 (те же методы, дефолты, переменная окружения `ORBITFLARE_JETSTREAM_URL` и failover):

```rust theme={null}
use orbitflare_sdk::jetstream::v2::{JetstreamClientBuilder, TransactionFilter};

let client = JetstreamClientBuilder::new()
    .url("http://ny.jetstream.orbitflare.com")
    .fallback_url("http://fra.jetstream.orbitflare.com")
    .build()?;
```

### Построение типизированных фильтров

Каждый фильтр — это `TransactionFilter` с выбранным клиентом id (через `.with_id()`). Этот id возвращается в каждой совпавшей транзакции и в подтверждении валидации фильтра, так что вы можете сопоставлять совпадения и позже удалить фильтр.

```rust theme={null}
let filter = TransactionFilter::new()
    .account_include(["6EF8rrecthR5Dkzon8Nwu78hRvfCKubJ14M5uBEwF6P"])
    .include_enrichment(true)
    .with_id("pumpfun");
```

Фильтр должен задать хотя бы одно из `account_include`, `account_exclude` или `account_required`; пустые фильтры (совпадающие со всем) отклоняются. `include_enrichment` необязателен (по умолчанию выключен) и применяется ко всей подписке: если хотя бы один активный фильтр его включает, каждая получаемая транзакция обогащается.

### Подписка на транзакции

```rust theme={null}
use orbitflare_sdk::proto::jetstream::v2::subscribe_transactions_response::Payload;

let mut stream = client.subscribe_transactions(vec![filter]);

while let Some(resp) = stream.next().await {
    let resp = resp?;
    match resp.payload {
        Some(Payload::Transaction(ft)) => {
            // resp.sequence - monotonic per-stream sequence number
            // ft.filter_ids - which of your filters this tx matched
            if let Some(tx) = &ft.transaction {
                // tx.slot, tx.signature, tx.account_keys, tx.instructions
                // enrichment (only when include_enrichment is set): tx.fee_payer,
                //   tx.program_ids, tx.compute_unit_price, tx.compute_limit,
                //   tx.loaded_writable_addresses, tx.loaded_readonly_addresses
                println!("seq={} slot={} cu_price={}", resp.sequence, tx.slot, tx.compute_unit_price);
            }
        }
        Some(Payload::FilterValidation(r)) => {
            // one per filter after each add/remove
            println!("filter {} accepted={} {}", r.filter_id, r.accepted, r.rejection_reason);
        }
        Some(Payload::Heartbeat(hb)) => {
            // hb.server_ts_ms - server wall-clock, ms since epoch
            let _ = hb;
        }
        Some(Payload::Pong(p)) => {
            // p.ping_id
            let _ = p;
        }
        None => {}
    }
}
```

### Управление фильтрами на активном потоке

Возьмите handle из потока и добавляйте или удаляйте фильтры без переподключения. Удаления ссылаются на id фильтра, который вы задали через `.with_id()`.

```rust theme={null}
let handle = stream.handle();

handle.add_filters(vec![
    TransactionFilter::new()
        .account_include(["675kPX9MHTjS2zt1qfr1NYHuzeLXfQM9H24wFSUt1Mp8"])
        .with_id("raydium"),
])?;

handle.remove_filters(vec!["pumpfun".to_string()])?;
```

### События жизненного цикла слота

`subscribe_slots()` — это отдельный серверный поток, он не принимает фильтров.

```rust theme={null}
let mut slots = client.subscribe_slots();

while let Some(event) = slots.next().await {
    let event = event?;
    // event.slot, event.status() (a SlotStatus enum: alive / complete / dead),
    // event.parent_slot, event.current_leader, event.sequence
    println!("slot={} status={:?}", event.slot, event.status());
}
```

### Проверки состояния

`get_version()` и `ping()` — унарные вызовы с тем же failover, что и у потоков.

```rust theme={null}
println!("version={}", client.get_version().await?);
println!("ping={}", client.ping(7).await?);
```

Полную спецификацию protobuf, форматы сообщений, поля обогащения и статусы слотов см. в [справочнике протокола JetStream v2](/ru/data-streaming/jetstream-v2-reference).
