ch32rv 0.12.0

Flashing and debugging tool for WCH CH32 RISC-V microcontrollers
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
//! en: DMI-backed runtime I/O sources shared by `monitor`, `run` and `arduino monitor`
//! (docs/cli.ja.md §4.5): `dmdata` (the ch32fun/minichlink DM data0/data1 mailbox), `dmseq` (the
//! same mailbox with sequence numbers and a CRC-8 - OEP `target.console` framing 2) and `rtt`
//! (a SEGGER-format ring buffer in target RAM). Both are bidirectional: one [`DmiSource::poll`]
//! drains the target's output and hands over pending host input. Also here: the output [`Sink`]
//! (raw bytes on stdout, or `output` NDJSON events on stderr under `--json`, where stdout is
//! reserved for the result envelope) and the reader thread that feeds host input (stdin or a
//! socket) into a channel.
//!
//! `rtt` reads and writes target RAM over the Debug Module, which needs the hart halted, so a
//! poll briefly halts the core (unless it already is - a `run` servicing semihosting leaves a
//! halted core alone). `dmdata` only touches the DM data registers and never halts.
//!
//! ja: `monitor`/`run`/`arduino monitor` が共用する DMI 経由の実行時 I/O source(dmdata / rtt)。
//! どちらも双方向で、1 回の [`DmiSource::poll`] で target 出力を汲み host 入力を渡す。出力先
//! [`Sink`](生 stdout。`--json` 時は stdout を envelope に譲り stderr の `output` event)と、
//! host 入力(stdin / socket)を channel に流す reader thread もここに置く。`rtt` は RAM の
//! 読み書きに halt が要るので poll ごとに一瞬 halt する(既に halt 中ならそのまま)。

use std::io::{Read, Write};
use std::sync::mpsc::{self, Receiver};
use std::time::{Duration, Instant};

use ch32rv_contract::event::Event;
use ch32rv_contract::policy::MonitorSource;
use ch32rv_contract::{ErrorKind, Warning};
use ch32rv_dmi::{DmSeq, DmiError, dmseq};

use crate::args::Cli;
use crate::session::Session;

// ---- RTT control block layout (SEGGER format; ArduinoCore-CH32 SerialRTT publishes the same) ----

/// All CH32 parts map SRAM at this base.
const RTT_RAM_BASE: u32 = 0x2000_0000;
/// The control-block id string the target publishes once its RTT channel is up.
const RTT_MAGIC: &[u8] = b"SEGGER RTT";
/// Sanity cap on a ring-buffer size read out of RAM (reject a half-initialized / garbage block).
const RTT_MAX_BUF: u32 = 0x1_0000;
/// Sanity cap on the channel counts in the header.
const RTT_MAX_CHANNELS: u32 = 16;
/// Scan length when the target's SRAM size is unknown (the `_SEGGER_RTT` block lives in early .bss).
const RTT_DEFAULT_SCAN: u32 = 8 * 1024;
/// The WCH-Link's bulk read rejects/times-out on a very large single region, so read the scan
/// window in transfers this size (8 KiB is proven to work well within the transport timeout).
const RTT_READ_CHUNK: u32 = 8192;
/// Header: id[16] + max_up(4) + max_down(4); then `max_up` up descriptors, then the down ones.
const RTT_HEADER_LEN: u32 = 24;
/// One ring descriptor: name(+0) buffer(+4) size(+8) write_off(+12) read_off(+16) flags(+20).
const RTT_DESC_LEN: u32 = 24;
const RTT_DESC_WRITE_OFF: u32 = 12;
const RTT_DESC_READ_OFF: u32 = 16;

fn le32(b: &[u8], off: usize) -> u32 {
    u32::from_le_bytes([b[off], b[off + 1], b[off + 2], b[off + 3]])
}

/// A control block found in a RAM snapshot: its offset and channel counts.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct RttHeader {
    offset: usize,
    max_up: u32,
    max_down: u32,
}

/// The live fields of one ring descriptor.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct RingDesc {
    buffer: u32,
    size: u32,
    wr: u32,
    rd: u32,
}

impl RingDesc {
    /// Parse and validate a 24-byte descriptor (buffer in RAM, sane size, offsets in range).
    fn parse(d: &[u8]) -> Option<Self> {
        if d.len() < RTT_DESC_LEN as usize {
            return None;
        }
        let desc = RingDesc {
            buffer: le32(d, 4),
            size: le32(d, 8),
            wr: le32(d, 12),
            rd: le32(d, 16),
        };
        let ok = desc.size > 0
            && desc.size <= RTT_MAX_BUF
            && desc.wr < desc.size
            && desc.rd < desc.size
            && desc.buffer >= RTT_RAM_BASE;
        ok.then_some(desc)
    }

