Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 19 additions & 12 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@ import { buildProviderTalkerLookups } from './nmea0183TalkerGroups'
import { pipedProviders } from './pipedproviders'
import { EventsActorId, WithWrappedEmitter, wrapEmitter } from './events'
import { StalenessEnforcer } from './staleness'
import { ProviderStatusEmitter } from './providerStatusEmitter'
import { Zones } from './zones'
import checkNodeVersion from './version'
import helmet from 'helmet'
Expand Down Expand Up @@ -105,6 +106,8 @@ function cloneDelta(delta: any): any {
}
}

const PROVIDER_STATUS_REFRESH_INTERVAL_MS = 5 * 1000

class Server {
app: ServerApp &
WithConfig &
Expand All @@ -120,6 +123,7 @@ class Server {
// Pending sourceRef migration timers; cleared on stop() so a deferred
// migration scheduled before a restart cannot fire on a torn-down app.
pendingSourceRefMigrations?: Set<NodeJS.Timeout>
private providerStatusEmitter: ProviderStatusEmitter

constructor(opts: { securityConfig: SecurityConfig }) {
checkNodeVersion()
Expand Down Expand Up @@ -225,6 +229,15 @@ class Server {
}
Object.assign(app, pluginManager)

const providerStatusEmitter = new ProviderStatusEmitter(() => {
app.emit('serverevent', {
type: 'PROVIDERSTATUS',
from: 'signalk-server',
data: app.getProviderStatus()
})
})
this.providerStatusEmitter = providerStatusEmitter

app.setPluginStatus = (providerId: string, statusMessage: string) => {
doSetProviderStatus(providerId, statusMessage, 'status', 'plugin')
}
Expand Down Expand Up @@ -269,11 +282,7 @@ class Server {

status.message = statusMessage

app.emit('serverevent', {
type: 'PROVIDERSTATUS',
from: 'signalk-server',
data: app.getProviderStatus()
})
providerStatusEmitter.request()
}

app.getProviderStatus = () => {
Expand Down Expand Up @@ -670,13 +679,10 @@ class Server {
app.stalenessEnforcer = undefined
}
app.intervals.push(
setInterval(() => {
app.emit('serverevent', {
type: 'PROVIDERSTATUS',
from: 'signalk-server',
data: app.getProviderStatus()
})
}, 5 * 1000)
setInterval(
() => self.providerStatusEmitter.emitNow(),
PROVIDER_STATUS_REFRESH_INTERVAL_MS
)
)

function serverUpgradeIsAvailable(err: any, newVersion?: string) {
Expand Down Expand Up @@ -806,6 +812,7 @@ class Server {
}

async stop(cb?: () => void) {
this.providerStatusEmitter.cancel()
if (!this.app.started) {
return this
}
Expand Down
49 changes: 49 additions & 0 deletions src/providerStatusEmitter.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
/**
* Coalesces PROVIDERSTATUS serverevents.
*
* setPluginStatus / setProviderStatus each carry the full status list, and
* serverevents are written to every admin-UI WebSocket without
* backpressure, so a plugin calling them in a tight loop can push megabytes
* per second into each socket. request() emits at most once per
* PROVIDER_STATUS_MIN_INTERVAL_MS: the first request after a quiet interval
* emits promptly, and every further request before the interval has passed
* folds into a single emit scheduled for when it does. The admin UI thus
* sees the latest status within one interval however chatty the caller is.
*/

export const PROVIDER_STATUS_MIN_INTERVAL_MS = 1000

export class ProviderStatusEmitter {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this is not a ProviderStatusEmitter - it does not emit provider statuses, but is used to emit them. to me it is ThrottledCaller - it calls the given function but throttles the calls as described. imho naming should reflect that and min interval come from callee.

private timer: NodeJS.Timeout | undefined
private lastEmitAt = 0

constructor(
private readonly emit: () => void,
private readonly minIntervalMs = PROVIDER_STATUS_MIN_INTERVAL_MS
) {}

request(): void {
if (this.timer) return
const wait = Math.max(0, this.lastEmitAt + this.minIntervalMs - Date.now())
this.timer = setTimeout(() => {
this.timer = undefined
this.emitNow()
}, wait)
// A pending status refresh must never keep a stopping server alive.
this.timer.unref?.()
}

/** Emit immediately, dropping any pending request it would duplicate. */
emitNow(): void {
this.cancel()
this.lastEmitAt = Date.now()
this.emit()
}

cancel(): void {
if (this.timer) {
clearTimeout(this.timer)
this.timer = undefined
}
}
}
74 changes: 74 additions & 0 deletions test/providerStatusEmitter.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
import { expect } from 'chai'
import { ProviderStatusEmitter } from '../dist/providerStatusEmitter'

const INTERVAL_MS = 40
// Long enough for a zero-delay timer to fire, short against INTERVAL_MS.
const SETTLE_MS = 5
const BURST_SIZE = 1000

const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms))

function burst(emitter: ProviderStatusEmitter) {
for (let i = 0; i < BURST_SIZE; i++) {
emitter.request()
}
}

describe('ProviderStatusEmitter', function () {
it('folds a burst into one prompt emit and one more when the interval expires', async function () {
let emits = 0
const emitter = new ProviderStatusEmitter(() => emits++, INTERVAL_MS)

burst(emitter)
await sleep(SETTLE_MS)
expect(emits).to.equal(1)

burst(emitter)
await sleep(INTERVAL_MS / 2)
expect(emits).to.equal(1)
await sleep(INTERVAL_MS)
expect(emits).to.equal(2)

emitter.cancel()
})

it('emits promptly when the last emit is older than the interval', async function () {
let emits = 0
const emitter = new ProviderStatusEmitter(() => emits++, INTERVAL_MS)

emitter.request()
await sleep(SETTLE_MS)
expect(emits).to.equal(1)

await sleep(INTERVAL_MS + SETTLE_MS)
emitter.request()
await sleep(SETTLE_MS)
expect(emits).to.equal(2)

emitter.cancel()
})

it('emitNow drops a pending request instead of emitting twice', async function () {
let emits = 0
const emitter = new ProviderStatusEmitter(() => emits++, INTERVAL_MS)

emitter.emitNow()
emitter.request()
expect(emits).to.equal(1)
emitter.emitNow()
expect(emits).to.equal(2)
await sleep(INTERVAL_MS + SETTLE_MS)
expect(emits).to.equal(2)
})

it('cancel discards a pending request', async function () {
let emits = 0
const emitter = new ProviderStatusEmitter(() => emits++, INTERVAL_MS)

emitter.emitNow()
emitter.request()
emitter.cancel()
await sleep(INTERVAL_MS + SETTLE_MS)
expect(emits).to.equal(1)
})
})
Loading