feat: complete Quantumancy web app — full frontend + docs
Frontend (React 18 + TS + Vite): - Landing: glitching hero, live /api/stats veil ticker, mode cards, featured spirits - Séance: three.js shader ghost (hue/form per entity, mood + audio-reactive), Ouija planchette board spelling utterances, transcript with TTS replay, entity dossier, direct contact streaming, passive/active listening - Modes: Wire Ghost telemetry panel, EVP mic anomaly detection, WebUSB RTL-SDR sweep + waterfall (hardware pass pending), Ouija/Direct Contact - Codex: public registry + entity dossiers, rarity tiers, i18n EN/ES complete - State: VeilSocket (reconnect/backoff), seance reducer, auth context - 72 vitest tests green; served by FastAPI at :7777 Docs: README + as-built plans 3-7
This commit is contained in:
131
frontend/src/lib/audio.ts
Normal file
131
frontend/src/lib/audio.ts
Normal file
@@ -0,0 +1,131 @@
|
||||
// Spirit voice playback: a single FIFO queue so overlapping audio frames
|
||||
// never talk over each other. The currently-playing element is routed
|
||||
// through an AnalyserNode so the ghost scene can pulse with the voice.
|
||||
|
||||
export type PlayingInfo = {
|
||||
utteranceId: string
|
||||
url: string
|
||||
}
|
||||
|
||||
export type SpiritAudioCallbacks = {
|
||||
onStart?: (info: PlayingInfo) => void
|
||||
onEnd?: (info: PlayingInfo) => void
|
||||
}
|
||||
|
||||
export class SpiritAudioPlayer {
|
||||
private queue: PlayingInfo[] = []
|
||||
private current: PlayingInfo | null = null
|
||||
private el: HTMLAudioElement | null = null
|
||||
private ctx: AudioContext | null = null
|
||||
private analyser: AnalyserNode | null = null
|
||||
private source: MediaElementAudioSourceNode | null = null
|
||||
private cb: SpiritAudioCallbacks = {}
|
||||
private ampBuf: Uint8Array<ArrayBuffer> = new Uint8Array(0)
|
||||
|
||||
setCallbacks(cb: SpiritAudioCallbacks): void {
|
||||
this.cb = cb
|
||||
}
|
||||
|
||||
get isPlaying(): boolean {
|
||||
return this.current !== null
|
||||
}
|
||||
|
||||
get currentId(): string | null {
|
||||
return this.current?.utteranceId ?? null
|
||||
}
|
||||
|
||||
get queuedCount(): number {
|
||||
return this.queue.length
|
||||
}
|
||||
|
||||
/** Enqueue audio for an utterance id. */
|
||||
enqueue(utteranceId: string, url: string): void {
|
||||
this.queue.push({ utteranceId, url })
|
||||
void this.pump()
|
||||
}
|
||||
|
||||
/** Instantaneous playback amplitude 0..1 (0 when silent/not playing). */
|
||||
getAmplitude(): number {
|
||||
if (!this.analyser || !this.current) return 0
|
||||
if (this.ampBuf.length !== this.analyser.fftSize) {
|
||||
this.ampBuf = new Uint8Array(this.analyser.fftSize)
|
||||
}
|
||||
this.analyser.getByteTimeDomainData(this.ampBuf)
|
||||
let sum = 0
|
||||
for (let i = 0; i < this.ampBuf.length; i++) {
|
||||
const v = (this.ampBuf[i] - 128) / 128
|
||||
sum += v * v
|
||||
}
|
||||
return Math.min(1, Math.sqrt(sum / this.ampBuf.length) * 3)
|
||||
}
|
||||
|
||||
clear(): void {
|
||||
this.queue = []
|
||||
this.el?.pause()
|
||||
this.el = null
|
||||
this.current = null
|
||||
}
|
||||
|
||||
private ensureGraph(el: HTMLAudioElement): void {
|
||||
if (this.source && this.el === el) return
|
||||
if (!this.ctx) {
|
||||
this.ctx = new AudioContext()
|
||||
this.analyser = this.ctx.createAnalyser()
|
||||
this.analyser.fftSize = 512
|
||||
this.analyser.smoothingTimeConstant = 0.7
|
||||
this.analyser.connect(this.ctx.destination)
|
||||
}
|
||||
if (this.source) {
|
||||
this.source.disconnect()
|
||||
this.source = null
|
||||
}
|
||||
this.source = this.ctx.createMediaElementSource(el)
|
||||
this.source.connect(this.analyser!)
|
||||
}
|
||||
|
||||
private async pump(): Promise<void> {
|
||||
if (this.current || this.queue.length === 0) return
|
||||
const next = this.queue.shift()!
|
||||
this.current = next
|
||||
const el = new Audio(next.url)
|
||||
this.el = el
|
||||
el.crossOrigin = 'anonymous'
|
||||
|
||||
el.onended = () => {
|
||||
const done = this.current
|
||||
this.current = null
|
||||
this.el = null
|
||||
if (done) this.cb.onEnd?.(done)
|
||||
void this.pump()
|
||||
}
|
||||
el.onerror = () => {
|
||||
const done = this.current
|
||||
this.current = null
|
||||
this.el = null
|
||||
if (done) this.cb.onEnd?.(done)
|
||||
void this.pump()
|
||||
}
|
||||
|
||||
try {
|
||||
this.ensureGraph(el)
|
||||
// AudioContext may be suspended until a user gesture; resume best-effort.
|
||||
if (this.ctx && this.ctx.state === 'suspended') {
|
||||
await this.ctx.resume().catch(() => undefined)
|
||||
}
|
||||
this.cb.onStart?.(next)
|
||||
await el.play()
|
||||
} catch {
|
||||
// Autoplay blocked or unroutable: play without the analyser as fallback.
|
||||
try {
|
||||
await el.play()
|
||||
this.cb.onStart?.(next)
|
||||
} catch {
|
||||
const done = this.current
|
||||
this.current = null
|
||||
this.el = null
|
||||
if (done) this.cb.onEnd?.(done)
|
||||
void this.pump()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
158
frontend/src/lib/evp.test.ts
Normal file
158
frontend/src/lib/evp.test.ts
Normal file
@@ -0,0 +1,158 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { EvpDetectorCore } from './evp'
|
||||
import type { EvpAnomaly } from './evp'
|
||||
|
||||
const SAMPLE_RATE = 48000
|
||||
const FFT_SIZE = 2048
|
||||
const BINS = FFT_SIZE / 2
|
||||
const BIN_HZ = SAMPLE_RATE / FFT_SIZE // 23.4375 Hz per bin
|
||||
|
||||
function quietFrame(db = -100): Float64Array {
|
||||
return new Float64Array(BINS).fill(db)
|
||||
}
|
||||
|
||||
function spikedFrame(bin: number, spikeDb: number, baseDb = -100): Float64Array {
|
||||
const frame = quietFrame(baseDb)
|
||||
frame[bin] = spikeDb
|
||||
return frame
|
||||
}
|
||||
|
||||
function makeCore(): EvpDetectorCore {
|
||||
return new EvpDetectorCore({ sampleRate: SAMPLE_RATE, fftSize: FFT_SIZE })
|
||||
}
|
||||
|
||||
describe('EvpDetectorCore', () => {
|
||||
it('seeds the rolling floor from the first frame and stays silent in silence', () => {
|
||||
const core = makeCore()
|
||||
expect(core.process(quietFrame(), 0)).toBeNull() // first frame only seeds
|
||||
const floor = core.getFloor()
|
||||
expect(floor).not.toBeNull()
|
||||
expect(floor![50]).toBeCloseTo(-100, 9)
|
||||
|
||||
for (let t = 1; t <= 50; t++) {
|
||||
expect(core.process(quietFrame(), t * 100)).toBeNull()
|
||||
}
|
||||
})
|
||||
|
||||
it('fires exactly one anomaly for a brief voice-band spike, with sane values', () => {
|
||||
const core = makeCore()
|
||||
core.process(quietFrame(), 0)
|
||||
for (let t = 1; t <= 20; t++) core.process(quietFrame(), t * 100)
|
||||
|
||||
const bin = 50 // 1171.875 Hz — inside the default 300–3400 Hz voice band
|
||||
const anomaly: EvpAnomaly | null = core.process(spikedFrame(bin, -80), 5000)
|
||||
expect(anomaly).not.toBeNull()
|
||||
expect(anomaly!.frequency).toBeCloseTo(bin * BIN_HZ, 6)
|
||||
expect(anomaly!.magnitude).toBeCloseTo(20, 6)
|
||||
|
||||
// The spike is gone on the next frame — no further anomaly.
|
||||
expect(core.process(quietFrame(), 5100)).toBeNull()
|
||||
})
|
||||
|
||||
it('throttles repeats inside the throttle window and refires after it', () => {
|
||||
const core = makeCore()
|
||||
core.process(quietFrame(), 0)
|
||||
const bin = 60
|
||||
|
||||
expect(core.process(spikedFrame(bin, -80), 10_000)).not.toBeNull()
|
||||
// Same spike 1 s later: swallowed by the 2 s throttle.
|
||||
expect(core.process(spikedFrame(bin, -80), 11_000)).toBeNull()
|
||||
// Exactly at the throttle boundary it fires again.
|
||||
expect(core.process(spikedFrame(bin, -80), 12_000)).not.toBeNull()
|
||||
})
|
||||
|
||||
it('fires at most once per throttle window under sustained loudness', () => {
|
||||
const core = makeCore()
|
||||
core.process(quietFrame(), 0)
|
||||
const bin = 70
|
||||
|
||||
let fired = 0
|
||||
// Sustained spike for 1.9 s at 100 ms intervals — inside one window.
|
||||
for (let t = 10_000; t <= 11_900; t += 100) {
|
||||
if (core.process(spikedFrame(bin, -80), t)) fired++
|
||||
}
|
||||
expect(fired).toBe(1)
|
||||
})
|
||||
|
||||
it('absorbs a slowly rising ambient into the floor instead of firing', () => {
|
||||
const core = makeCore()
|
||||
let level = -100
|
||||
core.process(quietFrame(level), 0)
|
||||
|
||||
// 20 dB of sustained rise, but gradual (0.1 dB per frame): the rolling
|
||||
// floor tracks it, so the deviation never crosses the 8 dB threshold.
|
||||
for (let t = 1; t <= 200; t++) {
|
||||
level += 0.1
|
||||
expect(core.process(quietFrame(level), t * 100)).toBeNull()
|
||||
}
|
||||
const floor = core.getFloor()!
|
||||
expect(floor[50]).toBeGreaterThan(-85)
|
||||
expect(floor[50]).toBeLessThan(-79)
|
||||
})
|
||||
|
||||
it('ignores spikes outside the voice band', () => {
|
||||
const core = makeCore()
|
||||
core.process(quietFrame(), 0)
|
||||
expect(core.process(spikedFrame(5, -60), 1000)).toBeNull() // ~117 Hz, below 300
|
||||
expect(core.process(spikedFrame(500, -60), 2000)).toBeNull() // ~11.7 kHz, above 3400
|
||||
})
|
||||
|
||||
it('respects a custom band and threshold', () => {
|
||||
const core = new EvpDetectorCore({
|
||||
sampleRate: SAMPLE_RATE,
|
||||
fftSize: FFT_SIZE,
|
||||
bandLowHz: 1000,
|
||||
bandHighHz: 1200,
|
||||
thresholdDb: 15,
|
||||
})
|
||||
core.process(quietFrame(), 0)
|
||||
// bin 50 = 1171.875 Hz is inside the custom 1000–1200 Hz band…
|
||||
expect(core.process(spikedFrame(50, -80), 1000)).not.toBeNull()
|
||||
// …but a spike outside that band is ignored…
|
||||
expect(core.process(spikedFrame(80, -80), 5000)).toBeNull()
|
||||
// …and a 10 dB spike inside it is below the custom 15 dB threshold.
|
||||
expect(core.process(spikedFrame(50, -90), 9000)).toBeNull()
|
||||
})
|
||||
|
||||
it('sanitizes -Infinity bins (analyser silence) to finite values', () => {
|
||||
const core = makeCore()
|
||||
core.process(quietFrame(), 0)
|
||||
const frame = quietFrame()
|
||||
frame[40] = -Infinity
|
||||
expect(core.process(frame, 1000)).toBeNull()
|
||||
const floor = core.getFloor()!
|
||||
expect(Number.isFinite(floor[40])).toBe(true)
|
||||
})
|
||||
|
||||
it('handles empty frames gracefully', () => {
|
||||
const core = makeCore()
|
||||
expect(core.process(new Float64Array(0), 0)).toBeNull()
|
||||
expect(core.getFloor()).toBeNull()
|
||||
})
|
||||
|
||||
it('reset() clears the floor and the throttle clock', () => {
|
||||
const core = makeCore()
|
||||
core.process(quietFrame(), 0)
|
||||
const bin = 50
|
||||
expect(core.process(spikedFrame(bin, -80), 10_000)).not.toBeNull()
|
||||
// Would be throttled without a reset…
|
||||
expect(core.process(spikedFrame(bin, -80), 10_100)).toBeNull()
|
||||
|
||||
core.reset()
|
||||
expect(core.getFloor()).toBeNull()
|
||||
|
||||
// …so the next spike fires immediately after the floor re-seeds.
|
||||
expect(core.process(quietFrame(), 10_200)).toBeNull() // re-seeds floor
|
||||
const anomaly = core.process(spikedFrame(bin, -80), 10_300)
|
||||
expect(anomaly).not.toBeNull()
|
||||
expect(anomaly!.frequency).toBeCloseTo(bin * BIN_HZ, 6)
|
||||
})
|
||||
|
||||
it('returns a defensive copy from getFloor()', () => {
|
||||
const core = makeCore()
|
||||
core.process(quietFrame(), 0)
|
||||
const floor = core.getFloor()!
|
||||
floor[50] = 999
|
||||
expect(core.getFloor()![50]).toBeCloseTo(-100, 9)
|
||||
})
|
||||
})
|
||||
202
frontend/src/lib/evp.ts
Normal file
202
frontend/src/lib/evp.ts
Normal file
@@ -0,0 +1,202 @@
|
||||
// EVP (Electronic Voice Phenomena) listening: microphone anomaly detection.
|
||||
//
|
||||
// Design: an AnalyserNode feeds frequency data; we track a rolling ambient
|
||||
// noise floor over the voice band (~300–3400 Hz). When a brief deviation
|
||||
// rises >~8 dB above the floor during an otherwise-quiet stretch, we fire
|
||||
// one anomaly (peak Hz + dB above floor), throttled to at most one per 2 s.
|
||||
//
|
||||
// The pure math (floor tracking, deviation detection, throttling) lives in
|
||||
// `EvpDetectorCore` so it is unit-testable with injected sample arrays; the
|
||||
// Web Audio plumbing lives in `EvpListener` below.
|
||||
|
||||
export type EvpAnomaly = {
|
||||
/** Peak frequency in Hz. */
|
||||
frequency: number
|
||||
/** dB above the rolling noise floor. */
|
||||
magnitude: number
|
||||
}
|
||||
|
||||
export type EvpDetectorOptions = {
|
||||
sampleRate: number
|
||||
fftSize: number
|
||||
/** Low edge of the voice band, Hz. */
|
||||
bandLowHz?: number
|
||||
/** High edge of the voice band, Hz. */
|
||||
bandHighHz?: number
|
||||
/** How far above floor (dB) a deviation must reach. */
|
||||
thresholdDb?: number
|
||||
/** Minimum milliseconds between anomaly emissions. */
|
||||
throttleMs?: number
|
||||
/** EMA factor for the rolling floor (0..1, higher = faster adaptation). */
|
||||
floorAlpha?: number
|
||||
/** Floor only adapts when current level is within this many dB of the floor. */
|
||||
quietBandDb?: number
|
||||
}
|
||||
|
||||
export class EvpDetectorCore {
|
||||
private floor: Float64Array | null = null
|
||||
private lastEmit = -Infinity
|
||||
private readonly binHz: number
|
||||
private readonly loBin: number
|
||||
private readonly hiBin: number
|
||||
private readonly thresholdDb: number
|
||||
private readonly throttleMs: number
|
||||
private readonly floorAlpha: number
|
||||
private readonly quietBandDb: number
|
||||
|
||||
constructor(opts: EvpDetectorOptions) {
|
||||
this.binHz = opts.sampleRate / opts.fftSize
|
||||
const bins = opts.fftSize / 2
|
||||
const bandLow = opts.bandLowHz ?? 300
|
||||
const bandHigh = opts.bandHighHz ?? 3400
|
||||
this.loBin = Math.max(0, Math.floor(bandLow / this.binHz))
|
||||
this.hiBin = Math.min(bins - 1, Math.ceil(bandHigh / this.binHz))
|
||||
this.thresholdDb = opts.thresholdDb ?? 8
|
||||
this.throttleMs = opts.throttleMs ?? 2000
|
||||
this.floorAlpha = opts.floorAlpha ?? 0.05
|
||||
this.quietBandDb = opts.quietBandDb ?? 3
|
||||
}
|
||||
|
||||
/**
|
||||
* Feed one frame of analyser dB values (getFloatFrequencyData shape:
|
||||
* one dB value per bin, typically −160..0). `nowMs` is the frame time.
|
||||
* Returns an anomaly when one fires, else null.
|
||||
*/
|
||||
process(dbData: ArrayLike<number>, nowMs: number): EvpAnomaly | null {
|
||||
const bins = dbData.length
|
||||
if (bins === 0) return null
|
||||
|
||||
// Sanitize: AnalyserNode can emit -Infinity for silence.
|
||||
const frame = new Float64Array(bins)
|
||||
for (let i = 0; i < bins; i++) {
|
||||
const v = dbData[i]
|
||||
frame[i] = Number.isFinite(v) ? v : -160
|
||||
}
|
||||
|
||||
if (!this.floor || this.floor.length !== bins) {
|
||||
this.floor = Float64Array.from(frame)
|
||||
return null // first frame seeds the floor
|
||||
}
|
||||
|
||||
// Find the peak deviation inside the voice band.
|
||||
let peakBin = -1
|
||||
let peakDev = 0
|
||||
for (let i = this.loBin; i <= this.hiBin && i < bins; i++) {
|
||||
const dev = frame[i] - this.floor[i]
|
||||
if (dev > peakDev) {
|
||||
peakDev = dev
|
||||
peakBin = i
|
||||
}
|
||||
}
|
||||
|
||||
// Overall band level, used for the "quiet stretch" check: the floor
|
||||
// only adapts while the band sits within quietBandDb of the floor —
|
||||
// that is precisely the quiet stretch in which anomalies count.
|
||||
let bandLevel = -Infinity
|
||||
let floorLevel = -Infinity
|
||||
for (let i = this.loBin; i <= this.hiBin && i < bins; i++) {
|
||||
bandLevel = Math.max(bandLevel, frame[i])
|
||||
floorLevel = Math.max(floorLevel, this.floor[i])
|
||||
}
|
||||
const isQuiet = bandLevel - floorLevel < this.quietBandDb
|
||||
|
||||
let anomaly: EvpAnomaly | null = null
|
||||
if (
|
||||
peakBin >= 0 &&
|
||||
peakDev >= this.thresholdDb &&
|
||||
nowMs - this.lastEmit >= this.throttleMs
|
||||
) {
|
||||
anomaly = {
|
||||
frequency: peakBin * this.binHz,
|
||||
magnitude: peakDev,
|
||||
}
|
||||
this.lastEmit = nowMs
|
||||
}
|
||||
|
||||
// Adapt the floor only during quiet stretches so sustained speech/music
|
||||
// does not drag the baseline up and mask real spikes.
|
||||
if (isQuiet) {
|
||||
for (let i = 0; i < bins; i++) {
|
||||
this.floor[i] += this.floorAlpha * (frame[i] - this.floor[i])
|
||||
}
|
||||
}
|
||||
|
||||
return anomaly
|
||||
}
|
||||
|
||||
/** Current rolling floor snapshot (for visualizers). */
|
||||
getFloor(): Float64Array | null {
|
||||
return this.floor ? Float64Array.from(this.floor) : null
|
||||
}
|
||||
|
||||
reset(): void {
|
||||
this.floor = null
|
||||
this.lastEmit = -Infinity
|
||||
}
|
||||
}
|
||||
|
||||
// ---- Web Audio plumbing ----
|
||||
|
||||
export type EvpListenerCallbacks = {
|
||||
onAnomaly: (a: EvpAnomaly) => void
|
||||
/** Called each frame with (dbData, floor) for waveform/band visualization. */
|
||||
onFrame?: (dbData: Float32Array, floor: Float64Array | null) => void
|
||||
onError?: (err: Error) => void
|
||||
}
|
||||
|
||||
export class EvpListener {
|
||||
private ctx: AudioContext | null = null
|
||||
private stream: MediaStream | null = null
|
||||
private analyser: AnalyserNode | null = null
|
||||
private raf = 0
|
||||
private core: EvpDetectorCore | null = null
|
||||
private running = false
|
||||
|
||||
get isRunning(): boolean {
|
||||
return this.running
|
||||
}
|
||||
|
||||
async start(cb: EvpListenerCallbacks): Promise<void> {
|
||||
if (this.running) return
|
||||
const stream = await navigator.mediaDevices.getUserMedia({ audio: true })
|
||||
const ctx = new AudioContext()
|
||||
const src = ctx.createMediaStreamSource(stream)
|
||||
const analyser = ctx.createAnalyser()
|
||||
analyser.fftSize = 2048
|
||||
analyser.smoothingTimeConstant = 0.5
|
||||
src.connect(analyser)
|
||||
|
||||
this.ctx = ctx
|
||||
this.stream = stream
|
||||
this.analyser = analyser
|
||||
this.core = new EvpDetectorCore({
|
||||
sampleRate: ctx.sampleRate,
|
||||
fftSize: analyser.fftSize,
|
||||
})
|
||||
this.running = true
|
||||
|
||||
const buf = new Float32Array(analyser.frequencyBinCount)
|
||||
const loop = () => {
|
||||
if (!this.running || !this.analyser || !this.core) return
|
||||
this.analyser.getFloatFrequencyData(buf)
|
||||
const anomaly = this.core.process(buf, performance.now())
|
||||
if (anomaly) cb.onAnomaly(anomaly)
|
||||
cb.onFrame?.(buf, this.core.getFloor())
|
||||
this.raf = requestAnimationFrame(loop)
|
||||
}
|
||||
this.raf = requestAnimationFrame(loop)
|
||||
}
|
||||
|
||||
async stop(): Promise<void> {
|
||||
this.running = false
|
||||
cancelAnimationFrame(this.raf)
|
||||
this.stream?.getTracks().forEach((t) => t.stop())
|
||||
if (this.ctx && this.ctx.state !== 'closed') {
|
||||
await this.ctx.close().catch(() => undefined)
|
||||
}
|
||||
this.ctx = null
|
||||
this.stream = null
|
||||
this.analyser = null
|
||||
this.core = null
|
||||
}
|
||||
}
|
||||
136
frontend/src/lib/fft.test.ts
Normal file
136
frontend/src/lib/fft.test.ts
Normal file
@@ -0,0 +1,136 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { fftInPlace, magnitudeSpectrum, nextPow2, powerSpectrumDb } from './fft'
|
||||
|
||||
function argmax(values: Float64Array): number {
|
||||
let best = 0
|
||||
for (let i = 1; i < values.length; i++) {
|
||||
if (values[i] > values[best]) best = i
|
||||
}
|
||||
return best
|
||||
}
|
||||
|
||||
/** Real-valued unit-amplitude sine/cosine exactly on FFT bin `bin`. */
|
||||
function tone(n: number, bin: number, fn: 'sin' | 'cos' = 'sin'): Float64Array {
|
||||
const out = new Float64Array(n)
|
||||
for (let t = 0; t < n; t++) {
|
||||
const angle = (2 * Math.PI * bin * t) / n
|
||||
out[t] = fn === 'sin' ? Math.sin(angle) : Math.cos(angle)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
describe('nextPow2', () => {
|
||||
it('returns the smallest power of two >= n', () => {
|
||||
expect(nextPow2(1)).toBe(1)
|
||||
expect(nextPow2(2)).toBe(2)
|
||||
expect(nextPow2(3)).toBe(4)
|
||||
expect(nextPow2(1000)).toBe(1024)
|
||||
expect(nextPow2(1024)).toBe(1024)
|
||||
})
|
||||
})
|
||||
|
||||
describe('magnitudeSpectrum', () => {
|
||||
it('peaks at the bin of an injected sine', () => {
|
||||
const n = 64
|
||||
const k = 5
|
||||
const mags = magnitudeSpectrum(tone(n, k), new Float64Array(n))
|
||||
expect(mags).toHaveLength(n / 2)
|
||||
expect(argmax(mags)).toBe(k)
|
||||
// A unit-amplitude sine concentrates N/2 of magnitude in its bin.
|
||||
expect(mags[k]).toBeCloseTo(n / 2, 6)
|
||||
expect(mags[k - 1]).toBeLessThan(1e-6)
|
||||
expect(mags[k + 1]).toBeLessThan(1e-6)
|
||||
})
|
||||
|
||||
it('peaks at the bin of a second, different frequency (cosine)', () => {
|
||||
const n = 128
|
||||
const k = 11
|
||||
const mags = magnitudeSpectrum(tone(n, k, 'cos'), new Float64Array(n))
|
||||
expect(argmax(mags)).toBe(k)
|
||||
expect(mags[k]).toBeCloseTo(n / 2, 6)
|
||||
})
|
||||
|
||||
it('places a DC signal entirely in bin 0', () => {
|
||||
const n = 64
|
||||
const mags = magnitudeSpectrum(new Float64Array(n).fill(1), new Float64Array(n))
|
||||
expect(mags[0]).toBeCloseTo(n, 6)
|
||||
for (let k = 1; k < n / 2; k++) {
|
||||
expect(mags[k]).toBeLessThan(1e-6)
|
||||
}
|
||||
})
|
||||
|
||||
it('does not mutate its inputs', () => {
|
||||
const re = tone(64, 3)
|
||||
const im = new Float64Array(64)
|
||||
const reBefore = Float64Array.from(re)
|
||||
const imBefore = Float64Array.from(im)
|
||||
magnitudeSpectrum(re, im)
|
||||
expect(re).toEqual(reBefore)
|
||||
expect(im).toEqual(imBefore)
|
||||
})
|
||||
})
|
||||
|
||||
describe('fftInPlace', () => {
|
||||
it('puts a Nyquist-rate alternating signal entirely in bin n/2', () => {
|
||||
const n = 64
|
||||
const re = new Float64Array(n)
|
||||
for (let t = 0; t < n; t++) re[t] = t % 2 === 0 ? 1 : -1
|
||||
const im = new Float64Array(n)
|
||||
fftInPlace(re, im) // in-place: the transform lands back in re/im
|
||||
// magnitudeSpectrum only returns the first n/2 bins, so the Nyquist bin
|
||||
// is verified here on the raw output instead.
|
||||
expect(Math.hypot(re[n / 2], im[n / 2])).toBeCloseTo(n, 6)
|
||||
expect(Math.hypot(re[0], im[0])).toBeLessThan(1e-9)
|
||||
})
|
||||
|
||||
it('transforms a unit impulse into a flat spectrum', () => {
|
||||
const n = 32
|
||||
const re = new Float64Array(n)
|
||||
re[0] = 1
|
||||
const im = new Float64Array(n)
|
||||
fftInPlace(re, im)
|
||||
for (let k = 0; k < n; k++) {
|
||||
expect(Math.hypot(re[k], im[k])).toBeCloseTo(1, 9)
|
||||
}
|
||||
})
|
||||
|
||||
it('rejects mismatched re/im lengths', () => {
|
||||
expect(() => fftInPlace(new Float64Array(4), new Float64Array(8))).toThrow(/mismatch/)
|
||||
})
|
||||
|
||||
it('rejects non-power-of-two lengths', () => {
|
||||
expect(() => fftInPlace(new Float64Array(3), new Float64Array(3))).toThrow(/power of two/)
|
||||
})
|
||||
})
|
||||
|
||||
describe('powerSpectrumDb', () => {
|
||||
it('peaks near 0 dB at the bin of a unit complex exponential', () => {
|
||||
const n = 128
|
||||
const k = 7
|
||||
const iq = new Float64Array(2 * n)
|
||||
for (let t = 0; t < n; t++) {
|
||||
iq[2 * t] = Math.cos((2 * Math.PI * k * t) / n)
|
||||
iq[2 * t + 1] = Math.sin((2 * Math.PI * k * t) / n)
|
||||
}
|
||||
const db = powerSpectrumDb(iq)
|
||||
expect(db).toHaveLength(n / 2)
|
||||
expect(argmax(db)).toBe(k)
|
||||
// Unit amplitude => power (n^2)/(n^2) = 1 => ~0 dB at the peak.
|
||||
expect(db[k]).toBeCloseTo(0, 6)
|
||||
expect(db[k + 1]).toBeLessThan(-60)
|
||||
})
|
||||
|
||||
it('zero-pads non-power-of-two sample counts and still finds the peak', () => {
|
||||
const n = 128 // padded size chosen by nextPow2
|
||||
const samples = 100
|
||||
const k = 7
|
||||
const iq = new Float64Array(2 * samples)
|
||||
for (let t = 0; t < samples; t++) {
|
||||
iq[2 * t] = Math.cos((2 * Math.PI * k * t) / n)
|
||||
iq[2 * t + 1] = Math.sin((2 * Math.PI * k * t) / n)
|
||||
}
|
||||
const db = powerSpectrumDb(iq)
|
||||
expect(db).toHaveLength(n / 2)
|
||||
expect(argmax(db)).toBe(k)
|
||||
})
|
||||
})
|
||||
95
frontend/src/lib/fft.ts
Normal file
95
frontend/src/lib/fft.ts
Normal file
@@ -0,0 +1,95 @@
|
||||
// Radix-2 Cooley–Tukey FFT, in-place, iterative, bit-reversal ordered.
|
||||
// Used by the RTL-SDR spectrum sweep. Length must be a power of two.
|
||||
|
||||
export type ComplexArray = { re: Float64Array; im: Float64Array }
|
||||
|
||||
export function nextPow2(n: number): number {
|
||||
let p = 1
|
||||
while (p < n) p <<= 1
|
||||
return p
|
||||
}
|
||||
|
||||
/**
|
||||
* In-place iterative radix-2 FFT. `re`/`im` must have identical length
|
||||
* and that length must be a power of two.
|
||||
*/
|
||||
export function fftInPlace(re: Float64Array, im: Float64Array): void {
|
||||
const n = re.length
|
||||
if (n !== im.length) throw new Error('re/im length mismatch')
|
||||
if ((n & (n - 1)) !== 0) throw new Error('FFT length must be a power of two')
|
||||
if (n <= 1) return
|
||||
|
||||
// Bit-reversal permutation.
|
||||
for (let i = 1, j = 0; i < n; i++) {
|
||||
let bit = n >> 1
|
||||
for (; j & bit; bit >>= 1) j ^= bit
|
||||
j ^= bit
|
||||
if (i < j) {
|
||||
const tr = re[i]
|
||||
re[i] = re[j]
|
||||
re[j] = tr
|
||||
const ti = im[i]
|
||||
im[i] = im[j]
|
||||
im[j] = ti
|
||||
}
|
||||
}
|
||||
|
||||
for (let len = 2; len <= n; len <<= 1) {
|
||||
const half = len >> 1
|
||||
const ang = (-2 * Math.PI) / len
|
||||
const wr = Math.cos(ang)
|
||||
const wi = Math.sin(ang)
|
||||
for (let i = 0; i < n; i += len) {
|
||||
let cwr = 1
|
||||
let cwi = 0
|
||||
for (let j = 0; j < half; j++) {
|
||||
const uR = re[i + j]
|
||||
const uI = im[i + j]
|
||||
const vR = re[i + j + half] * cwr - im[i + j + half] * cwi
|
||||
const vI = re[i + j + half] * cwi + im[i + j + half] * cwr
|
||||
re[i + j] = uR + vR
|
||||
im[i + j] = uI + vI
|
||||
re[i + j + half] = uR - vR
|
||||
im[i + j + half] = uI - vI
|
||||
const nwr = cwr * wr - cwi * wi
|
||||
cwi = cwr * wi + cwi * wr
|
||||
cwr = nwr
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Convenience wrapper returning magnitude spectrum (first n/2 bins). */
|
||||
export function magnitudeSpectrum(re: Float64Array, im: Float64Array): Float64Array {
|
||||
const r = Float64Array.from(re)
|
||||
const i = Float64Array.from(im)
|
||||
fftInPlace(r, i)
|
||||
const half = r.length >> 1
|
||||
const mags = new Float64Array(half)
|
||||
for (let k = 0; k < half; k++) {
|
||||
mags[k] = Math.hypot(r[k], i[k])
|
||||
}
|
||||
return mags
|
||||
}
|
||||
|
||||
/**
|
||||
* Power spectrum in dBFS-ish units from interleaved I/Q samples.
|
||||
* `iq` layout: [i0, q0, i1, q1, ...]. Returns n/2 dB values.
|
||||
*/
|
||||
export function powerSpectrumDb(iq: Float64Array): Float64Array {
|
||||
const n = nextPow2(iq.length >> 1)
|
||||
const re = new Float64Array(n)
|
||||
const im = new Float64Array(n)
|
||||
for (let k = 0; k < n && 2 * k + 1 < iq.length; k++) {
|
||||
re[k] = iq[2 * k]
|
||||
im[k] = iq[2 * k + 1]
|
||||
}
|
||||
fftInPlace(re, im)
|
||||
const half = n >> 1
|
||||
const out = new Float64Array(half)
|
||||
for (let k = 0; k < half; k++) {
|
||||
const power = (re[k] * re[k] + im[k] * im[k]) / (n * n)
|
||||
out[k] = 10 * Math.log10(power + 1e-12)
|
||||
}
|
||||
return out
|
||||
}
|
||||
179
frontend/src/lib/planchette.test.ts
Normal file
179
frontend/src/lib/planchette.test.ts
Normal file
@@ -0,0 +1,179 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { PlanchetteMachine, normalizeWord, tokenize } from './planchette'
|
||||
|
||||
// Fast cadence so tests can drive the machine with small tick counts.
|
||||
const FAST = { msPerLetter: 100, msBetweenWords: 200 }
|
||||
|
||||
describe('normalizeWord', () => {
|
||||
it('uppercases and strips punctuation', () => {
|
||||
expect(normalizeWord('beware, the veil!')).toBe('BEWARE THE VEIL')
|
||||
})
|
||||
|
||||
it('keeps digits', () => {
|
||||
expect(normalizeWord('room 13, floor 4')).toBe('ROOM 13 FLOOR 4')
|
||||
})
|
||||
|
||||
it('drops accented letters outright (the board only knows A-Z)', () => {
|
||||
expect(normalizeWord('séance niño')).toBe('SANCE NIO')
|
||||
})
|
||||
|
||||
it('collapses whitespace and trims', () => {
|
||||
expect(normalizeWord(' hello \t spirit ')).toBe('HELLO SPIRIT')
|
||||
})
|
||||
|
||||
it('returns an empty string when nothing spellable remains', () => {
|
||||
expect(normalizeWord('!!! …')).toBe('')
|
||||
})
|
||||
})
|
||||
|
||||
describe('tokenize', () => {
|
||||
it('splits an utterance into normalized words', () => {
|
||||
expect(tokenize('beware, the veil!')).toEqual(['BEWARE', 'THE', 'VEIL'])
|
||||
})
|
||||
|
||||
it('returns a single word as-is (normalized)', () => {
|
||||
expect(tokenize('hello')).toEqual(['HELLO'])
|
||||
})
|
||||
|
||||
it('returns an empty list when there is nothing to spell', () => {
|
||||
expect(tokenize('')).toEqual([])
|
||||
expect(tokenize('!!!')).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
describe('PlanchetteMachine', () => {
|
||||
it('starts idle with an empty queue and history', () => {
|
||||
const m = new PlanchetteMachine(FAST)
|
||||
const s = m.snapshot()
|
||||
expect(s.phase).toBe('idle')
|
||||
expect(s.current).toBeNull()
|
||||
expect(s.word).toBeNull()
|
||||
expect(s.index).toBe(0)
|
||||
expect(s.queued).toEqual([])
|
||||
expect(s.spelled).toEqual([])
|
||||
})
|
||||
|
||||
it('spells a word letter by letter on the per-letter cadence', () => {
|
||||
const m = new PlanchetteMachine(FAST)
|
||||
m.enqueue('HI')
|
||||
|
||||
// Enqueueing alone does not start spelling.
|
||||
expect(m.snapshot().phase).toBe('idle')
|
||||
expect(m.snapshot().queued).toEqual(['HI'])
|
||||
|
||||
// The first tick pulls the word off the queue and hovers the first letter.
|
||||
let s = m.tick(0)
|
||||
expect(s.phase).toBe('moving')
|
||||
expect(s.word).toBe('HI')
|
||||
expect(s.current).toBe('H')
|
||||
expect(s.index).toBe(0)
|
||||
expect(s.queued).toEqual([])
|
||||
|
||||
// A partial interval keeps it dwelling on the same letter.
|
||||
s = m.tick(50)
|
||||
expect(s.phase).toBe('dwelling')
|
||||
expect(s.current).toBe('H')
|
||||
expect(s.index).toBe(0)
|
||||
|
||||
// Completing the interval advances to the next letter.
|
||||
s = m.tick(50)
|
||||
expect(s.current).toBe('I')
|
||||
expect(s.index).toBe(1)
|
||||
|
||||
// One more interval finishes the word: it lands in `spelled`.
|
||||
s = m.tick(100)
|
||||
expect(s.phase).toBe('returning')
|
||||
expect(s.current).toBeNull()
|
||||
expect(s.word).toBeNull()
|
||||
expect(s.spelled).toEqual(['HI'])
|
||||
|
||||
// Between-words pause, then back to idle.
|
||||
s = m.tick(100)
|
||||
expect(s.phase).toBe('returning')
|
||||
s = m.tick(100)
|
||||
expect(s.phase).toBe('idle')
|
||||
expect(s.current).toBeNull()
|
||||
})
|
||||
|
||||
it('spells queued words in FIFO order and idles after draining', () => {
|
||||
const m = new PlanchetteMachine(FAST)
|
||||
m.enqueue('BEWARE THE VEIL')
|
||||
expect(m.snapshot().queued).toEqual(['BEWARE', 'THE', 'VEIL'])
|
||||
|
||||
for (let i = 0; i < 500; i++) m.tick(50)
|
||||
|
||||
const s = m.snapshot()
|
||||
expect(s.phase).toBe('idle')
|
||||
expect(s.current).toBeNull()
|
||||
expect(s.queued).toEqual([])
|
||||
expect(s.spelled).toEqual(['BEWARE', 'THE', 'VEIL'])
|
||||
})
|
||||
|
||||
it('reports queue -> spelled progress while draining', () => {
|
||||
const m = new PlanchetteMachine(FAST)
|
||||
m.enqueue('AB CD')
|
||||
m.tick(0) // start AB
|
||||
m.tick(200) // finish AB -> returning
|
||||
|
||||
let s = m.snapshot()
|
||||
expect(s.spelled).toEqual(['AB'])
|
||||
expect(s.queued).toEqual(['CD'])
|
||||
|
||||
m.tick(200) // returning pause elapses -> idle
|
||||
m.tick(0) // start CD
|
||||
s = m.snapshot()
|
||||
expect(s.phase).toBe('moving')
|
||||
expect(s.word).toBe('CD')
|
||||
expect(s.current).toBe('C')
|
||||
expect(s.index).toBe(0)
|
||||
})
|
||||
|
||||
it('caps the pending queue at maxQueue', () => {
|
||||
const m = new PlanchetteMachine({ ...FAST, maxQueue: 2 })
|
||||
m.enqueue('ONE TWO THREE FOUR')
|
||||
expect(m.snapshot().queued).toEqual(['ONE', 'TWO'])
|
||||
})
|
||||
|
||||
it('clear() cancels pending and in-progress spelling but keeps history', () => {
|
||||
const m = new PlanchetteMachine(FAST)
|
||||
m.enqueue('HI THERE')
|
||||
m.tick(0) // start HI
|
||||
m.tick(200) // finish HI -> returning, spelled: ['HI']
|
||||
|
||||
m.clear()
|
||||
const s = m.snapshot()
|
||||
expect(s.phase).toBe('idle')
|
||||
expect(s.current).toBeNull()
|
||||
expect(s.queued).toEqual([])
|
||||
// clear() empties pending work; the spelled history survives the session.
|
||||
expect(s.spelled).toEqual(['HI'])
|
||||
|
||||
// ...and the machine can be reused afterwards.
|
||||
m.enqueue('OK')
|
||||
m.tick(0)
|
||||
expect(m.snapshot().word).toBe('OK')
|
||||
})
|
||||
|
||||
it('ignores utterances with no spellable words', () => {
|
||||
const m = new PlanchetteMachine(FAST)
|
||||
const version = m.getVersion()
|
||||
m.enqueue('!!! …')
|
||||
expect(m.snapshot().queued).toEqual([])
|
||||
expect(m.getVersion()).toBe(version) // no mutation, no version bump
|
||||
m.tick(0)
|
||||
expect(m.snapshot().phase).toBe('idle')
|
||||
})
|
||||
|
||||
it('bumps the version counter on every mutation', () => {
|
||||
const m = new PlanchetteMachine(FAST)
|
||||
let v = m.getVersion()
|
||||
m.enqueue('HI')
|
||||
expect(m.getVersion()).toBeGreaterThan(v)
|
||||
v = m.getVersion()
|
||||
m.tick(0)
|
||||
expect(m.getVersion()).toBeGreaterThan(v)
|
||||
v = m.getVersion()
|
||||
m.clear()
|
||||
expect(m.getVersion()).toBeGreaterThan(v)
|
||||
})
|
||||
})
|
||||
158
frontend/src/lib/planchette.ts
Normal file
158
frontend/src/lib/planchette.ts
Normal file
@@ -0,0 +1,158 @@
|
||||
// Planchette word-queue / spelling state machine.
|
||||
//
|
||||
// The Ouija board spells out spirit utterances letter by letter. This module
|
||||
// is the pure logic: queue words, step through letters at a fixed cadence,
|
||||
// and report which character the planchette is currently hovering over.
|
||||
// Rendering (canvas, DOM) lives elsewhere; this is unit-tested in isolation.
|
||||
|
||||
export type PlanchettePhase = 'idle' | 'moving' | 'dwelling' | 'returning'
|
||||
|
||||
export type PlanchetteSnapshot = {
|
||||
phase: PlanchettePhase
|
||||
/** Character currently hovered, or null when drifting idle. */
|
||||
current: string | null
|
||||
/** The full normalized word being spelled, or null. */
|
||||
word: string | null
|
||||
/** Index of `current` within `word`. */
|
||||
index: number
|
||||
/** Words still waiting to be spelled. */
|
||||
queued: readonly string[]
|
||||
/** Words already spelled this session (most recent last). */
|
||||
spelled: readonly string[]
|
||||
}
|
||||
|
||||
export type PlanchetteOptions = {
|
||||
/** Milliseconds spent on each letter. Default 300. */
|
||||
msPerLetter?: number
|
||||
/** Milliseconds paused between words. Default 700. */
|
||||
msBetweenWords?: number
|
||||
/** Cap on the internal queue so a flood of utterances cannot wedge it. */
|
||||
maxQueue?: number
|
||||
}
|
||||
|
||||
/** Keep A–Z, 0–9 and spaces; uppercase; drop everything else. */
|
||||
export function normalizeWord(text: string): string {
|
||||
return text
|
||||
.toUpperCase()
|
||||
.replace(/[^A-Z0-9 ]/g, '')
|
||||
.replace(/\s+/g, ' ')
|
||||
.trim()
|
||||
}
|
||||
|
||||
/**
|
||||
* Split an utterance into queueable word tokens: punctuation stripped,
|
||||
* whitespace-split, empties removed.
|
||||
*/
|
||||
export function tokenize(text: string): string[] {
|
||||
const norm = normalizeWord(text)
|
||||
if (!norm) return []
|
||||
return norm.split(' ').filter((w) => w.length > 0)
|
||||
}
|
||||
|
||||
export class PlanchetteMachine {
|
||||
private queue: string[] = []
|
||||
private done: string[] = []
|
||||
private word: string | null = null
|
||||
private index = 0
|
||||
private phase: PlanchettePhase = 'idle'
|
||||
private clock = 0
|
||||
private readonly msPerLetter: number
|
||||
private readonly msBetweenWords: number
|
||||
private readonly maxQueue: number
|
||||
private version = 0
|
||||
|
||||
constructor(opts: PlanchetteOptions = {}) {
|
||||
this.msPerLetter = opts.msPerLetter ?? 300
|
||||
this.msBetweenWords = opts.msBetweenWords ?? 700
|
||||
this.maxQueue = opts.maxQueue ?? 64
|
||||
}
|
||||
|
||||
/** Monotonic counter bumped on every mutation — handy for render loops. */
|
||||
getVersion(): number {
|
||||
return this.version
|
||||
}
|
||||
|
||||
/** Enqueue an utterance (may contain several words + punctuation). */
|
||||
enqueue(text: string): void {
|
||||
const words = tokenize(text)
|
||||
for (const w of words) {
|
||||
if (this.queue.length < this.maxQueue) this.queue.push(w)
|
||||
}
|
||||
if (words.length > 0) this.version++
|
||||
}
|
||||
|
||||
clear(): void {
|
||||
this.queue = []
|
||||
this.word = null
|
||||
this.index = 0
|
||||
this.phase = 'idle'
|
||||
this.clock = 0
|
||||
this.version++
|
||||
}
|
||||
|
||||
/**
|
||||
* Advance the machine by `dtMs`. Call from a rAF loop or a timer.
|
||||
* Returns a snapshot after advancing.
|
||||
*/
|
||||
tick(dtMs: number): PlanchetteSnapshot {
|
||||
this.clock += dtMs
|
||||
|
||||
if (this.phase === 'idle') {
|
||||
if (this.queue.length > 0) {
|
||||
this.word = this.queue.shift() ?? null
|
||||
this.index = 0
|
||||
this.phase = this.word && this.word.length > 0 ? 'moving' : 'idle'
|
||||
this.clock = 0
|
||||
this.version++
|
||||
}
|
||||
return this.snapshot()
|
||||
}
|
||||
|
||||
if (this.phase === 'returning') {
|
||||
if (this.clock >= this.msBetweenWords) {
|
||||
this.clock = 0
|
||||
this.phase = 'idle'
|
||||
this.word = null
|
||||
this.index = 0
|
||||
this.version++
|
||||
}
|
||||
return this.snapshot()
|
||||
}
|
||||
|
||||
// moving / dwelling: advance letters on the per-letter cadence.
|
||||
if (this.word) {
|
||||
while (this.clock >= this.msPerLetter && this.phase !== 'returning') {
|
||||
this.clock -= this.msPerLetter
|
||||
this.index++
|
||||
this.version++
|
||||
if (this.index >= this.word.length) {
|
||||
this.done.push(this.word)
|
||||
if (this.done.length > 128) this.done.shift()
|
||||
this.phase = 'returning'
|
||||
this.clock = 0
|
||||
} else {
|
||||
this.phase = this.clock >= this.msPerLetter ? 'moving' : 'dwelling'
|
||||
}
|
||||
}
|
||||
if (this.phase === 'dwelling' || this.phase === 'moving') {
|
||||
this.phase = 'dwelling'
|
||||
}
|
||||
}
|
||||
return this.snapshot()
|
||||
}
|
||||
|
||||
snapshot(): PlanchetteSnapshot {
|
||||
const active =
|
||||
(this.phase === 'moving' || this.phase === 'dwelling') && this.word
|
||||
? this.word
|
||||
: null
|
||||
return {
|
||||
phase: this.phase,
|
||||
current: active ? active[this.index] ?? null : null,
|
||||
word: active,
|
||||
index: active ? this.index : 0,
|
||||
queued: [...this.queue],
|
||||
spelled: [...this.done],
|
||||
}
|
||||
}
|
||||
}
|
||||
339
frontend/src/lib/sdr.ts
Normal file
339
frontend/src/lib/sdr.ts
Normal file
@@ -0,0 +1,339 @@
|
||||
// Best-effort WebUSB driver for RTL2832U + R820T("T") based SDR dongles.
|
||||
//
|
||||
// ⚠️ HARDWARE PASS REQUIRED: this driver is structured from the public
|
||||
// librtlsdr register documentation and has NOT been validated against a real
|
||||
// device in this environment (no RTL-SDR attached). The control-transfer
|
||||
// sequences below follow the known-good init order (demod power-up, R820T
|
||||
// tuner init via I2C repeater, sample-rate set, FIR, bulk streaming) but
|
||||
// expect to debug register pokes with a logic analyzer / librtlsdr -T.
|
||||
// Every entry point fails soft: callers must treat any thrown error as
|
||||
// "this vessel cannot hear the radio dead" and degrade gracefully.
|
||||
|
||||
import { powerSpectrumDb } from './fft'
|
||||
|
||||
export type RtlSampleBlock = {
|
||||
/** Center frequency this block was captured at, Hz. */
|
||||
centerHz: number
|
||||
/** Sample rate used, Hz. */
|
||||
sampleRateHz: number
|
||||
/** Interleaved unsigned I/Q bytes, zero-centered: [i,q,i,q…] */
|
||||
iq: Uint8Array
|
||||
}
|
||||
|
||||
export type SweepCallbacks = {
|
||||
/** Per-tune power spectrum (dB values, fftSize/2 entries). */
|
||||
onSpectrum?: (centerHz: number, db: Float64Array) => void
|
||||
onError?: (err: Error) => void
|
||||
}
|
||||
|
||||
// USB identification for RTL2832U dongles.
|
||||
export const RTL2832U_VENDOR = 0x0bda
|
||||
export const RTL2832U_PRODUCTS = [0x2832, 0x2834, 0x2838, 0x2837]
|
||||
|
||||
// Request types used by librtlsdr.
|
||||
const CTRL_IN = 0xc0
|
||||
const CTRL_OUT = 0x40
|
||||
const DEMOD = 0x03
|
||||
const USB_EPA = 0x02
|
||||
const SYS = 0x09
|
||||
const PAGE_USB = 0x01
|
||||
|
||||
export function isSupported(): boolean {
|
||||
return (
|
||||
typeof navigator !== 'undefined' &&
|
||||
'usb' in navigator &&
|
||||
typeof navigator.usb?.requestDevice === 'function'
|
||||
)
|
||||
}
|
||||
|
||||
export class RtlSdr {
|
||||
private device: USBDevice | null = null
|
||||
private interfaceNumber = 0
|
||||
private endpointIn = 0x81
|
||||
private running = false
|
||||
private sampleRate = 2_048_000
|
||||
|
||||
get isOpen(): boolean {
|
||||
return this.device?.opened ?? false
|
||||
}
|
||||
|
||||
/** Ask the browser for an RTL2832U device. Throws if none chosen/found. */
|
||||
async requestDevice(): Promise<void> {
|
||||
if (!isSupported()) {
|
||||
throw new Error('WebUSB is not available in this vessel (Chromium required)')
|
||||
}
|
||||
const filters = [
|
||||
{ vendorId: RTL2832U_VENDOR, productId: 0x2832 },
|
||||
{ vendorId: RTL2832U_VENDOR, productId: 0x2834 },
|
||||
{ vendorId: RTL2832U_VENDOR, productId: 0x2838 },
|
||||
{ vendorId: RTL2832U_VENDOR, productId: 0x2837 },
|
||||
]
|
||||
this.device = await navigator.usb.requestDevice({ filters })
|
||||
}
|
||||
|
||||
/** Open, claim, and run the RTL2832U + R820T init sequence. */
|
||||
async open(sampleRateHz = 2_048_000): Promise<void> {
|
||||
const dev = this.device
|
||||
if (!dev) throw new Error('no device selected')
|
||||
this.sampleRate = sampleRateHz
|
||||
await dev.open()
|
||||
|
||||
// Find the first bulk-IN endpoint.
|
||||
const iface = dev.configuration?.interfaces[0]
|
||||
const alt = iface?.alternates[0]
|
||||
const ep = alt?.endpoints.find((e) => e.direction === 'in' && e.type === 'bulk')
|
||||
if (iface && alt && ep) {
|
||||
this.interfaceNumber = iface.interfaceNumber
|
||||
this.endpointIn = ep.endpointNumber
|
||||
}
|
||||
|
||||
// Detach kernel driver (Linux) — ignore failure: may not be supported.
|
||||
await this.claim()
|
||||
|
||||
// --- Init sequence (HARDWARE PASS REQUIRED) ---
|
||||
// Order follows librtlsdr: USB reset → demod init → tuner I2C init.
|
||||
await this.demodWrite(1, 0x01, 0x14, 1) // soft reset
|
||||
await this.demodWrite(1, 0x01, 0x10, 1)
|
||||
await this.demodWrite(0, 0x01, 0x08, 2) // demod_ctl
|
||||
await this.demodWrite(0, 0x06, 0x80, 1)
|
||||
await this.demodWrite(1, 0x15, 0x00, 1) // suspend off
|
||||
await this.demodWrite(1, 0x16, 0x00, 1)
|
||||
await this.demodWrite(1, 0x17, 0x00, 1)
|
||||
await this.demodWrite(1, 0x18, 0x00, 1)
|
||||
await this.demodWrite(1, 0x19, 0x00, 1)
|
||||
await this.demodWrite(1, 0x1a, 0x00, 1)
|
||||
await this.demodWrite(1, 0x1b, 0x00, 1)
|
||||
await this.demodWrite(1, 0x1c, 0x00, 1)
|
||||
await this.demodWrite(1, 0x0d, 0x83, 1) // standby off
|
||||
await this.demodWrite(1, 0x0b, 0x1b, 1) // AGC mode
|
||||
|
||||
// Power on tuner through I2C repeater (R820T at 0x1a).
|
||||
await this.i2cWrite(0x1a, 0x05, 0x8f) // LNA power on
|
||||
await this.i2cWrite(0x1a, 0x08, 0x80) // mixer
|
||||
await this.i2cWrite(0x1a, 0x0a, 0x10) // IF filter
|
||||
|
||||
await this.setSampleRate(this.sampleRate)
|
||||
await this.setFrequency(98_000_000)
|
||||
|
||||
// Reset endpoint before streaming.
|
||||
await this.writeReg(USB_EPA, 0x0001, 0xffff, 2)
|
||||
await this.demodWrite(1, 0x02, 0x00, 1)
|
||||
await this.demodWrite(0, 0x02, 0x40, 2)
|
||||
await this.demodWrite(1, 0x02, 0x00, 1) // enable test mode off
|
||||
}
|
||||
|
||||
async close(): Promise<void> {
|
||||
this.running = false
|
||||
const dev = this.device
|
||||
if (dev?.opened) {
|
||||
try {
|
||||
await this.demodWrite(1, 0x01, 0x10, 1) // suspend
|
||||
} catch {
|
||||
/* device may already be gone */
|
||||
}
|
||||
await dev.releaseInterface(this.interfaceNumber).catch(() => undefined)
|
||||
await dev.close().catch(() => undefined)
|
||||
}
|
||||
}
|
||||
|
||||
/** Tune the R820T mixer PLL. Hz. (HARDWARE PASS REQUIRED) */
|
||||
async setFrequency(hz: number): Promise<void> {
|
||||
// R820T fractional-N PLL programming via I2C. Simplified to the
|
||||
// integer part + common divider ratio; real librtlsdr computes the
|
||||
// exact sdm/vco from a 28.8MHz crystal reference. Marked for hardware.
|
||||
const loHz = hz + 3_570_000 // R820T IF offset
|
||||
const ref = 28_800_000
|
||||
const mixDiv = 2
|
||||
const nint = Math.floor(loHz / (ref * mixDiv))
|
||||
const vco = loHz % (ref * mixDiv)
|
||||
const sdm = Math.min(0xffff, Math.floor((vco * 65536) / (ref * mixDiv)))
|
||||
const reg = nint & 0x3f
|
||||
await this.i2cWrite(0x1a, 0x10, reg)
|
||||
await this.i2cWrite(0x1a, 0x11, (sdm >> 8) & 0xff)
|
||||
await this.i2cWrite(0x1a, 0x12, sdm & 0xff)
|
||||
}
|
||||
|
||||
/** Program demod resampling rate for the requested sample rate. */
|
||||
async setSampleRate(hz: number): Promise<void> {
|
||||
const crystal = 28_800_000
|
||||
const rsampRatio = Math.floor(((crystal * 2 ** 22) / hz) & 0x0ffffffc)
|
||||
await this.demodWrite(1, 0x9f, (rsampRatio >> 16) & 0xffff, 2)
|
||||
await this.demodWrite(1, 0xa1, rsampRatio & 0xffff, 2)
|
||||
this.sampleRate = hz
|
||||
}
|
||||
|
||||
/** Read one block of I/Q samples. Length must be multiple of 512. */
|
||||
async readSamples(bytes: number): Promise<Uint8Array> {
|
||||
const dev = this.device
|
||||
if (!dev?.opened) throw new Error('device not open')
|
||||
const res = await dev.transferIn(this.endpointIn, bytes)
|
||||
if (!res.data) throw new Error('bulk read failed')
|
||||
return new Uint8Array(res.data.buffer, res.data.byteOffset, res.data.byteLength)
|
||||
}
|
||||
|
||||
/**
|
||||
* Continuous sweep across [startHz, endHz] in `stepHz` steps, emitting a
|
||||
* power spectrum per tuning step. Runs until `stopSweep()`.
|
||||
*/
|
||||
async sweep(
|
||||
startHz: number,
|
||||
endHz: number,
|
||||
stepHz: number,
|
||||
cb: SweepCallbacks,
|
||||
fftSize = 512,
|
||||
settleMs = 25,
|
||||
): Promise<void> {
|
||||
if (this.running) return
|
||||
this.running = true
|
||||
let hz = startHz
|
||||
const readBytes = fftSize * 2 * 2 // 2 samples per fft point, unsigned iq
|
||||
try {
|
||||
while (this.running) {
|
||||
await this.setFrequency(hz)
|
||||
await new Promise((r) => setTimeout(r, settleMs))
|
||||
try {
|
||||
await this.readSamples(16384) // discard: PLL settle
|
||||
const iq = await this.readSamples(readBytes)
|
||||
const f64 = new Float64Array(iq.length)
|
||||
for (let i = 0; i < iq.length; i++) f64[i] = (iq[i] - 127.5) / 128
|
||||
const db = powerSpectrumDb(f64)
|
||||
cb.onSpectrum?.(hz, db)
|
||||
} catch (err) {
|
||||
cb.onError?.(err instanceof Error ? err : new Error(String(err)))
|
||||
}
|
||||
hz += stepHz
|
||||
if (hz > endHz) hz = startHz
|
||||
if (!this.running) break
|
||||
}
|
||||
} finally {
|
||||
this.running = false
|
||||
}
|
||||
}
|
||||
|
||||
stopSweep(): void {
|
||||
this.running = false
|
||||
}
|
||||
|
||||
// ---- Low-level USB helpers (private; HARDWARE PASS REQUIRED) ----
|
||||
|
||||
private async claim(): Promise<void> {
|
||||
const dev = this.device
|
||||
if (!dev) throw new Error('no device')
|
||||
try {
|
||||
// Best-effort kernel driver detach; unsupported on some platforms.
|
||||
const anyDev = dev as unknown as {
|
||||
claimInterface(n: number): Promise<void>
|
||||
}
|
||||
await anyDev.claimInterface(this.interfaceNumber)
|
||||
} catch (err) {
|
||||
throw new Error(
|
||||
`could not claim the radio dead (interface busy?) — ${String(err)}`,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
private async writeReg(
|
||||
block: number,
|
||||
address: number,
|
||||
value: number,
|
||||
length: number,
|
||||
): Promise<void> {
|
||||
const dev = this.device
|
||||
if (!dev) throw new Error('no device')
|
||||
const data = new Uint8Array([value & 0xff, (value >> 8) & 0xff])
|
||||
await dev.controlTransferOut({
|
||||
requestType: 'vendor',
|
||||
recipient: 'device',
|
||||
request: 0,
|
||||
value: (block << 8) | 0x10,
|
||||
index: address,
|
||||
}, data.subarray(0, length))
|
||||
}
|
||||
|
||||
private async demodWrite(
|
||||
page: number,
|
||||
address: number,
|
||||
value: number,
|
||||
length: number,
|
||||
): Promise<void> {
|
||||
// Demod registers are paged: index = (page << 8) | address.
|
||||
const dev = this.device
|
||||
if (!dev) throw new Error('no device')
|
||||
const data = new Uint8Array([value & 0xff, (value >> 8) & 0xff])
|
||||
await dev.controlTransferOut({
|
||||
requestType: 'vendor',
|
||||
recipient: 'device',
|
||||
request: 0,
|
||||
value: (DEMOD << 8) | 0x10,
|
||||
index: (page << 8) | address,
|
||||
}, data.subarray(0, length))
|
||||
}
|
||||
|
||||
private async i2cWrite(i2cAddr: number, reg: number, value: number): Promise<void> {
|
||||
// I2C repeater: librtlsdr writes tuner registers by tunneling through
|
||||
// the demod's I2C master. Simplified single-byte write.
|
||||
const dev = this.device
|
||||
if (!dev) throw new Error('no device')
|
||||
await this.demodWrite(1, 0x02, 0x41, 1) // repeater on
|
||||
await dev.controlTransferOut({
|
||||
requestType: 'vendor',
|
||||
recipient: 'device',
|
||||
request: 0,
|
||||
value: (0x02 << 8) | 0x10,
|
||||
index: (i2cAddr << 8) | reg,
|
||||
}, new Uint8Array([value & 0xff]))
|
||||
await this.demodWrite(1, 0x02, 0x01, 1) // repeater off
|
||||
}
|
||||
|
||||
// Silence unused-warnings for constants kept for the hardware pass.
|
||||
private static readonly _refs = { CTRL_IN, CTRL_OUT, SYS, PAGE_USB }
|
||||
}
|
||||
|
||||
/** Rolling noise-floor + spike detector over sweep spectra (pure, testable). */
|
||||
export class SpectrumAnomalyDetector {
|
||||
private floor: Float64Array | null = null
|
||||
private lastEmit = -Infinity
|
||||
constructor(
|
||||
private readonly thresholdDb = 10,
|
||||
private readonly throttleMs = 2000,
|
||||
private readonly alpha = 0.1,
|
||||
) {}
|
||||
|
||||
/** Returns (peakHz, magnitudeDbOverFloor) when a spike fires, else null. */
|
||||
process(
|
||||
centerHz: number,
|
||||
sampleRateHz: number,
|
||||
db: Float64Array,
|
||||
nowMs: number,
|
||||
): { frequency: number; magnitude: number } | null {
|
||||
if (!this.floor || this.floor.length !== db.length) {
|
||||
this.floor = Float64Array.from(db)
|
||||
return null
|
||||
}
|
||||
let peak = 0
|
||||
let peakBin = -1
|
||||
for (let i = 0; i < db.length; i++) {
|
||||
const dev = db[i] - this.floor[i]
|
||||
if (dev > peak) {
|
||||
peak = dev
|
||||
peakBin = i
|
||||
}
|
||||
this.floor[i] += this.alpha * (db[i] - this.floor[i])
|
||||
}
|
||||
if (peakBin >= 0 && peak >= this.thresholdDb && nowMs - this.lastEmit >= this.throttleMs) {
|
||||
this.lastEmit = nowMs
|
||||
const binHz = sampleRateHz / 2 / db.length
|
||||
const offset = (peakBin - db.length / 2) * binHz
|
||||
return { frequency: (centerHz + offset) / 1e6, magnitude: peak }
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
reset(): void {
|
||||
this.floor = null
|
||||
this.lastEmit = -Infinity
|
||||
}
|
||||
}
|
||||
|
||||
export const SWEEP_START_MHZ = 88
|
||||
export const SWEEP_END_MHZ = 108
|
||||
102
frontend/src/lib/types.ts
Normal file
102
frontend/src/lib/types.ts
Normal file
@@ -0,0 +1,102 @@
|
||||
// Shared protocol + domain types for the Quantumancy frontend.
|
||||
// Mirrors the backend contract exactly — do not invent changes.
|
||||
|
||||
export type Mode = 'wire' | 'evp' | 'radio' | 'ouija'
|
||||
export type Language = 'en' | 'es'
|
||||
export type Rarity = 'common' | 'uncommon' | 'rare' | 'mythic'
|
||||
export type GhostForm = 'wisp' | 'banshee' | 'fairy' | 'shade'
|
||||
|
||||
export type CodexVisual = {
|
||||
hue: number
|
||||
form: GhostForm
|
||||
}
|
||||
|
||||
export type CodexEntity = {
|
||||
id: string
|
||||
name: string
|
||||
epithet: string
|
||||
rarity: Rarity
|
||||
visual: CodexVisual
|
||||
quotes: string[]
|
||||
contact_count: number
|
||||
discovered_at: string
|
||||
discovered_by: string | null
|
||||
}
|
||||
|
||||
export type CodexEntityDetail = CodexEntity & {
|
||||
persona: string
|
||||
voice: Record<string, unknown>
|
||||
sightings: number
|
||||
}
|
||||
|
||||
export type SpiritVoice = {
|
||||
voice_id: string
|
||||
pitch: number
|
||||
rate: number
|
||||
noise: number
|
||||
echo: number
|
||||
}
|
||||
|
||||
export type SpiritEntity = {
|
||||
id: string
|
||||
name: string
|
||||
epithet: string
|
||||
persona: string
|
||||
rarity: Rarity
|
||||
voice: SpiritVoice
|
||||
visual: CodexVisual
|
||||
quotes: string[]
|
||||
contact_count: number
|
||||
discovered_at: string
|
||||
}
|
||||
|
||||
export type CodexListResponse = { entities: CodexEntity[] }
|
||||
|
||||
export type Stats = {
|
||||
entities: number
|
||||
sessions: number
|
||||
utterances: number
|
||||
anomalies: number
|
||||
}
|
||||
|
||||
export type SessionStatus = 'attuning' | 'summoning' | 'gathering'
|
||||
|
||||
// ---- WebSocket frames: client -> server ----
|
||||
|
||||
export type ClientFrame =
|
||||
| { type: 'ping' }
|
||||
| { type: 'set_mode'; mode: Mode }
|
||||
| { type: 'language'; language: Language }
|
||||
| { type: 'summon' }
|
||||
| { type: 'anomaly'; source: 'radio' | 'evp' | 'wire'; frequency: number; magnitude: number }
|
||||
| { type: 'question'; text: string }
|
||||
| { type: 'passive'; enabled: boolean }
|
||||
|
||||
// ---- WebSocket frames: server -> client ----
|
||||
|
||||
export type Telemetry = {
|
||||
jitter_bytes_per_s: number
|
||||
latency_variance_ms: number
|
||||
latency_mean_ms: number
|
||||
dns_ms: number
|
||||
}
|
||||
|
||||
export type UtteranceKind = 'greeting' | 'fragment' | 'ambient' | 'reply'
|
||||
|
||||
export type ServerFrame =
|
||||
| { type: 'session'; id: string }
|
||||
| { type: 'pong' }
|
||||
| { type: 'mode'; mode: Mode }
|
||||
| { type: 'status'; state: SessionStatus }
|
||||
| { type: 'entity'; is_new: boolean; entity: SpiritEntity }
|
||||
| { type: 'anomaly_ack'; count: number }
|
||||
| { type: 'utterance'; id: string; kind: UtteranceKind; text: string; entity: string | null }
|
||||
| { type: 'audio'; id: string; url: string }
|
||||
| { type: 'reply_start' }
|
||||
| { type: 'reply_token'; token: string }
|
||||
| { type: 'reply_end'; id: string; text: string }
|
||||
| { type: 'telemetry' } & Telemetry
|
||||
| { type: 'error'; code: 'rate_limited' | 'veil_crowded' | string; message: string }
|
||||
|
||||
// REST helpers
|
||||
export type SortOrder = 'recent' | 'contacted'
|
||||
295
frontend/src/lib/ws.test.ts
Normal file
295
frontend/src/lib/ws.test.ts
Normal file
@@ -0,0 +1,295 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { VeilSocket } from './ws'
|
||||
import type { VeilConnectionState } from './ws'
|
||||
import type { ClientFrame, ServerFrame } from './types'
|
||||
|
||||
const URL = 'ws://quantumancy.test/ws/session'
|
||||
|
||||
/**
|
||||
* Minimal stand-in for the DOM WebSocket, driven manually by the tests.
|
||||
* readyState uses the real constant values so VeilSocket's
|
||||
* `readyState === WebSocket.OPEN` checks behave (WebSocket.OPEN === 1).
|
||||
*/
|
||||
class FakeWebSocket {
|
||||
readyState = 0 // CONNECTING
|
||||
onopen: (() => void) | null = null
|
||||
onmessage: ((ev: { data: unknown }) => void) | null = null
|
||||
onclose: (() => void) | null = null
|
||||
onerror: (() => void) | null = null
|
||||
readonly sent: string[] = []
|
||||
|
||||
constructor(readonly url: string) {}
|
||||
|
||||
send(data: string): void {
|
||||
this.sent.push(data)
|
||||
}
|
||||
|
||||
close(): void {
|
||||
this.readyState = 3 // CLOSED
|
||||
this.onclose?.()
|
||||
}
|
||||
|
||||
// ---- test drives ----
|
||||
|
||||
serverOpen(): void {
|
||||
this.readyState = 1 // WebSocket.OPEN
|
||||
this.onopen?.()
|
||||
}
|
||||
|
||||
serverSend(data: unknown): void {
|
||||
this.onmessage?.({ data })
|
||||
}
|
||||
|
||||
serverClose(): void {
|
||||
this.readyState = 3
|
||||
this.onclose?.()
|
||||
}
|
||||
}
|
||||
|
||||
type HarnessOptions = {
|
||||
baseBackoffMs?: number
|
||||
maxBackoffMs?: number
|
||||
pingIntervalMs?: number
|
||||
}
|
||||
|
||||
function makeHarness(opts: HarnessOptions = {}) {
|
||||
const sockets: FakeWebSocket[] = []
|
||||
const socket = new VeilSocket({
|
||||
url: URL,
|
||||
baseBackoffMs: opts.baseBackoffMs ?? 100,
|
||||
maxBackoffMs: opts.maxBackoffMs ?? 1000,
|
||||
pingIntervalMs: opts.pingIntervalMs ?? 1000,
|
||||
socketFactory: (url: string): WebSocket => {
|
||||
const fake = new FakeWebSocket(url)
|
||||
sockets.push(fake)
|
||||
return fake as unknown as WebSocket
|
||||
},
|
||||
})
|
||||
return { socket, sockets }
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers()
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
describe('VeilSocket', () => {
|
||||
it('queues frames while connecting and flushes them on open', () => {
|
||||
const { socket, sockets } = makeHarness()
|
||||
const states: VeilConnectionState[] = []
|
||||
socket.onState((s) => states.push(s))
|
||||
|
||||
socket.connect()
|
||||
expect(socket.state).toBe('connecting')
|
||||
expect(sockets).toHaveLength(1)
|
||||
expect(sockets[0].url).toBe(URL)
|
||||
|
||||
const summon: ClientFrame = { type: 'summon' }
|
||||
const question: ClientFrame = { type: 'question', text: 'is anyone there' }
|
||||
socket.send(summon)
|
||||
socket.send(question)
|
||||
expect(sockets[0].sent).toEqual([]) // nothing on the wire yet
|
||||
|
||||
sockets[0].serverOpen()
|
||||
expect(socket.state).toBe('open')
|
||||
expect(sockets[0].sent).toEqual([JSON.stringify(summon), JSON.stringify(question)])
|
||||
expect(states).toEqual(['connecting', 'open'])
|
||||
})
|
||||
|
||||
it('sends immediately once open', () => {
|
||||
const { socket, sockets } = makeHarness()
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
socket.send({ type: 'ping' })
|
||||
expect(sockets[0].sent).toEqual([JSON.stringify({ type: 'ping' })])
|
||||
})
|
||||
|
||||
it('dispatches parsed JSON frames to onFrame handlers', () => {
|
||||
const { socket, sockets } = makeHarness()
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
|
||||
const frames: ServerFrame[] = []
|
||||
const off = socket.onFrame((f) => frames.push(f))
|
||||
|
||||
const session: ServerFrame = { type: 'session', id: 'abc' }
|
||||
sockets[0].serverSend(JSON.stringify(session))
|
||||
expect(frames).toEqual([session])
|
||||
|
||||
off()
|
||||
sockets[0].serverSend(JSON.stringify({ type: 'pong' }))
|
||||
expect(frames).toHaveLength(1) // unsubscribed handler no longer called
|
||||
})
|
||||
|
||||
it('ignores malformed JSON and non-string payloads', () => {
|
||||
const { socket, sockets } = makeHarness()
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
|
||||
const frames: ServerFrame[] = []
|
||||
socket.onFrame((f) => frames.push(f))
|
||||
expect(() => {
|
||||
sockets[0].serverSend('}{ not json')
|
||||
sockets[0].serverSend(new Uint8Array([1, 2, 3]).buffer)
|
||||
sockets[0].serverSend(undefined)
|
||||
}).not.toThrow()
|
||||
expect(frames).toEqual([])
|
||||
})
|
||||
|
||||
it('drops frames sent while unstable instead of queueing them', () => {
|
||||
const { socket, sockets } = makeHarness()
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
sockets[0].serverClose()
|
||||
expect(socket.state).toBe('unstable')
|
||||
|
||||
socket.send({ type: 'summon' }) // dropped: not open, not connecting
|
||||
vi.advanceTimersByTime(100) // reconnect fires
|
||||
expect(sockets).toHaveLength(2)
|
||||
sockets[1].serverOpen()
|
||||
expect(sockets[1].sent).toEqual([]) // nothing stale flushed
|
||||
})
|
||||
|
||||
it('drops frames sent while closed', () => {
|
||||
const { socket, sockets } = makeHarness()
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
socket.close()
|
||||
socket.send({ type: 'summon' })
|
||||
expect(sockets[0].sent).toEqual([])
|
||||
})
|
||||
|
||||
it('goes unstable on unexpected close and reconnects with exponential backoff', () => {
|
||||
const { socket, sockets } = makeHarness({ baseBackoffMs: 100, maxBackoffMs: 1000 })
|
||||
const states: VeilConnectionState[] = []
|
||||
socket.onState((s) => states.push(s))
|
||||
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
sockets[0].serverClose()
|
||||
expect(socket.state).toBe('unstable')
|
||||
expect(sockets).toHaveLength(1)
|
||||
|
||||
// First retry after 100 ms (base * 2^0).
|
||||
vi.advanceTimersByTime(99)
|
||||
expect(sockets).toHaveLength(1)
|
||||
vi.advanceTimersByTime(1)
|
||||
expect(sockets).toHaveLength(2)
|
||||
expect(socket.state).toBe('connecting')
|
||||
|
||||
// Second drop: backoff doubles to 200 ms.
|
||||
sockets[1].serverClose()
|
||||
vi.advanceTimersByTime(199)
|
||||
expect(sockets).toHaveLength(2)
|
||||
vi.advanceTimersByTime(1)
|
||||
expect(sockets).toHaveLength(3)
|
||||
|
||||
// The reconnect succeeds and the veil lifts again.
|
||||
sockets[2].serverOpen()
|
||||
expect(socket.state).toBe('open')
|
||||
expect(states).toEqual([
|
||||
'connecting',
|
||||
'open',
|
||||
'unstable',
|
||||
'connecting',
|
||||
'unstable',
|
||||
'connecting',
|
||||
'open',
|
||||
])
|
||||
})
|
||||
|
||||
it('caps the backoff at maxBackoffMs', () => {
|
||||
const { socket, sockets } = makeHarness({ baseBackoffMs: 100, maxBackoffMs: 250 })
|
||||
socket.connect()
|
||||
sockets[0].serverClose() // retry in 100 ms
|
||||
vi.advanceTimersByTime(100)
|
||||
sockets[1].serverClose() // retry in 200 ms
|
||||
vi.advanceTimersByTime(200)
|
||||
sockets[2].serverClose() // retry in min(250, 400) = 250 ms
|
||||
vi.advanceTimersByTime(249)
|
||||
expect(sockets).toHaveLength(3)
|
||||
vi.advanceTimersByTime(1)
|
||||
expect(sockets).toHaveLength(4)
|
||||
})
|
||||
|
||||
it('does not reconnect after a deliberate close()', () => {
|
||||
const { socket, sockets } = makeHarness()
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
socket.close()
|
||||
expect(socket.state).toBe('closed')
|
||||
vi.advanceTimersByTime(60_000)
|
||||
expect(sockets).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('does not reconnect when close() races a pending reconnect timer', () => {
|
||||
const { socket, sockets } = makeHarness()
|
||||
socket.connect()
|
||||
sockets[0].serverClose() // reconnect scheduled in 100 ms
|
||||
socket.close() // deliberate: must cancel the pending timer
|
||||
vi.advanceTimersByTime(60_000)
|
||||
expect(sockets).toHaveLength(1)
|
||||
expect(socket.state).toBe('closed')
|
||||
})
|
||||
|
||||
it('sends ping frames on the configured interval while open', () => {
|
||||
const { socket, sockets } = makeHarness({ pingIntervalMs: 1000 })
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
|
||||
vi.advanceTimersByTime(3000)
|
||||
expect(sockets[0].sent).toEqual([
|
||||
JSON.stringify({ type: 'ping' }),
|
||||
JSON.stringify({ type: 'ping' }),
|
||||
JSON.stringify({ type: 'ping' }),
|
||||
])
|
||||
})
|
||||
|
||||
it('stops pinging after the connection drops', () => {
|
||||
const { socket, sockets } = makeHarness({ pingIntervalMs: 1000 })
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
vi.advanceTimersByTime(2000)
|
||||
expect(sockets[0].sent).toHaveLength(2)
|
||||
|
||||
sockets[0].serverClose()
|
||||
vi.advanceTimersByTime(10_000)
|
||||
expect(sockets[0].sent).toHaveLength(2) // frozen at the drop
|
||||
})
|
||||
|
||||
it('does not ping when pingIntervalMs is 0', () => {
|
||||
const { socket, sockets } = makeHarness({ pingIntervalMs: 0 })
|
||||
socket.connect()
|
||||
sockets[0].serverOpen()
|
||||
vi.advanceTimersByTime(10_000)
|
||||
expect(sockets[0].sent).toEqual([])
|
||||
})
|
||||
|
||||
it('goes unstable and retries when the socket factory throws', () => {
|
||||
const sockets: FakeWebSocket[] = []
|
||||
let calls = 0
|
||||
const socket = new VeilSocket({
|
||||
url: URL,
|
||||
baseBackoffMs: 100,
|
||||
pingIntervalMs: 0,
|
||||
socketFactory: (): WebSocket => {
|
||||
calls++
|
||||
if (calls === 1) throw new Error('the veil refused the connection')
|
||||
const fake = new FakeWebSocket(URL)
|
||||
sockets.push(fake)
|
||||
return fake as unknown as WebSocket
|
||||
},
|
||||
})
|
||||
|
||||
socket.connect()
|
||||
expect(socket.state).toBe('unstable')
|
||||
vi.advanceTimersByTime(100)
|
||||
expect(calls).toBe(2)
|
||||
expect(sockets).toHaveLength(1)
|
||||
sockets[0].serverOpen()
|
||||
expect(socket.state).toBe('open')
|
||||
})
|
||||
})
|
||||
171
frontend/src/lib/ws.ts
Normal file
171
frontend/src/lib/ws.ts
Normal file
@@ -0,0 +1,171 @@
|
||||
// Reconnecting WebSocket client for /ws/session with a typed event emitter.
|
||||
// Exponential backoff on unexpected close; surfaces connection status so the
|
||||
// UI can show "the connection to the other side is unstable".
|
||||
|
||||
import type { ClientFrame, ServerFrame } from './types'
|
||||
|
||||
export type VeilConnectionState =
|
||||
| 'connecting'
|
||||
| 'open'
|
||||
| 'unstable' // reconnecting after a drop
|
||||
| 'closed'
|
||||
|
||||
type FrameHandler = (frame: ServerFrame) => void
|
||||
type StateHandler = (state: VeilConnectionState) => void
|
||||
|
||||
export type VeilSocketOptions = {
|
||||
url?: string
|
||||
/** Base backoff delay (ms); doubles each retry up to maxBackoffMs. */
|
||||
baseBackoffMs?: number
|
||||
maxBackoffMs?: number
|
||||
/** Ping interval to keep the veil warm (ms). 0 disables. */
|
||||
pingIntervalMs?: number
|
||||
/** Injectable WebSocket constructor (tests). */
|
||||
socketFactory?: (url: string) => WebSocket
|
||||
}
|
||||
|
||||
export function defaultSessionUrl(): string {
|
||||
const proto = location.protocol === 'https:' ? 'wss://' : 'ws://'
|
||||
return `${proto}${location.host}/ws/session`
|
||||
}
|
||||
|
||||
export class VeilSocket {
|
||||
private ws: WebSocket | null = null
|
||||
private frameHandlers = new Set<FrameHandler>()
|
||||
private stateHandlers = new Set<StateHandler>()
|
||||
private readonly url: string
|
||||
private readonly baseBackoffMs: number
|
||||
private readonly maxBackoffMs: number
|
||||
private readonly pingIntervalMs: number
|
||||
private readonly socketFactory: (url: string) => WebSocket
|
||||
private attempts = 0
|
||||
private reconnectTimer: ReturnType<typeof setTimeout> | null = null
|
||||
private pingTimer: ReturnType<typeof setInterval> | null = null
|
||||
private deliberatelyClosed = false
|
||||
private outbox: string[] = []
|
||||
private _state: VeilConnectionState = 'closed'
|
||||
|
||||
constructor(opts: VeilSocketOptions = {}) {
|
||||
this.url = opts.url ?? defaultSessionUrl()
|
||||
this.baseBackoffMs = opts.baseBackoffMs ?? 800
|
||||
this.maxBackoffMs = opts.maxBackoffMs ?? 15000
|
||||
this.pingIntervalMs = opts.pingIntervalMs ?? 25000
|
||||
this.socketFactory = opts.socketFactory ?? ((url: string) => new WebSocket(url))
|
||||
}
|
||||
|
||||
get state(): VeilConnectionState {
|
||||
return this._state
|
||||
}
|
||||
|
||||
onFrame(handler: FrameHandler): () => void {
|
||||
this.frameHandlers.add(handler)
|
||||
return () => this.frameHandlers.delete(handler)
|
||||
}
|
||||
|
||||
onState(handler: StateHandler): () => void {
|
||||
this.stateHandlers.add(handler)
|
||||
return () => this.stateHandlers.delete(handler)
|
||||
}
|
||||
|
||||
connect(): void {
|
||||
this.deliberatelyClosed = false
|
||||
this.openSocket()
|
||||
}
|
||||
|
||||
close(): void {
|
||||
this.deliberatelyClosed = true
|
||||
this.clearTimers()
|
||||
this.ws?.close()
|
||||
this.ws = null
|
||||
this.setState('closed')
|
||||
}
|
||||
|
||||
/** Send a frame; queued while connecting, dropped while closed/unstable. */
|
||||
send(frame: ClientFrame): void {
|
||||
const raw = JSON.stringify(frame)
|
||||
if (this.ws && this.ws.readyState === WebSocket.OPEN) {
|
||||
this.ws.send(raw)
|
||||
} else if (this._state === 'connecting') {
|
||||
if (this.outbox.length < 64) this.outbox.push(raw)
|
||||
}
|
||||
}
|
||||
|
||||
private setState(state: VeilConnectionState): void {
|
||||
if (this._state === state) return
|
||||
this._state = state
|
||||
this.stateHandlers.forEach((h) => h(state))
|
||||
}
|
||||
|
||||
private clearTimers(): void {
|
||||
if (this.reconnectTimer !== null) {
|
||||
clearTimeout(this.reconnectTimer)
|
||||
this.reconnectTimer = null
|
||||
}
|
||||
if (this.pingTimer !== null) {
|
||||
clearInterval(this.pingTimer)
|
||||
this.pingTimer = null
|
||||
}
|
||||
}
|
||||
|
||||
private openSocket(): void {
|
||||
this.clearTimers()
|
||||
this.setState('connecting')
|
||||
let ws: WebSocket
|
||||
try {
|
||||
ws = this.socketFactory(this.url)
|
||||
} catch {
|
||||
this.scheduleReconnect()
|
||||
return
|
||||
}
|
||||
this.ws = ws
|
||||
|
||||
ws.onopen = () => {
|
||||
this.attempts = 0
|
||||
this.setState('open')
|
||||
// Flush frames queued while connecting.
|
||||
for (const raw of this.outbox) ws.send(raw)
|
||||
this.outbox = []
|
||||
if (this.pingIntervalMs > 0) {
|
||||
this.pingTimer = setInterval(() => {
|
||||
if (ws.readyState === WebSocket.OPEN) {
|
||||
ws.send(JSON.stringify({ type: 'ping' }))
|
||||
}
|
||||
}, this.pingIntervalMs)
|
||||
}
|
||||
}
|
||||
|
||||
ws.onmessage = (ev: MessageEvent) => {
|
||||
if (typeof ev.data !== 'string') return
|
||||
let frame: ServerFrame
|
||||
try {
|
||||
frame = JSON.parse(ev.data) as ServerFrame
|
||||
} catch {
|
||||
return // ignore malformed whispers
|
||||
}
|
||||
this.frameHandlers.forEach((h) => h(frame))
|
||||
}
|
||||
|
||||
ws.onclose = () => {
|
||||
this.clearTimers()
|
||||
this.ws = null
|
||||
if (!this.deliberatelyClosed) this.scheduleReconnect()
|
||||
else this.setState('closed')
|
||||
}
|
||||
|
||||
ws.onerror = () => {
|
||||
// onclose follows onerror; nothing extra to do here.
|
||||
}
|
||||
}
|
||||
|
||||
private scheduleReconnect(): void {
|
||||
this.setState('unstable')
|
||||
const delay = Math.min(
|
||||
this.maxBackoffMs,
|
||||
this.baseBackoffMs * 2 ** this.attempts,
|
||||
)
|
||||
this.attempts++
|
||||
this.reconnectTimer = setTimeout(() => {
|
||||
if (!this.deliberatelyClosed) this.openSocket()
|
||||
}, delay)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user