【免费下载链接】rocketride-server
High-performance AI pipeline engine with a C++ core and 50+ Python-extensible nodes. Build, debug, and scale LLM workflows with 13+ model providers, 8+ vector databases, and agent orchestration, all from your IDE. Includes VS Code extension, TypeScript/Python SDKs, and Docker deployment.
media_inspect是 RocketRide(rocketride-server 仓库nodes/下的 Python 扩展节点)中用于处理流式音频/视频的专业节点:它接收以 BEGIN/WRITE/END 三段式到达的媒体流,在输入对象关闭时完成处理,支持元数据探测(probe)、静音/RMS 响度测量(levels)与按时间点抽取 JPEG 静帧(stills)三种操作模式。读完本文,你可以理解该节点的流生命周期与隔离机制、完整配置参数及其在源码中的钳制规则、三种mode的请求结构与返回格式,并学会用仓库自带的 example.pipe 把它接入 webhook → parse → 处理 → 响应的真实管道。
节点定位:只处理流,不碰存储
根据 README 的定义,media_inspect的核心行为边界非常明确:
- 接收 BEGIN/WRITE/END 流式媒体,输入对象关闭时才真正处理;
- 临时文件仅私有于当前对象,在“处理完成、失败或取消清理”三种路径之后统一释放;
- 节点从不调用账号存储、不抓取 URL、不保留持久源缓存——结果持久化必须由下游节点完成。
这个“只过手、不落地”的契约在 IInstance.py 中可以得到印证:MediaInstance类注释写明 “Own temporary files until closing; downstream nodes own persistence”,其close()方法(L218-L228)会关闭所有进行中的流文件、释放资源并销毁Workspace,且可被重复调用——这正是失败/取消路径下也能干净释放临时盘的实现保障。
输入校验的源码细节
在 instance.py 的_receive方法中,BEGIN 阶段会依次检查:
- 同一 lane 上存在未结束的旧流 → 抛
Previous media stream did not finish; - 流描述符必须是 JSON 对象;
name取自描述符的name或resource_name; - 重名拒绝:与已接收或进行中的流重名即失败;
- 每对象最多 128 个输入流(
len(self._inputs) + len(self._active) >= 128时拒绝); - 名称以保留目录
outputs/开头会被拒绝; - 声明的
size若为非负整数且超过剩余字节预算,在 BEGIN 即失败。
WRITE 阶段则在每次写入前再次检查剩余预算(L124-L127):单块超出剩余max_input_mb预算、或累计字节超过声明size,都会立即报错;END 时若收到 0 字节或实际字节数与声明不符,则报Empty or truncated media stream。这说明 README 中“输入限制在每块写入前强制执行(即使未声明 size)”不是纸面约定,而是逐块落地的检查。
Lane 拓扑:两条输入 lane,三类输出
README 给出的 lane 矩阵如下,配置时可原样用于核对管道连线:
| Lane in | Lane out | 说明 |
|---|---|---|
video | text | JSON 测量/报告(stills 模式下附带图片名标记) |
video | answers | 结构化结果 manifest |
video | image | 命名的 JPEG 流 |
audio | text | JSON 测量/报告(stills 模式下附带图片名标记) |
audio | answers | 结构化结果 manifest |
audio | image | 命名的 JPEG 流 |
从源码看,输出分发逻辑在 instance.py:closing()会强制要求至少连接了text或answers消费者(否则抛Connect a text or answers consumer),而imagelane 仅在 stills 模式使用且需要独立连接图片 sink。answers输出通过Answer(expectJson=True)承载完整结果对象,text输出则是json.dumps后的同一份 payload(L200-L208)。
配置参数:profile 覆盖 + 硬钳制
README 指出:所有覆盖项写在所选 profile 块中;chunk_kb控制输出块大小(钳制到 64–8192 KB);输入字节增量落盘到临时磁盘;event_type命名进度事件,且进度不写入任何文件。
IGlobal.py 定义了默认值并在beginGlobal中做类型转换与钳制:
DEFAULTS = { 'request': '{}', 'event_type': 'media_inspect', 'chunk_kb': 1024, 'max_input_mb': 16384, 'silence_db': -40.0, 'silence_ms': 700, 'peak_ms': 100, 'scene_threshold': 0.35, }注意源码中的钳制比 README 正文更精确:chunk_bytes = max(64, min(8192, chunk_kb)) * 1024,max_input_mb钳制到 1–1048576,peak_ms最小 10,silence_ms最小 1。profile 的读取由 config.py 的load_node_config完成:管线 profile 值按默认值的类型(int/float/bool/str)逐键覆盖,非法数值静默回退默认值。
完整 Schema 表(与 README 生成区 一致):
| Field | Type | Description | Default |
|---|---|---|---|
media_inspect.chunk_kb | number | 流式输出块大小(KB),钳制 64–8192 | 1024 |
media_inspect.event_type | string | 本管道发出的进度事件名 | "media_inspect" |
media_inspect.max_input_mb | number | 每对象累计输入上限(MiB),含全部资产;钳制 1–1048576,未声明 size 也强制 | 16384 |
media_inspect.peak_ms | number | 单个响度桶长度(ms),最小 10 | 100 |
media_inspect.profile | string | 配置 profile 名 | "default" |
media_inspect.request | string | 操作请求(JSON 字符串):probe、levels(range、ranges、scan、scan_scenes)或 stills(stills 数组) | "{}" |
media_inspect.scene_threshold | number | 镜头切换需跨越的场景分数阈值 | 0.35 |
media_inspect.silence_db | number | 静音阈值(dB) | -40 |
media_inspect.silence_ms | number | 值得报告的最短静音时长(ms) | 700 |
进度事件通过 instance.py 的_status()以 SSE 形式推给引擎监控(monitorSSE),携带schema_version、node、stage及自由字段(如stills、scanning、stills_done等 stage);投递失败只记 debug 日志,不影响媒体输出。
request:一个 JSON 字符串决定操作模式
request是编码为字符串的 JSON 对象,mode选择操作,input_name可选地指定接收到的媒体描述符名。默认是probe;源必须以媒体流形式到达——单独的存储路径不构成输入。README 特别强调:顶层的write_to、probe_to、status_to、report_to会被拒绝。
源码中对应 instance.py 的_request():
LEGACY_DESTINATIONS = {'write_to', 'probe_to', 'status_to', 'report_to'} # ... if LEGACY_DESTINATIONS.intersection(request): raise ValueError('Storage destinations are unsupported; connect a downstream sink') mode = request.get('mode', self.modes[0]) if mode not in self.modes: raise ValueError('mode must be one of: ' + ', '.join(self.modes))modes = ('probe', 'levels', 'stills'),非法 mode 直接报错。源解析规则(closing,L151-L161)为:给了input_name就按名取;只有一条音视频流则自动使用;多条则必须显式指定,否则报Multiple media streams require input_name。另外_process(L176-L178)明确约束:本处理器每个对象只接受恰好 1 条媒体流。
probe:一次拿到全部元数据
probe返回时长、显示/编码尺寸、像素宽高比、编解码器、帧率、音频声道与采样率等。返回结构由 media.py 的probe_payload组装:
{ "schema_version": 1, "kind": "media_probe", "mode": "probe", "duration_ms": 123456, "width": 853, "height": 480, "sar": 1.185185, "coded_width": 720, "coded_height": 480, "fps": 29.97, "has_video": true, "has_audio": true, "size": 12345678, "video_codec": "h264", "audio_codec": "aac", "audio_channels": 2, "audio_sample_rate": 48000, "container": "mov,mp4,m4a,3gp,3g2,mj2", "context": { "mode": "probe", "source": "source.media" } }值得注意的实现细节:width/height是显示(方形像素)尺寸——各向异性像素源(DV、HDV、4:3 广播母带,如 720×480 @ SAR 32:27)会先经display_dims换算,sar与coded_width/coded_height随行携带,调用方无需二次探测即可区分方形与非方形源(media.py)。探测主路径用 PyAV(probe(),L323-L374);PyAV 缺失时回退到解析ffmpeg -i输出(_probe_ffmpeg,L258-L320,子进程超时 60 秒)——这正是 README 所说“元数据子进程最大 60 秒超时”的实现。
levels:静音、RMS 峰值与镜头切换
levels在一次请求中做三类全区间扫描加一类按需测量(levels.py 模块 docstring 有完整说明):
range:start_ms-end_ms(也接受[start, end]形式,见 parse_range);省略时默认全时长;ranges:JSON 数组[[a, b], …](ms),只对指定窗口测 RMS(dBFS),每个窗口一次短解码,不触碰文件其余部分;scan: "no":跳过全区间扫描,此时不产生任何解码;scan_scenes: "no":只跳过镜头切换扫描(场景检测是全扫描中最贵的一步);silence_db、silence_ms、peak_ms、scene_threshold均可在 request 内覆盖同名 profile 参数(IInstance.py)。
测量管线(scan_levels):
- silences:区间先解码成 16 kHz 单声道 PCM WAV(
ANALYSIS_RATE = 16000),再跑 ffmpegsilencedetect(keep_pause_ms=0,报告原始静音而非“可安全剪掉的静音”); - peaks:同一 WAV 按
peak_ms分桶计算 RMS,归一化到最响桶(0..1),peaks[i]对应start_ms + i*peak_ms起的桶——绘图与“候选剪辑点两侧是否有声音”的判断都靠它; - scenes:对缩放到 320 px 宽的降采样解码跑场景分数,超过
scene_threshold的时间戳即镜头切换点; - ranges:
measure_ranges逐窗口测 dBFS,测不到的窗口返回None而不是编造数字。
所有时间都在源时钟上:扫描先 seek 到区间起点,报告前把区间偏移加回。无画面的区间不报 scenes,无声的区间不报 silences/peaks——“这是回答,不是失败”。scan: no与scan_scenes: no对长录制的成本差异在源码注释中被量化为“场景检测约每分钟对应 10 分钟画面”的量级差,是长文件上“秒级与分钟级”的区别。
stills:按时间点抽 JPEG 并走 image lane
stills模式下stills是[{id, t_ms, crop, width}]数组,crop为显示像素下的[x, y, width, height];quality设置 JPEG 量化值(钳制 2–31,缺省或 0 用默认值 2,见 IInstance.py 的max(2, min(31, quality))与media_lib.STILL_QUALITY = 2)。README 给的示例可直接使用:
{"mode":"stills","stills":[{"id":"poster","t_ms":1000,"width":320}]}行为要点(IInstance.py 的_stills与 media.py):
- 逐条容错:无效条目只进
failed列表,其余条目继续执行——“一条坏 crop 不再拖垮整批静帧”; - crop 自动钳位:
crop_window(media.py L73-L88)把越界框收回画面内部,x/y钳到[0, w-cw]/[0, h-ch],宽高至少 1 px; - id 必须是不含路径分隔符的普通文件名(
still_entry校验/、\、.、..与控制字符),因为它直接作为输出流名; - 无视频轨的文件该条计入
failed; - 输出必须接 image sink:
_image_listener(L104-L109)检测hasListener('image'),没有监听者时抛Connect an image sink for extracted stills; - 命名标记机制:每张图在 image lane 上以独立流(BEGIN/WRITE/END)发出,而名字先在 text lane 上写一行标记再发图。源码注释解释了原因:response 节点会把一次运行写入 image lane 的全部字节聚合进同一个输入对象名下的 text 缓冲,image lane 的描述符名到不了读 text 的调用方——所以标记行写在 text 上是“名字能活下来”的唯一途径。这是理解该节点与下游 response 节点配合的关键细节。
单帧抽取cut_still的时间精度规则在 picture_pass 中有专门论述:媒体元素在时刻 t 显示的是pts <= t的帧,而 ffmpeg-ss t返回pts >= t的帧——因此 seek 回退一个网格步长、用fps滤镜以start_time锚定并round=down,配合settb=1/1000000与 0.1 ms 的FRAME_EPSILON_S,使输出帧与“播放到该毫秒时你看到的帧”逐帧一致。多帧 JPEG 流(image2pipe/mjpeg)按 SOI/EOI 标记精确切分(split_jpegs,L94-L109),因为熵编码数据中 0xFF 后必跟 0x00 或重启标记,切分点不会误判。
超时与资源边界
README 的“Stream lifecycle and limits”一节列出的边界在源码中逐项可查:
- FFmpeg 编码超时:环境变量
ROCKETRIDE_MEDIA_FFMPEG_TIMEOUT限制每次 FFmpeg 编码,默认 3600 秒,允许大于 0 到 86400;非法值显式失败而非静默回退。实现见 media.py:非数值、<=0、>86400或非有限数都抛ValueError。测试 test_review_regressions.py 中test_ffmpeg_timeout_is_bounded_and_reported与test_timeout_configuration_rejects_unbounded_values正是对这条约束的回归验证; - 元数据子进程:最大 60 秒超时(
_probe_ffmpeg的timeout=60); max_input_mb:默认 16384 MiB(16 GiB),钳制 1–1048576 MiB;在 BEGIN(声明 size 超出剩余预算)与每个 WRITE 块前双重执行。README 特别澄清:该入口限制不约束解码/中间输出大小;- 流数量:每对象最多 128 条输入流;本处理器每对象只接受恰好 1 条媒体流;
- 失败面:空流、截断流、超声明大小的流、重名流均显式失败;
- 异常终止:进程被强杀可能留下 scratch 文件,由宿主清理;
- 浏览器中转:节点契约不需要浏览器下载/再上传中继。README 同时说明:标准
filestore_source目前拒绝保存超过 100 MiB 的文件,直接 webhook 上传可绕开该源端限制;这些节点不修改标准 source/sink 的实现。持久的进度、恢复与输出命名属于管线/应用编排层的职责,不在本节点内。
示例管线:webhook → parse → media_inspect → response + 图片 sink
example.pipe 给出了完整可运行的拓扑(运行前需按部署环境配置各 sink 目标):
{ "id": "media", "provider": "media_inspect", "config": { "profile": "default", "default": { "request": "{\"mode\": \"probe\"}" } }, "input": [ { "from": "parse", "lane": "video" }, { "from": "parse", "lane": "audio" } ] }节点链为webhook(source)→parse(消费tagslane)→media_inspect(消费video/audiolane)→response_answers(消费answerslane),另有一条media_inspect的imagelane 接到filestoresink(targetDir: "media-results/image"、onConflict: "unique")。默认probe请求返回元数据;把request改成stills后,image sink 即开始落盘静帧。README 同时说明:该节点不加载任何语音模型。
测试与依赖
验证层位于 nodes/test/media_inspect:
- test_streams.py:流生命周期与对象隔离;
- test_live.py:真实 FFmpeg 操作与在线引擎集成;
- test_review_regressions.py:非法 ranges 永不退化为全量输入、range 输入形式、FFmpeg 超时钳制、stills 质量钳制(
test_still_quality_is_bounded)等回归项。
依赖(requirements.txt,与 README 生成区一致):
| 依赖 | 版本约束 |
|---|---|
av(PyAV) | >=18,<19 |
imageio-ffmpeg | >=0.6,<0.7 |
numpy | >=2,<3 |
FFmpeg 可执行文件的解析顺序在 media.py:环境变量MEDIA_TOOLKIT_FFMPEG优先,其次imageio_ffmpeg内置二进制,最后退回 PATH 上的ffmpeg。
小结
media_inspect的设计可以概括为三句话:流进流出、对象级隔离(临时文件随对象生死,持久化交给下游 sink);三种模式各司其职(probe 拿元数据、levels 做静音/RMS/场景测量、stills 抽帧走 image lane);所有边界显式失败(超时钳制、字节预算逐块检查、crop 钳位、无效静帧逐条降级)。结合 IInstance.py、levels.py 与 example.pipe,你可以在 RocketRide 管道中放心地把它作为媒体分析的前置节点,并用nodes/test/media_inspect下的测试验证任何定制改动。
【免费下载链接】rocketride-server
High-performance AI pipeline engine with a C++ core and 50+ Python-extensible nodes. Build, debug, and scale LLM workflows with 13+ model providers, 8+ vector databases, and agent orchestration, all from your IDE. Includes VS Code extension, TypeScript/Python SDKs, and Docker deployment.
相关推荐
3步掌握Seraphine智能助手:你的英雄联盟排位赛专属数据分析解决方案
3步掌握Seraphine智能助手:你的英雄联盟排位赛专属数据分析解决方案 你是否曾在英雄联盟排位赛中遇到过这样的困境?BP阶段手忙脚乱,既要分析对手战绩,又要
RocketRide Gemini Vision 节点:在 LLM 流水线中实现图像分析、OCR 与视频帧理解
RocketRide Gemini Vision 节点:在 LLM 流水线中实现图像分析、OCR 与视频帧理解 RocketRide 的 llm_vision_
librealsense C API 深度测距实战:基于 rs-distance 示例解析深度帧流与中心点距离测量
librealsense C API 深度测距实战:基于 rs distance 示例解析深度帧流与中心点距离测量 导读 本文围绕 librealsense(R
智能硬件音视频计算机视觉
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考