    fn used(&self) -> u32 {
        if self.wr >= self.rd {
            self.wr - self.rd
        } else {
            self.size - self.rd + self.wr
        }
    }
}

/// Find a SEGGER RTT control block in a RAM snapshot: the magic, a header with sane channel
/// counts, and an up[0] descriptor that validates. Validating rejects a stray copy of the magic
/// that lives in a ring buffer's own contents, not in a real block.
fn find_control_block(snap: &[u8]) -> Option<RttHeader> {
    let mut from = 0usize;
    while let Some(rel) = snap[from..]
        .windows(RTT_MAGIC.len())
        .position(|w| w == RTT_MAGIC)
    {
        let pos = from + rel;
        let hdr_end = pos + RTT_HEADER_LEN as usize;
        if hdr_end + RTT_DESC_LEN as usize <= snap.len() {
            let max_up = le32(snap, pos + 16);
            let max_down = le32(snap, pos + 20);
            if (1..=RTT_MAX_CHANNELS).contains(&max_up)
                && max_down <= RTT_MAX_CHANNELS
                && RingDesc::parse(&snap[hdr_end..]).is_some()
            {
                return Some(RttHeader {
                    offset: pos,
                    max_up,
                    max_down,
                });
            }
        }
        from = pos + 1;
    }
    None
}

// ---- sources ----

/// Target addresses of the RTT ring descriptors this session streams (channel 0 each way).
#[derive(Debug, Clone, Copy)]
pub(crate) struct RttChannels {
    up: u32,
    down: Option<u32>,
}

/// An opened runtime output source over the Debug Module.
pub(crate) enum DmiSource {
    Dmdata,
    Dmseq(Box<DmseqStream>),
    Rtt(RttChannels),
}

/// en: How long to read nothing that is a dmseq frame before saying there is no dmseq console on
/// this target. The spec leaves the delay to the host; long enough that a target which simply has
/// not printed yet (it still posts empty frames) is never reported, short enough to answer
/// "why is nothing coming out?" while the user is still looking.
/// ja: 「この target に dmseq console が無い」と報告するまでの時間(仕様は host 任せ)。まだ印字
/// していないだけの target(空フレームは出している)を誤報しない長さで、利用者が見ている間に
/// 「何も出ない」に答えられる短さ。
const DMSEQ_NO_CONSOLE_AFTER: Duration = Duration::from_secs(3);

/// One `dmseq` session: the framing state plus what the stream should tell the user about it.
pub(crate) struct DmseqStream {
    session: DmSeq,
    opened: Instant,
    no_console_reported: bool,
    notices: Vec<(&'static str, String)>,
}

pub(crate) enum OpenError {
    /// The source is not a DMI one (uart/sdi go through the probe's CDC port).
    NotDmi,
    /// No RTT control block appeared in the scanned RAM window.
    NoControlBlock {
        scan_len: u32,
    },
    Dmi(DmiError),
}

impl std::fmt::Display for OpenError {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            OpenError::NotDmi => f.write_str("not a DMI source"),
            OpenError::NoControlBlock { scan_len } => write!(
                f,
                "no SEGGER RTT control block in the first {scan_len} bytes of RAM (flash an RTT sketch first)"
            ),
            OpenError::Dmi(e) => write!(f, "{e}"),
        }
    }
}

impl From<DmiError> for OpenError {
    fn from(e: DmiError) -> Self {
        OpenError::Dmi(e)
    }
}

fn transport(e: ch32rv_wchlink::WchLinkError) -> DmiError {
    DmiError::Transport(e.to_string())
}

/// Read `len` bytes of target memory into one buffer, chunked so each transfer stays small. Stops
/// early (returning what it has) if a chunk fails, so a short read still lets the scan try.
fn read_region(session: &mut Session, base: u32, len: u32) -> Vec<u8> {
    let mut buf = Vec::with_capacity(len as usize);
    let mut off = 0u32;
    while off < len {
        let want = RTT_READ_CHUNK.min(len - off);
        match session.link().read_mem(base + off, want) {
            Ok(mut chunk) => {
                buf.append(&mut chunk);
                off += want;
            }
            Err(_) => break,
        }
    }
    buf
}

