feat(media-streaming): deepen swarm, manifest, and protection APIs
Chunk scheduler cancel, helper swarm viewer registry, bandwidth reset, stream manifest layers, telemetry averages, tree leaves, access tokens. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
@@ -101,6 +101,23 @@ class HyperP2PBandwidthAggregator extends EventEmitter {
|
||||
return n
|
||||
}
|
||||
|
||||
maxSeq () {
|
||||
let m = -1
|
||||
for (const c of this._chunks.values()) if (c.seq > m) m = c.seq
|
||||
return m
|
||||
}
|
||||
|
||||
reset () {
|
||||
const n = this._chunks.size
|
||||
this._chunks.clear()
|
||||
this._sources.clear()
|
||||
this._nextSeq = 0
|
||||
this._stats.ingested = 0
|
||||
this._stats.duplicates = 0
|
||||
this._stats.bytes = 0
|
||||
return n
|
||||
}
|
||||
|
||||
getStats () {
|
||||
return mediaStats(this._stats, PROTOCOL, {
|
||||
streamId: this.streamId,
|
||||
|
||||
@@ -101,6 +101,17 @@ class HyperP2PChunkSchedulerMedia extends EventEmitter {
|
||||
return n
|
||||
}
|
||||
|
||||
cancelSeq (seq) {
|
||||
const idx = this._queue.findIndex((e) => e.chunk.seq === seq)
|
||||
if (idx < 0) return false
|
||||
this._queue.splice(idx, 1)
|
||||
return true
|
||||
}
|
||||
|
||||
maxPriorityEntry () {
|
||||
return this._queue.length ? { ...this._queue[0] } : null
|
||||
}
|
||||
|
||||
getStats () {
|
||||
return mediaStats(this._stats, PROTOCOL, { pending: this._queue.length })
|
||||
}
|
||||
|
||||
@@ -76,6 +76,11 @@ class HyperP2PContentProtection extends EventEmitter {
|
||||
return had
|
||||
}
|
||||
|
||||
encryptBatch (payloads, keyId = null) {
|
||||
if (!Array.isArray(payloads)) throw new Error('payloads array required')
|
||||
return payloads.map((p) => this.encryptSegment(p, keyId))
|
||||
}
|
||||
|
||||
getStats () {
|
||||
return mediaStats(this._stats, PROTOCOL, { activeKeyId: this._activeKeyId, keys: this._keys.size })
|
||||
}
|
||||
|
||||
@@ -93,6 +93,18 @@ class HyperP2PHelperSwarmCoordinator extends EventEmitter {
|
||||
return this.autoAssign(viewerId)
|
||||
}
|
||||
|
||||
listViewers () {
|
||||
return [...this._assignments.keys()]
|
||||
}
|
||||
|
||||
removeHelper (peerId) {
|
||||
return this._helpers.delete(assertPeerId(peerId))
|
||||
}
|
||||
|
||||
assignmentCount () {
|
||||
return this._assignments.size
|
||||
}
|
||||
|
||||
_gossip (payload) {
|
||||
if (this._peerMsgs) gossipSend(this, payload)
|
||||
}
|
||||
|
||||
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"name": "hyper-p2p-helper-swarm-coordinator",
|
||||
"version": "0.0.0-scaffold",
|
||||
"version": "0.3.1",
|
||||
"lockfileVersion": 3,
|
||||
"requires": true,
|
||||
"packages": {
|
||||
"": {
|
||||
"name": "hyper-p2p-helper-swarm-coordinator",
|
||||
"version": "0.0.0-scaffold",
|
||||
"version": "0.3.1",
|
||||
"license": "Apache-2.0",
|
||||
"dependencies": {
|
||||
"b4a": "^1.6.7",
|
||||
|
||||
@@ -139,6 +139,15 @@ class HyperP2PMediaTreeOrchestrator extends EventEmitter {
|
||||
return path
|
||||
}
|
||||
|
||||
leaves () {
|
||||
return [...this._nodes.values()].filter((n) => n.children.length === 0).map((n) => n.peerId)
|
||||
}
|
||||
|
||||
setRoot (peerId) {
|
||||
this._root = assertPeerId(peerId)
|
||||
return this._root
|
||||
}
|
||||
|
||||
_gossip (payload) {
|
||||
if (!this._peerMsgs) return
|
||||
gossipSend(this, payload)
|
||||
|
||||
@@ -84,6 +84,15 @@ class HyperP2PStreamAccessControl extends EventEmitter {
|
||||
return [...this._tokens.values()].filter((c) => c.streamId === sid && c.expiresAt > now)
|
||||
}
|
||||
|
||||
tokensExpiringWithin (ms = 3600000) {
|
||||
const cutoff = Date.now() + ms
|
||||
return [...this._tokens.values()].filter((c) => c.expiresAt <= cutoff)
|
||||
}
|
||||
|
||||
activeTokenCount () {
|
||||
return this.listActive().length
|
||||
}
|
||||
|
||||
getStats () {
|
||||
return mediaStats(this._stats, PROTOCOL, { active: this.listActive().length })
|
||||
}
|
||||
|
||||
@@ -85,6 +85,16 @@ class HyperP2PStreamManifest extends EventEmitter {
|
||||
return m ? m.layers.length : 0
|
||||
}
|
||||
|
||||
removeManifest (streamId) {
|
||||
return this._manifests.delete(assertStreamId(streamId))
|
||||
}
|
||||
|
||||
highestBitrateLayer (streamId) {
|
||||
const m = this.getManifest(streamId)
|
||||
if (!m || !m.layers.length) return null
|
||||
return m.layers.reduce((a, b) => ((a.bitrate || 0) >= (b.bitrate || 0) ? a : b))
|
||||
}
|
||||
|
||||
getStats () {
|
||||
return mediaStats(this._stats, PROTOCOL, { streams: this._manifests.size })
|
||||
}
|
||||
|
||||
@@ -96,6 +96,25 @@ class HyperP2PStreamTelemetry extends EventEmitter {
|
||||
return ranked.sort((a, b) => b.stalls - a.stalls).slice(0, limit)
|
||||
}
|
||||
|
||||
avgBitrate () {
|
||||
let sum = 0
|
||||
let n = 0
|
||||
for (const hist of this._sessions.values()) {
|
||||
const last = hist[hist.length - 1]
|
||||
if (last) {
|
||||
sum += last.bitrate
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n ? Math.floor(sum / n) : 0
|
||||
}
|
||||
|
||||
clearAll () {
|
||||
const n = this._sessions.size
|
||||
this._sessions.clear()
|
||||
return n
|
||||
}
|
||||
|
||||
getStats () {
|
||||
return mediaStats(this._stats, PROTOCOL, { sessions: this._sessions.size })
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user