> ## 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 客户端

> 以 gRPC 流形式接收 OrbitFlare 解码碎片，支持交易与账户过滤器。

## 安装

```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 客户端](/cn/sdk/rust-grpc) 相同 — 默认与行为一致。

## 编写 YAML 配置

JetStream 支持交易与账户过滤器。不包含槽位、区块或承诺 — 这些为 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 客户端](/cn/sdk/rust-grpc) 完全一致。

## 完整示例

监听 Raydium AMM  swap 并打印每条交易的签名与指令数量：

```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`，因此你可以检测丢失的消息。
* **可选的富化** — 按交易请求手续费付款方、程序 id、计算单元价格、计算上限、解析后的地址表地址等。
* **槽位生命周期事件** — 一条独立的服务端流，推送槽位的 alive/complete/dead 事件。

v2 客户端位于 `jetstream::v2` 之下。构建器与 v1 完全一致（相同的方法、默认值、`ORBITFLARE_JETSTREAM_URL` 环境变量与故障转移）：

```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()?;
```

### 构建类型化过滤器

每个过滤器都是一个带有客户端选定 id（通过 `.with_id()`）的 `TransactionFilter`。该 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 => {}
    }
}
```

### 在活动流上管理过滤器

从流中取得一个句柄，即可在无需重连的情况下添加或移除过滤器。移除时引用你通过 `.with_id()` 指定的过滤器 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()` 是一元调用，具有与流相同的故障转移。

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

完整的 protobuf 规范、消息结构、富化字段与槽位状态，请参见 [JetStream v2 协议参考](/cn/data-streaming/jetstream-v2-reference)。