impl DmiSource {
    /// en: Open `source` on an attached session. `rtt` scans RAM for the control block (the
    /// target has usually been running since power-on, so begin() has published it; a few
    /// retries cover a very early attach) and leaves the core HALTED - the caller resumes it.
    /// `dmdata` needs no setup. A block with more than one channel per direction is reported as
    /// a warning: only channel 0 is streamed.
    /// ja: attach 済み session 上で source を開く。`rtt` は RAM を走査して control block を見つけ、
    /// core を halt のまま返す(resume は呼び出し側)。channel が複数ある block は警告して 0 番だけ流す。
    pub(crate) fn open(
        session: &mut Session,
        source: MonitorSource,
        warnings: &mut Vec<Warning>,
    ) -> Result<Self, OpenError> {
        match source {
            MonitorSource::Dmdata => {
                // Empty the mailbox before the core runs: a word left in data0 by attach/flash
                // reads as host input on the target (see `DebugModule::dmdata_clear`).
                session.dm().dmdata_clear()?;
                Ok(DmiSource::Dmdata)
            }
            MonitorSource::Dmseq => {
                // en: Nothing to set up, and in particular the mailbox is NOT cleared the way
                // `dmdata` clears it. dmseq's ownership rule is that the host writes DATA0 only
                // while bit 7 is set, and a conforming target reads a zero word as silence, not as
                // an answer - so clearing could only destroy a frame that was already posted and
                // make the target wait out its timeout (up to 1 s) before posting it again, with a
                // spurious "output was dropped" warning on the way. What a stale mailbox needs
                // instead is already in the framing: whatever the probe's attach left in DATA0
                // (measured: 0xe339e339 on a CH32V203, 0xffffffff elsewhere) fails the target's
                // answer check, so the target posts its frame again by itself. Verified on a
                // CH32V203 with a target left latched for 12 s: the session syncs and streams both
                // ways without the host writing anything first.
                // ja: 準備は無く、特に **`dmdata` のような mailbox クリアはしない**。dmseq の所有権
                // 規則では host は bit 7 が 1 のときしか DATA0 を書かず、規格どおりの target は 0 の
                // word を「沈黙」と読む(答えとは見ない)。したがってクリアは、既に出ているフレームを
                // 壊して target の timeout(最大 1 秒)を待たせ、偽の「出力が捨てられた」警告まで出す
                // だけになる。古い mailbox の始末は framing 側にある: probe の attach が DATA0 に
                // 残す値(実測 CH32V203 は 0xe339e339、他は 0xffffffff)は target の答え検査を通らない
                // ので、target が自分でフレームを出し直す。12 秒放置して latch させた CH32V203 で
                // 確認済み(host が先に何も書かなくても同期し、双方向に流れる)。
                Ok(DmiSource::Dmseq(Box::new(DmseqStream {
                    session: DmSeq::new(),
                    opened: Instant::now(),
                    no_console_reported: false,
                    notices: Vec::new(),
                })))
            }
            MonitorSource::Rtt => Self::open_rtt(session, warnings),
            MonitorSource::Uart | MonitorSource::Sdi | MonitorSource::FixtureUart => {
                Err(OpenError::NotDmi)
            }
        }
    }

    fn open_rtt(session: &mut Session, warnings: &mut Vec<Warning>) -> Result<Self, OpenError> {
        // How much RAM to scan for the control block: the target's SRAM (from the DB) or a default.
        let scan_len = {
            let db = session.db();
            match db.resolve_by_chip_id(session.attach.chip_id) {
                ch32rv_target::Resolution::Sku(s) if s.sram_bytes > 0 => {
                    s.sram_bytes.min(64 * 1024)
                }
                _ => RTT_DEFAULT_SCAN,
            }
        };
        let mut found = None;
        for _ in 0..10 {
            session.dm().halt()?;
            let snap = read_region(session, RTT_RAM_BASE, scan_len);
            if let Some(h) = find_control_block(&snap) {
                found = Some(h);
                break;
            }
            session.dm().resume()?;
            std::thread::sleep(Duration::from_millis(100));
        }
        let Some(h) = found else {
            return Err(OpenError::NoControlBlock { scan_len });
        };
        if h.max_up > 1 || h.max_down > 1 {
            warnings.push(Warning {
                code: "rtt-channels".to_owned(),
                msg: format!(
                    "the RTT control block has {} up / {} down channels; only channel 0 is streamed",
                    h.max_up, h.max_down
                ),
            });
        }
        let base = RTT_RAM_BASE + h.offset as u32;
        let up = base + RTT_HEADER_LEN;
        let down = (h.max_down >= 1).then_some(up + RTT_DESC_LEN * h.max_up);
        Ok(DmiSource::Rtt(RttChannels { up, down }))
    }

