hisui 2025.3.3

Recording Composition Tool Hisui
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
use std::collections::VecDeque;
use std::sync::{Arc, OnceLock, mpsc};

use orfail::OrFail;
use shiguredo_mp4::boxes::SampleEntry;

use crate::{
    encoder::VideoEncoderOptions,
    types::CodecName,
    video::{VideoFormat, VideoFrame},
    video_av1, video_h264, video_h265,
};

/// callback で受け取った圧縮フレームと、user_data として渡した入力側メタデータ
#[derive(Debug)]
struct EncodedFrameWithMeta {
    data: Vec<u8>,
    keyframe: bool,
    input_frame: VideoFrame,
}

/// callback スレッドが参照する遅延確定コンテキスト。
///
/// AV1 の `av1_sequence_header` は `Encoder::new` 後の `get_sequence_params()` でしか得られないが、
/// `Encoder::new` は handler を消費する API のため、handler にはあらかじめ `OnceLock` だけを
/// キャプチャさせておき、`Encoder::new` の直後に `HandlerContext` を確定して `set()` する。
struct HandlerContext {
    /// AV1 のキーフレームに Sequence Header OBU を付与するためのバイト列。
    /// H.264 / H.265 では意味を持たないため `None`
    av1_sequence_header: Option<Vec<u8>>,
}

type HandlerContextSlot = Arc<OnceLock<HandlerContext>>;

/// callback スレッドから呼ばれる Handler をラップした型
type NvcodecHandler = shiguredo_nvcodec::FnEncodeHandler<VideoFrame, shiguredo_nvcodec::Error>;

/// callback スレッドから main スレッドへ送るメッセージ。
///
/// Ok / Err を混ぜて単一 channel で流し、`drain_rx()` で受信側が振り分ける。
/// エラー発生後は上位が即停止する前提のため、エラーメッセージがキュー内のどこにあっても構わない
type EncodeMessage = Result<EncodedFrameWithMeta, orfail::Failure>;

#[derive(Debug)]
pub struct NvcodecEncoder {
    inner: shiguredo_nvcodec::Encoder<NvcodecHandler>,
    rx: mpsc::Receiver<EncodeMessage>,
    /// `rx` から `drain_rx()` で吸い上げた成功フレームを保持する
    ok_buffer: VecDeque<EncodedFrameWithMeta>,
    /// `rx` から吸い上げた初回エラー。以降のエラーは破棄する
    pending_error: Option<orfail::Failure>,
    encoded_format: VideoFormat,
    /// 最初の出力フレームに sample_entry を載せるために保持する。一度 take() したら以降は None のまま
    sample_entry: Option<SampleEntry>,
    /// `shiguredo_nvcodec::Encoder` 内でまだ drain されていない frame 数の推定値。
    /// `encode()` の呼び出しで +1、`inner.flush()` の完了で 0 にリセットする。
    /// `next_encoded_frame()` の pop は NVENC の外 (`rx` から吸い上げた buffer) の話なのでここでは触らない。
    in_flight: usize,
    /// `encode()` 送信直前に許容する `in_flight` の上限値。
    ///
    /// `shiguredo_nvcodec::Encoder` の worker は内部で `n_encoder_buffer = frame_interval_p + 3` で
    /// 送信可能な最大 in-flight を管理し、超えた frame を `"encoder buffer is full"` で捨てる。
    /// この上限に達する前に `inner.flush()` で drain する必要があるため、上限は `n_encoder_buffer - 1`
    /// (= `frame_interval_p + 2`) にする。`EncoderConfig::frame_interval_p` は encoder ごとに固定なので、
    /// build 時に算出して保持する
    in_flight_limit: usize,
}

impl NvcodecEncoder {
    pub fn new_h264(options: &VideoEncoderOptions) -> orfail::Result<Self> {
        let width = options.width.get();
        let height = options.height.get();
        log::debug!("create nvcodec(H264) encoder: {}x{}", width, height);

        let mut config = options.encode_params.nvcodec_h264.clone();
        Self::override_config_size_and_bitrate(&mut config, options);
        log::debug!("nvcodec h264 encoder config: {config:?}");

        Self::build_encoder(config, VideoFormat::H264, |seq_params| {
            let entry =
                video_h264::h264_sample_entry_from_annexb(width, height, &seq_params).or_fail()?;
            Ok((
                entry,
                HandlerContext {
                    av1_sequence_header: None,
                },
            ))
        })
    }

