The site now ingests the sidecar's live WebSocket feed and persists it to its own MariaDB, and re-broadcasts curated events to browsers over SSE. - schema: shard_events (append-only notable-kind log, sha1 dedupe_key + INSERT IGNORE for idempotent reconnect backfill), shard_online (current players, upsert/refresh/remove), shard_economy (gold-supply series), shard_houses (per-house decay stage + derived is_idoc). - model/shardEvents + model/shardState: the .db.js/.model.js split; writes take camelCase event data, reads are shaped; online upsert uses COALESCE so a partial char.vitals refresh never blanks login fields. - utils/shardIngest: single dispatcher routing each kind to state writes and/or the event log, then the broadcaster. High-frequency kinds (char.vitals, economy.supply) update state only. A changed server.hello bootId clears the stale online roster. Deps are injected for unit testing. - utils/uoLinkSocket: the server's first outbound WS client (ws dep). Verifies the ws.hello protocol, backfills via /history + /economy on every (re)connect (dedupe handles overlap), reconnects with capped backoff, and mirrors connection state into uo_link_config. Self-guards: only connects when the integration is enabled with a token. - utils/shardBroadcast: SSE fan-out with public (safe kinds only) and admin (all) channels, keepalive pings, per-client cleanup. - server.js: start the ingest socket on boot (no-op until configured) and stop it + close SSE streams on graceful shutdown. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011qPmpmVH1xGCiZoz9m9vW3
196 lines
6.0 KiB
JavaScript
196 lines
6.0 KiB
JavaScript
// ── uo-link WebSocket ingest client ────────────────────────────────────────
|
|
//
|
|
// Long-lived client that connects to the sidecar's push-only WS feed and pumps
|
|
// every frame through the ingest dispatcher. This is the server's first
|
|
// outbound WebSocket. Lifecycle:
|
|
// • start() — connect if the config is enabled and has a token; verify the
|
|
// ws.hello protocol; backfill missed events via /history on every
|
|
// (re)connect (INSERT IGNORE dedupes the overlap); reconnect with
|
|
// capped backoff.
|
|
// • stop() — close the socket and stop reconnecting (graceful shutdown).
|
|
// Connection state is mirrored into uo_link_config (plugin_connected / status /
|
|
// last_event_at) so the admin panel and public status endpoint have live data.
|
|
|
|
const WebSocket = require('ws')
|
|
|
|
const uoLinkConfig = require('../model/uoLinkConfig/uoLinkConfig.model')
|
|
const uoLinkClient = require('./uoLinkClient')
|
|
const shardIngest = require('./shardIngest')
|
|
const log = require('./logger')('uo-link-socket')
|
|
|
|
const BACKOFF_MIN_MS = 1000
|
|
const BACKOFF_MAX_MS = 30000
|
|
const BACKFILL_LIMIT = 500
|
|
|
|
let ws = null
|
|
let reconnectTimer = null
|
|
let backoff = BACKOFF_MIN_MS
|
|
let running = false // set by start()/stop(); guards auto-reconnect
|
|
let helloSeen = false
|
|
|
|
const state = {
|
|
connected: false,
|
|
lastEventAt: null,
|
|
lastConnectedAt: null,
|
|
reconnects: 0,
|
|
protocol: null,
|
|
}
|
|
|
|
function buildUrl(wsUrl, token) {
|
|
const sep = wsUrl.includes('?') ? '&' : '?'
|
|
return token ? `${wsUrl}${sep}token=${encodeURIComponent(token)}` : wsUrl
|
|
}
|
|
|
|
// Pull recent events from the sidecar's own store and replay them through the
|
|
// dispatcher (fromBackfill = no SSE re-broadcast). dedupe_key + INSERT IGNORE
|
|
// make this idempotent, so overlap with what we already stored is harmless.
|
|
async function backfill() {
|
|
try {
|
|
const hist = await uoLinkClient.getHistory({ limit: BACKFILL_LIMIT })
|
|
if (hist.ok && hist.data && Array.isArray(hist.data.events)) {
|
|
// History is newest-first; replay oldest-first so latest-wins state (e.g.
|
|
// house.decay stage) settles correctly.
|
|
const events = [...hist.data.events].reverse()
|
|
for (const ev of events) await shardIngest.ingest(ev, { fromBackfill: true })
|
|
log.info('backfilled events from /history', { count: events.length })
|
|
}
|
|
const eco = await uoLinkClient.getEconomy(200)
|
|
if (eco.ok && eco.data && Array.isArray(eco.data.series)) {
|
|
const series = [...eco.data.series].reverse()
|
|
for (const ev of series) await shardIngest.ingest(ev, { fromBackfill: true })
|
|
}
|
|
} catch (err) {
|
|
log.warn('backfill failed (continuing on live feed)', { message: err.message })
|
|
}
|
|
}
|
|
|
|
function scheduleReconnect() {
|
|
if (!running) return
|
|
clearTimeout(reconnectTimer)
|
|
reconnectTimer = setTimeout(connect, backoff)
|
|
log.info(`reconnecting in ${backoff}ms`)
|
|
backoff = Math.min(backoff * 2, BACKOFF_MAX_MS)
|
|
}
|
|
|
|
async function connect() {
|
|
if (!running) return
|
|
let config
|
|
try {
|
|
config = await uoLinkConfig.getWithToken()
|
|
} catch (err) {
|
|
log.error('could not read uo-link config', err)
|
|
scheduleReconnect()
|
|
return
|
|
}
|
|
if (!config || !config.enabled || !config.wsUrl || !config.token) {
|
|
log.info('uo-link WS not started (disabled or missing url/token)')
|
|
running = false
|
|
return
|
|
}
|
|
|
|
state.protocol = config.protocol || 1
|
|
helloSeen = false
|
|
const url = buildUrl(config.wsUrl, config.token)
|
|
|
|
try {
|
|
ws = new WebSocket(url)
|
|
} catch (err) {
|
|
log.error('failed to open WS', err)
|
|
scheduleReconnect()
|
|
return
|
|
}
|
|
|
|
ws.on('open', async () => {
|
|
log.info('uo-link WS connected')
|
|
state.connected = true
|
|
state.lastConnectedAt = Date.now()
|
|
backoff = BACKOFF_MIN_MS
|
|
await uoLinkConfig.recordStatus({ status: 'connected', statusDetail: null, pluginConnected: true }).catch(() => {})
|
|
await backfill()
|
|
})
|
|
|
|
ws.on('message', async (raw) => {
|
|
let event
|
|
try {
|
|
event = JSON.parse(raw.toString())
|
|
} catch {
|
|
log.warn('dropping non-JSON WS frame')
|
|
return
|
|
}
|
|
|
|
if (event.kind === 'ws.hello') {
|
|
helloSeen = true
|
|
if (event.protocol && event.protocol !== state.protocol) {
|
|
log.error('uo-link protocol mismatch on ws.hello — closing', {
|
|
expected: state.protocol,
|
|
got: event.protocol,
|
|
})
|
|
await uoLinkConfig
|
|
.recordStatus({ status: 'error', statusDetail: `protocol mismatch: expected ${state.protocol}, got ${event.protocol}` })
|
|
.catch(() => {})
|
|
running = false
|
|
try {
|
|
ws.close()
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
return
|
|
}
|
|
if (event.kind === 'pong') return // sidecar heartbeat — ignore
|
|
|
|
state.lastEventAt = Number.isFinite(event.t) ? event.t : Date.now()
|
|
await shardIngest.ingest(event)
|
|
})
|
|
|
|
ws.on('close', async () => {
|
|
state.connected = false
|
|
if (running) state.reconnects += 1
|
|
log.warn('uo-link WS closed')
|
|
await uoLinkConfig
|
|
.recordStatus({ status: running ? 'reconnecting' : 'disconnected', pluginConnected: false })
|
|
.catch(() => {})
|
|
ws = null
|
|
scheduleReconnect()
|
|
})
|
|
|
|
ws.on('error', (err) => {
|
|
log.warn('uo-link WS error', { message: err.message })
|
|
// 'close' fires after 'error'; reconnect is scheduled there.
|
|
})
|
|
}
|
|
|
|
// Begin (or restart) the WS client. Idempotent — a running client is stopped
|
|
// first so a config save can re-point it at a new URL/token.
|
|
async function start() {
|
|
stop()
|
|
running = true
|
|
backoff = BACKOFF_MIN_MS
|
|
await connect()
|
|
}
|
|
|
|
// Stop the client and cancel any pending reconnect. Called on shutdown and
|
|
// before a restart.
|
|
function stop() {
|
|
running = false
|
|
clearTimeout(reconnectTimer)
|
|
reconnectTimer = null
|
|
if (ws) {
|
|
try {
|
|
ws.removeAllListeners()
|
|
ws.close()
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
ws = null
|
|
}
|
|
state.connected = false
|
|
}
|
|
|
|
// Ingestion stats for the admin panel.
|
|
function getState() {
|
|
return { ...state, running }
|
|
}
|
|
|
|
module.exports = { start, stop, getState }
|