    /// en: Anything the source wants said to the user since the last poll (dmseq: the target's
    /// timeout, "no console here"). Taken, not peeked, so each notice is reported once.
    /// ja: 前回以降に source が利用者へ伝えたいこと(dmseq の TO / console 不在)。取り出したら消える。
    pub(crate) fn take_notices(&mut self) -> Vec<(&'static str, String)> {
        match self {
            DmiSource::Dmseq(st) => std::mem::take(&mut st.notices),
            DmiSource::Dmdata | DmiSource::Rtt(_) => Vec::new(),
        }
    }

    pub(crate) fn name(&self) -> &'static str {
        match self {
            DmiSource::Dmdata => "dmdata",
            DmiSource::Dmseq(_) => "dmseq",
            DmiSource::Rtt(_) => "rtt",
        }
    }

    /// How long to wait between polls that returned nothing.
    pub(crate) fn idle(&self) -> Duration {
        match self {
            DmiSource::Dmdata => Duration::from_millis(2),
            // The target waits 1 s for an answer once it has had one, and drops what it writes
            // after that, so the poll interval is what decides whether output survives.
            DmiSource::Dmseq(_) => Duration::from_millis(2),
            DmiSource::Rtt(_) => Duration::from_millis(50),
        }
    }

    /// One-line description for the human "monitor: ..." banner.
    pub(crate) fn describe(&self) -> String {
        match self {
            DmiSource::Dmdata => "dmdata (DMI mailbox, core runs)".to_owned(),
            DmiSource::Dmseq(_) => "dmseq (DMI mailbox, sequenced + CRC, core runs)".to_owned(),
            DmiSource::Rtt(ch) => format!(
                "rtt (RAM ring @ 0x{:08x}, core briefly halts per poll)",
                ch.up - RTT_HEADER_LEN
            ),
        }
    }

    /// en: One exchange: returns the target's output since the last poll and removes from the
    /// front of `input` whatever was handed to the target (dmdata: up to 3 bytes per target
    /// frame; rtt: as much as the down ring has room for).
    /// ja: 1 回の交換。前回以降の target 出力を返し、target へ渡せた分だけ `input` の先頭を消す。
    pub(crate) fn poll(
        &mut self,
        session: &mut Session,
        input: &mut Vec<u8>,
    ) -> Result<Vec<u8>, DmiError> {
        match self {
            DmiSource::Dmdata => {
                let r = session.dm().dmdata_poll(&input[..input.len().min(3)])?;
                input.drain(..r.sent);
                Ok(r.received)
            }
            DmiSource::Dmseq(st) => {
                let take = input.len().min(dmseq::MAX_HOST_PAYLOAD);
                let r = session.dm().dmseq_poll(&mut st.session, &input[..take])?;
                input.drain(..r.sent);
                if r.timed_out {
                    st.notices.push((
                        "dmseq-target-timeout",
                        "the target stopped waiting for an answer; output written while nobody \
                         answered was dropped"
                            .to_owned(),
                    ));
                }
                if !st.session.synced()
                    && !st.no_console_reported
                    && st.opened.elapsed() >= DMSEQ_NO_CONSOLE_AFTER
                {
                    st.no_console_reported = true;
                    st.notices.push((
                        "dmseq-no-console",
                        format!(
                            "no dmseq frame in {} s ({} word(s) rejected): this target may not \
                             have a dmseq console (SerialDMSeq) - `--source dmdata` speaks the \
                             older framing",
                            DMSEQ_NO_CONSOLE_AFTER.as_secs(),
                            st.session.invalid
                        ),
                    ));
                }
                Ok(r.received)
            }
            DmiSource::Rtt(ch) => {
                let ch = *ch;
                // RAM access needs the hart halted. Leave an already-halted core (a semihosting
                // stop being serviced by `run`) alone; otherwise halt for the exchange and resume.
                let was_halted = session.dm().is_halted()?;
                if !was_halted {
                    session.dm().halt()?;
                }
                let result = rtt_exchange(session, ch, input);
                if !was_halted {
                    let _ = session.dm().resume();
                }
                result
            }
        }
    }
}

