コンテンツにスキップ

送信処理

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 と障害時の処理へ進む.

Client は 2 つのスレッドに分かれている.

2 つのスレッド

RT スレッドは EtherCAT の周期でひたすら送受信を繰り返しており, ユーザースレッドはそこへコマンドを渡して結果を待つ. 以下では, デバイス 2 台に PerDevice のフレームを 1 つ送る場合を例に, 送信から受信までの処理を追う.

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
}

FrameFrames の中身を借りているだけなので, そのまま 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 で待つ.

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 は終わりであり, ユーザーは ResponseFutureawait して待つ.

以下は RT スレッド側の処理である.

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_frameSlot の内容を 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 が SeqCmd, 残りがペイロードである. 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 つである.

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 のフレームがもう一度流れる.

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 では 0
for (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 フレームは ackdata の 2 byte である. ack はそのデバイスが最後に処理し終えた SEQ を表す.

ack_offset は 例えば, デバイスが一個前のSEQを返した場合は 255 となり, ack_offset > span で弾かれる.

デバイス 2 台の両方が ACK を返すと, acked0b11 になる.

送信直後 acked = 0b00
dev1 のみ ACK acked = 0b10
両方 ACK acked = 0b11 = all_acked → 完了

all_acked は全デバイスのビットが立った値である.

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 が起こされる.

stop-and-wait の 1 フレーム

send が返す ResponseFuture をすぐに await せず貯めれば, 応答を待たずに次のフレームを送れる. このとき pending に複数の Inflight が並ぶ.

pending: ┌────────┬────────┬────────┐
│ seq=3 │ seq=4 │ seq=5 │
│ acked │ acked │ acked │
└────────┴────────┴────────┘
front back

stage_newnext_seq を 1 つずつ配って push_back するので, pending の SEQ は連番の昇順になる.

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 しか受理しない.

route_ackstake(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=41 採用
古い front=5, ack=4255 無視
異常 front=3, span=2, ack=96 無視

上記の 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 を掃除し, それぞれについて 待っている ResponseFuturecause() を返し, Slot をプールへ戻す.

teardownによる後処理は正常系ではほぼ意味を持たない. 一方リンクエラーでは RT スレッドだけが死に, Client は生き残る. このとき後処理を怠ると, デッドロックが発生する可能性がある. SlotPool::acquire は取得した Semaphore の許可を forget するため, in-flight の Slotを回収しないと, プールが枯渇し以降のsendacquire` で永久にブロックする可能性がある.

また, エラー種別が具体的になるというメリットもある.

障害時は stage_tx が新規コマンドの投入より回復を優先する.

優先度 条件 動作
1 reset_remaining > 0 Cmd::Reset をブロードキャストし続ける.
2 resync.active pending の先頭を再送する.
3 held_exclusive あり pending が空になるまで待って投入する.
4 通常 cmd_rx.try_recv() で新規を 1 件取る.
Normal
│ 先頭が timeout_cycles 経過
Resync ──── pending が空 ────▶ Normal
│ (この間, 先頭を毎サイクル再送)
│ max_resync_rounds 回超過
Reset
│ SEQ を 0 から振り直し
Resync (Reset 後)
│ 再度 max_resync_rounds 回超過
GiveUp ──── 全 pending を Timeout で失敗 ────▶ Normal

pending が空になれば, どの状態からでも Normal へ戻る.

状態の判定は pending の先頭の age (応答を待っているサイクル数) で行う. 既定値 (timeout_cycles = 10, max_resync_rounds = 8, reset_resend_cycles = 2) では次のようになる.

サイクル 状態 出来事
〜10 Normal 先頭の age を数える.
10 Normal → Resync agetimeout_cycles に到達. 再送モードへ入る.
20〜90 Resync 10 サイクルごとに rounds が増える. その間は先頭を再送.
90 Resync → Reset roundsmax_resync_rounds に到達. Cmd::Reset を送る.
92 Reset → Resync Reset 送出が終わり, SEQ を 0 から振り直す.
102〜172 Resync 再び 10 サイクルごとに rounds が増える.
172 Resync → GiveUp Reset 済みなので今度は諦め, 全 pending をタイムアウトさせて Normal へ戻る.

各サイクル数は 3 つの設定値から決まる.

  • 再送モードに入るまで: agetimeout_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 で回復する.

route_acks で 1 本でも完了すれば on_ack_progress が呼ばれ, roundsreset_tried と先頭の age が戻る.

一方で active は保持し, pending が空になるまで再送モードを維持し, 溜まっているフレームを 1 本ずつ流し切ってから通常運転へ戻る. stage_tx の優先順位で示したとおり, その間は新規コマンドを受け付けない.

再送を繰り返しても進まない場合, デバイスの expected_seq が想定外の値になっている (例えばデバイスが再起動した等) 可能性が高い. この場合は再送では合わせられないので, Cmd::Reset で SEQ をリセットする.

デバイス側は Cmd::Reset を受けると次の状態になる.

  • expected_seq を 0 に戻す
  • ack0xFF にする
  • 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 = 0xFFdistance_from255 になり, spanの判定で弾かれる.)

rx_valid = false はデバイスがEtherCATバスに参加していない状態であり, 再送しても Cmd::Reset を送っても届かないのでどうしようもない. そのため, stale_limit サイクル待って一括でタイムアウトさせる.

この間は route_acksadvance_head も呼ばれないため, 先頭の age は進まない.

fail_pending_timeout は全 ResponseFutureError::Timeout を返して Slot をプールへ戻す. RT ループ自体は継続し, 次のサイクルから新しいコマンドを受け付ける.

DatagramBuilderbuild は「送れば適用される」前提で 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)
}
}