    pub fn new_h265(options: &VideoEncoderOptions) -> orfail::Result<Self> {
        let width = options.width.get();
        let height = options.height.get();
        log::debug!("create nvcodec(H265) encoder: {}x{}", width, height);

        let mut config = options.encode_params.nvcodec_h265.clone();
        Self::override_config_size_and_bitrate(&mut config, options);
        log::debug!("nvcodec h265 encoder config: {config:?}");

        let frame_rate = options.frame_rate;
        Self::build_encoder(config, VideoFormat::H265, move |seq_params| {
            let entry = video_h265::h265_sample_entry_from_annexb(
                width,
                height,
                frame_rate,
                &seq_params,
            )
            .or_fail()?;
            Ok((
                entry,
                HandlerContext {
                    av1_sequence_header: None,
                },
            ))
        })
    }

    pub fn new_av1(options: &VideoEncoderOptions) -> orfail::Result<Self> {
        let width = options.width;
        let height = options.height;
        log::debug!(
            "create nvcodec(AV1) encoder: {}x{}",
            width.get(),
            height.get()
        );

        let mut config = options.encode_params.nvcodec_av1.clone();
        Self::override_config_size_and_bitrate(&mut config, options);
        log::debug!("nvcodec av1 encoder config: {config:?}");

        Self::build_encoder(config, VideoFormat::Av1, move |seq_params| {
            // NVENC SDK 13.0 のドキュメント (https://docs.nvidia.com/video-technologies/video-codec-sdk/13.0/nvenc-video-encoder-api-prog-guide/index.html#retrieving-sequence-parameters)
            // には以下の記載がある:
            //   "By default, SPS/PPS and Sequence Header OBU data will be attached to every IDR frame and Key frame for H.264/HEVC and AV1 respectively."
            //
            // しかし実際には、AV1 の場合、最初のキーフレームにのみ Sequence Header OBU が付与され、
            // 二番目以降のキーフレームには含まれない。これにより、二番目以降のキーフレームからシークすると、
            // デコーダが解像度やプロファイルなどの情報を取得できず、映像が再生できない問題が発生する。
            //
            // そのため、ここで Sequence Header OBU を get_sequence_params() で取得して保持しておき、
            // キーフレームのエンコード時に Sequence Header が含まれていない場合は明示的に付与する。
            let entry = video_av1::av1_sample_entry(width, height, &seq_params);
            Ok((
                entry,
                HandlerContext {
                    av1_sequence_header: Some(seq_params),
                },
            ))
        })
    }

    /// codec 別 config を受けて、handler 準備 -> Encoder::new -> sample_entry 確定までの共通シーケンスを実行する
    fn build_encoder(
        config: shiguredo_nvcodec::EncoderConfig,
        encoded_format: VideoFormat,
        make_context: impl FnOnce(Vec<u8>) -> orfail::Result<(SampleEntry, HandlerContext)>,
    ) -> orfail::Result<Self> {
        let context_slot: HandlerContextSlot = Arc::new(OnceLock::new());
        let (tx, rx) = mpsc::channel();

        let in_flight_limit = config.frame_interval_p as usize + 2;

        let handler = build_handler(tx, context_slot.clone(), encoded_format);
        let inner = shiguredo_nvcodec::Encoder::new(config, handler).or_fail()?;

        let seq_params = inner.get_sequence_params().or_fail()?;
        let (sample_entry, context) = make_context(seq_params).or_fail()?;
        context_slot
            .set(context)
            .ok()
            .expect("BUG: HandlerContext must not be set before Encoder::new returns");

        Ok(Self {
            inner,
            rx,
            ok_buffer: VecDeque::new(),
            pending_error: None,
            encoded_format,
            sample_entry: Some(sample_entry),
            in_flight: 0,
            in_flight_limit,
        })
    }

    fn override_config_size_and_bitrate(
        config: &mut shiguredo_nvcodec::EncoderConfig,
        options: &VideoEncoderOptions,
    ) {
        config.width = options.width.get() as u32;
        config.height = options.height.get() as u32;
        config.framerate_num = options.frame_rate.numerator.get() as u32;
        config.framerate_den = options.frame_rate.denumerator.get() as u32;
        config.average_bitrate = Some(options.bitrate as u32);
    }