/// Drain up[0] (acknowledging by moving its read offset) and push host input into down[0].
fn rtt_exchange(
    session: &mut Session,
    ch: RttChannels,
    input: &mut Vec<u8>,
) -> Result<Vec<u8>, DmiError> {
    let raw = session
        .link()
        .read_mem(ch.up, RTT_DESC_LEN)
        .map_err(transport)?;
    let Some(up) = RingDesc::parse(&raw) else {
        // Not ready or garbage: let the target run and retry next poll.
        return Ok(Vec::new());
    };
    let mut out = Vec::new();
    if up.wr != up.rd {
        let link = session.link();
        if up.wr > up.rd {
            out = link
                .read_mem(up.buffer + up.rd, up.wr - up.rd)
                .map_err(transport)?;
        } else {
            // Wrapped: [rd, size) then [0, wr).
            out = link
                .read_mem(up.buffer + up.rd, up.size - up.rd)
                .map_err(transport)?;
            if up.wr > 0 {
                out.extend(link.read_mem(up.buffer, up.wr).map_err(transport)?);
            }
        }
        // Tell the target we drained: up[0].read_off = write_off.
        session.dm().write_mem32(ch.up + RTT_DESC_READ_OFF, up.wr)?;
    }
    if let Some(down_addr) = ch.down
        && !input.is_empty()
    {
        let raw = session
            .link()
            .read_mem(down_addr, RTT_DESC_LEN)
            .map_err(transport)?;
        if let Some(down) = RingDesc::parse(&raw) {
            // A ring spends one slot telling full from empty.
            let room = (down.size - 1 - down.used()) as usize;
            let n = input.len().min(room);
            if n > 0 {
                let first = n.min((down.size - down.wr) as usize);
                let mut dm = session.dm();
                dm.write_mem(down.buffer + down.wr, &input[..first])?;
                if n > first {
                    dm.write_mem(down.buffer, &input[first..n])?;
                }
                // Publish the offset only after the bytes are in place (the target reads it).
                dm.write_mem32(
                    down_addr + RTT_DESC_WRITE_OFF,
                    (down.wr + n as u32) % down.size,
                )?;
                input.drain(..n);
            }
        }
    }
    Ok(out)
}

/// Hand the source's notices to the sink. Called from every streaming loop (`monitor` / `run`).
pub(crate) fn report_notices(src: &mut DmiSource, sink: &mut Sink) {
    for (code, msg) in src.take_notices() {
        sink.warn(code, &msg);
    }
}

/// Exit code class for a DMI failure while streaming (docs/cli.ja.md §3.6: 40 either way, the
/// JSON `kind` tells a true timeout from a failed transfer).
pub(crate) fn dmi_error_kind(e: &DmiError) -> ErrorKind {
    match e {
        DmiError::Timeout => ErrorKind::TransportTimeout,
        _ => ErrorKind::TransferFailed,
    }
}

// ---- output sink ----

/// en: Where the target's output goes. Human mode: raw bytes straight to stdout. `--json`: one
/// `output` NDJSON event per chunk on stderr, split on UTF-8 character boundaries so a
/// multi-byte character straddling two 7-byte dmdata frames is not mangled.
/// ja: target 出力の行き先。human は生 byte を stdout へ。`--json` は stderr へ `output` event
/// (UTF-8 文字境界で分割。dmdata の 7 byte frame をまたぐ多バイト文字を壊さない)。
pub(crate) enum Sink {
    Raw,
    Json {
        source: &'static str,
        pending: Vec<u8>,
    },
}

impl Sink {
    pub(crate) fn new(cli: &Cli, source: &'static str) -> Self {
        if cli.json {
            Sink::Json {
                source,
                pending: Vec::new(),
            }
        } else {
            Sink::Raw
        }
    }

    pub(crate) fn write(&mut self, bytes: &[u8]) {
        if bytes.is_empty() {
            return;
        }
        match self {
            Sink::Raw => {
                let mut out = std::io::stdout().lock();
                let _ = out.write_all(bytes);
                let _ = out.flush();
            }
            Sink::Json { source, pending } => {
                pending.extend_from_slice(bytes);
                let text = take_text(pending);
                emit_output(source, text);
            }
        }
    }

