239 lines
6.6 KiB
JavaScript
239 lines
6.6 KiB
JavaScript
import {
|
|
setupBareOsMeshdropChannel,
|
|
PROTOCOL_MESHDROP_CHANNEL_NAME,
|
|
BARE_OS_MESHDROP_WIRE_SCHEMA_VERSION,
|
|
bareOsProtMuxMeshdropChannelEnabled
|
|
} from 'bare-os-protocol'
|
|
|
|
/**
|
|
* @typedef {{ chan: import('protomux').Channel, mux: import('protomux').Protomux, socket: any, id: string | null, meshdropChan?: import('protomux').Channel | null }} SwarmPeer
|
|
*/
|
|
|
|
/**
|
|
* @param {Record<string, string | undefined>} [env]
|
|
*/
|
|
export function bareOsMeshdropMuxEnabled(env = globalThis.process?.env) {
|
|
return bareOsProtMuxMeshdropChannelEnabled(env || {})
|
|
}
|
|
|
|
/**
|
|
* @param {import('./swarm-disk.js').SwarmDisk} disk
|
|
* @param {Record<string, string | undefined>} [env]
|
|
*/
|
|
export function ensureDiskBareOsMeshdropTransport(disk, env = {}) {
|
|
if (!disk || !bareOsMeshdropMuxEnabled(globalThis.process?.env)) return
|
|
const merged = {
|
|
.../** @type {Record<string, string | undefined>} */ (
|
|
globalThis.process?.env || {}
|
|
),
|
|
...env
|
|
}
|
|
disk.bareOsMeshdropService = createBareOsMeshdropService({ env: merged })
|
|
}
|
|
|
|
/**
|
|
* @param {Record<string, string | undefined>} env
|
|
*/
|
|
function meshdropHistoryMax(env) {
|
|
const raw = String(env.BARE_OS_MESHDROP_HISTORY_MAX ?? '').trim()
|
|
const n = raw ? Number.parseInt(raw, 10) : NaN
|
|
if (Number.isFinite(n) && n >= 32 && n <= 10000) return n
|
|
return 2048
|
|
}
|
|
|
|
/**
|
|
* @param {Record<string, string | undefined>} env
|
|
*/
|
|
function meshdropPayloadMaxBytes(env) {
|
|
const raw = String(env.BARE_OS_MESHDROP_MAX_PAYLOAD_BYTES ?? '').trim()
|
|
const n = raw ? Number.parseInt(raw, 10) : NaN
|
|
if (Number.isFinite(n) && n >= 512 && n <= 512 * 1024) return n
|
|
return 64 * 1024
|
|
}
|
|
|
|
export function createBareOsMeshdropService(opts = {}) {
|
|
const env = opts.env || globalThis.process?.env || {}
|
|
const historyMax = meshdropHistoryMax(
|
|
/** @type {Record<string, string | undefined>} */ (env)
|
|
)
|
|
const payloadMaxBytes = meshdropPayloadMaxBytes(
|
|
/** @type {Record<string, string | undefined>} */ (env)
|
|
)
|
|
/** @type {Set<(ev: Record<string, unknown>) => void>} */
|
|
const subscribers = new Set()
|
|
/** @type {Array<Record<string, unknown>>} */
|
|
const history = []
|
|
/** @type {Map<string, number>} */
|
|
const seenFrame = new Map()
|
|
const metrics = {
|
|
rxEnvelope: 0,
|
|
txEnvelope: 0,
|
|
droppedPayload: 0,
|
|
deduped: 0
|
|
}
|
|
|
|
function trimSeen() {
|
|
const now = Date.now()
|
|
for (const [k, exp] of seenFrame) {
|
|
if (exp < now) seenFrame.delete(k)
|
|
}
|
|
}
|
|
|
|
function pushHistory(rec) {
|
|
history.push(rec)
|
|
while (history.length > historyMax) history.shift()
|
|
}
|
|
|
|
/**
|
|
* @param {Record<string, unknown>} frame
|
|
*/
|
|
function frameId(frame) {
|
|
const p = frame && typeof frame.payload === 'object' ? frame.payload : {}
|
|
const transferId = typeof p.transferId === 'string' ? p.transferId : ''
|
|
const offerId = typeof p.offerId === 'string' ? p.offerId : ''
|
|
const idx = typeof p.index === 'number' ? p.index : -1
|
|
const kind = typeof frame.kind === 'string' ? frame.kind : 'unknown'
|
|
const frameUid = typeof p.frameId === 'string' ? p.frameId : ''
|
|
return frameUid || `${transferId}:${offerId}:${kind}:${idx}`
|
|
}
|
|
|
|
/**
|
|
* @param {import('./swarm-disk.js').SwarmDisk} disk
|
|
* @param {SwarmPeer} fromPeer
|
|
* @param {Record<string, unknown>} frame
|
|
*/
|
|
function relayEnvelope(disk, fromPeer, frame) {
|
|
for (const p of disk.peers) {
|
|
if (p === fromPeer) continue
|
|
const ch = p.meshdropChan
|
|
if (!ch || !ch.messages || !ch.messages[0]) continue
|
|
try {
|
|
ch.messages[0].send(frame)
|
|
metrics.txEnvelope++
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* @param {import('./swarm-disk.js').SwarmDisk} disk
|
|
* @param {SwarmPeer} fromPeer
|
|
* @param {Record<string, unknown>} frame
|
|
*/
|
|
function ingestEnvelope(disk, fromPeer, frame) {
|
|
if (!frame || typeof frame !== 'object') return
|
|
const encodedLen = JSON.stringify(frame).length
|
|
if (encodedLen > payloadMaxBytes) {
|
|
metrics.droppedPayload++
|
|
return
|
|
}
|
|
trimSeen()
|
|
const fid = frameId(frame)
|
|
if (fid) {
|
|
if (seenFrame.has(fid)) {
|
|
metrics.deduped++
|
|
return
|
|
}
|
|
seenFrame.set(fid, Date.now() + 120_000)
|
|
}
|
|
metrics.rxEnvelope++
|
|
const rec = {
|
|
...frame,
|
|
fromPeerKey: fromPeer.id || '',
|
|
receivedAtMs: Date.now()
|
|
}
|
|
pushHistory(rec)
|
|
for (const fn of subscribers) {
|
|
try {
|
|
fn(rec)
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
relayEnvelope(disk, fromPeer, frame)
|
|
}
|
|
|
|
return {
|
|
PROTOCOL_MESHDROP_CHANNEL_NAME,
|
|
metrics,
|
|
history() {
|
|
return [...history]
|
|
},
|
|
subscribe(fn) {
|
|
subscribers.add(fn)
|
|
return () => subscribers.delete(fn)
|
|
},
|
|
/**
|
|
* @param {import('./swarm-disk.js').SwarmDisk} disk
|
|
* @param {import('protomux').Protomux} mux
|
|
* @param {any} _socket
|
|
* @param {SwarmPeer} peer
|
|
*/
|
|
pairOnMux(disk, mux, _socket, peer) {
|
|
setupBareOsMeshdropChannel(mux, {
|
|
onEnvelope(m, _ch) {
|
|
try {
|
|
disk.protomuxMeshdropChannelRxTotal =
|
|
(disk.protomuxMeshdropChannelRxTotal || 0) + 1
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
ingestEnvelope(disk, peer, m)
|
|
},
|
|
onChannelOpened(chan) {
|
|
peer.meshdropChan = chan
|
|
mux.stream?.once?.('close', () => {
|
|
peer.meshdropChan = null
|
|
})
|
|
}
|
|
})
|
|
},
|
|
/**
|
|
* @param {import('./swarm-disk.js').SwarmDisk} disk
|
|
* @param {Record<string, unknown>} frame
|
|
* @param {{ sender?: string }} [meta]
|
|
*/
|
|
broadcastLocal(disk, frame, meta = {}) {
|
|
const safeFrame = {
|
|
schemaVersion: BARE_OS_MESHDROP_WIRE_SCHEMA_VERSION,
|
|
sender: String(meta.sender || env.USER || 'local'),
|
|
tsMs: Date.now(),
|
|
...frame
|
|
}
|
|
const encodedLen = JSON.stringify(safeFrame).length
|
|
if (encodedLen > payloadMaxBytes) {
|
|
metrics.droppedPayload++
|
|
return { ok: false, reason: 'payload_too_large' }
|
|
}
|
|
pushHistory({ ...safeFrame, local: true, receivedAtMs: Date.now() })
|
|
for (const fn of subscribers) {
|
|
try {
|
|
fn({ ...safeFrame, local: true })
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
for (const p of disk.peers) {
|
|
const ch = p.meshdropChan
|
|
if (!ch || !ch.messages || !ch.messages[0]) continue
|
|
try {
|
|
ch.messages[0].send(safeFrame)
|
|
metrics.txEnvelope++
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
return { ok: true }
|
|
},
|
|
snapshotMetrics() {
|
|
return {
|
|
...metrics,
|
|
historyMax,
|
|
payloadMaxBytes,
|
|
protocol: PROTOCOL_MESHDROP_CHANNEL_NAME
|
|
}
|
|
}
|
|
}
|
|
}
|