    pub fn encode(&mut self, frame: &VideoFrame) -> orfail::Result<()> {
        self.take_pending_error()?;

        (frame.format == VideoFormat::I420).or_fail()?;

        // 上限に達したら inner.flush() で全 in-flight を drain してから次の送信に進む
        // (batched flush 方式)。shiguredo_nvcodec は送信直前の in-flight が n_encoder_buffer 未満で
        // ないと "encoder buffer is full" で frame を捨てるため、送信前に必ず上限以下に抑える。
        if self.in_flight >= self.in_flight_limit {
            self.inner.flush().or_fail()?;
            self.in_flight = 0;
        }

        // I420 から NV12 への変換
        let width = frame.width;
        let height = frame.height;
        let (y_plane, u_plane, v_plane) = frame.as_yuv_planes().or_fail()?;

        // NV12 用のバッファを確保
        let y_size = width * height;
        let uv_width = width.div_ceil(2);
        let uv_height = height.div_ceil(2);
        let uv_size = uv_width * uv_height * 2; // U と V が交互に配置されているため
        let total_size = y_size + uv_size;

        let mut nv12_data = vec![0u8; total_size];
        let (nv12_y, nv12_uv) = nv12_data.split_at_mut(y_size);

        // libyuv を使って I420 から NV12 に変換
        let src = shiguredo_libyuv::I420Planes {
            y: y_plane,
            y_stride: width,
            u: u_plane,
            u_stride: uv_width,
            v: v_plane,
            v_stride: uv_width,
        };

        let mut dst = shiguredo_libyuv::Nv12PlanesMut {
            y: nv12_y,
            y_stride: width,
            uv: nv12_uv,
            uv_stride: width,
        };

        let size = shiguredo_libyuv::ImageSize::new(width, height);
        shiguredo_libyuv::i420_to_nv12(&src, &mut dst, size).or_fail()?;

        // callback スレッドで input_frame のメタデータを復元するために軽量な to_stripped() を渡す
        let encode_options = shiguredo_nvcodec::EncodeOptions {
            force_intra: false,
            force_idr: false,
            output_spspps: false,
        };
        self.inner
            .encode(&nv12_data, &encode_options, frame.to_stripped())
            .or_fail()?;
        // encode() は job_tx.send() で fire-and-forget なので、NVENC 内 pending が +1 された分を
        // 追跡する。次回 encode() の冒頭で LIMIT を超えていれば flush で drain する。
        self.in_flight += 1;
        Ok(())
    }

    pub fn finish(&mut self) -> orfail::Result<()> {
        // encode() は上限到達時のみ flush する batched flush 方式なので、
        // finish() 時点では in_flight が 0〜in_flight_limit の範囲を取りうる。
        // 残 in-flight を drain して EOS 前の callback を発火させる。
        self.inner.flush().or_fail()?;
        self.in_flight = 0;
        self.take_pending_error()?;
        Ok(())
    }

    /// callback スレッドから受け取ったメッセージを `ok_buffer` / `pending_error` に振り分ける
    fn drain_rx(&mut self) {
        while let Ok(msg) = self.rx.try_recv() {
            match msg {
                Ok(encoded) => self.ok_buffer.push_back(encoded),
                Err(err) => {
                    // 初回エラーだけ保持、以降は破棄
                    self.pending_error.get_or_insert(err);
                }
            }
        }
    }

    /// callback スレッドで発生したエラーがあれば取り出して返す
    fn take_pending_error(&mut self) -> orfail::Result<()> {
        self.drain_rx();
        if let Some(err) = self.pending_error.take() {
            return Err(err);
        }
        Ok(())
    }

    pub fn next_encoded_frame(&mut self) -> Option<VideoFrame> {
        self.drain_rx();
        let encoded = self.ok_buffer.pop_front()?;
        Some(VideoFrame {
            source_id: encoded.input_frame.source_id.clone(),
            data: encoded.data,
            format: self.encoded_format,
            keyframe: encoded.keyframe,
            width: encoded.input_frame.width,
            height: encoded.input_frame.height,
            timestamp: encoded.input_frame.timestamp,
            duration: encoded.input_frame.duration,
            sample_entry: self.sample_entry.take(),
        })
    }

    pub fn codec(&self) -> CodecName {
        self.encoded_format.codec_name().expect("infallible")
    }
}