    /// en: Tell the user something about the stream itself (not target output): human mode gets
    /// the same `warning[code]: msg` line the rest of the CLI uses, `--json` a `warn` event on
    /// stderr, where the runtime output already goes.
    /// ja: stream 自体についての通知(target 出力ではない)。human は CLI 共通の
    /// `warning[code]: msg`、`--json` は stderr の `warn` event。
    pub(crate) fn warn(&mut self, code: &str, msg: &str) {
        match self {
            Sink::Raw => eprintln!("warning[{code}]: {msg}"),
            Sink::Json { .. } => {
                let ev = Event::Warn {
                    code: code.to_owned(),
                    msg: msg.to_owned(),
                };
                if let Ok(line) = serde_json::to_string(&ev) {
                    let mut err = std::io::stderr().lock();
                    let _ = writeln!(err, "{line}");
                }
            }
        }
    }

    /// Flush an incomplete trailing sequence (lossily) when the stream ends.
    pub(crate) fn finish(&mut self) {
        if let Sink::Json { source, pending } = self
            && !pending.is_empty()
        {
            let text = String::from_utf8_lossy(pending).into_owned();
            pending.clear();
            emit_output(source, text);
        }
    }
}

fn emit_output(source: &str, data: String) {
    if data.is_empty() {
        return;
    }
    let ev = Event::Output {
        source: source.to_owned(),
        data,
    };
    if let Ok(line) = serde_json::to_string(&ev) {
        let mut err = std::io::stderr().lock();
        let _ = writeln!(err, "{line}");
    }
}

/// Take the longest run of complete UTF-8 characters off the front of `buf`, replacing invalid
/// sequences with U+FFFD. An incomplete multi-byte sequence at the end (at most 3 bytes) stays in
/// `buf` for the next chunk to complete.
fn take_text(buf: &mut Vec<u8>) -> String {
    let mut out = String::new();
    let mut i = 0usize;
    loop {
        match std::str::from_utf8(&buf[i..]) {
            Ok(s) => {
                out.push_str(s);
                i = buf.len();
                break;
            }
            Err(e) => {
                let ok = e.valid_up_to();
                out.push_str(&String::from_utf8_lossy(&buf[i..i + ok]));
                i += ok;
                match e.error_len() {
                    Some(bad) => {
                        out.push('\u{FFFD}');
                        i += bad;
                    }
                    // Incomplete sequence at the end: keep it for the next chunk.
                    None => break,
                }
            }
        }
    }
    buf.drain(..i);
    out
}

// ---- host input ----

/// en: Read `r` (stdin, a socket) on its own thread and hand chunks to the returned channel;
/// ends quietly at EOF or error. The main loop pulls from it between polls so a blocking read
/// never stalls the target exchange.
/// ja: `r`(stdin / socket)を別 thread で読み、chunk を channel に渡す。EOF/エラーで静かに終わる。
pub(crate) fn spawn_reader(mut r: impl Read + Send + 'static) -> Receiver<Vec<u8>> {
    let (tx, rx) = mpsc::channel();
    std::thread::spawn(move || {
        let mut buf = [0u8; 256];
        loop {
            match r.read(&mut buf) {
                Ok(0) => break,
                Ok(n) => {
                    if tx.send(buf[..n].to_vec()).is_err() {
                        break;
                    }
                }
                Err(e) if e.kind() == std::io::ErrorKind::Interrupted => {}
                Err(_) => break,
            }
        }
    });
    rx
}

/// Move everything the reader has produced so far into `into` (non-blocking).
pub(crate) fn drain_input(rx: &Receiver<Vec<u8>>, into: &mut Vec<u8>) {
    while let Ok(chunk) = rx.try_recv() {
        into.extend_from_slice(&chunk);
    }
}

