package com.inspect.nvr.hik.service; import org.junit.Test; import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; public class HikFfmpegH264PublisherTest { @Test public void usesPipeInputAndPublishesToSessionRtmpUrl() { CapturingProcessFactory processFactory = new CapturingProcessFactory(new TestProcess()); HikFfmpegH264Publisher publisher = new HikFfmpegH264Publisher( processFactory, "C:/tools/ffmpeg/bin/ffmpeg.exe", "rtmp://127.0.0.1:1935", 4); publisher.start("a1-1"); assertTrue(processFactory.command().contains("pipe:0")); assertTrue(processFactory.command().contains("rtmp://127.0.0.1:1935/playback/a1-1")); publisher.close(); } @Test public void doesNotUseArrivalWallClockForRawH264Timestamps() { CapturingProcessFactory processFactory = new CapturingProcessFactory(new TestProcess()); HikFfmpegH264Publisher publisher = new HikFfmpegH264Publisher( processFactory, "ffmpeg", "rtmp://127.0.0.1:1935", 4); publisher.start("a1-1"); assertFalse(processFactory.command().contains("-use_wallclock_as_timestamps")); publisher.close(); } @Test public void forwardsRawH264WithoutRequiringDecoderSupport() { CapturingProcessFactory processFactory = new CapturingProcessFactory(new TestProcess()); HikFfmpegH264Publisher publisher = new HikFfmpegH264Publisher( processFactory, "ffmpeg", "rtmp://127.0.0.1:1935", 4); publisher.start("a1-1"); // 原码流转发兼容海康部分通道使用的H264数据分区帧。 assertTrue(processFactory.command().contains("copy")); assertTrue(processFactory.command().contains("-framerate")); assertTrue(processFactory.command().contains("25")); assertFalse(processFactory.command().contains("libx264")); publisher.close(); } @Test public void scalesRawH264FrameRateForNativeFastPlayback() { CapturingProcessFactory processFactory = new CapturingProcessFactory(new TestProcess()); HikFfmpegH264Publisher publisher = new HikFfmpegH264Publisher( processFactory, "ffmpeg", "rtmp://127.0.0.1:1935", 4); publisher.start("a1-4", 4); int frameRateOption = processFactory.command().indexOf("-framerate"); // 四倍速NVR帧必须按100fps生成时间戳,避免HLS窗口被浏览器快速耗尽。 assertEquals("100", processFactory.command().get(frameRateOption + 1)); int timestampFilterOption = processFactory.command().indexOf("-bsf:v"); // 原始H264没有PTS/DTS,流复制时必须显式给每个包补100fps连续时间戳。 assertEquals( "setts=pts=N/(100*TB):dts=N/(100*TB):duration=1/(100*TB)", processFactory.command().get(timestampFilterOption + 1)); publisher.close(); } /** * 验证历史回放仅保留 FFmpeg 错误且关闭进度输出,避免媒体包调试内容进入 Docker 日志。 */ @Test public void limitsFfmpegOutputToErrorsWithoutProgressStats() { CapturingProcessFactory processFactory = new CapturingProcessFactory(new TestProcess()); HikFfmpegH264Publisher publisher = new HikFfmpegH264Publisher( processFactory, "ffmpeg", "rtmp://127.0.0.1:1935", 4); // 启动发布器以捕获实际传给 FFmpeg 的命令行。 publisher.start("a1-1"); // 日志等级必须为 error,确保 RTMP 逐包调试和十六进制转储不会写入容器日志。 int logLevelOption = processFactory.command().indexOf("-loglevel"); assertEquals("error", processFactory.command().get(logLevelOption + 1)); // 长时间回放还需关闭进度行,避免无错误时仍持续输出统计信息。 assertTrue(processFactory.command().contains("-nostats")); assertFalse(processFactory.command().contains("debug")); publisher.close(); } @Test public void preservesIFrameWhenQueueIsFull() throws Exception { BlockingOutputStream output = new BlockingOutputStream(); CapturingProcessFactory processFactory = new CapturingProcessFactory(new TestProcess(output)); HikFfmpegH264Publisher publisher = new HikFfmpegH264Publisher( processFactory, "ffmpeg", "rtmp://127.0.0.1:1935", 2); publisher.start("a1-1"); assertTrue(publisher.offer(frameOfType(2))); assertTrue(output.awaitFirstWrite()); assertTrue(publisher.offer(frameOfType(2))); assertTrue(publisher.offer(frameOfType(3))); assertTrue(publisher.offer(frameOfType(1))); assertTrue(publisher.hasQueuedIFrame()); output.release(); publisher.close(); } private HikPlaybackFrame frameOfType(int packetType) { return new HikPlaybackFrame(packetType, new byte[]{0, 0, 0, 1, 9}); } private static final class CapturingProcessFactory extends HikFfmpegProcessFactory { private final Process process; private List command = Collections.emptyList(); private CapturingProcessFactory(Process process) { this.process = process; } @Override public Process start(List command) { this.command = new ArrayList<>(command); return process; } private List command() { return command; } } private static final class TestProcess extends Process { private final OutputStream output; private TestProcess() { this(new ByteArrayOutputStream()); } private TestProcess(OutputStream output) { this.output = output; } @Override public OutputStream getOutputStream() { return output; } @Override public InputStream getInputStream() { return new ByteArrayInputStream(new byte[0]); } @Override public InputStream getErrorStream() { return new ByteArrayInputStream(new byte[0]); } @Override public int waitFor() { return 0; } @Override public boolean waitFor(long timeout, TimeUnit unit) { return true; } @Override public int exitValue() { return 0; } @Override public void destroy() { } @Override public Process destroyForcibly() { return this; } } private static final class BlockingOutputStream extends OutputStream { private final CountDownLatch firstWrite = new CountDownLatch(1); private final CountDownLatch release = new CountDownLatch(1); @Override public void write(int value) throws IOException { block(); } @Override public void write(byte[] bytes, int offset, int length) throws IOException { block(); } @Override public void close() { release(); } private void block() throws IOException { firstWrite.countDown(); try { if (!release.await(5, TimeUnit.SECONDS)) { throw new IOException("test writer timed out"); } } catch (InterruptedException exception) { Thread.currentThread().interrupt(); throw new IOException("test writer interrupted", exception); } } private boolean awaitFirstWrite() throws InterruptedException { return firstWrite.await(2, TimeUnit.SECONDS); } private void release() { release.countDown(); } } }