/// shiguredo_nvcodec::Encoder が消費する handler を構築する
fn build_handler(
    tx: mpsc::Sender<EncodeMessage>,
    context_slot: HandlerContextSlot,
    encoded_format: VideoFormat,
) -> NvcodecHandler {
    shiguredo_nvcodec::FnEncodeHandler::new(move |result| {
        handle_encode_callback(&tx, &context_slot, encoded_format, result);
    })
}

/// callback スレッドから呼ばれるコールバック本体
fn handle_encode_callback(
    tx: &mpsc::Sender<EncodeMessage>,
    context_slot: &OnceLock<HandlerContext>,
    encoded_format: VideoFormat,
    result: std::result::Result<
        shiguredo_nvcodec::EncodedFrame<VideoFrame>,
        shiguredo_nvcodec::Error,
    >,
) {
    let message = match result {
        Ok(encoded_frame) => {
            let context = context_slot
                .get()
                .expect("BUG: HandlerContext must be set before first encode() call");
            let keyframe = matches!(
                encoded_frame.picture_type(),
                shiguredo_nvcodec::PictureType::I | shiguredo_nvcodec::PictureType::Idr
            );
            let (data, input_frame) = encoded_frame.into_parts();
            convert_encoded_data(encoded_format, data, keyframe, context).map(|frame_data| {
                EncodedFrameWithMeta {
                    data: frame_data,
                    keyframe,
                    input_frame,
                }
            })
        }
        Err(err) => Err(orfail::Failure::new(format!("nvcodec encode error: {err}"))),
    };
    // Receiver drop 時のみ失敗するが、そのタイミングでは Encoder も drop 目前なので無視
    let _ = tx.send(message);
}

/// エンコード出力フレームを VideoFormat に応じて MP4 に載せる形式へ変換する
fn convert_encoded_data(
    encoded_format: VideoFormat,
    data: Vec<u8>,
    keyframe: bool,
    context: &HandlerContext,
) -> orfail::Result<Vec<u8>> {
    if encoded_format == VideoFormat::Av1 {
        // encoded_format == Av1 の分岐に入るのは new_av1 経路のみ。そこでは
        // make_context が Some(seq_params) を確定して slot に set() する契約
        let seq_header = context
            .av1_sequence_header
            .as_deref()
            .expect("BUG: AV1 encoder must have av1_sequence_header set");
        Ok(prepend_av1_sequence_header_if_needed(
            data, keyframe, seq_header,
        ))
    } else {
        convert_annexb_to_mp4(&data)
    }
}

/// AV1 のキーフレームで Sequence Header OBU が欠落している場合のみ、先頭に付与して返す
fn prepend_av1_sequence_header_if_needed(
    data: Vec<u8>,
    keyframe: bool,
    seq_header: &[u8],
) -> Vec<u8> {
    if !keyframe || has_sequence_header(&data) {
        return data;
    }
    log::debug!(
        "prepending Sequence Header OBU to AV1 keyframe (seq_header: {} bytes, frame: {} bytes)",
        seq_header.len(),
        data.len()
    );
    let mut new_data = Vec::with_capacity(seq_header.len() + data.len());
    new_data.extend_from_slice(seq_header);
    new_data.extend_from_slice(&data);
    new_data
}

/// AV1 ペイロードの先頭に Sequence Header OBU が含まれているかチェック
fn has_sequence_header(data: &[u8]) -> bool {
    // 最低限 OBU Header の 1 バイトが必要
    if data.is_empty() {
        return false;
    }

    // 先頭の OBU Header を解析
    // obu_header のビット構成:
    //   - bit 0: obu_forbidden_bit (常に0)
    //   - bit 1-4: obu_type
    //   - bit 5: obu_extension_flag
    //   - bit 6: obu_has_size_field
    //   - bit 7: obu_reserved_1bit
    let obu_header = data[0];
    let obu_has_extension = (obu_header & 0b0010_0000) != 0;

    // OBU Extension が存在する場合は 2 バイト目も必要
    if obu_has_extension && data.len() < 2 {
        return false;
    }

    let obu_type = (obu_header >> 3) & 0x0F;

    // 先頭が Sequence Header (type=1) なら true
    obu_type == 1
}

