Skip to content
This is the development version of the documentation. It may change before the next release. See 0.6.x for the latest release.

送信処理

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

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