This commit is contained in:
@@ -39,7 +39,7 @@ P2P: [Browser] ◄──────► [Browser]
|
|||||||
- **HRPC** - Remote procedure calls with streaming support
|
- **HRPC** - Remote procedure calls with streaming support
|
||||||
|
|
||||||
### Optional capabilities
|
### Optional capabilities
|
||||||
- **Media** - Host-side image/video via `bare-media` + `bare-ffmpeg` (`BridgeSwarm.media.*`) — **included in the default host**; see [docs/CAPABILITIES.md](docs/CAPABILITIES.md) and [docs/DEFAULT-MODULES.md](docs/DEFAULT-MODULES.md)
|
- **Media** - Host-side image/video via `bare-media` + `bare-ffmpeg` (`BridgeSwarm.media.*`) — batch jobs **and live VP9 encode** — included in the default host; see [docs/CAPABILITIES.md](docs/CAPABILITIES.md)
|
||||||
|
|
||||||
## Quick Start
|
## Quick Start
|
||||||
|
|
||||||
|
|||||||
@@ -28,7 +28,7 @@ These are the main APIs your page uses. For host request types and event payload
|
|||||||
|
|
||||||
- **`BridgeSwarm.capabilities`** — Bare capability packs on the host (`media` is default). `list()`, `has(pack)`, `call(pack, cmd, payload)`, `on(event, fn)`. See [CAPABILITIES.md](CAPABILITIES.md) / [DEFAULT-MODULES.md](DEFAULT-MODULES.md).
|
- **`BridgeSwarm.capabilities`** — Bare capability packs on the host (`media` is default). `list()`, `has(pack)`, `call(pack, cmd, payload)`, `on(event, fn)`. See [CAPABILITIES.md](CAPABILITIES.md) / [DEFAULT-MODULES.md](DEFAULT-MODULES.md).
|
||||||
|
|
||||||
- **`BridgeSwarm.media.*`** — Media pack wrappers (`info`, `imageTransform`, `extractFrame`, `transcode`, `cancel`, `writeInput`, `readOutput`).
|
- **`BridgeSwarm.media.*`** — Media pack: batch (`info`, `imageTransform`, `extractFrame`, `transcode`, …) and live encode (`encodeStart`, `encodePushFrame`, `encodeStop`, `liveSession`, `attachLiveReceiver`). See [CAPABILITIES.md](CAPABILITIES.md).
|
||||||
---
|
---
|
||||||
|
|
||||||
## Host request types and events
|
## Host request types and events
|
||||||
@@ -145,14 +145,15 @@ Full details and examples: [DATA-API.md](DATA-API.md).
|
|||||||
| `capabilities.has` | `{ pack }` | `{ ok, pack, has }` |
|
| `capabilities.has` | `{ pack }` | `{ ok, pack, has }` |
|
||||||
| `capability` | `{ pack, cmd, payload }` | pack-specific |
|
| `capability` | `{ pack, cmd, payload }` | pack-specific |
|
||||||
| `media.info` / `media.imageTransform` / `media.extractFrame` / `media.transcode` / `media.cancel` / … | command payload | see [CAPABILITIES.md](CAPABILITIES.md) |
|
| `media.info` / `media.imageTransform` / `media.extractFrame` / `media.transcode` / `media.cancel` / … | command payload | see [CAPABILITIES.md](CAPABILITIES.md) |
|
||||||
|
| `media.encodeStart` / `encodePushFrame` / `encodePush` / `encodeSubscribe` / `encodeStop` | live session payload | `{ ok, sessionId, … }` + `cap-chunk` segments |
|
||||||
|
|
||||||
### Capability events
|
### Capability events
|
||||||
|
|
||||||
| Event | Payload |
|
| Event | Payload |
|
||||||
|-------|---------|
|
|-------|---------|
|
||||||
| `cap-chunk` | `{ pack, jobId, index?, data?, progress?, bytes? }` |
|
| `cap-chunk` | `{ pack, jobId/sessionId, kind?, seq?, index?, data?, progress?, bytes? }` |
|
||||||
| `cap-end` | `{ pack, jobId, … }` |
|
| `cap-end` | `{ pack, jobId/sessionId, … }` |
|
||||||
| `cap-error` | `{ pack, jobId, message }` |
|
| `cap-error` | `{ pack, jobId/sessionId, message }` |
|
||||||
|
|
||||||
These events are broadcast to all subscribed tabs (no `swarmId` filter).
|
These events are broadcast to all subscribed tabs (no `swarmId` filter).
|
||||||
|
|
||||||
|
|||||||
+28
-8
@@ -11,23 +11,30 @@ await BridgeSwarm.capabilities.has('media')
|
|||||||
// Generic dispatch
|
// Generic dispatch
|
||||||
await BridgeSwarm.capabilities.call('media', 'info', { path: '…' })
|
await BridgeSwarm.capabilities.call('media', 'info', { path: '…' })
|
||||||
|
|
||||||
// Thin wrappers
|
// Batch wrappers
|
||||||
await BridgeSwarm.media.info({ dataBase64, filename })
|
await BridgeSwarm.media.info({ dataBase64, filename })
|
||||||
await BridgeSwarm.media.imageTransform({ dataBase64, maxWidth: 640, mimetype: 'image/webp' })
|
await BridgeSwarm.media.imageTransform({ dataBase64, maxWidth: 640, mimetype: 'image/webp' })
|
||||||
await BridgeSwarm.media.extractFrame({ dataBase64, frameIndex: 0 })
|
await BridgeSwarm.media.extractFrame({ dataBase64, frameIndex: 0 })
|
||||||
await BridgeSwarm.media.transcode({ dataBase64, format: 'webm' }, { onProgress })
|
await BridgeSwarm.media.transcode({ dataBase64, format: 'webm' }, { onProgress })
|
||||||
await BridgeSwarm.media.cancel(jobId)
|
await BridgeSwarm.media.cancel(jobId)
|
||||||
|
|
||||||
|
// Live encode (host bare-ffmpeg → MSE)
|
||||||
|
const live = BridgeSwarm.media.liveSession(videoEl)
|
||||||
|
await live.start({ width: 640, height: 360, fps: 10, ingest: 'frames' })
|
||||||
|
await live.pushFrame(canvas)
|
||||||
|
await live.subscribe({ swarm: { connIds: [conn.connId] } })
|
||||||
|
await live.stop()
|
||||||
```
|
```
|
||||||
|
|
||||||
Streaming results use host events (under the ~1 MB native-messaging limit):
|
Streaming results use host events (under the ~1 MB native-messaging limit):
|
||||||
|
|
||||||
| Event | Payload |
|
| Event | Payload |
|
||||||
|-------|---------|
|
|-------|---------|
|
||||||
| `cap-chunk` | `{ pack, jobId, index?, data? (base64), progress?, bytes? }` |
|
| `cap-chunk` | `{ pack, jobId/sessionId, kind?, index?, data? (base64), progress?, bytes? }` |
|
||||||
| `cap-end` | `{ pack, jobId, chunks?, path?, mimetype?, … }` |
|
| `cap-end` | `{ pack, jobId/sessionId, chunks?, path?, mimetype?, … }` |
|
||||||
| `cap-error` | `{ pack, jobId, message }` |
|
| `cap-error` | `{ pack, jobId/sessionId, message }` |
|
||||||
|
|
||||||
`BridgeSwarm.media.*` helpers assign a `jobId`, listen for these events, and resolve with assembled `dataBase64` (or a host-relative `path` for large transcodes).
|
Live sessions emit `cap-chunk` with `kind: 'segment'` (WebM fragments for MSE) or `kind: 'drop'` under backpressure.
|
||||||
|
|
||||||
```js
|
```js
|
||||||
const off = BridgeSwarm.capabilities.on('cap-chunk', (p) => console.log(p))
|
const off = BridgeSwarm.capabilities.on('cap-chunk', (p) => console.log(p))
|
||||||
@@ -53,16 +60,29 @@ Shipped in the standard `bridge-swarm-host` artifact (`bare-media` + `bare-ffmpe
|
|||||||
| `media.imageTransform` | Decode → optional resize/crop → encode (webp/jpeg/png); stream base64 |
|
| `media.imageTransform` | Decode → optional resize/crop → encode (webp/jpeg/png); stream base64 |
|
||||||
| `media.extractFrame` | Extract one video frame → encode; stream base64 |
|
| `media.extractFrame` | Extract one video frame → encode; stream base64 |
|
||||||
| `media.transcode` | Async job to webm/mp4/mkv (VP9+Opus where supported); progress events |
|
| `media.transcode` | Async job to webm/mp4/mkv (VP9+Opus where supported); progress events |
|
||||||
| `media.cancel` | Cancel in-flight job |
|
| `media.encodeStart` | Start live VP9/WebM session (`ingest: frames\|segments`, egress page/swarm/file) |
|
||||||
|
| `media.encodePushFrame` | Push JPEG/WebP frame (base64) into live session |
|
||||||
|
| `media.encodePush` | Push MediaRecorder timeslice or frame alias |
|
||||||
|
| `media.encodeSubscribe` | Update fan-out `connIds` / page egress mid-session |
|
||||||
|
| `media.encodeStop` | Flush encoder, finalize archive, `cap-end` |
|
||||||
|
| `media.cancel` | Cancel batch job or live session |
|
||||||
| `media.writeInput` / `media.readOutput` | Chunked upload / download under `cap-jobs/` |
|
| `media.writeInput` / `media.readOutput` | Chunked upload / download under `cap-jobs/` |
|
||||||
|
|
||||||
After installing or updating the host on macOS, run `npm run repair:macos` (or `npm run install:capability:media`, which now aliases repair) so new `.bare` addons are extracted and codesigned. Fully quit the browser, then:
|
### Live encode notes
|
||||||
|
|
||||||
|
- Chrome NMH ~1 MB/message → ingest **compressed JPEG frames** (or WebM slices), not raw RGBA.
|
||||||
|
- Defaults for demos: **640×360 @ 10fps**; soft max ~720p @ 15fps; session max duration 10 minutes.
|
||||||
|
- Egress: `page` (MSE segments), `swarm` (BSML-framed binary on Hyperswarm sockets), `file` (`cap-jobs/<sessionId>/live.webm`).
|
||||||
|
- Peer fan-out writes on the **host socket** (good for bitrates); the receiving page still reads via NMH `conn.on('data')` unless you stay host-side.
|
||||||
|
- Helpers: `BridgeSwarm.media.liveSession(videoEl)` and `BridgeSwarm.media.attachLiveReceiver(conn, videoEl)`.
|
||||||
|
|
||||||
|
After installing or updating the host on macOS, run `npm run repair:macos` so `.bare` addons (including ffmpeg) are extracted and codesigned. Fully quit the browser, then:
|
||||||
|
|
||||||
```js
|
```js
|
||||||
await BridgeSwarm.capabilities.has('media') // true
|
await BridgeSwarm.capabilities.has('media') // true
|
||||||
```
|
```
|
||||||
|
|
||||||
Demo: [`examples/media-demo/`](../examples/media-demo/) (`npm run examples` → `/media-demo/`).
|
Demos: [`examples/live-encode/`](../examples/live-encode/), [`examples/media-demo/`](../examples/media-demo/), [`examples/clip-studio/`](../examples/clip-studio/).
|
||||||
|
|
||||||
## Bundled for upcoming packs
|
## Bundled for upcoming packs
|
||||||
|
|
||||||
|
|||||||
@@ -30,7 +30,7 @@ Curated modules shipped in the **default** BridgeSwarm native host. Not every `b
|
|||||||
| Module | Why |
|
| Module | Why |
|
||||||
|--------|-----|
|
|--------|-----|
|
||||||
| `bare-media` | High-level image/video API for `BridgeSwarm.media.*` |
|
| `bare-media` | High-level image/video API for `BridgeSwarm.media.*` |
|
||||||
| `bare-ffmpeg` | Transcode / frame extract (pulled by bare-media; listed explicitly) |
|
| `bare-ffmpeg` | Batch transcode **and** live VP9/WebM encode sessions (`encodeStart` / `encodePushFrame`) |
|
||||||
| Image codecs | Via bare-media: `bare-jpeg`, `bare-png`, `bare-webp`, `bare-gif`, `bare-heif`, `bare-bmp`, `bare-ico`, `bare-tiff`, `bare-svg`, `bare-image-resample`, `bare-exif` |
|
| Image codecs | Via bare-media: `bare-jpeg`, `bare-png`, `bare-webp`, `bare-gif`, `bare-heif`, `bare-bmp`, `bare-ico`, `bare-tiff`, `bare-svg`, `bare-image-resample`, `bare-exif` |
|
||||||
|
|
||||||
### Local power (bundled for capability packs)
|
### Local power (bundled for capability packs)
|
||||||
|
|||||||
+12
-8
@@ -26,6 +26,7 @@ Shared chrome lives in [`shared/`](shared/) (`theme.css`, `chrome.css`, `boot.js
|
|||||||
| Screenshare | http://127.0.0.1:4173/screenshare/ |
|
| Screenshare | http://127.0.0.1:4173/screenshare/ |
|
||||||
| Data API | http://127.0.0.1:4173/data-demo/ |
|
| Data API | http://127.0.0.1:4173/data-demo/ |
|
||||||
| Auto-replicate | http://127.0.0.1:4173/sync-demo/ |
|
| Auto-replicate | http://127.0.0.1:4173/sync-demo/ |
|
||||||
|
| Live Encode | http://127.0.0.1:4173/live-encode/ |
|
||||||
| Media Studio | http://127.0.0.1:4173/media-demo/ |
|
| Media Studio | http://127.0.0.1:4173/media-demo/ |
|
||||||
| Clip Studio | http://127.0.0.1:4173/clip-studio/ |
|
| Clip Studio | http://127.0.0.1:4173/clip-studio/ |
|
||||||
| SDK Demo | http://127.0.0.1:4173/sdk-demo/ |
|
| SDK Demo | http://127.0.0.1:4173/sdk-demo/ |
|
||||||
@@ -50,6 +51,17 @@ Collaborative canvas with color/size tools.
|
|||||||
### Screenshare (`screenshare/`)
|
### Screenshare (`screenshare/`)
|
||||||
Live **WebRTC** video; BridgeSwarm is signaling only (not host ffmpeg).
|
Live **WebRTC** video; BridgeSwarm is signaling only (not host ffmpeg).
|
||||||
|
|
||||||
|
## Media
|
||||||
|
|
||||||
|
### Live Encode (`live-encode/`)
|
||||||
|
Host **bare-ffmpeg** live VP9/WebM: capture → JPEG frames → MSE preview of the **host** encode, optional peer fan-out.
|
||||||
|
|
||||||
|
### Media Studio (`media-demo/`)
|
||||||
|
Probe, WebP transform, extract frame, batch transcode (`BridgeSwarm.media.*`).
|
||||||
|
|
||||||
|
### Clip Studio (`clip-studio/`)
|
||||||
|
Record screen/camera → nearline batch `transcode` / poster (not continuous live encode).
|
||||||
|
|
||||||
## Data
|
## Data
|
||||||
|
|
||||||
### Data API (`data-demo/`)
|
### Data API (`data-demo/`)
|
||||||
@@ -58,14 +70,6 @@ Hypercore / Hyperbee / Hyperdrive / Autobase / Hyperdb via `BridgeSwarm.request`
|
|||||||
### Auto-replicate (`sync-demo/`)
|
### Auto-replicate (`sync-demo/`)
|
||||||
`setAutoReplicate` Hypercore sync across peers. Note: those sockets are taken over for replication.
|
`setAutoReplicate` Hypercore sync across peers. Note: those sockets are taken over for replication.
|
||||||
|
|
||||||
## Media
|
|
||||||
|
|
||||||
### Media Studio (`media-demo/`)
|
|
||||||
Probe, WebP transform, extract frame, transcode (`BridgeSwarm.media.*`).
|
|
||||||
|
|
||||||
### Clip Studio (`clip-studio/`)
|
|
||||||
Record screen/camera → upload → host `bare-ffmpeg` batch transcode / poster frame (nearline, not live encode).
|
|
||||||
|
|
||||||
## Platform
|
## Platform
|
||||||
|
|
||||||
### SDK Demo (`sdk-demo/`)
|
### SDK Demo (`sdk-demo/`)
|
||||||
|
|||||||
@@ -21,8 +21,11 @@
|
|||||||
</p>
|
</p>
|
||||||
|
|
||||||
<div class="bs-banner bs-banner--info">
|
<div class="bs-banner bs-banner--info">
|
||||||
Capture → <code>MediaRecorder</code> → <code>media.writeInput</code> →
|
<strong>Nearline batch</strong>: Capture → <code>MediaRecorder</code> →
|
||||||
<code>transcode</code> / <code>extractFrame</code>. For live WebRTC, use
|
<code>transcode</code> / <code>extractFrame</code>.
|
||||||
|
For continuous host <code>bare-ffmpeg</code> encode, use
|
||||||
|
<a href="../live-encode/" style="color:var(--accent-2)">Live Encode</a>.
|
||||||
|
For browser WebRTC A/V, use
|
||||||
<a href="../screenshare/" style="color:var(--accent-2)">Screenshare</a>.
|
<a href="../screenshare/" style="color:var(--accent-2)">Screenshare</a>.
|
||||||
</div>
|
</div>
|
||||||
<div id="missing" class="bs-banner bs-banner--warn" hidden>
|
<div id="missing" class="bs-banner bs-banner--warn" hidden>
|
||||||
|
|||||||
+7
-2
@@ -70,14 +70,19 @@
|
|||||||
|
|
||||||
<div class="bs-category">Media</div>
|
<div class="bs-category">Media</div>
|
||||||
<div class="bs-grid">
|
<div class="bs-grid">
|
||||||
|
<a class="bs-card" href="./live-encode/">
|
||||||
|
<h2>Live Encode</h2>
|
||||||
|
<p>Real host bare-ffmpeg: capture → VP9/WebM live mux → MSE preview (+ peer fan-out).</p>
|
||||||
|
<span class="bs-path">/live-encode/</span>
|
||||||
|
</a>
|
||||||
<a class="bs-card" href="./media-demo/">
|
<a class="bs-card" href="./media-demo/">
|
||||||
<h2>Media Studio</h2>
|
<h2>Media Studio</h2>
|
||||||
<p>Probe, WebP transform, extract frame, and transcode via host ffmpeg.</p>
|
<p>Probe, WebP transform, extract frame, and batch transcode via host ffmpeg.</p>
|
||||||
<span class="bs-path">/media-demo/</span>
|
<span class="bs-path">/media-demo/</span>
|
||||||
</a>
|
</a>
|
||||||
<a class="bs-card" href="./clip-studio/">
|
<a class="bs-card" href="./clip-studio/">
|
||||||
<h2>Clip Studio</h2>
|
<h2>Clip Studio</h2>
|
||||||
<p>Record screen/camera, then batch-transcode with bare-ffmpeg (nearline).</p>
|
<p>Record a short clip, then batch-transcode with bare-ffmpeg (nearline).</p>
|
||||||
<span class="bs-path">/clip-studio/</span>
|
<span class="bs-path">/clip-studio/</span>
|
||||||
</a>
|
</a>
|
||||||
</div>
|
</div>
|
||||||
|
|||||||
@@ -0,0 +1,220 @@
|
|||||||
|
(function () {
|
||||||
|
var logEl = document.getElementById('log');
|
||||||
|
var statusEl = document.getElementById('status');
|
||||||
|
var statsEl = document.getElementById('stats');
|
||||||
|
var localVideo = document.getElementById('local');
|
||||||
|
var encodedVideo = document.getElementById('encoded');
|
||||||
|
var remoteVideo = document.getElementById('remote');
|
||||||
|
var livePill = document.getElementById('livePill');
|
||||||
|
|
||||||
|
var captureStream = null;
|
||||||
|
var pump = null;
|
||||||
|
var live = null;
|
||||||
|
var swarm = null;
|
||||||
|
var connections = [];
|
||||||
|
var receivers = [];
|
||||||
|
|
||||||
|
function log(msg, type) {
|
||||||
|
BridgeSwarmExamples.logLine(logEl, msg, type);
|
||||||
|
}
|
||||||
|
function setStatus(msg, isError) {
|
||||||
|
statusEl.textContent = msg;
|
||||||
|
statusEl.className = 'bs-status' + (isError ? ' error' : '');
|
||||||
|
}
|
||||||
|
function setStats(s) {
|
||||||
|
statsEl.textContent =
|
||||||
|
'frames in: ' + (s.framesIn || 0) +
|
||||||
|
' · encoded: ' + (s.framesEncoded || 0) +
|
||||||
|
' · out: ' + Math.round((s.bytesOut || 0) / 1024) + ' KB' +
|
||||||
|
' · drops: ' + (s.drops || 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
function stopCaptureOnly() {
|
||||||
|
if (pump) {
|
||||||
|
pump.stop();
|
||||||
|
pump = null;
|
||||||
|
}
|
||||||
|
if (captureStream) {
|
||||||
|
captureStream.getTracks().forEach(function (t) { t.stop(); });
|
||||||
|
captureStream = null;
|
||||||
|
}
|
||||||
|
localVideo.srcObject = null;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function beginCapture(kind) {
|
||||||
|
stopCaptureOnly();
|
||||||
|
captureStream = await BridgeSwarmLiveCapture.startCapture(kind);
|
||||||
|
localVideo.srcObject = captureStream;
|
||||||
|
document.getElementById('btnStart').disabled = false;
|
||||||
|
log('Capture ready (' + kind + ')', 'val');
|
||||||
|
setStatus('Capture ready. Start live encode when you want host ffmpeg.');
|
||||||
|
captureStream.getVideoTracks()[0].addEventListener('ended', function () {
|
||||||
|
stopAll();
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
async function startEncode() {
|
||||||
|
if (!captureStream || live) return;
|
||||||
|
var width = Number(document.getElementById('width').value) || 640;
|
||||||
|
var height = Number(document.getElementById('height').value) || 360;
|
||||||
|
var fps = Number(document.getElementById('fps').value) || 10;
|
||||||
|
var crf = Number(document.getElementById('crf').value) || 34;
|
||||||
|
|
||||||
|
live = window.BridgeSwarm.media.liveSession(encodedVideo, {
|
||||||
|
onDrop: function (p) {
|
||||||
|
log('egress drop (backpressure) total=' + (p.drops || '?'), 'sys');
|
||||||
|
},
|
||||||
|
onStats: function (p) {
|
||||||
|
if (p.framesIn != null) setStats(p);
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
var started = await live.start({
|
||||||
|
width: width,
|
||||||
|
height: height,
|
||||||
|
fps: fps,
|
||||||
|
crf: crf,
|
||||||
|
ingest: 'frames',
|
||||||
|
egress: { page: true, file: true, swarm: { connIds: fanoutConnIds() } },
|
||||||
|
});
|
||||||
|
log('Live session ' + started.sessionId, 'val');
|
||||||
|
livePill.textContent = 'LIVE encode';
|
||||||
|
livePill.className = 'bs-pill bs-pill--live';
|
||||||
|
document.getElementById('btnStart').disabled = true;
|
||||||
|
document.getElementById('btnStop').disabled = false;
|
||||||
|
setStatus('Encoding on host bare-ffmpeg…');
|
||||||
|
|
||||||
|
pump = BridgeSwarmLiveCapture.startFramePump(localVideo, {
|
||||||
|
fps: fps,
|
||||||
|
maxWidth: width,
|
||||||
|
maxHeight: height,
|
||||||
|
onFrame: function (canvas) {
|
||||||
|
return live.pushFrame(canvas, 0.7).then(function (r) {
|
||||||
|
setStats(r);
|
||||||
|
}).catch(function (e) {
|
||||||
|
log(e.message, 'err');
|
||||||
|
});
|
||||||
|
},
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
function fanoutConnIds() {
|
||||||
|
if (!document.getElementById('fanout').checked) return [];
|
||||||
|
return connections.map(function (c) { return c.conn.connId; });
|
||||||
|
}
|
||||||
|
|
||||||
|
async function syncFanout() {
|
||||||
|
if (!live || !live.sessionId) return;
|
||||||
|
try {
|
||||||
|
await live.subscribe({ swarm: { connIds: fanoutConnIds() } });
|
||||||
|
log('Fan-out connIds=' + fanoutConnIds().length, 'sys');
|
||||||
|
} catch (e) {
|
||||||
|
log(e.message, 'err');
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function stopAll() {
|
||||||
|
if (pump) {
|
||||||
|
pump.stop();
|
||||||
|
pump = null;
|
||||||
|
}
|
||||||
|
if (live) {
|
||||||
|
try {
|
||||||
|
var r = await live.stop();
|
||||||
|
if (r && r.path) log('Archive ' + r.path, 'val');
|
||||||
|
} catch (e) {
|
||||||
|
log(e.message, 'err');
|
||||||
|
}
|
||||||
|
live = null;
|
||||||
|
}
|
||||||
|
stopCaptureOnly();
|
||||||
|
livePill.textContent = '';
|
||||||
|
livePill.className = 'bs-pill';
|
||||||
|
document.getElementById('btnStart').disabled = true;
|
||||||
|
document.getElementById('btnStop').disabled = true;
|
||||||
|
setStatus('Stopped.');
|
||||||
|
}
|
||||||
|
|
||||||
|
function attachReceiver(conn) {
|
||||||
|
var recv = window.BridgeSwarm.media.attachLiveReceiver(conn, remoteVideo, {
|
||||||
|
onSegment: function () {
|
||||||
|
document.getElementById('recvStage').classList.add('has-stream');
|
||||||
|
},
|
||||||
|
});
|
||||||
|
receivers.push(recv);
|
||||||
|
conn.on('end', function () {
|
||||||
|
recv.destroy();
|
||||||
|
receivers = receivers.filter(function (r) { return r !== recv; });
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
BridgeSwarmExamples.waitForBridgeSwarm().then(async function () {
|
||||||
|
var has = await window.BridgeSwarm.capabilities.has('media');
|
||||||
|
if (!has) {
|
||||||
|
document.getElementById('missing').hidden = false;
|
||||||
|
setStatus('Media pack missing', true);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
setStatus('Ready. Capture screen or camera, then start live encode.');
|
||||||
|
document.getElementById('btnJoin').disabled = false;
|
||||||
|
log('Media pack ready (live encode)', 'val');
|
||||||
|
|
||||||
|
document.getElementById('btnScreen').onclick = function () {
|
||||||
|
beginCapture('screen').catch(function (e) { log(e.message, 'err'); });
|
||||||
|
};
|
||||||
|
document.getElementById('btnCamera').onclick = function () {
|
||||||
|
beginCapture('camera').catch(function (e) { log(e.message, 'err'); });
|
||||||
|
};
|
||||||
|
document.getElementById('btnStart').onclick = function () {
|
||||||
|
startEncode().catch(function (e) {
|
||||||
|
setStatus(e.message, true);
|
||||||
|
log(e.message, 'err');
|
||||||
|
});
|
||||||
|
};
|
||||||
|
document.getElementById('btnStop').onclick = function () { stopAll(); };
|
||||||
|
document.getElementById('fanout').onchange = function () { syncFanout(); };
|
||||||
|
|
||||||
|
document.getElementById('btnJoin').onclick = async function () {
|
||||||
|
if (swarm) return;
|
||||||
|
try {
|
||||||
|
var topic = document.getElementById('topic').value.trim() || 'bridge-swarm-live-encode';
|
||||||
|
swarm = new window.BridgeSwarm({ appName: 'bridge-swarm-live-encode' });
|
||||||
|
swarm.on('connection', function (conn, peerInfo) {
|
||||||
|
connections.push({ conn: conn, peerInfo: peerInfo });
|
||||||
|
document.getElementById('peers').textContent = 'Peers: ' + connections.length;
|
||||||
|
log('Peer ' + (peerInfo.publicKey || '').slice(0, 16) + '…', 'peer');
|
||||||
|
attachReceiver(conn);
|
||||||
|
syncFanout();
|
||||||
|
conn.on('end', function () {
|
||||||
|
connections = connections.filter(function (c) { return c.conn !== conn; });
|
||||||
|
document.getElementById('peers').textContent = 'Peers: ' + connections.length;
|
||||||
|
syncFanout();
|
||||||
|
});
|
||||||
|
});
|
||||||
|
await swarm.join(topic);
|
||||||
|
document.getElementById('btnJoin').disabled = true;
|
||||||
|
document.getElementById('btnLeave').disabled = false;
|
||||||
|
log('Joined ' + topic, 'sys');
|
||||||
|
} catch (e) {
|
||||||
|
log(e.message, 'err');
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
document.getElementById('btnLeave').onclick = async function () {
|
||||||
|
receivers.forEach(function (r) { try { r.destroy(); } catch (_) {} });
|
||||||
|
receivers = [];
|
||||||
|
if (swarm) {
|
||||||
|
try { await swarm.destroy(); } catch (_) {}
|
||||||
|
}
|
||||||
|
swarm = null;
|
||||||
|
connections = [];
|
||||||
|
document.getElementById('peers').textContent = 'Peers: 0';
|
||||||
|
document.getElementById('btnJoin').disabled = false;
|
||||||
|
document.getElementById('btnLeave').disabled = true;
|
||||||
|
await syncFanout();
|
||||||
|
log('Left topic', 'sys');
|
||||||
|
};
|
||||||
|
}).catch(function (e) {
|
||||||
|
setStatus(e.message, true);
|
||||||
|
});
|
||||||
|
})();
|
||||||
@@ -0,0 +1,90 @@
|
|||||||
|
<!DOCTYPE html>
|
||||||
|
<html lang="en">
|
||||||
|
<head>
|
||||||
|
<meta charset="UTF-8">
|
||||||
|
<meta name="viewport" content="width=device-width, initial-scale=1">
|
||||||
|
<title>BridgeSwarm – Live Encode</title>
|
||||||
|
<link rel="stylesheet" href="../shared/theme.css">
|
||||||
|
<link rel="stylesheet" href="../shared/chrome.css">
|
||||||
|
<link rel="stylesheet" href="style.css">
|
||||||
|
<script src="../shared/boot.js"></script>
|
||||||
|
<script src="../shared/live-capture.js"></script>
|
||||||
|
<script>if (BridgeSwarmExamples.guardFileProtocol('live-encode/')) { /* blocked */ }</script>
|
||||||
|
</head>
|
||||||
|
<body>
|
||||||
|
<div class="bs-page bs-page--wide">
|
||||||
|
<a class="bs-back" href="../">← Examples</a>
|
||||||
|
<p class="bs-brand">BridgeSwarm</p>
|
||||||
|
<h1 class="bs-title">Live Encode</h1>
|
||||||
|
<p class="bs-lede">
|
||||||
|
Real host <code>bare-ffmpeg</code> encode: capture → JPEG frames over NMH → VP9/WebM live mux → MSE preview.
|
||||||
|
Not WebRTC — the video you see is what the host encoded.
|
||||||
|
</p>
|
||||||
|
|
||||||
|
<div class="bs-banner bs-banner--info">
|
||||||
|
Default <strong>640×360 @ 10fps</strong> keeps NMH happy. For pure browser P2P A/V without host encode, see
|
||||||
|
<a href="../screenshare/" style="color:var(--accent-2)">Screenshare</a>.
|
||||||
|
</div>
|
||||||
|
<div id="missing" class="bs-banner bs-banner--warn" hidden>
|
||||||
|
Media / live encode unavailable. Rebuild the host and run <code>npm run repair:macos</code> if needed.
|
||||||
|
</div>
|
||||||
|
|
||||||
|
<div class="bs-section">
|
||||||
|
<h2>1. Capture & encode</h2>
|
||||||
|
<p class="bs-status" id="status">Waiting for extension…</p>
|
||||||
|
<div class="bs-row">
|
||||||
|
<button class="primary" id="btnScreen">Screen</button>
|
||||||
|
<button class="secondary" id="btnCamera">Camera</button>
|
||||||
|
<button class="primary" id="btnStart" disabled>Start live encode</button>
|
||||||
|
<button class="danger" id="btnStop" disabled>Stop</button>
|
||||||
|
<span id="livePill" class="bs-pill"></span>
|
||||||
|
</div>
|
||||||
|
<div class="bs-row">
|
||||||
|
<label class="ctrl">Width <input type="number" id="width" value="640" min="160" max="1280"></label>
|
||||||
|
<label class="ctrl">Height <input type="number" id="height" value="360" min="90" max="720"></label>
|
||||||
|
<label class="ctrl">FPS <input type="number" id="fps" value="10" min="1" max="15"></label>
|
||||||
|
<label class="ctrl">CRF <input type="number" id="crf" value="34" min="20" max="50"></label>
|
||||||
|
</div>
|
||||||
|
<div class="bs-meta" id="stats">frames in: 0 · encoded: 0 · out: 0 KB · drops: 0</div>
|
||||||
|
</div>
|
||||||
|
|
||||||
|
<div class="stages">
|
||||||
|
<div class="bs-section stage-col">
|
||||||
|
<h2>Local capture</h2>
|
||||||
|
<div class="bs-stage">
|
||||||
|
<video id="local" autoplay playsinline muted></video>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
<div class="bs-section stage-col">
|
||||||
|
<h2>Host encode (MSE)</h2>
|
||||||
|
<div class="bs-stage">
|
||||||
|
<video id="encoded" autoplay playsinline muted controls></video>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
|
||||||
|
<div class="bs-section">
|
||||||
|
<h2>2. Fan out to peers (optional)</h2>
|
||||||
|
<div class="bs-row">
|
||||||
|
<input type="text" id="topic" value="bridge-swarm-live-encode">
|
||||||
|
<button class="secondary" id="btnJoin" disabled>Join topic</button>
|
||||||
|
<button class="ghost" id="btnLeave" disabled>Leave</button>
|
||||||
|
</div>
|
||||||
|
<div class="bs-meta" id="peers">Peers: 0</div>
|
||||||
|
<div class="bs-row">
|
||||||
|
<label class="ctrl"><input type="checkbox" id="fanout"> Fan out host-encoded stream to peers</label>
|
||||||
|
</div>
|
||||||
|
<div class="bs-stage" id="recvStage">
|
||||||
|
<video id="remote" autoplay playsinline muted controls></video>
|
||||||
|
<p class="placeholder" id="recvPlaceholder">Join in a second tab with fan-out enabled on the broadcaster.</p>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
|
||||||
|
<div class="bs-section">
|
||||||
|
<h2>Log</h2>
|
||||||
|
<div class="bs-log" id="log"></div>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
<script src="app.js"></script>
|
||||||
|
</body>
|
||||||
|
</html>
|
||||||
@@ -0,0 +1,52 @@
|
|||||||
|
.bs-page code {
|
||||||
|
font-family: var(--font-mono);
|
||||||
|
font-size: 0.85em;
|
||||||
|
color: var(--accent-2);
|
||||||
|
}
|
||||||
|
.stages {
|
||||||
|
display: grid;
|
||||||
|
gap: 0.85rem;
|
||||||
|
grid-template-columns: 1fr 1fr;
|
||||||
|
}
|
||||||
|
@media (max-width: 800px) {
|
||||||
|
.stages { grid-template-columns: 1fr; }
|
||||||
|
}
|
||||||
|
.stage-col { margin-bottom: 0; }
|
||||||
|
.bs-stage video {
|
||||||
|
width: 100%;
|
||||||
|
max-height: 320px;
|
||||||
|
background: #0a0e14;
|
||||||
|
}
|
||||||
|
.ctrl {
|
||||||
|
display: inline-flex;
|
||||||
|
align-items: center;
|
||||||
|
gap: 0.35rem;
|
||||||
|
font-size: 0.85rem;
|
||||||
|
color: var(--muted);
|
||||||
|
}
|
||||||
|
.ctrl input[type="number"] {
|
||||||
|
width: 4.5rem;
|
||||||
|
flex: none;
|
||||||
|
}
|
||||||
|
.ctrl input[type="checkbox"] {
|
||||||
|
width: auto;
|
||||||
|
flex: none;
|
||||||
|
}
|
||||||
|
#recvStage {
|
||||||
|
position: relative;
|
||||||
|
min-height: 180px;
|
||||||
|
margin-top: 0.75rem;
|
||||||
|
}
|
||||||
|
#recvPlaceholder {
|
||||||
|
position: absolute;
|
||||||
|
inset: 0;
|
||||||
|
display: flex;
|
||||||
|
align-items: center;
|
||||||
|
justify-content: center;
|
||||||
|
margin: 0;
|
||||||
|
padding: 1rem;
|
||||||
|
text-align: center;
|
||||||
|
color: var(--muted);
|
||||||
|
pointer-events: none;
|
||||||
|
}
|
||||||
|
#recvStage.has-stream #recvPlaceholder { display: none; }
|
||||||
@@ -17,7 +17,10 @@
|
|||||||
<h1 class="bs-title">Screenshare</h1>
|
<h1 class="bs-title">Screenshare</h1>
|
||||||
<p class="bs-lede">Live WebRTC video. BridgeSwarm carries signaling only — not host ffmpeg.</p>
|
<p class="bs-lede">Live WebRTC video. BridgeSwarm carries signaling only — not host ffmpeg.</p>
|
||||||
<div class="bs-banner bs-banner--info">
|
<div class="bs-banner bs-banner--info">
|
||||||
For host <code>bare-ffmpeg</code> batch encode of a recorded clip, see
|
This demo is <strong>WebRTC</strong> (browser encode). For real host
|
||||||
|
<code>bare-ffmpeg</code> live encode → MSE, open
|
||||||
|
<a href="../live-encode/" style="color:var(--accent-2)">Live Encode</a>.
|
||||||
|
Nearline clip jobs:
|
||||||
<a href="../clip-studio/" style="color:var(--accent-2)">Clip Studio</a>.
|
<a href="../clip-studio/" style="color:var(--accent-2)">Clip Studio</a>.
|
||||||
</div>
|
</div>
|
||||||
<p class="bs-status" id="status">Waiting for extension…</p>
|
<p class="bs-status" id="status">Waiting for extension…</p>
|
||||||
|
|||||||
@@ -0,0 +1,90 @@
|
|||||||
|
/**
|
||||||
|
* Shared capture helpers for live-encode / clip-studio demos.
|
||||||
|
*/
|
||||||
|
(function (global) {
|
||||||
|
'use strict';
|
||||||
|
|
||||||
|
function startCapture(kind) {
|
||||||
|
if (kind === 'camera') {
|
||||||
|
return navigator.mediaDevices.getUserMedia({
|
||||||
|
video: { width: { ideal: 1280 }, height: { ideal: 720 } },
|
||||||
|
audio: false,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
return navigator.mediaDevices.getDisplayMedia({
|
||||||
|
video: { frameRate: { ideal: 30 } },
|
||||||
|
audio: false,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Pump video element frames onto a canvas at target FPS, calling onFrame(canvas).
|
||||||
|
* Returns { stop, canvas, setFps }.
|
||||||
|
*/
|
||||||
|
function startFramePump(videoEl, opts) {
|
||||||
|
opts = opts || {};
|
||||||
|
var fps = opts.fps || 10;
|
||||||
|
var maxW = opts.maxWidth || 640;
|
||||||
|
var maxH = opts.maxHeight || 360;
|
||||||
|
var canvas = opts.canvas || document.createElement('canvas');
|
||||||
|
var ctx = canvas.getContext('2d', { alpha: false });
|
||||||
|
var running = true;
|
||||||
|
var timer = null;
|
||||||
|
var lastPush = 0;
|
||||||
|
var inflight = false;
|
||||||
|
|
||||||
|
function sizeCanvas() {
|
||||||
|
var vw = videoEl.videoWidth || maxW;
|
||||||
|
var vh = videoEl.videoHeight || maxH;
|
||||||
|
var scale = Math.min(maxW / vw, maxH / vh, 1);
|
||||||
|
canvas.width = Math.max(2, Math.round(vw * scale / 2) * 2);
|
||||||
|
canvas.height = Math.max(2, Math.round(vh * scale / 2) * 2);
|
||||||
|
}
|
||||||
|
|
||||||
|
function tick() {
|
||||||
|
if (!running) return;
|
||||||
|
var interval = 1000 / fps;
|
||||||
|
var now = performance.now();
|
||||||
|
if (now - lastPush >= interval && !inflight && videoEl.readyState >= 2) {
|
||||||
|
lastPush = now;
|
||||||
|
if (!canvas.width || !canvas.height) sizeCanvas();
|
||||||
|
try {
|
||||||
|
ctx.drawImage(videoEl, 0, 0, canvas.width, canvas.height);
|
||||||
|
} catch (_) {}
|
||||||
|
if (typeof opts.onFrame === 'function') {
|
||||||
|
inflight = true;
|
||||||
|
Promise.resolve(opts.onFrame(canvas))
|
||||||
|
.catch(function () {})
|
||||||
|
.finally(function () { inflight = false; });
|
||||||
|
}
|
||||||
|
}
|
||||||
|
timer = requestAnimationFrame(tick);
|
||||||
|
}
|
||||||
|
|
||||||
|
videoEl.addEventListener('loadedmetadata', sizeCanvas);
|
||||||
|
if (videoEl.videoWidth) sizeCanvas();
|
||||||
|
timer = requestAnimationFrame(tick);
|
||||||
|
|
||||||
|
return {
|
||||||
|
canvas: canvas,
|
||||||
|
stop: function () {
|
||||||
|
running = false;
|
||||||
|
if (timer) cancelAnimationFrame(timer);
|
||||||
|
timer = null;
|
||||||
|
},
|
||||||
|
setFps: function (n) {
|
||||||
|
fps = Math.max(1, Math.min(30, Number(n) || fps));
|
||||||
|
},
|
||||||
|
setMaxSize: function (w, h) {
|
||||||
|
maxW = w;
|
||||||
|
maxH = h;
|
||||||
|
sizeCanvas();
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
global.BridgeSwarmLiveCapture = {
|
||||||
|
startCapture: startCapture,
|
||||||
|
startFramePump: startFramePump,
|
||||||
|
};
|
||||||
|
})(typeof window !== 'undefined' ? window : globalThis);
|
||||||
@@ -554,6 +554,286 @@
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function blobToBase64(blob) {
|
||||||
|
return new Promise(function (resolve, reject) {
|
||||||
|
const reader = new FileReader();
|
||||||
|
reader.onload = function () {
|
||||||
|
const s = String(reader.result || '');
|
||||||
|
const i = s.indexOf(',');
|
||||||
|
resolve(i >= 0 ? s.slice(i + 1) : s);
|
||||||
|
};
|
||||||
|
reader.onerror = reject;
|
||||||
|
reader.readAsDataURL(blob);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
function canvasToJpegBase64(canvas, quality) {
|
||||||
|
return new Promise(function (resolve, reject) {
|
||||||
|
if (!canvas.toBlob) {
|
||||||
|
const dataUrl = canvas.toDataURL('image/jpeg', quality == null ? 0.72 : quality);
|
||||||
|
const i = dataUrl.indexOf(',');
|
||||||
|
resolve(i >= 0 ? dataUrl.slice(i + 1) : dataUrl);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
canvas.toBlob(function (blob) {
|
||||||
|
if (!blob) {
|
||||||
|
reject(new Error('canvas.toBlob failed'));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
blobToBase64(blob).then(resolve, reject);
|
||||||
|
}, 'image/jpeg', quality == null ? 0.72 : quality);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Live encode session helper: host VP9/WebM → MSE preview.
|
||||||
|
* @param {HTMLVideoElement} videoEl
|
||||||
|
* @param {object} [opts]
|
||||||
|
*/
|
||||||
|
function createLiveSession(videoEl, opts) {
|
||||||
|
opts = opts || {};
|
||||||
|
let sessionId = null;
|
||||||
|
let mediaSource = null;
|
||||||
|
let sourceBuffer = null;
|
||||||
|
let queue = [];
|
||||||
|
let appending = false;
|
||||||
|
let closed = false;
|
||||||
|
const pendingBySeq = new Map();
|
||||||
|
const listeners = { onSegment: opts.onSegment, onStats: opts.onStats, onDrop: opts.onDrop };
|
||||||
|
|
||||||
|
function flushQueue() {
|
||||||
|
if (appending || !sourceBuffer || sourceBuffer.updating || !queue.length) return;
|
||||||
|
appending = true;
|
||||||
|
const buf = queue.shift();
|
||||||
|
try {
|
||||||
|
sourceBuffer.appendBuffer(buf);
|
||||||
|
} catch (err) {
|
||||||
|
appending = false;
|
||||||
|
if (listeners.onStats) listeners.onStats({ appendError: err.message });
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function onSourceBufferUpdateEnd() {
|
||||||
|
appending = false;
|
||||||
|
flushQueue();
|
||||||
|
}
|
||||||
|
|
||||||
|
function ensureMse() {
|
||||||
|
if (mediaSource) return Promise.resolve();
|
||||||
|
return new Promise(function (resolve, reject) {
|
||||||
|
mediaSource = new MediaSource();
|
||||||
|
videoEl.src = URL.createObjectURL(mediaSource);
|
||||||
|
mediaSource.addEventListener('sourceopen', function () {
|
||||||
|
try {
|
||||||
|
const mime = opts.mimeType || 'video/webm; codecs="vp9"';
|
||||||
|
sourceBuffer = mediaSource.addSourceBuffer(mime);
|
||||||
|
sourceBuffer.mode = 'sequence';
|
||||||
|
sourceBuffer.addEventListener('updateend', onSourceBufferUpdateEnd);
|
||||||
|
resolve();
|
||||||
|
} catch (err) {
|
||||||
|
reject(err);
|
||||||
|
}
|
||||||
|
}, { once: true });
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
function onCapEvent(e) {
|
||||||
|
const msg = e.detail;
|
||||||
|
if (!msg || msg.type !== 'event' || !sessionId) return;
|
||||||
|
const payload = msg.payload || {};
|
||||||
|
if (payload.sessionId !== sessionId && payload.jobId !== sessionId) return;
|
||||||
|
if (msg.event === 'cap-chunk') {
|
||||||
|
if (payload.kind === 'drop') {
|
||||||
|
if (listeners.onDrop) listeners.onDrop(payload);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (payload.kind === 'segment' && payload.data) {
|
||||||
|
const seq = payload.seq != null ? payload.seq : 0;
|
||||||
|
let entry = pendingBySeq.get(seq);
|
||||||
|
if (!entry) {
|
||||||
|
entry = { parts: [], expected: payload.parts || 1 };
|
||||||
|
pendingBySeq.set(seq, entry);
|
||||||
|
}
|
||||||
|
if (payload.parts) entry.expected = payload.parts;
|
||||||
|
entry.parts[payload.index != null ? payload.index : 0] = payload.data;
|
||||||
|
let ready = entry.parts.length >= entry.expected;
|
||||||
|
for (let i = 0; i < entry.expected; i++) {
|
||||||
|
if (entry.parts[i] == null) ready = false;
|
||||||
|
}
|
||||||
|
if (ready) {
|
||||||
|
pendingBySeq.delete(seq);
|
||||||
|
flushSegmentBase64(entry.parts.join(''));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (msg.event === 'cap-end' || msg.event === 'cap-error') {
|
||||||
|
if (listeners.onStats) listeners.onStats(payload);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function flushSegmentBase64(b64) {
|
||||||
|
const binary = atob(b64);
|
||||||
|
const bytes = new Uint8Array(binary.length);
|
||||||
|
for (let i = 0; i < binary.length; i++) bytes[i] = binary.charCodeAt(i);
|
||||||
|
if (listeners.onSegment) listeners.onSegment(bytes);
|
||||||
|
ensureMse().then(function () {
|
||||||
|
queue.push(bytes.buffer);
|
||||||
|
flushQueue();
|
||||||
|
}).catch(function (err) {
|
||||||
|
if (listeners.onStats) listeners.onStats({ mseError: err.message });
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
window.addEventListener('bridge-swarm-event', onCapEvent);
|
||||||
|
|
||||||
|
return {
|
||||||
|
start: function (startOpts) {
|
||||||
|
startOpts = Object.assign({
|
||||||
|
width: 640,
|
||||||
|
height: 360,
|
||||||
|
fps: 10,
|
||||||
|
ingest: 'frames',
|
||||||
|
egress: { page: true, file: true },
|
||||||
|
}, startOpts || {});
|
||||||
|
if (!startOpts.sessionId) {
|
||||||
|
startOpts.sessionId = 'live_' + Date.now() + '_' + Math.random().toString(36).slice(2, 8);
|
||||||
|
}
|
||||||
|
sessionId = startOpts.sessionId;
|
||||||
|
return BridgeSwarm.media.encodeStart(startOpts).then(function (r) {
|
||||||
|
sessionId = r.sessionId || sessionId;
|
||||||
|
return r;
|
||||||
|
});
|
||||||
|
},
|
||||||
|
pushFrame: function (canvasOrBlob, quality) {
|
||||||
|
if (!sessionId) return Promise.reject(new Error('session not started'));
|
||||||
|
const p = (canvasOrBlob instanceof HTMLCanvasElement || (canvasOrBlob && canvasOrBlob.tagName === 'CANVAS'))
|
||||||
|
? canvasToJpegBase64(canvasOrBlob, quality)
|
||||||
|
: (canvasOrBlob instanceof Blob
|
||||||
|
? blobToBase64(canvasOrBlob)
|
||||||
|
: Promise.resolve(canvasOrBlob));
|
||||||
|
return p.then(function (dataBase64) {
|
||||||
|
return BridgeSwarm.media.encodePushFrame({
|
||||||
|
sessionId: sessionId,
|
||||||
|
dataBase64: dataBase64,
|
||||||
|
mimetype: 'image/jpeg',
|
||||||
|
});
|
||||||
|
});
|
||||||
|
},
|
||||||
|
pushSegment: function (blob) {
|
||||||
|
if (!sessionId) return Promise.reject(new Error('session not started'));
|
||||||
|
return blobToBase64(blob).then(function (dataBase64) {
|
||||||
|
return BridgeSwarm.media.encodePush({
|
||||||
|
sessionId: sessionId,
|
||||||
|
dataBase64: dataBase64,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
},
|
||||||
|
subscribe: function (sub) {
|
||||||
|
if (!sessionId) return Promise.reject(new Error('session not started'));
|
||||||
|
return BridgeSwarm.media.encodeSubscribe(Object.assign({ sessionId: sessionId }, sub || {}));
|
||||||
|
},
|
||||||
|
stop: function () {
|
||||||
|
if (closed) return Promise.resolve();
|
||||||
|
closed = true;
|
||||||
|
window.removeEventListener('bridge-swarm-event', onCapEvent);
|
||||||
|
if (!sessionId) return Promise.resolve();
|
||||||
|
const id = sessionId;
|
||||||
|
sessionId = null;
|
||||||
|
return BridgeSwarm.media.encodeStop(id).finally(function () {
|
||||||
|
try {
|
||||||
|
if (mediaSource && mediaSource.readyState === 'open') mediaSource.endOfStream();
|
||||||
|
} catch (_) {}
|
||||||
|
});
|
||||||
|
},
|
||||||
|
get sessionId() { return sessionId; },
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Parse BSML-framed host media live packets from a swarm connection into MSE. */
|
||||||
|
function attachLiveReceiver(conn, videoEl, opts) {
|
||||||
|
opts = opts || {};
|
||||||
|
let buf = new Uint8Array(0);
|
||||||
|
const session = {
|
||||||
|
closed: false,
|
||||||
|
_queue: [],
|
||||||
|
_ms: null,
|
||||||
|
_sb: null,
|
||||||
|
_appending: false,
|
||||||
|
};
|
||||||
|
|
||||||
|
function concat(a, b) {
|
||||||
|
const out = new Uint8Array(a.length + b.length);
|
||||||
|
out.set(a, 0);
|
||||||
|
out.set(b, a.length);
|
||||||
|
return out;
|
||||||
|
}
|
||||||
|
|
||||||
|
function flush() {
|
||||||
|
if (session._appending || !session._sb || session._sb.updating || !session._queue.length) return;
|
||||||
|
session._appending = true;
|
||||||
|
try {
|
||||||
|
session._sb.appendBuffer(session._queue.shift());
|
||||||
|
} catch (_) {
|
||||||
|
session._appending = false;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function ensureMse() {
|
||||||
|
if (session._ms) return Promise.resolve();
|
||||||
|
return new Promise(function (resolve, reject) {
|
||||||
|
session._ms = new MediaSource();
|
||||||
|
videoEl.src = URL.createObjectURL(session._ms);
|
||||||
|
session._ms.addEventListener('sourceopen', function () {
|
||||||
|
try {
|
||||||
|
session._sb = session._ms.addSourceBuffer(opts.mimeType || 'video/webm; codecs="vp9"');
|
||||||
|
session._sb.mode = 'sequence';
|
||||||
|
session._sb.addEventListener('updateend', function () {
|
||||||
|
session._appending = false;
|
||||||
|
flush();
|
||||||
|
});
|
||||||
|
resolve();
|
||||||
|
} catch (err) {
|
||||||
|
reject(err);
|
||||||
|
}
|
||||||
|
}, { once: true });
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
function onData(data) {
|
||||||
|
const chunk = data instanceof Uint8Array ? data : new Uint8Array(data);
|
||||||
|
buf = concat(buf, chunk);
|
||||||
|
while (buf.length >= 8) {
|
||||||
|
if (buf[0] !== 0x42 || buf[1] !== 0x53 || buf[2] !== 0x4d || buf[3] !== 0x4c) {
|
||||||
|
// resync
|
||||||
|
buf = buf.subarray(1);
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
const len = (buf[4] << 24) | (buf[5] << 16) | (buf[6] << 8) | buf[7];
|
||||||
|
if (buf.length < 8 + len) break;
|
||||||
|
const payload = buf.subarray(8, 8 + len);
|
||||||
|
buf = buf.subarray(8 + len);
|
||||||
|
const copy = payload.slice();
|
||||||
|
ensureMse().then(function () {
|
||||||
|
session._queue.push(copy.buffer);
|
||||||
|
flush();
|
||||||
|
}).catch(function () {});
|
||||||
|
if (opts.onSegment) opts.onSegment(copy);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
conn.on('data', onData);
|
||||||
|
return {
|
||||||
|
destroy: function () {
|
||||||
|
session.closed = true;
|
||||||
|
conn.off('data', onData);
|
||||||
|
try {
|
||||||
|
if (session._ms && session._ms.readyState === 'open') session._ms.endOfStream();
|
||||||
|
} catch (_) {}
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
BridgeSwarm.capabilities = {
|
BridgeSwarm.capabilities = {
|
||||||
list: function (options) {
|
list: function (options) {
|
||||||
return BridgeSwarm.request('capabilities.list', {}, options).then(function (r) {
|
return BridgeSwarm.request('capabilities.list', {}, options).then(function (r) {
|
||||||
@@ -602,6 +882,24 @@
|
|||||||
readOutput: function (payload, options) {
|
readOutput: function (payload, options) {
|
||||||
return capabilityCall('media', 'readOutput', payload, options);
|
return capabilityCall('media', 'readOutput', payload, options);
|
||||||
},
|
},
|
||||||
|
encodeStart: function (payload, options) {
|
||||||
|
return capabilityCall('media', 'encodeStart', payload, options);
|
||||||
|
},
|
||||||
|
encodePushFrame: function (payload, options) {
|
||||||
|
return capabilityCall('media', 'encodePushFrame', payload, options);
|
||||||
|
},
|
||||||
|
encodePush: function (payload, options) {
|
||||||
|
return capabilityCall('media', 'encodePush', payload, options);
|
||||||
|
},
|
||||||
|
encodeSubscribe: function (payload, options) {
|
||||||
|
return capabilityCall('media', 'encodeSubscribe', payload, options);
|
||||||
|
},
|
||||||
|
encodeStop: function (sessionId, options) {
|
||||||
|
const payload = typeof sessionId === 'string' ? { sessionId: sessionId } : (sessionId || {});
|
||||||
|
return capabilityCall('media', 'encodeStop', payload, options);
|
||||||
|
},
|
||||||
|
liveSession: createLiveSession,
|
||||||
|
attachLiveReceiver: attachLiveReceiver,
|
||||||
};
|
};
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -0,0 +1,581 @@
|
|||||||
|
/**
|
||||||
|
* Live bare-ffmpeg encode sessions for the media capability pack.
|
||||||
|
* Ingest JPEG/WebP frames (or MediaRecorder WebM slices) → VP9/WebM live mux
|
||||||
|
* → page cap-chunk segments + optional Hyperswarm fanout + file archive.
|
||||||
|
*/
|
||||||
|
|
||||||
|
const path = require('bare-path');
|
||||||
|
const fs = require('bare-fs');
|
||||||
|
const b4a = require('b4a');
|
||||||
|
const { ensureJobDir, sanitizeId, getCapJobsRoot } = require('./paths.js');
|
||||||
|
|
||||||
|
const MAX_CHUNK_CHARS = 700000;
|
||||||
|
const MAX_DURATION_MS = 10 * 60 * 1000;
|
||||||
|
const MAX_PAGE_QUEUE = 8;
|
||||||
|
const MAGIC = Buffer.from('BSML'); // BridgeSwarm Media Live
|
||||||
|
|
||||||
|
/** @type {Map<string, object>} */
|
||||||
|
const sessions = new Map();
|
||||||
|
|
||||||
|
let nextSeq = 0;
|
||||||
|
|
||||||
|
function makeSessionId() {
|
||||||
|
return `live_${Date.now()}_${nextSeq++}`;
|
||||||
|
}
|
||||||
|
|
||||||
|
function base64ToBuffer(b64) {
|
||||||
|
return b4a.from(b64, 'base64');
|
||||||
|
}
|
||||||
|
|
||||||
|
function bufferToBase64(buf) {
|
||||||
|
return b4a.toString(buf, 'base64');
|
||||||
|
}
|
||||||
|
|
||||||
|
function frameSwarmPayload(buf) {
|
||||||
|
const header = Buffer.alloc(8);
|
||||||
|
MAGIC.copy(header, 0);
|
||||||
|
header.writeUInt32BE(buf.length, 4);
|
||||||
|
return Buffer.concat([header, buf]);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @param {object} bareMedia - { image }
|
||||||
|
* @param {object} ffmpeg - bare-ffmpeg module
|
||||||
|
*/
|
||||||
|
function createLiveEncodeCommands(bareMedia, ffmpeg) {
|
||||||
|
const { image } = bareMedia;
|
||||||
|
|
||||||
|
function emitSegment(session, chunk) {
|
||||||
|
const buf = b4a.from(chunk);
|
||||||
|
session.stats.bytesOut += buf.length;
|
||||||
|
|
||||||
|
if (session.fileStream) {
|
||||||
|
try {
|
||||||
|
session.fileStream.write(buf);
|
||||||
|
} catch (_) {}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (session.egress.swarm && session.egress.swarm.connIds && session.writeToConn) {
|
||||||
|
const framed = frameSwarmPayload(buf);
|
||||||
|
for (const connId of session.egress.swarm.connIds) {
|
||||||
|
try {
|
||||||
|
session.writeToConn(connId, framed);
|
||||||
|
} catch (_) {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (session.egress.page !== false) {
|
||||||
|
if (session.pageQueue >= MAX_PAGE_QUEUE) {
|
||||||
|
session.stats.drops += 1;
|
||||||
|
session.emit('cap-chunk', {
|
||||||
|
pack: 'media',
|
||||||
|
sessionId: session.id,
|
||||||
|
jobId: session.id,
|
||||||
|
kind: 'drop',
|
||||||
|
frames: 1,
|
||||||
|
drops: session.stats.drops,
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
session.pageQueue += 1;
|
||||||
|
if (session._segSeq == null) session._segSeq = 0;
|
||||||
|
const seq = session._segSeq++;
|
||||||
|
const b64 = bufferToBase64(buf);
|
||||||
|
let index = 0;
|
||||||
|
const parts = Math.ceil(b64.length / MAX_CHUNK_CHARS) || 1;
|
||||||
|
for (let offset = 0; offset < b64.length; offset += MAX_CHUNK_CHARS) {
|
||||||
|
session.emit('cap-chunk', {
|
||||||
|
pack: 'media',
|
||||||
|
sessionId: session.id,
|
||||||
|
jobId: session.id,
|
||||||
|
kind: 'segment',
|
||||||
|
seq,
|
||||||
|
index,
|
||||||
|
parts,
|
||||||
|
data: b64.slice(offset, offset + MAX_CHUNK_CHARS),
|
||||||
|
bytes: buf.length,
|
||||||
|
});
|
||||||
|
index++;
|
||||||
|
}
|
||||||
|
session.pageQueue = Math.max(0, session.pageQueue - 1);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function destroySession(session) {
|
||||||
|
if (!session || session.destroyed) return;
|
||||||
|
session.destroyed = true;
|
||||||
|
session.cancelled = true;
|
||||||
|
try {
|
||||||
|
if (session.encoder) {
|
||||||
|
try {
|
||||||
|
session.encoder.sendFrame(null);
|
||||||
|
const pkt = new ffmpeg.Packet();
|
||||||
|
while (session.encoder.receivePacket(pkt)) {
|
||||||
|
pkt.streamIndex = session.outStream.index;
|
||||||
|
try {
|
||||||
|
pkt.rescaleTimestamps(session.encoder.timeBase, session.outStream.timeBase);
|
||||||
|
} catch (_) {}
|
||||||
|
session.outputFormat.writeFrame(pkt);
|
||||||
|
pkt.unref();
|
||||||
|
}
|
||||||
|
} catch (_) {}
|
||||||
|
}
|
||||||
|
if (session.outputFormat) {
|
||||||
|
try {
|
||||||
|
session.outputFormat.writeTrailer();
|
||||||
|
} catch (_) {}
|
||||||
|
try {
|
||||||
|
session.outputFormat.destroy();
|
||||||
|
} catch (_) {}
|
||||||
|
}
|
||||||
|
} catch (_) {}
|
||||||
|
try {
|
||||||
|
if (session.encoder) session.encoder.destroy();
|
||||||
|
} catch (_) {}
|
||||||
|
try {
|
||||||
|
if (session.scaler) session.scaler.destroy();
|
||||||
|
} catch (_) {}
|
||||||
|
try {
|
||||||
|
if (session.rgbaFrame) session.rgbaFrame.destroy();
|
||||||
|
} catch (_) {}
|
||||||
|
try {
|
||||||
|
if (session.yuvFrame) session.yuvFrame.destroy();
|
||||||
|
} catch (_) {}
|
||||||
|
try {
|
||||||
|
if (session.fileStream) session.fileStream.end();
|
||||||
|
} catch (_) {}
|
||||||
|
sessions.delete(session.id);
|
||||||
|
}
|
||||||
|
|
||||||
|
async function openSession(opts, ctx) {
|
||||||
|
const width = Math.min(1280, Math.max(160, Number(opts.width) || 640));
|
||||||
|
const height = Math.min(720, Math.max(90, Number(opts.height) || 360));
|
||||||
|
const fps = Math.min(30, Math.max(1, Number(opts.fps) || 10));
|
||||||
|
const ingest = opts.ingest === 'segments' ? 'segments' : 'frames';
|
||||||
|
const egress = opts.egress || { page: true, file: true };
|
||||||
|
const sessionId = sanitizeId(opts.sessionId || makeSessionId());
|
||||||
|
|
||||||
|
if (sessions.has(sessionId)) {
|
||||||
|
throw new Error('session already exists');
|
||||||
|
}
|
||||||
|
|
||||||
|
ensureJobDir(sessionId);
|
||||||
|
const archiveRel = path.join(sessionId, 'live.webm');
|
||||||
|
const archivePath = path.join(getCapJobsRoot(), archiveRel);
|
||||||
|
let fileStream = null;
|
||||||
|
if (egress.file !== false) {
|
||||||
|
fileStream = fs.createWriteStream(archivePath);
|
||||||
|
}
|
||||||
|
|
||||||
|
const pendingChunks = [];
|
||||||
|
const sessionRef = { session: null };
|
||||||
|
const io = new ffmpeg.IOContext(64 * 1024, {
|
||||||
|
onwrite: (chunk) => {
|
||||||
|
const copy = b4a.from(chunk);
|
||||||
|
if (sessionRef.session) emitSegment(sessionRef.session, copy);
|
||||||
|
else pendingChunks.push(copy);
|
||||||
|
return chunk.length;
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
const outputFormat = new ffmpeg.OutputFormatContext('webm', io);
|
||||||
|
const outStream = outputFormat.createStream();
|
||||||
|
outStream.codecParameters.id = ffmpeg.constants.codecs.VP9;
|
||||||
|
outStream.codecParameters.type = ffmpeg.constants.mediaTypes.VIDEO;
|
||||||
|
outStream.codecParameters.width = width;
|
||||||
|
outStream.codecParameters.height = height;
|
||||||
|
outStream.codecParameters.format = ffmpeg.constants.pixelFormats.YUV420P;
|
||||||
|
outStream.timeBase = new ffmpeg.Rational(1, fps);
|
||||||
|
|
||||||
|
const encoder = outStream.encoder();
|
||||||
|
encoder.timeBase = outStream.timeBase;
|
||||||
|
try {
|
||||||
|
encoder.frameRate = new ffmpeg.Rational(fps, 1);
|
||||||
|
} catch (_) {
|
||||||
|
try {
|
||||||
|
encoder.framerate = new ffmpeg.Rational(fps, 1);
|
||||||
|
} catch (_) {}
|
||||||
|
}
|
||||||
|
encoder.gopSize = Math.max(fps, 10);
|
||||||
|
|
||||||
|
const encOpts = ffmpeg.Dictionary.from({
|
||||||
|
allow_sw: '1',
|
||||||
|
deadline: 'realtime',
|
||||||
|
'cpu-used': '6',
|
||||||
|
crf: String(opts.crf != null ? opts.crf : 34),
|
||||||
|
b: '0',
|
||||||
|
});
|
||||||
|
encoder.open(encOpts);
|
||||||
|
outStream.codecParameters.fromContext(encoder);
|
||||||
|
|
||||||
|
const muxOpts = ffmpeg.Dictionary.from({ live: '1' });
|
||||||
|
outputFormat.writeHeader(muxOpts);
|
||||||
|
|
||||||
|
const scaler = new ffmpeg.Scaler(
|
||||||
|
ffmpeg.constants.pixelFormats.RGBA,
|
||||||
|
width,
|
||||||
|
height,
|
||||||
|
ffmpeg.constants.pixelFormats.YUV420P,
|
||||||
|
width,
|
||||||
|
height
|
||||||
|
);
|
||||||
|
|
||||||
|
const rgbaFrame = new ffmpeg.Frame();
|
||||||
|
rgbaFrame.width = width;
|
||||||
|
rgbaFrame.height = height;
|
||||||
|
rgbaFrame.format = ffmpeg.constants.pixelFormats.RGBA;
|
||||||
|
rgbaFrame.alloc();
|
||||||
|
|
||||||
|
const yuvFrame = new ffmpeg.Frame();
|
||||||
|
yuvFrame.width = width;
|
||||||
|
yuvFrame.height = height;
|
||||||
|
yuvFrame.format = ffmpeg.constants.pixelFormats.YUV420P;
|
||||||
|
yuvFrame.alloc();
|
||||||
|
|
||||||
|
const session = {
|
||||||
|
id: sessionId,
|
||||||
|
width,
|
||||||
|
height,
|
||||||
|
fps,
|
||||||
|
ingest,
|
||||||
|
egress: {
|
||||||
|
page: egress.page !== false,
|
||||||
|
file: egress.file !== false,
|
||||||
|
swarm: egress.swarm || { connIds: [] },
|
||||||
|
},
|
||||||
|
emit: ctx.emit,
|
||||||
|
writeToConn: ctx.writeToConn,
|
||||||
|
outputFormat,
|
||||||
|
outStream,
|
||||||
|
encoder,
|
||||||
|
scaler,
|
||||||
|
rgbaFrame,
|
||||||
|
yuvFrame,
|
||||||
|
fileStream,
|
||||||
|
archiveRel,
|
||||||
|
archivePath,
|
||||||
|
pts: 0,
|
||||||
|
cancelled: false,
|
||||||
|
destroyed: false,
|
||||||
|
pageQueue: 0,
|
||||||
|
startedAt: Date.now(),
|
||||||
|
maxDurationMs: opts.maxDurationMs || MAX_DURATION_MS,
|
||||||
|
stats: { framesIn: 0, framesEncoded: 0, bytesOut: 0, drops: 0 },
|
||||||
|
segmentBuffer: Buffer.alloc(0),
|
||||||
|
packet: new ffmpeg.Packet(),
|
||||||
|
};
|
||||||
|
|
||||||
|
sessionRef.session = session;
|
||||||
|
sessions.set(sessionId, session);
|
||||||
|
|
||||||
|
for (const c of pendingChunks) emitSegment(session, c);
|
||||||
|
|
||||||
|
return session;
|
||||||
|
}
|
||||||
|
|
||||||
|
function encodeRgbaFrame(session, rgba) {
|
||||||
|
if (session.cancelled || session.destroyed) throw new Error('session cancelled');
|
||||||
|
if (Date.now() - session.startedAt > session.maxDurationMs) {
|
||||||
|
throw new Error('session max duration exceeded');
|
||||||
|
}
|
||||||
|
|
||||||
|
const srcW = rgba.width;
|
||||||
|
const srcH = rgba.height;
|
||||||
|
let data = rgba.data;
|
||||||
|
if (rgba.frames && rgba.frames[0]) {
|
||||||
|
data = rgba.frames[0].data || rgba.frames[0];
|
||||||
|
}
|
||||||
|
|
||||||
|
if (srcW !== session.width || srcH !== session.height || !session._scalerSrc) {
|
||||||
|
try {
|
||||||
|
if (session.scaler) session.scaler.destroy();
|
||||||
|
} catch (_) {}
|
||||||
|
session.scaler = new ffmpeg.Scaler(
|
||||||
|
ffmpeg.constants.pixelFormats.RGBA,
|
||||||
|
srcW,
|
||||||
|
srcH,
|
||||||
|
ffmpeg.constants.pixelFormats.YUV420P,
|
||||||
|
session.width,
|
||||||
|
session.height
|
||||||
|
);
|
||||||
|
session._scalerSrc = srcW + 'x' + srcH;
|
||||||
|
try {
|
||||||
|
session.rgbaFrame.destroy();
|
||||||
|
} catch (_) {}
|
||||||
|
session.rgbaFrame = new ffmpeg.Frame();
|
||||||
|
session.rgbaFrame.width = srcW;
|
||||||
|
session.rgbaFrame.height = srcH;
|
||||||
|
session.rgbaFrame.format = ffmpeg.constants.pixelFormats.RGBA;
|
||||||
|
session.rgbaFrame.alloc();
|
||||||
|
}
|
||||||
|
|
||||||
|
const imageHandle = new ffmpeg.Image(ffmpeg.constants.pixelFormats.RGBA, srcW, srcH);
|
||||||
|
const srcBuf = b4a.from(data);
|
||||||
|
srcBuf.copy(imageHandle.data);
|
||||||
|
imageHandle.fill(session.rgbaFrame);
|
||||||
|
|
||||||
|
session.scaler.scale(session.rgbaFrame, session.yuvFrame);
|
||||||
|
session.yuvFrame.pts = session.pts++;
|
||||||
|
session.stats.framesIn += 1;
|
||||||
|
|
||||||
|
if (!session.encoder.sendFrame(session.yuvFrame)) {
|
||||||
|
session.stats.drops += 1;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
while (session.encoder.receivePacket(session.packet)) {
|
||||||
|
session.packet.streamIndex = session.outStream.index;
|
||||||
|
try {
|
||||||
|
session.packet.rescaleTimestamps(session.encoder.timeBase, session.outStream.timeBase);
|
||||||
|
} catch (_) {}
|
||||||
|
session.outputFormat.writeFrame(session.packet);
|
||||||
|
session.stats.framesEncoded += 1;
|
||||||
|
session.packet.unref();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function pushFrame(session, payload) {
|
||||||
|
const buf = base64ToBuffer(payload.dataBase64 || '');
|
||||||
|
if (buf.length > 900000) throw new Error('frame too large for NMH (max ~900KB)');
|
||||||
|
const rgba = await image.decode(buf, { maxFrames: 1 });
|
||||||
|
encodeRgbaFrame(session, rgba);
|
||||||
|
}
|
||||||
|
|
||||||
|
async function pushSegment(session, payload) {
|
||||||
|
const buf = base64ToBuffer(payload.dataBase64 || '');
|
||||||
|
if (buf.length > 900000) throw new Error('segment too large for NMH');
|
||||||
|
session.segmentBuffer = Buffer.concat([session.segmentBuffer, buf]);
|
||||||
|
|
||||||
|
let offset = 0;
|
||||||
|
const size = () => session.segmentBuffer.length;
|
||||||
|
const io = new ffmpeg.IOContext(8192, {
|
||||||
|
onread: (buffer, requested) => {
|
||||||
|
const avail = size() - offset;
|
||||||
|
if (avail <= 0) return 0;
|
||||||
|
const n = Math.min(requested, avail);
|
||||||
|
session.segmentBuffer.copy(buffer, 0, offset, offset + n);
|
||||||
|
offset += n;
|
||||||
|
return n;
|
||||||
|
},
|
||||||
|
onseek: (o, whence) => {
|
||||||
|
if (whence === ffmpeg.constants.seek.SIZE) return size();
|
||||||
|
if (whence === ffmpeg.constants.seek.SET) offset = o;
|
||||||
|
else if (whence === ffmpeg.constants.seek.CUR) offset += o;
|
||||||
|
else if (whence === ffmpeg.constants.seek.END) offset = size() + o;
|
||||||
|
else return -1;
|
||||||
|
return offset;
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
let inputFormat = null;
|
||||||
|
try {
|
||||||
|
inputFormat = new ffmpeg.InputFormatContext(io);
|
||||||
|
const stream = inputFormat.getBestStream(ffmpeg.constants.mediaTypes.VIDEO);
|
||||||
|
if (!stream) return;
|
||||||
|
const decoder = stream.decoder();
|
||||||
|
decoder.open();
|
||||||
|
const pkt = new ffmpeg.Packet();
|
||||||
|
const frame = new ffmpeg.Frame();
|
||||||
|
let decoded = 0;
|
||||||
|
while (inputFormat.readFrame(pkt)) {
|
||||||
|
if (pkt.streamIndex !== stream.index) {
|
||||||
|
pkt.unref();
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if (decoder.sendPacket(pkt)) {
|
||||||
|
while (decoder.receiveFrame(frame)) {
|
||||||
|
const toRgba = new ffmpeg.Scaler(
|
||||||
|
frame.format,
|
||||||
|
frame.width,
|
||||||
|
frame.height,
|
||||||
|
ffmpeg.constants.pixelFormats.RGBA,
|
||||||
|
frame.width,
|
||||||
|
frame.height
|
||||||
|
);
|
||||||
|
const rgbaFrame = new ffmpeg.Frame();
|
||||||
|
rgbaFrame.width = frame.width;
|
||||||
|
rgbaFrame.height = frame.height;
|
||||||
|
rgbaFrame.format = ffmpeg.constants.pixelFormats.RGBA;
|
||||||
|
rgbaFrame.alloc();
|
||||||
|
toRgba.scale(frame, rgbaFrame);
|
||||||
|
const img = new ffmpeg.Image(
|
||||||
|
ffmpeg.constants.pixelFormats.RGBA,
|
||||||
|
rgbaFrame.width,
|
||||||
|
rgbaFrame.height
|
||||||
|
);
|
||||||
|
img.read(rgbaFrame);
|
||||||
|
encodeRgbaFrame(session, {
|
||||||
|
width: rgbaFrame.width,
|
||||||
|
height: rgbaFrame.height,
|
||||||
|
data: img.data,
|
||||||
|
});
|
||||||
|
try { toRgba.destroy(); } catch (_) {}
|
||||||
|
try { rgbaFrame.destroy(); } catch (_) {}
|
||||||
|
decoded += 1;
|
||||||
|
if (decoded >= 2) break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
pkt.unref();
|
||||||
|
if (decoded >= 2) break;
|
||||||
|
}
|
||||||
|
try { decoder.destroy(); } catch (_) {}
|
||||||
|
} catch (err) {
|
||||||
|
if (!/Invalid|End of file|Immediate exit|Invalid data/i.test(String(err.message || err))) {
|
||||||
|
throw err;
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
try {
|
||||||
|
if (inputFormat) inputFormat.destroy();
|
||||||
|
} catch (_) {}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (session.segmentBuffer.length > 8 * 1024 * 1024) {
|
||||||
|
session.segmentBuffer = session.segmentBuffer.subarray(session.segmentBuffer.length - 4 * 1024 * 1024);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
async encodeStart(ctx) {
|
||||||
|
const { payload, reply } = ctx;
|
||||||
|
try {
|
||||||
|
const session = await openSession(payload || {}, ctx);
|
||||||
|
reply({
|
||||||
|
ok: true,
|
||||||
|
sessionId: session.id,
|
||||||
|
jobId: session.id,
|
||||||
|
status: 'live',
|
||||||
|
width: session.width,
|
||||||
|
height: session.height,
|
||||||
|
fps: session.fps,
|
||||||
|
ingest: session.ingest,
|
||||||
|
format: 'webm',
|
||||||
|
archivePath: session.egress.file ? session.archiveRel : undefined,
|
||||||
|
});
|
||||||
|
} catch (err) {
|
||||||
|
reply({ ok: false, error: err.message });
|
||||||
|
}
|
||||||
|
},
|
||||||
|
|
||||||
|
async encodePushFrame(ctx) {
|
||||||
|
const { payload, reply } = ctx;
|
||||||
|
const session = sessions.get(payload.sessionId || payload.jobId);
|
||||||
|
if (!session) {
|
||||||
|
reply({ ok: false, error: 'session not found' });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
await pushFrame(session, payload);
|
||||||
|
reply({
|
||||||
|
ok: true,
|
||||||
|
sessionId: session.id,
|
||||||
|
framesIn: session.stats.framesIn,
|
||||||
|
framesEncoded: session.stats.framesEncoded,
|
||||||
|
bytesOut: session.stats.bytesOut,
|
||||||
|
drops: session.stats.drops,
|
||||||
|
});
|
||||||
|
} catch (err) {
|
||||||
|
reply({ ok: false, error: err.message, sessionId: session.id });
|
||||||
|
}
|
||||||
|
},
|
||||||
|
|
||||||
|
async encodePush(ctx) {
|
||||||
|
const { payload, reply } = ctx;
|
||||||
|
const session = sessions.get(payload.sessionId || payload.jobId);
|
||||||
|
if (!session) {
|
||||||
|
reply({ ok: false, error: 'session not found' });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
if (session.ingest === 'frames') {
|
||||||
|
await pushFrame(session, payload);
|
||||||
|
} else {
|
||||||
|
await pushSegment(session, payload);
|
||||||
|
}
|
||||||
|
reply({
|
||||||
|
ok: true,
|
||||||
|
sessionId: session.id,
|
||||||
|
framesIn: session.stats.framesIn,
|
||||||
|
framesEncoded: session.stats.framesEncoded,
|
||||||
|
bytesOut: session.stats.bytesOut,
|
||||||
|
drops: session.stats.drops,
|
||||||
|
});
|
||||||
|
} catch (err) {
|
||||||
|
reply({ ok: false, error: err.message, sessionId: session.id });
|
||||||
|
}
|
||||||
|
},
|
||||||
|
|
||||||
|
async encodeSubscribe(ctx) {
|
||||||
|
const { payload, reply } = ctx;
|
||||||
|
const session = sessions.get(payload.sessionId || payload.jobId);
|
||||||
|
if (!session) {
|
||||||
|
reply({ ok: false, error: 'session not found' });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (payload.page != null) session.egress.page = !!payload.page;
|
||||||
|
if (payload.file != null) session.egress.file = !!payload.file;
|
||||||
|
if (payload.swarm && Array.isArray(payload.swarm.connIds)) {
|
||||||
|
session.egress.swarm = { connIds: payload.swarm.connIds.slice() };
|
||||||
|
} else if (Array.isArray(payload.connIds)) {
|
||||||
|
session.egress.swarm = { connIds: payload.connIds.slice() };
|
||||||
|
}
|
||||||
|
reply({
|
||||||
|
ok: true,
|
||||||
|
sessionId: session.id,
|
||||||
|
egress: session.egress,
|
||||||
|
});
|
||||||
|
},
|
||||||
|
|
||||||
|
async encodeStop(ctx) {
|
||||||
|
const { payload, reply, emit } = ctx;
|
||||||
|
const id = payload.sessionId || payload.jobId;
|
||||||
|
const session = sessions.get(id);
|
||||||
|
if (!session) {
|
||||||
|
reply({ ok: false, error: 'session not found' });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
const stats = { ...session.stats };
|
||||||
|
const archiveRel = session.archiveRel;
|
||||||
|
destroySession(session);
|
||||||
|
emit('cap-end', {
|
||||||
|
pack: 'media',
|
||||||
|
sessionId: id,
|
||||||
|
jobId: id,
|
||||||
|
kind: 'live-end',
|
||||||
|
path: archiveRel,
|
||||||
|
...stats,
|
||||||
|
});
|
||||||
|
reply({
|
||||||
|
ok: true,
|
||||||
|
sessionId: id,
|
||||||
|
path: archiveRel,
|
||||||
|
...stats,
|
||||||
|
});
|
||||||
|
},
|
||||||
|
|
||||||
|
cancelSession(sessionId) {
|
||||||
|
const session = sessions.get(sessionId);
|
||||||
|
if (!session) return false;
|
||||||
|
session.emit('cap-error', {
|
||||||
|
pack: 'media',
|
||||||
|
sessionId,
|
||||||
|
jobId: sessionId,
|
||||||
|
message: 'cancelled',
|
||||||
|
});
|
||||||
|
destroySession(session);
|
||||||
|
return true;
|
||||||
|
},
|
||||||
|
|
||||||
|
cleanupAll() {
|
||||||
|
for (const id of [...sessions.keys()]) {
|
||||||
|
destroySession(sessions.get(id));
|
||||||
|
}
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = {
|
||||||
|
createLiveEncodeCommands,
|
||||||
|
sessions,
|
||||||
|
frameSwarmPayload,
|
||||||
|
MAGIC,
|
||||||
|
};
|
||||||
@@ -7,6 +7,7 @@ const path = require('bare-path');
|
|||||||
const fs = require('bare-fs');
|
const fs = require('bare-fs');
|
||||||
const b4a = require('b4a');
|
const b4a = require('b4a');
|
||||||
const { resolveAllowedPath, ensureJobDir, sanitizeId, getCapJobsRoot } = require('./paths.js');
|
const { resolveAllowedPath, ensureJobDir, sanitizeId, getCapJobsRoot } = require('./paths.js');
|
||||||
|
const { createLiveEncodeCommands } = require('./live-encode.js');
|
||||||
|
|
||||||
/** Soft limit so JSON + base64 stays under Chrome NMH ~1 MB */
|
/** Soft limit so JSON + base64 stays under Chrome NMH ~1 MB */
|
||||||
const MAX_CHUNK_CHARS = 700000;
|
const MAX_CHUNK_CHARS = 700000;
|
||||||
@@ -56,8 +57,9 @@ function emitBase64Chunks(emit, pack, jobId, buffer, extra) {
|
|||||||
return index;
|
return index;
|
||||||
}
|
}
|
||||||
|
|
||||||
function createMediaPack(bareMedia) {
|
function createMediaPack(bareMedia, ffmpeg) {
|
||||||
const { image, video } = bareMedia;
|
const { image, video } = bareMedia;
|
||||||
|
const live = ffmpeg ? createLiveEncodeCommands(bareMedia, ffmpeg) : null;
|
||||||
|
|
||||||
async function writeInputIfNeeded(payload, jobId) {
|
async function writeInputIfNeeded(payload, jobId) {
|
||||||
if (payload.dataBase64) {
|
if (payload.dataBase64) {
|
||||||
@@ -237,14 +239,15 @@ function createMediaPack(bareMedia) {
|
|||||||
|
|
||||||
async cancel(ctx) {
|
async cancel(ctx) {
|
||||||
const { payload, reply } = ctx;
|
const { payload, reply } = ctx;
|
||||||
const jobId = payload.jobId;
|
const jobId = payload.jobId || payload.sessionId;
|
||||||
if (!jobId) {
|
if (!jobId) {
|
||||||
reply({ ok: false, error: 'jobId required' });
|
reply({ ok: false, error: 'jobId required' });
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
const job = jobs.get(jobId);
|
const job = jobs.get(jobId);
|
||||||
if (job) job.cancelled = true;
|
if (job) job.cancelled = true;
|
||||||
reply({ ok: true, jobId, cancelled: !!job });
|
const liveCancelled = live ? live.cancelSession(jobId) : false;
|
||||||
|
reply({ ok: true, jobId, cancelled: !!(job || liveCancelled) });
|
||||||
},
|
},
|
||||||
|
|
||||||
async writeInput(ctx) {
|
async writeInput(ctx) {
|
||||||
@@ -279,12 +282,23 @@ function createMediaPack(bareMedia) {
|
|||||||
},
|
},
|
||||||
};
|
};
|
||||||
|
|
||||||
|
if (live) {
|
||||||
|
commands.encodeStart = live.encodeStart;
|
||||||
|
commands.encodePushFrame = live.encodePushFrame;
|
||||||
|
commands.encodePush = live.encodePush;
|
||||||
|
commands.encodeSubscribe = live.encodeSubscribe;
|
||||||
|
commands.encodeStop = live.encodeStop;
|
||||||
|
}
|
||||||
|
|
||||||
return {
|
return {
|
||||||
id: 'media',
|
id: 'media',
|
||||||
commands,
|
commands,
|
||||||
onLoad() {
|
onLoad() {
|
||||||
ensureJobDir('_ready');
|
ensureJobDir('_ready');
|
||||||
},
|
},
|
||||||
|
cleanup() {
|
||||||
|
if (live) live.cleanupAll();
|
||||||
|
},
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1205,6 +1205,16 @@ async function handleMessageAsync(send, msg) {
|
|||||||
reply,
|
reply,
|
||||||
emit,
|
emit,
|
||||||
storageRoot: getStorageRoot(),
|
storageRoot: getStorageRoot(),
|
||||||
|
writeToConn(connId, buf) {
|
||||||
|
const entry = connections.get(connId);
|
||||||
|
if (!entry || entry._closed) return false;
|
||||||
|
try {
|
||||||
|
entry.socket.write(buf);
|
||||||
|
return true;
|
||||||
|
} catch (_) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
},
|
||||||
});
|
});
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
reply({ ok: false, error: err.message });
|
reply({ ok: false, error: err.message });
|
||||||
@@ -1222,6 +1232,16 @@ async function handleMessageAsync(send, msg) {
|
|||||||
reply,
|
reply,
|
||||||
emit,
|
emit,
|
||||||
storageRoot: getStorageRoot(),
|
storageRoot: getStorageRoot(),
|
||||||
|
writeToConn(connId, buf) {
|
||||||
|
const entry = connections.get(connId);
|
||||||
|
if (!entry || entry._closed) return false;
|
||||||
|
try {
|
||||||
|
entry.socket.write(buf);
|
||||||
|
return true;
|
||||||
|
} catch (_) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
},
|
||||||
});
|
});
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
reply({ ok: false, error: err.message });
|
reply({ ok: false, error: err.message });
|
||||||
@@ -1244,6 +1264,10 @@ async function handleMessageAsync(send, msg) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
function cleanup() {
|
function cleanup() {
|
||||||
|
try {
|
||||||
|
const mediaPack = capabilities.getPack('media');
|
||||||
|
if (mediaPack && typeof mediaPack.cleanup === 'function') mediaPack.cleanup();
|
||||||
|
} catch (_) {}
|
||||||
for (const entry of connections.values()) {
|
for (const entry of connections.values()) {
|
||||||
try {
|
try {
|
||||||
entry._closed = true;
|
entry._closed = true;
|
||||||
|
|||||||
@@ -4,15 +4,16 @@
|
|||||||
*/
|
*/
|
||||||
|
|
||||||
import * as bareMedia from 'bare-media';
|
import * as bareMedia from 'bare-media';
|
||||||
|
import * as bareFfmpeg from 'bare-ffmpeg';
|
||||||
import registry from './capabilities/registry.js';
|
import registry from './capabilities/registry.js';
|
||||||
import mediaMod from './capabilities/media.js';
|
import mediaMod from './capabilities/media.js';
|
||||||
import { logErr } from './boot.mjs';
|
import { logErr } from './boot.mjs';
|
||||||
|
|
||||||
export function registerDefaultCapabilityPacks() {
|
export function registerDefaultCapabilityPacks() {
|
||||||
const pack = mediaMod.createMediaPack(bareMedia);
|
const pack = mediaMod.createMediaPack(bareMedia, bareFfmpeg);
|
||||||
registry.registerPack(pack);
|
registry.registerPack(pack);
|
||||||
if (typeof pack.onLoad === 'function') pack.onLoad();
|
if (typeof pack.onLoad === 'function') pack.onLoad();
|
||||||
logErr('capability pack registered: media');
|
logErr('capability pack registered: media (live encode enabled)');
|
||||||
return ['media'];
|
return ['media'];
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user