/// Annex B 形式から MP4 形式への変換
///
/// Annex B 形式: スタートコード (0x00000001 or 0x000001) + NALU データ
/// MP4 形式: サイズ (4バイト) + NALU データ
fn convert_annexb_to_mp4(annexb_data: &[u8]) -> orfail::Result<Vec<u8>> {
    let mut mp4_data = Vec::new();
    let mut pos = 0;

    while pos < annexb_data.len() {
        // スタートコードを探す (0x00000001 or 0x000001)
        let start_code_len =
            if pos + 4 <= annexb_data.len() && annexb_data[pos..pos + 4] == [0, 0, 0, 1] {
                4
            } else if pos + 3 <= annexb_data.len() && annexb_data[pos..pos + 3] == [0, 0, 1] {
                3
            } else if pos == 0 {
                return Err(orfail::Failure::new("No start code found at beginning"));
            } else {
                break;
            };

        pos += start_code_len;

        // 次のスタートコードまたはデータ終端を探す
        let nalu_start = pos;
        let mut nalu_end = annexb_data.len();

        for i in (pos + 3)..annexb_data.len() {
            if i + 4 <= annexb_data.len() && annexb_data[i..i + 4] == [0, 0, 0, 1] {
                nalu_end = i;
                break;
            }
            if i + 3 <= annexb_data.len() && annexb_data[i..i + 3] == [0, 0, 1] {
                nalu_end = i;
                break;
            }
        }

        let nalu_size = nalu_end - nalu_start;

        // MP4 形式: 4 バイトのサイズ + NALU データ
        mp4_data.extend_from_slice(&(nalu_size as u32).to_be_bytes());
        mp4_data.extend_from_slice(&annexb_data[nalu_start..nalu_end]);

        pos = nalu_end;
    }

    Ok(mp4_data)
}

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

    // AV1 OBU Header の bit 配置:
    //   - bit 7 (MSB): obu_forbidden_bit
    //   - bit 6-3:     obu_type (4 bits)
    //   - bit 2:       obu_extension_flag
    //   - bit 1:       obu_has_size_field
    //   - bit 0 (LSB): obu_reserved_1bit

    #[test]
    fn has_sequence_header_sh() {
        // obu_type = 1 (Sequence Header)
        assert!(has_sequence_header(&[0x08]));
    }

    #[test]
    fn has_sequence_header_td() {
        // obu_type = 2 (Temporal Delimiter)
        assert!(!has_sequence_header(&[0x10]));
    }

    #[test]
    fn has_sequence_header_frame_header() {
        // obu_type = 6 (Frame Header)
        assert!(!has_sequence_header(&[0x30]));
    }

    #[test]
    fn has_sequence_header_empty() {
        assert!(!has_sequence_header(&[]));
    }

    #[test]
    fn convert_annexb_to_mp4_4byte_start_code() {
        let annexb = [0, 0, 0, 1, 0xAA, 0xBB, 0xCC, 0xDD];
        let mp4 = convert_annexb_to_mp4(&annexb).expect("should succeed");
        assert_eq!(mp4, [0, 0, 0, 4, 0xAA, 0xBB, 0xCC, 0xDD]);
    }

    #[test]
    fn convert_annexb_to_mp4_3byte_start_code() {
        let annexb = [0, 0, 1, 0xAA, 0xBB, 0xCC];
        let mp4 = convert_annexb_to_mp4(&annexb).expect("should succeed");
        assert_eq!(mp4, [0, 0, 0, 3, 0xAA, 0xBB, 0xCC]);
    }

    #[test]
    fn convert_annexb_to_mp4_multiple_nalus() {
        let annexb = [
            0, 0, 0, 1, 0xAA, 0xBB, 0xCC, 0xDD, // NALU1
            0, 0, 0, 1, 0xEE, 0xFF, // NALU2
        ];
        let mp4 = convert_annexb_to_mp4(&annexb).expect("should succeed");
        assert_eq!(
            mp4,
            [
                0, 0, 0, 4, 0xAA, 0xBB, 0xCC, 0xDD, // NALU1
                0, 0, 0, 2, 0xEE, 0xFF, // NALU2
            ]
        );
    }

    #[test]
    fn convert_annexb_to_mp4_no_start_code_at_beginning() {
        let annexb = [0xAA, 0xBB];
        assert!(convert_annexb_to_mp4(&annexb).is_err());
    }

    #[test]
    fn convert_annexb_to_mp4_empty() {
        let annexb: [u8; 0] = [];
        let mp4 = convert_annexb_to_mp4(&annexb).expect("should succeed");
        assert!(mp4.is_empty());
    }
}