/// en: Stream `src` to `sink` and `input` to the target until `deadline` (None = until Ctrl-C).
/// The caller has resumed the core. Errors are DMI failures the caller turns into an exit code.
/// ja: `deadline` まで(None なら Ctrl-C まで)`src`→`sink`、`input`→target を流す。
pub(crate) fn stream(
    session: &mut Session,
    src: &mut DmiSource,
    sink: &mut Sink,
    input: &Receiver<Vec<u8>>,
    deadline: Option<Instant>,
) -> Result<(), DmiError> {
    let mut pending = Vec::new();
    loop {
        if let Some(dl) = deadline
            && Instant::now() >= dl
        {
            return Ok(());
        }
        drain_input(input, &mut pending);
        let out = src.poll(session, &mut pending)?;
        report_notices(src, sink);
        if out.is_empty() {
            std::thread::sleep(src.idle());
        } else {
            sink.write(&out);
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    /// Build a RAM snapshot with a control block at `at` (`max_up`/`max_down` channels) whose
    /// up[0] descriptor is valid and whose down[0], if any, sits after all the up descriptors.
    fn snapshot(at: usize, max_up: u32, max_down: u32, up0: RingDesc) -> Vec<u8> {
        let descs = (max_up + max_down) as usize;
        let mut ram = vec![0u8; at + 24 + 24 * descs + 8];
        ram[at..at + RTT_MAGIC.len()].copy_from_slice(RTT_MAGIC);
        ram[at + 16..at + 20].copy_from_slice(&max_up.to_le_bytes());
        ram[at + 20..at + 24].copy_from_slice(&max_down.to_le_bytes());
        let d = at + 24;
        ram[d + 4..d + 8].copy_from_slice(&up0.buffer.to_le_bytes());
        ram[d + 8..d + 12].copy_from_slice(&up0.size.to_le_bytes());
        ram[d + 12..d + 16].copy_from_slice(&up0.wr.to_le_bytes());
        ram[d + 16..d + 20].copy_from_slice(&up0.rd.to_le_bytes());
        ram
    }

    fn valid_up0() -> RingDesc {
        RingDesc {
            buffer: RTT_RAM_BASE + 0x100,
            size: 256,
            wr: 10,
            rd: 0,
        }
    }

    #[test]
    fn finds_valid_control_block_with_counts() {
        let ram = snapshot(64, 1, 1, valid_up0());
        assert_eq!(
            find_control_block(&ram),
            Some(RttHeader {
                offset: 64,
                max_up: 1,
                max_down: 1
            })
        );
    }

    #[test]
    fn down_descriptor_follows_all_up_descriptors() {
        // With 2 up channels, down[0] is the third descriptor after the header.
        let h = RttHeader {
            offset: 0,
            max_up: 2,
            max_down: 1,
        };
        let base = RTT_RAM_BASE + h.offset as u32;
        let up = base + RTT_HEADER_LEN;
        let down = up + RTT_DESC_LEN * h.max_up;
        assert_eq!(down, RTT_RAM_BASE + 24 + 48);
    }

    #[test]
    fn rejects_bogus_channel_counts() {
        // max_up = 0 (no up channel to stream) and an absurd count are both garbage.
        assert_eq!(find_control_block(&snapshot(0, 0, 1, valid_up0())), None);
        assert_eq!(find_control_block(&snapshot(0, 1000, 1, valid_up0())), None);
    }

    #[test]
    fn skips_magic_with_bogus_descriptor() {
        // A stray "SEGGER RTT" in buffer contents: the descriptor after it is garbage (size huge),
        // so it must not be mistaken for a real block.
        let bogus = RingDesc {
            buffer: 0,
            size: 0xFFFF_FFFF,
            wr: 0,
            rd: 0,
        };
        assert_eq!(find_control_block(&snapshot(64, 1, 1, bogus)), None);
        assert_eq!(find_control_block(&[0u8; 64]), None);
    }

    #[test]
    fn ring_used_handles_wrap() {
        let d = RingDesc {
            buffer: RTT_RAM_BASE,
            size: 16,
            wr: 2,
            rd: 14,
        };
        assert_eq!(d.used(), 4);
        assert_eq!(RingDesc { wr: 9, rd: 3, ..d }.used(), 6);
    }

    #[test]
    fn take_text_keeps_incomplete_tail_across_chunks() {
        // "こ" is e3 81 93; a 7-byte dmdata frame can end after e3 81.
        let mut buf = b"ab\xe3\x81".to_vec();
        assert_eq!(take_text(&mut buf), "ab");
        assert_eq!(buf, b"\xe3\x81");
        buf.extend_from_slice(b"\x93!");
        assert_eq!(take_text(&mut buf), "こ!");
        assert!(buf.is_empty());
    }

    #[test]
    fn take_text_replaces_invalid_bytes() {
        let mut buf = b"x\xffy".to_vec();
        assert_eq!(take_text(&mut buf), "x\u{FFFD}y");
        assert!(buf.is_empty());
    }

    #[test]
    fn le32_reads_little_endian() {
        assert_eq!(le32(&[0x78, 0x56, 0x34, 0x12], 0), 0x1234_5678);
    }
}