送信処理
This content is not available in your language yet.
DatagramBuilder が生成した Frame を実際にデバイスへ送り届ける処理について述べる.
実装は crates/autd3-rs/src/client/ にある.
Frames は複数のフレームを持ち, 送信ループは Client ではなくユーザー側が書く.
for frame in &frames { client.send_checked(frame).await?;}これは 1 フレームずつ応答を待つ stop-and-wait である.
send を使って ResponseFuture を貯めれば, 応答を待たずに次を送る streaming になる.
以下ではまず stop-and-wait が 1 フレームを送り終えるまでを順に追い, そのあと streaming と障害時の処理へ進む.
stop-and-wait
Section titled “stop-and-wait”Client は 2 つのスレッドに分かれている.
RT スレッドは EtherCAT の周期でひたすら送受信を繰り返しており, ユーザースレッドはそこへコマンドを渡して結果を待つ.
以下では, デバイス 2 台に PerDevice のフレームを 1 つ送る場合を例に, 送信から受信までの処理を追う.
① Slot を借りる
Section titled “① Slot を借りる”pub async fn send(&self, frame: Frame<'_>) -> Result<ResponseFuture, Error> { match frame.distribution() { Distribution::Broadcast => self.send_broadcast(&frame.datagrams()[0]).await, Distribution::PerDevice => self.send_datagrams(frame.datagrams()).await, }}PerDevice なので send_datagrams へ進む.
async fn send_datagrams(&self, datagrams: &[Datagram]) -> Result<ResponseFuture, Error> { if datagrams.len() != self.num_devices { /* エラー */ } let mut slot = self.pool.acquire().await; // ① slot.reset(Distribution::PerDevice); for (device, datagram) in datagrams.iter().enumerate() { slot.payload_mut(device).copy_from_slice(&datagram.payload); slot.set_cmd(device, datagram.cmd); } self.dispatch(slot, false).await}Frame は Frames の中身を借りているだけなので, そのまま RT スレッドへ渡すことはできない.
そこで SlotPool から Slot を借り, ペイロードをそこへ写す.
struct Slot { num_devices: usize, dist: Distribution, payload: Box<[u8]>, // num_devices * PAYLOAD_BYTES cmds: Box<[Cmd]>, data: Box<[u8]>, // 各デバイスからの応答バイト}Slot は送信内容と応答の両方を持つ, 1 回の送信に対応する構造体である.
なお, アロケーションを回避するため, Slot は毎回確保するのではなくプールから借りる.
pub(super) async fn acquire(&self) -> Slot { self.permits.acquire().await.expect("...").forget(); self.free.lock().unwrap_or_else(PoisonError::into_inner).pop().expect("...")}Semaphore の許可を取ってから free list から 1 つ取り出す.
空きが無ければ .await で待つ.
② RT スレッドへ渡す
Section titled “② RT スレッドへ渡す”async fn dispatch(&self, slot: Slot, exclusive: bool) -> Result<ResponseFuture, Error> { let (response_tx, response_rx) = self.completions.channel(); if let Err(e) = self.cmd_tx.send(CmdMessage { frame: slot, response_tx, exclusive }).await { self.pool.release(e.0.frame); return Err(Error::RtClosed); } Ok(response_rx)}完了通知用のチャネルを 1 組作り, Slot と送信側 (response_tx) をまとめて RT スレッドへ送る.
受信側 (response_rx) が ResponseFuture としてユーザーへ返る.
ここで send は終わりであり, ユーザーは ResponseFuture を await して待つ.
以下は RT スレッド側の処理である.
③ TX バッファへ書く
Section titled “③ TX バッファへ書く”RT スレッドは毎サイクル stage_tx で「今回何を送るか」を決める.
新規コマンドがあれば stage_new が呼ばれる.
fn stage_tx(&mut self) -> StageOutcome { /* エラーや障害時の処理は省略 */ let exclusive_inflight = /* 受信データを使うタイプのフレームを送信中の場合. 後述. */; if !exclusive_inflight && self.pending.len() < self.config.max_inflight.get() { match self.cmd_rx.try_recv() { Ok(msg) if msg.exclusive && !self.pending.is_empty() => { self.held_exclusive = Some(msg); } Ok(msg) => self.stage_new(msg), Err(mpsc::error::TryRecvError::Empty) => {} Err(mpsc::error::TryRecvError::Disconnected) => return StageOutcome::Disconnected, } } StageOutcome::Staged}
fn stage_new(&mut self, msg: CmdMessage) { let seq = self.next_seq; self.next_seq = self.next_seq.next(); stage_frame(seq, &msg.frame, &mut self.tx_bufs); self.pending.push_back(Inflight { seq, frame: msg.frame, acked: 0, age: 0, exclusive: msg.exclusive, response_tx: msg.response_tx, });}try_recv は 1 サイクルにつき 1 回しか呼ばれないため, 新規 SEQ の投入は最大 1 サイクルに 1 件である.
stage_frame が Slot の内容を TX バッファへ展開する.
fn stage_frame(seq: Seq, frame: &Slot, tx_bufs: &mut [[u8; TX_FRAME_BYTES]]) { for (device, buf) in tx_bufs.iter_mut().enumerate() { buf[0] = seq.get(); buf[1] = frame.cmd_for(device).as_u8(); buf[2..].copy_from_slice(frame.payload_for(device)); }}TX フレームの先頭 2 byte が Seq と Cmd, 残りがペイロードである.
Seq はこの送信に割り当てられた 8 bit の通し番号で, 周期的に 0 → 1 → … → 255 → 0 とラップする.
送信中のコマンドは pending に積まれる.
struct Inflight { seq: Seq, frame: Slot, acked: u128, // ACK を返したデバイスのビットマスク age: u32, // 応答を待っているサイクル数 exclusive: bool, response_tx: CompletionSender,}stop-and-wait では pending の要素は常に高々 1 つである.
④ 1 サイクル送受信する
Section titled “④ 1 サイクル送受信する”let CycleOutcome { rx_valid } = self.link.cycle(&self.tx_bufs, &mut self.rx_bufs)?;Link::cycle が 1 サイクル分の PDO 送受信を行う.
デバイスは受け取ったフレームを処理し, 結果を次のサイクルの RX に載せる. したがって, サイクル N で送ったフレームの ACK が読めるのは早くてもサイクル N+1 である.
サイクル N+1 の stage_tx では新規コマンドが無いので, TX バッファは前サイクルのまま送出される.
つまり同じ SEQ のフレームがもう一度流れる.
⑤ ACK を集める
Section titled “⑤ ACK を集める”rx_valid が真 (つまり正しくデータが受信できている) なら route_acks が呼ばれる.
rx_valid をどう返すかは Link の実装に依存するが, 通常は, WKC が期待値と一致しているかどうかで判定する. false になるのは例えばSafe-OPに落ちた場合などである.
let front_seq = pending.front().seq;let span = pending.back().seq.distance_from(front_seq) as usize; // stop-and-wait では 0for (device, rx_buf) in rx_bufs.iter().enumerate() { let rx = RxFrame::parse(rx_buf); let ack_offset = rx.ack.distance_from(front_seq) as usize; if ack_offset > span { continue; } // まだ届いていないデバイスは飛ばす let bit = 1u128 << device; for entry in pending.iter_mut().take(ack_offset + 1) { if entry.acked & bit == 0 { entry.acked |= bit; entry.frame.record_data(device, rx.data); // 応答バイトを Slot へ記録 } }}stop-and-wait ではpendingの要素は 1 つなので, front() == back() である.
RX フレームは ack と data の 2 byte である.
ack はそのデバイスが最後に処理し終えた SEQ を表す.
ack_offset は 例えば, デバイスが一個前のSEQを返した場合は 255 となり, ack_offset > span で弾かれる.
デバイス 2 台の両方が ACK を返すと, acked は 0b11 になる.
送信直後 acked = 0b00dev1 のみ ACK acked = 0b10両方 ACK acked = 0b11 = all_acked → 完了all_acked は全デバイスのビットが立った値である.
⑥ 完了を通知する
Section titled “⑥ 完了を通知する”while pending.front().is_some_and(|e| e.acked == all_acked) { let entry = pending.pop_front().expect("just checked"); entry.response_tx.send(Ok(Response::from_slice(entry.frame.data()))); self.pool.release(entry.frame); progressed = true;}Slot に記録した応答バイトを Response へ写して送信側へ渡し, Slot をプールへ返す.
これでユーザースレッドの ResponseFuture が起こされる.
streaming
Section titled “streaming”send が返す ResponseFuture をすぐに await せず貯めれば, 応答を待たずに次のフレームを送れる.
このとき pending に複数の Inflight が並ぶ.
pending: ┌────────┬────────┬────────┐ │ seq=3 │ seq=4 │ seq=5 │ │ acked │ acked │ acked │ └────────┴────────┴────────┘ front backstage_new が next_seq を 1 つずつ配って push_back するので, pending の SEQ は連番の昇順になる.
デバイス側の受理条件
Section titled “デバイス側の受理条件”if in_frame.seq == expected_seq { expected_seq = seq + 1; let data = dispatch(cmd, payload); tx = pack_tx(in_frame.seq, data); // ack = 今処理した SEQ, data = 結果} else { bump(Telemetry::SeqMismatch); // tx は据え置き}デバイスは順番通りの SEQ しか受理しない.
累積 ACK の展開
Section titled “累積 ACK の展開”route_acks の take(ack_offset + 1) は, 最新の ACK に従って, 古いFrameまで遡って処理する.
for entry in pending.iter_mut().take(ack_offset + 1) { .. }例として, デバイス 2 台, pending = [seq3, seq4, seq5] で, dev0 が 5 まで進み dev1 が 3 で足踏みしている場合を追う.
front = 3, back = 5, span = 2
dev0: ack=5 → offset = 2 ≤ 2 → take(3) → seq3,4,5 に bit0 を立てるdev1: ack=3 → offset = 0 ≤ 2 → take(1) → seq3 のみに bit1 を立てる
[seq3: 0b11] [seq4: 0b01] [seq5: 0b01] ↑ 完了
pending = [seq4: 0b01] [seq5: 0b01]次のサイクルで dev1 が ack=5 を返すと, seq4 と seq5 の両方が一度に埋まる.
dev1 は seq4 の ACK を報告していないが, ack=5 から遡って処理されたという扱いにする.
なお, 完了処理は front から順番に行うので, ユーザーから見た完了順序は送信順序と一致する.
ack_offset > span の 比較で, 2 パターンの不正な ACK を弾いている.
| 状況 | 計算 | 結果 |
|---|---|---|
| 正常 | front=3, ack=4 → 1 |
採用 |
| 古い | front=5, ack=4 → 255 |
無視 |
| 異常 | front=3, span=2, ack=9 → 6 |
無視 |
exclusive
Section titled “exclusive”上記の ACK 処理の影響で ACK が 2 つ以上飛んだ場合, 中間のエントリには後続フレームの応答バイトが記録される.
通常はこれで問題ないが, デバイスからの受信データを使うコマンドではこれが問題になる.
そのため, 受信データが必要なコマンド (例えばファームウェアのバージョンを読むコマンドなど) は, 送信中のフレームが完了するまで次のコマンドを送らないようにする.
ここまでは正常系だった. RT ループは毎サイクル, 送受信の結果に応じて 3 つの経路へ分岐する.
let mut link_error = None;loop { if closed { break } if matches!(stage_tx(), StageOutcome::Disconnected) { break } let rx_valid = match self.link.cycle(&self.tx_bufs, &mut self.rx_bufs) { Ok(CycleOutcome { rx_valid }) => rx_valid, Err(e) => { link_error = Some(format!("link cycle failed: {e}")); break } }; if reset_remaining > 0 { advance_reset_phase() } else if rx_valid { handle_healthy() } // 正常系 = route_acks else { handle_stale() }}teardown(link_error.as_deref());異常は主に3つのパターンに分かれる.
| 判定 | 意味 | 担当 | |
|---|---|---|---|
| リンク断 | cycle() が Err |
ソケット自体の異常. | teardown |
| RX 不信頼 | rx_valid == false |
WKC 不一致 / 全デバイスが OP でない. | handle_stale |
| ACK 停滞 | rx_valid だが先頭が進まない |
フレーム喪失, デバイス停止. | ResyncState |
Link 側の異常は Client 側ではどうすることもできないので, ループを抜けて RT スレッドを終了させる.
終了処理は理由によらず teardown に一本化されている.
fn teardown(&mut self, link_error: Option<&str>) { let cause = || link_error.map_or(Error::RtClosed, |msg| Error::Link(msg.to_owned())); if let Some(msg) = self.held_exclusive.take() { /* cause() を返して Slot を戻す */ } for entry in self.pending.drain(..) { /* 同上 */ } while let Ok(msg) = self.cmd_rx.try_recv() { /* 同上 */ }}pending / held_exclusive / チャネルに残った CmdMessage を掃除し, それぞれについて
待っている ResponseFuture に cause() を返し, Slot をプールへ戻す.
teardownによる後処理は正常系ではほぼ意味を持たない.
一方リンクエラーでは RT スレッドだけが死に, Client は生き残る.
このとき後処理を怠ると, デッドロックが発生する可能性がある.
SlotPool::acquire は取得した Semaphore の許可を forget するため, in-flight の Slotを回収しないと, プールが枯渇し以降のsendがacquire` で永久にブロックする可能性がある.
また, エラー種別が具体的になるというメリットもある.
stage_tx の優先順位
Section titled “stage_tx の優先順位”障害時は stage_tx が新規コマンドの投入より回復を優先する.
| 優先度 | 条件 | 動作 |
|---|---|---|
| 1 | reset_remaining > 0 |
Cmd::Reset をブロードキャストし続ける. |
| 2 | resync.active |
pending の先頭を再送する. |
| 3 | held_exclusive あり |
pending が空になるまで待って投入する. |
| 4 | 通常 | cmd_rx.try_recv() で新規を 1 件取る. |
ACK 停滞
Section titled “ACK 停滞” Normal │ 先頭が timeout_cycles 経過 ▼ Resync ──── pending が空 ────▶ Normal │ (この間, 先頭を毎サイクル再送) │ │ max_resync_rounds 回超過 ▼ Reset │ SEQ を 0 から振り直し ▼ Resync (Reset 後) │ 再度 max_resync_rounds 回超過 ▼ GiveUp ──── 全 pending を Timeout で失敗 ────▶ Normalpending が空になれば, どの状態からでも Normal へ戻る.
状態の判定は pending の先頭の age (応答を待っているサイクル数) で行う.
既定値 (timeout_cycles = 10, max_resync_rounds = 8, reset_resend_cycles = 2) では次のようになる.
| サイクル | 状態 | 出来事 |
|---|---|---|
| 〜10 | Normal | 先頭の age を数える. |
| 10 | Normal → Resync | age が timeout_cycles に到達. 再送モードへ入る. |
| 20〜90 | Resync | 10 サイクルごとに rounds が増える. その間は先頭を再送. |
| 90 | Resync → Reset | rounds が max_resync_rounds に到達. Cmd::Reset を送る. |
| 92 | Reset → Resync | Reset 送出が終わり, SEQ を 0 から振り直す. |
| 102〜172 | Resync | 再び 10 サイクルごとに rounds が増える. |
| 172 | Resync → GiveUp | Reset 済みなので今度は諦め, 全 pending をタイムアウトさせて Normal へ戻る. |
各サイクル数は 3 つの設定値から決まる.
- 再送モードに入るまで:
ageがtimeout_cyclesに達する =10サイクル. - Resync → Reset: 再送モードに入った後は
timeout_cyclesごとにroundsが 1 増え,max_resync_rounds回で Reset へ移る. つまりtimeout_cycles × max_resync_rounds=10 × 8=80サイクルで, 到達は10 + 80 =90. - Reset の送出:
reset_resend_cycles=2サイクルCmd::Resetを送り, 振り直す.90 + 2 =92. - Reset → GiveUp: 振り直し後も同じで,
timeout_cycles × max_resync_rounds=80サイクル後に GiveUp.92 + 80 =172.
まとめると, 決着までのサイクル数は次式になる.
timeout_cycles × (1 + 2 × max_resync_rounds) + reset_resend_cycles = 10 × (1 + 2 × 8) + 2 = 172よって 1 ms サイクルなら最悪でも約 172 ms で回復する.
進捗があれば状態をリセット
Section titled “進捗があれば状態をリセット”route_acks で 1 本でも完了すれば on_ack_progress が呼ばれ, rounds と reset_tried と先頭の age が戻る.
一方で active は保持し, pending が空になるまで再送モードを維持し, 溜まっているフレームを 1 本ずつ流し切ってから通常運転へ戻る.
stage_tx の優先順位で示したとおり, その間は新規コマンドを受け付けない.
SEQ のリセット
Section titled “SEQ のリセット”再送を繰り返しても進まない場合, デバイスの expected_seq が想定外の値になっている (例えばデバイスが再起動した等) 可能性が高い.
この場合は再送では合わせられないので, Cmd::Reset で SEQ をリセットする.
デバイス側は Cmd::Reset を受けると次の状態になる.
expected_seqを 0 に戻すackを0xFFにする- FIFO に積まれた未処理フレームを破棄する
Reset は FIFO を無視して即座に発行される.
ホスト側は Reset 送出が終わったところで pending の SEQ を 0 から振り直す.
fn advance_reset_phase(&mut self) { self.reset_remaining -= 1; if self.reset_remaining == 0 { let mut seq = Seq::ZERO; for entry in &mut self.pending { entry.seq = seq; entry.acked = 0; entry.age = 0; seq = seq.next(); } self.next_seq = seq; self.resync.active = !self.pending.is_empty(); self.resync.rounds = 0; }}reset_remaining > 0 の間は handle_healthy を呼ばないので, この期間の RX は ACK 処理に使わない.
(仮に使われたとしても, 振り直し後の front_seq = 0 に対して ack = 0xFF は distance_from が 255 になり, spanの判定で弾かれる.)
rx_valid = false はデバイスがEtherCATバスに参加していない状態であり, 再送しても Cmd::Reset を送っても届かないのでどうしようもない.
そのため, stale_limit サイクル待って一括でタイムアウトさせる.
この間は route_acks も advance_head も呼ばれないため, 先頭の age は進まない.
タイムアウト後
Section titled “タイムアウト後”fail_pending_timeout は全 ResponseFuture に Error::Timeout を返して Slot をプールへ戻す.
RT ループ自体は継続し, 次のサイクルから新しいコマンドを受け付ける.
mirror の Desync
Section titled “mirror の Desync”DatagramBuilder の build は「送れば適用される」前提で mirror を更新している.
送信が失敗した場合はその前提が崩れるため, mirror を Desynced にして以降の検証を止める必要がある.
特にタイムアウトは「適用されたかどうか自体が分からない」状態なので, 検証を続けることはできない.
この無効化は ResponseFuture が担う.
pub struct ResponseFuture { completion: Option<Arc<Completion>>, pool: Arc<CompletionPool>, mirror: Option<MirrorHandle>,}
impl ResponseFuture { fn finish(&mut self, result: Result<Response, Error>) -> Poll<Result<Response, Error>> { if result.is_err() && let Some(mirror) = self.mirror.take() { mirror.desync(); } Poll::Ready(result) }}