Usage — total data dir, Corestore, memory rings, warm HyperDB sample
CI / test (push) Successful in 1m2s
Release rolling / release (push) Has been cancelled

Hot memory — 15m → 24h (or custom points)
Warm memory — downsampled ring + samples-per-bucket
Disk history — auto-rotate by age (1 day → 1 year, forever, or custom days)
Soft disk budget + max warm points
Auto prune on a configurable hour interval
Dry-run prune / Prune now / Save to agent (admin role)
This commit is contained in:
Raven Scott
2026-07-19 12:55:13 -04:00
parent b5b2907496
commit c3b4240e47
26 changed files with 1610 additions and 17 deletions
+5
View File
@@ -24,6 +24,11 @@ PEARDATA_DEFAULT_ROLE=viewer
# PEARDATA_TIER0_POINTS=3600 # PEARDATA_TIER0_POINTS=3600
# PEARDATA_TIER1_POINTS=1440 # PEARDATA_TIER1_POINTS=1440
# PEARDATA_TIER1_EVERY=60 # PEARDATA_TIER1_EVERY=60
# Seeds for Data Manager (also persisted as data/retention.json after first boot):
# PEARDATA_WARM_RETENTION_MS=604800000
# PEARDATA_WARM_MAX_BYTES=0
# PEARDATA_WARM_MAX_POINTS=0
# PEARDATA_AUTO_PRUNE=1
# ── HyperDB (warm history + metadata + linked sync) ────────── # ── HyperDB (warm history + metadata + linked sync) ──────────
# PEARDATA_HYPERDB=1 # PEARDATA_HYPERDB=1
+15 -1
View File
@@ -31,6 +31,7 @@ import {
renderFleetCards, renderFleetCards,
} from './ui/fleet.js' } from './ui/fleet.js'
import { createLogsView } from './ui/logs.js' import { createLogsView } from './ui/logs.js'
import { createDataManager } from './ui/data-manager.js'
import { formatMib, formatKilobitsPerSec } from './shared/format.js' import { formatMib, formatKilobitsPerSec } from './shared/format.js'
const $ = (id) => document.getElementById(id) const $ = (id) => document.getElementById(id)
@@ -243,6 +244,13 @@ const logsView = createLogsView({
}, },
}) })
const dataManager = createDataManager({
$,
manager,
getRole: () => currentRole,
log: (msg) => log(msg),
})
function seriesMax() { function seriesMax() {
// Overview sparks follow Charts' active window when available (capped for density) // Overview sparks follow Charts' active window when available (capped for density)
try { try {
@@ -1003,6 +1011,7 @@ async function refreshMeta() {
const id = getClientIdentity() const id = getClientIdentity()
els.connMeta.textContent = `you ${id.publicKeyHex.slice(0, 12)}… · ${auth.role} · ${auth.authMode}` els.connMeta.textContent = `you ${id.publicKeyHex.slice(0, 12)}… · ${auth.role} · ${auth.authMode}`
logsView.syncSourceUi() logsView.syncSourceUi()
dataManager.syncGate()
renderPeers() renderPeers()
renderFleetStrip(fleet, manager.list()) renderFleetStrip(fleet, manager.list())
await populateExploreCharts() await populateExploreCharts()
@@ -1146,7 +1155,11 @@ function showView(name) {
} }
if (name === 'logs') logsView.enter() if (name === 'logs') logsView.enter()
if (name === 'fleet') loadFleetView() if (name === 'fleet') loadFleetView()
if (name === 'settings') syncSettingsUi() if (name === 'settings') {
syncSettingsUi()
const activeTab = document.querySelector('#settings-tabs .settings-tab.active')
if (activeTab?.dataset.settingsTab === 'data') dataManager.enter()
}
requestAnimationFrame(() => redrawAll()) requestAnimationFrame(() => redrawAll())
} }
@@ -1187,6 +1200,7 @@ document.querySelectorAll('#settings-tabs .settings-tab').forEach((tab) => {
document.querySelectorAll('.settings-panel').forEach((panel) => { document.querySelectorAll('.settings-panel').forEach((panel) => {
panel.classList.toggle('active', panel.dataset.settingsPanel === tab.dataset.settingsTab) panel.classList.toggle('active', panel.dataset.settingsPanel === tab.dataset.settingsTab)
}) })
if (tab.dataset.settingsTab === 'data') dataManager.enter()
}) })
}) })
+8 -2
View File
@@ -57,9 +57,15 @@ Permissions: directory `0700`. Do **not** commit `data/` or `.env`.
| Variable | Default | Description | | Variable | Default | Description |
|----------|---------|-------------| |----------|---------|-------------|
| `PEARDATA_SAMPLE_MS` | `1000` | Collector interval | | `PEARDATA_SAMPLE_MS` | `1000` | Collector interval |
| `PEARDATA_TIER0_POINTS` | `3600` | High-res ring size (~1h @ 1s) | | `PEARDATA_TIER0_POINTS` | `3600` | High-res ring size (~1h @ 1s); seeds Data Manager |
| `PEARDATA_TIER1_POINTS` | `1440` | Downsampled ring size | | `PEARDATA_TIER1_POINTS` | `1440` | Downsampled ring size; seeds Data Manager |
| `PEARDATA_TIER1_EVERY` | `60` | Samples per tier1 average (also HyperDB warm flush) | | `PEARDATA_TIER1_EVERY` | `60` | Samples per tier1 average (also HyperDB warm flush) |
| `PEARDATA_WARM_RETENTION_MS` | `604800000` (7d) | Soft age limit for HyperDB warm points (seed) |
| `PEARDATA_WARM_MAX_BYTES` | `0` | Soft corestore budget in bytes (`0` = unlimited) |
| `PEARDATA_WARM_MAX_POINTS` | `0` | Cap on warm points (`0` = unlimited) |
| `PEARDATA_AUTO_PRUNE` | on | `0` disables scheduled prune on boot |
Live retention is edited in the desktop **Settings → Data** tab (admin) and persisted to `$PEARDATA_DATA_DIR/retention.json`. Env values only seed the file on first boot.
--- ---
+10 -1
View File
@@ -108,9 +108,18 @@ Journal is **enabled by default** on Linux (`PEARDATA_JOURNAL=0` to disable). In
| Method | Role | Known jobs | | Method | Role | Known jobs |
|--------|------|------------| |--------|------|------------|
| `listJobs` | viewer | — | | `listJobs` | viewer | — |
| `runJob` | operator | `collectOnce`, `snapshot`, `gcBuffers` | | `runJob` | operator | `collectOnce`, `snapshot`, `gcBuffers` (prune) |
| `cancelJob` | operator | By job id | | `cancelJob` | operator | By job id |
### Data Manager (retention / storage)
| Method | Role | Notes |
|--------|------|-------|
| `getStorageInfo` | viewer | Disk usage, memory rings, warm sample stats, effective retention |
| `getRetentionConfig` | viewer | Current policy (`retention.json`) |
| `setRetentionConfig` | admin | Update hot/warm rings, age rotate, disk budget, auto-prune |
| `pruneHistory` | admin | `{ dryRun?, beforeMs? }` — delete warm points past retention |
### HyperDB / peer links ### HyperDB / peer links
| Method | Role | Notes | | Method | Role | Notes |
+1 -1
View File
@@ -45,7 +45,7 @@ curl -s http://127.0.0.1:18888/api/v3/health | jq
| GET | `/api/v3/db` | HyperDB public/discovery keys + collections | | GET | `/api/v3/db` | HyperDB public/discovery keys + collections |
| GET | `/api/v3/versions` | Agent / protocol / API versions | | GET | `/api/v3/versions` | Agent / protocol / API versions |
| GET | `/api/v3/me` | Anonymous REST identity note | | GET | `/api/v3/me` | Anonymous REST identity note |
| GET | `/api/v3/settings` | Runtime knobs | | GET | `/api/v3/settings` | Runtime knobs + retention + storage usage summary |
| GET | `/api/v3/config` | alias of settings | | GET | `/api/v3/config` | alias of settings |
| GET | `/health`, `/api/v1/health`, `/api/v3/health` | Aggregate health | | GET | `/health`, `/api/v1/health`, `/api/v3/health` | Aggregate health |
+2 -1
View File
@@ -47,7 +47,8 @@ Template rebrand, protocol, collector, memory store, anomalies, REST, desktop MV
| RPC peer-link + `getDbInfo` / REST `/api/v3/db` | Done | | RPC peer-link + `getDbInfo` / REST `/api/v3/db` | Done |
| Query fallback: memory miss / long window → HyperDB warm | Done | | Query fallback: memory miss / long window → HyperDB warm | Done |
| Hyperswarm mesh (`PEARDATA_SWARM=1`) | Done (opt-in) | | Hyperswarm mesh (`PEARDATA_SWARM=1`) | Done (opt-in) |
| Configurable retention knobs (tier sizes) | Done (`PEARDATA_TIER*`) | | Configurable retention knobs (tier sizes) | Done (`PEARDATA_TIER*` + Data Manager RPC/UI) |
| Data Manager (usage, age rotate, auto-prune) | Done (`retention.json`, Settings → Data) |
| Documented warm history across restart (operator M4) | Done (`store-hyperdb-fallback` + STORAGE doc) | | Documented warm history across restart (operator M4) | Done (`store-hyperdb-fallback` + STORAGE doc) |
| Export snapshot job → JSON / Prometheus push | Done (`exportSnapshot`, `prometheusPush`) | | Export snapshot job → JSON / Prometheus push | Done (`exportSnapshot`, `prometheusPush`) |
| Autobase multi-writer parents | Later (Phase 2c) | | Autobase multi-writer parents | Later (Phase 2c) |
+5
View File
@@ -129,6 +129,11 @@ Banner fields on agent start:
| `PEARDATA_SWARM` | off | `1` enables Hyperswarm `store.replicate` | | `PEARDATA_SWARM` | off | `1` enables Hyperswarm `store.replicate` |
| `PEARDATA_TIER1_EVERY` | `60` | Samples per warm bucket (~60s @ 1Hz) | | `PEARDATA_TIER1_EVERY` | `60` | Samples per warm bucket (~60s @ 1Hz) |
| `PEARDATA_TIER1_POINTS` | `1440` | In-memory tier1 ring (HyperDB keeps longer) | | `PEARDATA_TIER1_POINTS` | `1440` | In-memory tier1 ring (HyperDB keeps longer) |
| `PEARDATA_WARM_RETENTION_MS` | `7d` | Seeds Data Manager age rotate |
| `PEARDATA_WARM_MAX_BYTES` / `_POINTS` | `0` | Soft caps; `0` = unlimited |
| `PEARDATA_AUTO_PRUNE` | on | Scheduled prune after boot |
**Data Manager:** desktop **Settings → Data** (admin) calls `setRetentionConfig` / `pruneHistory`. Policy file: `$PEARDATA_DATA_DIR/retention.json`. Auto-prune runs on an interval and deletes HyperDB `@peardata/metric-point` rows older than the configured window (and enforces optional point/disk caps). Note: Hypercore deletes are logical; reclaiming physical disk may require compaction over time.
--- ---
+104 -1
View File
@@ -398,12 +398,13 @@
<div> <div>
<p class="dash-kicker">Preferences</p> <p class="dash-kicker">Preferences</p>
<h1 class="dash-title">Settings</h1> <h1 class="dash-title">Settings</h1>
<p class="page-subtitle">Appearance, charts, and desktop behaviour</p> <p class="page-subtitle">Appearance, charts, data retention, and desktop behaviour</p>
</div> </div>
</header> </header>
<div class="settings-subtabs" id="settings-tabs" role="tablist"> <div class="settings-subtabs" id="settings-tabs" role="tablist">
<button type="button" class="settings-tab active" data-settings-tab="appearance">Appearance</button> <button type="button" class="settings-tab active" data-settings-tab="appearance">Appearance</button>
<button type="button" class="settings-tab" data-settings-tab="charts">Charts</button> <button type="button" class="settings-tab" data-settings-tab="charts">Charts</button>
<button type="button" class="settings-tab" data-settings-tab="data">Data</button>
<button type="button" class="settings-tab" data-settings-tab="connections">Connections</button> <button type="button" class="settings-tab" data-settings-tab="connections">Connections</button>
<button type="button" class="settings-tab" data-settings-tab="notifications">Notifications</button> <button type="button" class="settings-tab" data-settings-tab="notifications">Notifications</button>
<button type="button" class="settings-tab" data-settings-tab="about">About</button> <button type="button" class="settings-tab" data-settings-tab="about">About</button>
@@ -452,6 +453,108 @@
</div> </div>
</div> </div>
<div class="settings-panel" data-settings-panel="data">
<div class="dash-card settings-section settings-section-wide">
<h3>Data Manager</h3>
<p class="hint">
Configure how much history the connected agent keeps in memory and on disk.
Changes apply live on the agent and persist under its data directory.
</p>
<p id="dm-admin-gate" class="dm-gate hint" role="status"></p>
<div class="dm-usage" aria-live="polite">
<div class="dm-usage-card">
<span class="stat-label">Total data</span>
<strong id="dm-usage-total"></strong>
</div>
<div class="dm-usage-card">
<span class="stat-label">Corestore</span>
<strong id="dm-usage-corestore"></strong>
</div>
<div class="dm-usage-card">
<span class="stat-label">Memory rings</span>
<strong id="dm-usage-memory"></strong>
</div>
<div class="dm-usage-card">
<span class="stat-label">Warm history</span>
<strong id="dm-usage-warm"></strong>
</div>
</div>
<p id="dm-usage-meta" class="muted settings-hint"></p>
<h3>Hot memory (1s)</h3>
<label>
Retention window
<select id="dm-hot-preset"></select>
</label>
<label>
Points (custom)
<input id="dm-hot-points" type="number" min="60" max="604800" step="60" value="3600" />
</label>
<p class="muted settings-hint">High-resolution samples kept in RAM. Larger windows use more memory.</p>
<h3>Warm memory (downsampled)</h3>
<label>
In-memory warm ring
<select id="dm-warm-mem-preset"></select>
</label>
<label>
Points (custom)
<input id="dm-warm-mem-points" type="number" min="10" max="525600" step="10" value="1440" />
</label>
<label>
Samples per warm bucket
<input id="dm-tier1-every" type="number" min="5" max="3600" step="1" value="60" />
</label>
<p class="muted settings-hint">Every N hot samples are averaged into one warm point (~60s at default).</p>
<h3>Disk history (HyperDB)</h3>
<label class="check-row">
<input type="checkbox" id="dm-warm-enabled" checked />
Auto-rotate warm history by age
</label>
<label>
Keep history for
<select id="dm-warm-preset"></select>
</label>
<label id="dm-warm-custom-wrap" class="hidden">
Custom days
<input id="dm-warm-custom-days" type="number" min="1" max="3650" step="1" value="7" />
</label>
<label>
Soft disk budget
<select id="dm-disk-preset"></select>
</label>
<label id="dm-disk-custom-wrap" class="hidden">
Custom budget (MB)
<input id="dm-disk-custom-mb" type="number" min="64" max="1048576" step="64" value="1024" />
</label>
<label>
Max warm points (0 = unlimited)
<input id="dm-warm-max-points" type="number" min="0" max="100000000" step="1000" value="0" />
</label>
<h3>Auto prune</h3>
<label class="check-row">
<input type="checkbox" id="dm-auto-prune" checked />
Run automatic prune on a schedule
</label>
<label>
Prune interval (hours)
<input id="dm-prune-interval-hours" type="number" min="1" max="168" step="1" value="1" />
</label>
<p id="dm-last-prune" class="muted settings-hint">Last prune: —</p>
<div class="dm-actions">
<button type="button" id="dm-refresh" class="ghost">Refresh</button>
<button type="button" id="dm-prune-dry" class="ghost">Dry-run prune</button>
<button type="button" id="dm-prune-now" class="ghost danger-btn">Prune now</button>
<button type="button" id="dm-save" class="primary">Save to agent</button>
</div>
<p id="dm-status" class="dm-status muted" role="status"></p>
</div>
</div>
<div class="settings-panel" data-settings-panel="connections"> <div class="settings-panel" data-settings-panel="connections">
<div class="dash-card settings-section"> <div class="dash-card settings-section">
<h3>Multi-agent</h3> <h3>Multi-agent</h3>
+2
View File
@@ -15,6 +15,8 @@ const MUTATING = new Set([
'exportSnapshot', 'exportSnapshot',
'linkPeer', 'linkPeer',
'unlinkPeer', 'unlinkPeer',
'setRetentionConfig',
'pruneHistory',
'handshake', 'handshake',
]) ])
+141
View File
@@ -282,6 +282,147 @@ export class PearDataModel extends ReadyResource {
})) }))
} }
/**
* Sample warm-point stats for Data Manager (bounded scan).
* @param {{ limit?: number }} [opts]
*/
async sampleMetricPointsStats(opts = {}) {
await this.ready()
const limit = Math.min(opts.limit || 50_000, 100_000)
const rows = await this.db
.find(
'@peardata/metric-point',
{
gte: { chart: '', ts: 0 },
lte: { chart: '\uffff', ts: Number.MAX_SAFE_INTEGER },
},
{ limit }
)
.toArray()
const charts = new Set()
let oldestTs = null
let newestTs = null
for (const r of rows) {
charts.add(r.chart)
if (oldestTs == null || r.ts < oldestTs) oldestTs = r.ts
if (newestTs == null || r.ts > newestTs) newestTs = r.ts
}
return {
pointsSampled: rows.length,
chartsSampled: charts.size,
oldestTs,
newestTs,
sampleCapped: rows.length >= limit,
}
}
/**
* Delete warm points with ts < beforeMs for the given charts.
* @param {number} beforeMs
* @param {{ charts?: string[], dryRun?: boolean, limitPerChart?: number }} [opts]
*/
async deleteMetricPointsBefore(beforeMs, opts = {}) {
await this.ready()
const cutoff = Number(beforeMs)
if (!Number.isFinite(cutoff) || cutoff <= 0) {
return { deleted: 0, scanned: 0 }
}
const dryRun = Boolean(opts.dryRun)
const limitPerChart = Math.min(opts.limitPerChart || 20_000, 50_000)
const charts = opts.charts?.length ? opts.charts : null
let deleted = 0
let scanned = 0
/** @type {string[]} */
let chartList = charts
if (!chartList) {
const sample = await this.db
.find(
'@peardata/metric-point',
{
gte: { chart: '', ts: 0 },
lte: { chart: '\uffff', ts: Number.MAX_SAFE_INTEGER },
},
{ limit: 50_000 }
)
.toArray()
chartList = [...new Set(sample.map((r) => r.chart))]
}
for (const chart of chartList) {
const rows = await this.db
.find(
'@peardata/metric-point',
{
gte: { chart: String(chart), ts: 0 },
lte: { chart: String(chart), ts: cutoff - 1 },
},
{ limit: limitPerChart }
)
.toArray()
scanned += rows.length
if (!rows.length) continue
if (dryRun) {
deleted += rows.length
continue
}
const tx = await this.db.exclusiveTransaction()
try {
for (const r of rows) {
await tx.delete('@peardata/metric-point', { chart: r.chart, ts: r.ts })
deleted++
}
await tx.flush()
} catch (err) {
await tx.close().catch(() => {})
throw err
}
}
return { deleted, scanned }
}
/**
* If total sampled points exceed maxPoints, delete oldest first.
* @param {number} maxPoints
* @param {{ dryRun?: boolean, charts?: string[] }} [opts]
*/
async enforceMetricPointCap(maxPoints, opts = {}) {
await this.ready()
const max = Math.floor(Number(maxPoints))
if (!Number.isFinite(max) || max <= 0) return { deleted: 0, scanned: 0 }
const limit = 100_000
const rows = await this.db
.find(
'@peardata/metric-point',
{
gte: { chart: '', ts: 0 },
lte: { chart: '\uffff', ts: Number.MAX_SAFE_INTEGER },
},
{ limit }
)
.toArray()
const scanned = rows.length
if (rows.length <= max) return { deleted: 0, scanned }
rows.sort((a, b) => a.ts - b.ts)
const over = rows.length - max
const victims = rows.slice(0, over)
if (opts.dryRun) return { deleted: victims.length, scanned }
const tx = await this.db.exclusiveTransaction()
try {
for (const r of victims) {
await tx.delete('@peardata/metric-point', { chart: r.chart, ts: r.ts })
}
await tx.flush()
} catch (err) {
await tx.close().catch(() => {})
throw err
}
return { deleted: victims.length, scanned }
}
// ── jobs ────────────────────────────────────────────────── // ── jobs ──────────────────────────────────────────────────
async putJob(job) { async putJob(job) {
+14
View File
@@ -43,6 +43,12 @@ import {
unsubscribeAnomalies, unsubscribeAnomalies,
} from '../services/subscriptions.js' } from '../services/subscriptions.js'
import { getJobs, knownJobNames } from '../services/jobs.js' import { getJobs, knownJobNames } from '../services/jobs.js'
import {
getRetentionConfig,
setRetentionConfig,
getStorageInfo,
pruneHistory,
} from '../services/retention.js'
import { formatAllMetrics } from '../rest/formatters.js' import { formatAllMetrics } from '../rest/formatters.js'
import { getDb } from '../db/index.js' import { getDb } from '../db/index.js'
import { attachLinkedPeer, isSwarmEnabled } from '../db/replicate.js' import { attachLinkedPeer, isSwarmEnabled } from '../db/replicate.js'
@@ -261,6 +267,14 @@ export function registerMonitorHandlers(session) {
session.respond('runJob', async (args) => getJobs().run(args.name, args.args || {})) session.respond('runJob', async (args) => getJobs().run(args.name, args.args || {}))
session.respond('cancelJob', async (args) => getJobs().cancel(args.id)) session.respond('cancelJob', async (args) => getJobs().cancel(args.id))
session.respond('getStorageInfo', async () => getStorageInfo())
session.respond('getRetentionConfig', async () => ({ config: getRetentionConfig() }))
session.respond('setRetentionConfig', async (args) => ({
success: true,
config: setRetentionConfig(args || {}),
}))
session.respond('pruneHistory', async (args) => pruneHistory(args || {}))
session.respond('mintInvite', async (args) => { session.respond('mintInvite', async (args) => {
const role = args.role || Roles.operator const role = args.role || Roles.operator
const ttlMs = args.ttlMs === undefined ? null : args.ttlMs const ttlMs = args.ttlMs === undefined ? null : args.ttlMs
+22 -1
View File
@@ -303,11 +303,32 @@ export async function handleRest(pathname, query) {
return json(snapshot) return json(snapshot)
} }
if (path === '/api/v3/settings' || path === '/api/v3/config') { if (path === '/api/v3/settings' || path === '/api/v3/config') {
let retention = null
let storage = null
try {
const { getRetentionConfig, getStorageInfo } = await import('../services/retention.js')
retention = getRetentionConfig()
storage = await getStorageInfo()
} catch {
// ignore
}
return json({ return json({
sample_ms: Number(process.env.PEARDATA_SAMPLE_MS) || 1000, sample_ms: Number(process.env.PEARDATA_SAMPLE_MS) || 1000,
rest_host: process.env.PEARDATA_REST_HOST || '127.0.0.1', rest_host: process.env.PEARDATA_REST_HOST || '127.0.0.1',
rest_port: Number(process.env.PEARDATA_REST_PORT) || 18888, rest_port: Number(process.env.PEARDATA_REST_PORT) || 18888,
tier0_points: Number(process.env.PEARDATA_TIER0_POINTS) || 3600, tier0_points: retention?.tier0Points ?? (Number(process.env.PEARDATA_TIER0_POINTS) || 3600),
tier1_points: retention?.tier1Points ?? (Number(process.env.PEARDATA_TIER1_POINTS) || 1440),
tier1_every: retention?.tier1Every ?? (Number(process.env.PEARDATA_TIER1_EVERY) || 60),
retention,
storage: storage
? {
dataDir: storage.dataDir,
dataDirBytes: storage.usage?.dataDirBytes,
corestoreBytes: storage.usage?.corestoreBytes,
memory: storage.memory,
warm: storage.warm,
}
: null,
}) })
} }
if (path === '/api/v3/stream_path') { if (path === '/api/v3/stream_path') {
+8
View File
@@ -24,6 +24,7 @@ import { startPipeline } from './pipeline.js'
import { startRestServer } from './rest/http-server.js' import { startRestServer } from './rest/http-server.js'
import { getCollector } from './services/collector.js' import { getCollector } from './services/collector.js'
import { flushWarmPending } from './services/warm-flush.js' import { flushWarmPending } from './services/warm-flush.js'
import { loadRetentionConfig, startAutoPrune, stopAutoPrune } from './services/retention.js'
import logger from './utils/logger.js' import logger from './utils/logger.js'
const bootStarted = Date.now() const bootStarted = Date.now()
@@ -68,7 +69,9 @@ if (db) {
} }
await startReplication() await startReplication()
loadRetentionConfig()
startPipeline() startPipeline()
startAutoPrune()
const restServer = startRestServer() const restServer = startRestServer()
let restTunnelInfo = null let restTunnelInfo = null
@@ -206,6 +209,11 @@ async function shutdown() {
} catch { } catch {
// ignore // ignore
} }
try {
stopAutoPrune()
} catch {
// ignore
}
try { try {
await flushWarmPending() await flushWarmPending()
} catch { } catch {
+8 -3
View File
@@ -29,9 +29,14 @@ const JOB_HANDLERS = {
minPoints: args.minPoints != null ? Number(args.minPoints) : undefined, minPoints: args.minPoints != null ? Number(args.minPoints) : undefined,
}) })
}, },
gcBuffers: async () => { gcBuffers: async (args = {}) => {
// ring buffers self-trim; placeholder for future disk GC const { pruneHistory } = await import('./retention.js')
return { ok: true } const result = await pruneHistory({
dryRun: Boolean(args.dryRun),
beforeMs: args.beforeMs != null ? Number(args.beforeMs) : undefined,
force: true,
})
return { ok: result.success !== false, ...result }
}, },
} }
+370
View File
@@ -0,0 +1,370 @@
/**
* Agent data manager — retention policy, auto-prune, storage usage.
*
* Persists under $PEARDATA_DATA_DIR/retention.json (survives restart).
* Env PEARDATA_TIER* seeds defaults on first boot when no file exists.
*/
import fs from 'fs'
import path from 'path'
import {
defaultRetentionConfig,
normalizeRetentionConfig,
retentionCutoffMs,
} from '../../shared/retention.js'
import { getDataDir, getDb, isHyperDbEnabled } from '../db/index.js'
import { getStore } from './store.js'
import { CHART_BY_ID } from '../../shared/metrics.js'
import logger from '../utils/logger.js'
const log = logger.child('retention')
/** @type {import('../../shared/retention.js').RetentionConfig} */
let config = defaultRetentionConfig()
/** @type {ReturnType<typeof setInterval>|null} */
let pruneTimer = null
let pruneRunning = false
function envInt(name, fallback) {
const n = Number(process.env[name])
return Number.isFinite(n) && n > 0 ? n : fallback
}
function configPath() {
return path.join(getDataDir(), 'retention.json')
}
function seedFromEnv() {
const d = defaultRetentionConfig()
const warmMsEnv = process.env.PEARDATA_WARM_RETENTION_MS
const warmRetentionMs = warmMsEnv != null ? envInt('PEARDATA_WARM_RETENTION_MS', d.warmRetentionMs) : d.warmRetentionMs
return normalizeRetentionConfig({
...d,
tier0Points: envInt('PEARDATA_TIER0_POINTS', d.tier0Points),
tier1Points: envInt('PEARDATA_TIER1_POINTS', d.tier1Points),
tier1Every: envInt('PEARDATA_TIER1_EVERY', d.tier1Every),
warmRetentionMs,
warmRetentionPreset: warmMsEnv != null ? 'custom' : d.warmRetentionPreset,
warmMaxBytes: Number(process.env.PEARDATA_WARM_MAX_BYTES) > 0 ? Number(process.env.PEARDATA_WARM_MAX_BYTES) : 0,
warmMaxPoints: Number(process.env.PEARDATA_WARM_MAX_POINTS) > 0 ? Number(process.env.PEARDATA_WARM_MAX_POINTS) : 0,
autoPruneEnabled: process.env.PEARDATA_AUTO_PRUNE === '0' ? false : d.autoPruneEnabled,
})
}
function save() {
try {
const dir = getDataDir()
if (!fs.existsSync(dir)) fs.mkdirSync(dir, { recursive: true, mode: 0o700 })
fs.writeFileSync(configPath(), JSON.stringify(config, null, 2), { mode: 0o600 })
} catch (err) {
log.warn('Failed to save retention config', { error: err.message })
}
}
/**
* Load persisted config (or env seeds) and apply to MetricStore.
*/
export function loadRetentionConfig() {
try {
const p = configPath()
if (fs.existsSync(p)) {
const raw = JSON.parse(fs.readFileSync(p, 'utf8'))
config = normalizeRetentionConfig(raw)
log.info('Loaded retention config', {
tier0: config.tier0Points,
tier1: config.tier1Points,
warmPreset: config.warmRetentionPreset,
autoPrune: config.autoPruneEnabled,
})
} else {
config = seedFromEnv()
save()
log.info('Seeded retention config from defaults/env')
}
} catch (err) {
log.warn('Failed to load retention config; using defaults', { error: err.message })
config = seedFromEnv()
}
applyRetentionToStore()
return getRetentionConfig()
}
export function getRetentionConfig() {
return { ...config }
}
/**
* @param {Partial<import('../../shared/retention.js').RetentionConfig>} patch
*/
export function setRetentionConfig(patch = {}) {
const next = normalizeRetentionConfig({
...config,
...patch,
// preserve prune stats unless explicitly cleared
lastPruneAt: patch.lastPruneAt !== undefined ? patch.lastPruneAt : config.lastPruneAt,
lastPruneDeleted:
patch.lastPruneDeleted !== undefined ? patch.lastPruneDeleted : config.lastPruneDeleted,
lastPruneBytesFreed:
patch.lastPruneBytesFreed !== undefined
? patch.lastPruneBytesFreed
: config.lastPruneBytesFreed,
})
config = next
save()
applyRetentionToStore()
restartAutoPrune()
return getRetentionConfig()
}
export function applyRetentionToStore() {
try {
getStore().setRetention({
tier0Max: config.tier0Points,
tier1Max: config.tier1Points,
tier1Every: config.tier1Every,
})
} catch (err) {
log.warn('Failed to apply retention to store', { error: err.message })
}
}
/**
* Recursive directory size (bytes).
* @param {string} dir
* @param {number} [budget] max entries walked
*/
export function directorySizeBytes(dir, budget = 250_000) {
let total = 0
let walked = 0
/** @type {string[]} */
const stack = [dir]
while (stack.length && walked < budget) {
const cur = stack.pop()
let entries
try {
entries = fs.readdirSync(cur, { withFileTypes: true })
} catch {
continue
}
for (const ent of entries) {
walked++
if (walked >= budget) break
const full = path.join(cur, ent.name)
try {
if (ent.isDirectory()) stack.push(full)
else if (ent.isFile() || ent.isSymbolicLink()) {
const st = fs.statSync(full)
total += st.size || 0
}
} catch {
// ignore races
}
}
}
return { bytes: total, truncated: walked >= budget, walked }
}
/**
* Snapshot of agent storage usage + effective retention.
*/
export async function getStorageInfo() {
const dataDir = getDataDir()
const corestorePath = path.join(dataDir, 'corestore')
const mem = getStore().memoryStats()
let dataDirSize = { bytes: 0, truncated: false, walked: 0 }
let corestoreSize = { bytes: 0, truncated: false, walked: 0 }
try {
if (fs.existsSync(dataDir)) dataDirSize = directorySizeBytes(dataDir)
} catch {
// ignore
}
try {
if (fs.existsSync(corestorePath)) corestoreSize = directorySizeBytes(corestorePath)
} catch {
// ignore
}
const db = getDb()
let warm = {
enabled: isHyperDbEnabled(),
open: Boolean(db),
pointsSampled: 0,
chartsSampled: 0,
oldestTs: null,
newestTs: null,
}
if (db) {
try {
const sample = await db.sampleMetricPointsStats({ limit: 50_000 })
warm = { ...warm, ...sample }
} catch (err) {
warm.error = err.message
}
}
const otherFiles = []
for (const name of ['retention.json', 'peer-policy.json', 'audit.log']) {
const p = path.join(dataDir, name)
try {
if (fs.existsSync(p)) {
const st = fs.statSync(p)
otherFiles.push({ name, bytes: st.size })
}
} catch {
// ignore
}
}
return {
dataDir,
hyperDbEnabled: isHyperDbEnabled(),
usage: {
dataDirBytes: dataDirSize.bytes,
corestoreBytes: corestoreSize.bytes,
dataDirTruncated: dataDirSize.truncated,
corestoreTruncated: corestoreSize.truncated,
otherFiles,
},
memory: mem,
warm,
retention: getRetentionConfig(),
measuredAt: Date.now(),
}
}
/**
* Prune warm HyperDB history by age / max points / soft disk budget.
* @param {{ dryRun?: boolean, beforeMs?: number|null, force?: boolean }} [opts]
*/
export async function pruneHistory(opts = {}) {
if (pruneRunning) return { success: false, error: 'prune already running' }
pruneRunning = true
const dryRun = Boolean(opts.dryRun)
const started = Date.now()
try {
const db = getDb()
if (!db) {
return {
success: true,
dryRun,
deleted: 0,
skipped: true,
reason: 'hyperdb disabled or not open',
elapsedMs: Date.now() - started,
}
}
const beforeMs =
opts.beforeMs != null
? Number(opts.beforeMs)
: retentionCutoffMs(config, started)
const chartIds = new Set([...CHART_BY_ID.keys(), ...getStore().series.keys()])
let deleted = 0
let scanned = 0
if (beforeMs != null && Number.isFinite(beforeMs) && beforeMs > 0) {
const age = await db.deleteMetricPointsBefore(beforeMs, {
charts: [...chartIds],
dryRun,
limitPerChart: 20_000,
})
deleted += age.deleted
scanned += age.scanned
}
// Cap total warm points (delete oldest across charts)
if (config.warmMaxPoints > 0) {
const cap = await db.enforceMetricPointCap(config.warmMaxPoints, {
dryRun,
charts: [...chartIds],
})
deleted += cap.deleted
scanned += cap.scanned
}
// Soft disk budget — if corestore over budget, age-prune more aggressively
if (config.warmMaxBytes > 0 && !dryRun) {
const { bytes } = directorySizeBytes(path.join(getDataDir(), 'corestore'))
if (bytes > config.warmMaxBytes) {
// Drop another 25% of retention window (or 1 day floor)
const extraMs = Math.max(
24 * 60 * 60 * 1000,
Math.floor((config.warmRetentionMs || 7 * 24 * 60 * 60 * 1000) * 0.25)
)
const tighter = (beforeMs != null ? beforeMs : started) - extraMs
if (tighter > 0) {
const extra = await db.deleteMetricPointsBefore(tighter, {
charts: [...chartIds],
dryRun: false,
limitPerChart: 20_000,
})
deleted += extra.deleted
scanned += extra.scanned
}
}
}
// Memory rings already self-trim via setRetention; force trim
getStore().trimToRetention()
const approxFreed = deleted * 180 // rough JSON row estimate
if (!dryRun && deleted > 0) {
config = normalizeRetentionConfig({
...config,
lastPruneAt: started,
lastPruneDeleted: deleted,
lastPruneBytesFreed: approxFreed,
})
save()
}
log.info(dryRun ? 'Dry-run prune' : 'Pruned warm history', {
deleted,
scanned,
beforeMs,
})
return {
success: true,
dryRun,
deleted,
scanned,
beforeMs: beforeMs ?? null,
approxBytesFreed: approxFreed,
elapsedMs: Date.now() - started,
retention: getRetentionConfig(),
}
} catch (err) {
log.warn('Prune failed', { error: err.message })
return { success: false, error: err.message, dryRun, elapsedMs: Date.now() - started }
} finally {
pruneRunning = false
}
}
export function startAutoPrune() {
stopAutoPrune()
if (!config.autoPruneEnabled) return
const interval = config.autoPruneIntervalMs || 60 * 60 * 1000
pruneTimer = setInterval(() => {
pruneHistory({ force: true }).catch(() => {})
}, interval)
if (typeof pruneTimer.unref === 'function') pruneTimer.unref()
// Kick once shortly after boot (let HyperDB settle)
setTimeout(() => {
pruneHistory({ force: true }).catch(() => {})
}, 30_000).unref?.()
log.info('Auto-prune scheduled', { intervalMs: interval })
}
export function stopAutoPrune() {
if (pruneTimer) {
clearInterval(pruneTimer)
pruneTimer = null
}
}
function restartAutoPrune() {
startAutoPrune()
}
+65
View File
@@ -26,6 +26,71 @@ export class MetricStore extends EventEmitter {
this.series = new Map() this.series = new Map()
} }
/**
* Live-update retention knobs (from Data Manager) and trim rings.
* @param {{ tier0Max?: number, tier1Max?: number, tier1Every?: number }} opts
*/
setRetention(opts = {}) {
if (opts.tier0Max != null && Number.isFinite(Number(opts.tier0Max)) && opts.tier0Max > 0) {
this.tier0Max = Math.floor(Number(opts.tier0Max))
}
if (opts.tier1Max != null && Number.isFinite(Number(opts.tier1Max)) && opts.tier1Max > 0) {
this.tier1Max = Math.floor(Number(opts.tier1Max))
}
if (opts.tier1Every != null && Number.isFinite(Number(opts.tier1Every)) && opts.tier1Every > 0) {
this.tier1Every = Math.floor(Number(opts.tier1Every))
}
this.trimToRetention()
return {
tier0Max: this.tier0Max,
tier1Max: this.tier1Max,
tier1Every: this.tier1Every,
}
}
/** Trim all series to current tier maxima. */
trimToRetention() {
for (const entry of this.series.values()) {
if (entry.points.length > this.tier0Max) {
entry.points.splice(0, entry.points.length - this.tier0Max)
}
if (entry.tier1.length > this.tier1Max) {
entry.tier1.splice(0, entry.tier1.length - this.tier1Max)
}
}
}
/** Approximate in-memory footprint for Data Manager. */
memoryStats() {
let tier0Points = 0
let tier1Points = 0
let charts = 0
let oldestTs = null
let newestTs = null
for (const entry of this.series.values()) {
charts++
tier0Points += entry.points.length
tier1Points += entry.tier1.length
const first = entry.points[0]?.ts
const last = entry.points[entry.points.length - 1]?.ts
if (first != null && (oldestTs == null || first < oldestTs)) oldestTs = first
if (last != null && (newestTs == null || last > newestTs)) newestTs = last
}
// Rough: ~48 bytes overhead + ~8 per numeric dim (assume ~6 dims) + JSON-ish
const approxBytes = charts * 256 + (tier0Points + tier1Points) * 96
return {
charts,
tier0Points,
tier1Points,
tier0Max: this.tier0Max,
tier1Max: this.tier1Max,
tier1Every: this.tier1Every,
approxBytes,
oldestTs,
newestTs,
}
}
/** /**
* @param {Array<{ chart: string, context: string, ts: number, values: Record<string, number|null> }>} batch * @param {Array<{ chart: string, context: string, ts: number, values: Record<string, number|null> }>} batch
*/ */
+13
View File
@@ -65,3 +65,16 @@ export function formatKilobitsPerSec(kbps) {
const { n, unit } = scaleUnit(bps, ['b/s', 'kb/s', 'Mb/s', 'Gb/s', 'Tb/s'], 1000) const { n, unit } = scaleUnit(bps, ['b/s', 'kb/s', 'Mb/s', 'Gb/s', 'Tb/s'], 1000)
return `${trimFixed(n, autoDigits(n))} ${unit}` return `${trimFixed(n, autoDigits(n))} ${unit}`
} }
/**
* Format raw bytes → B / KiB / MiB / GiB / TiB / PiB.
* @param {number|null|undefined} bytes
* @returns {string}
*/
export function formatBytes(bytes) {
if (bytes == null || !Number.isFinite(Number(bytes))) return '—'
const n0 = Number(bytes)
if (n0 === 0) return '0 B'
const { n, unit } = scaleUnit(n0, ['B', 'KiB', 'MiB', 'GiB', 'TiB', 'PiB'], 1024)
return `${trimFixed(n, autoDigits(n))} ${unit}`
}
+6
View File
@@ -91,6 +91,12 @@ export const MethodRoles = Object.freeze({
linkPeer: Roles.admin, linkPeer: Roles.admin,
unlinkPeer: Roles.admin, unlinkPeer: Roles.admin,
// Data Manager — retention + storage
getStorageInfo: Roles.viewer,
getRetentionConfig: Roles.viewer,
setRetentionConfig: Roles.admin,
pruneHistory: Roles.admin,
// parent / fleet aggregation // parent / fleet aggregation
getFleetHealth: Roles.viewer, getFleetHealth: Roles.viewer,
listChildPeers: Roles.viewer, listChildPeers: Roles.viewer,
+180
View File
@@ -0,0 +1,180 @@
/**
* Retention presets + config normalization (client + agent).
*/
/** @typedef {{
* tier0Points: number,
* tier1Points: number,
* tier1Every: number,
* warmRetentionEnabled: boolean,
* warmRetentionPreset: string,
* warmRetentionMs: number,
* warmMaxBytes: number,
* warmMaxPoints: number,
* autoPruneEnabled: boolean,
* autoPruneIntervalMs: number,
* lastPruneAt: number|null,
* lastPruneDeleted: number,
* lastPruneBytesFreed: number,
* }} RetentionConfig */
export const RETENTION_PRESETS = Object.freeze({
'1d': { id: '1d', label: '1 day', ms: 1 * 24 * 60 * 60 * 1000 },
'1w': { id: '1w', label: '1 week', ms: 7 * 24 * 60 * 60 * 1000 },
'2w': { id: '2w', label: '2 weeks', ms: 14 * 24 * 60 * 60 * 1000 },
'1m': { id: '1m', label: '1 month', ms: 30 * 24 * 60 * 60 * 1000 },
'3m': { id: '3m', label: '3 months', ms: 90 * 24 * 60 * 60 * 1000 },
'6m': { id: '6m', label: '6 months', ms: 180 * 24 * 60 * 60 * 1000 },
'1y': { id: '1y', label: '1 year', ms: 365 * 24 * 60 * 60 * 1000 },
forever: { id: 'forever', label: 'Forever', ms: 0 },
custom: { id: 'custom', label: 'Custom', ms: null },
})
/** Hot (1s) ring size presets → points. */
export const HOT_PRESETS = Object.freeze({
'15m': { id: '15m', label: '15 minutes', points: 900 },
'1h': { id: '1h', label: '1 hour', points: 3600 },
'6h': { id: '6h', label: '6 hours', points: 21600 },
'12h': { id: '12h', label: '12 hours', points: 43200 },
'24h': { id: '24h', label: '24 hours', points: 86400 },
custom: { id: 'custom', label: 'Custom', points: null },
})
/** In-memory warm ring presets (@ ~1m resolution). */
export const WARM_MEM_PRESETS = Object.freeze({
'1h': { id: '1h', label: '1 hour', points: 60 },
'6h': { id: '6h', label: '6 hours', points: 360 },
'24h': { id: '24h', label: '24 hours', points: 1440 },
'7d': { id: '7d', label: '7 days', points: 10080 },
custom: { id: 'custom', label: 'Custom', points: null },
})
/** Soft disk budget presets (bytes). 0 = unlimited. */
export const DISK_BUDGET_PRESETS = Object.freeze({
unlimited: { id: 'unlimited', label: 'Unlimited', bytes: 0 },
'256mb': { id: '256mb', label: '256 MB', bytes: 256 * 1024 * 1024 },
'512mb': { id: '512mb', label: '512 MB', bytes: 512 * 1024 * 1024 },
'1gb': { id: '1gb', label: '1 GB', bytes: 1024 * 1024 * 1024 },
'5gb': { id: '5gb', label: '5 GB', bytes: 5 * 1024 * 1024 * 1024 },
'20gb': { id: '20gb', label: '20 GB', bytes: 20 * 1024 * 1024 * 1024 },
custom: { id: 'custom', label: 'Custom', bytes: null },
})
/**
* @returns {RetentionConfig}
*/
export function defaultRetentionConfig() {
return {
tier0Points: 3600,
tier1Points: 1440,
tier1Every: 60,
warmRetentionEnabled: true,
warmRetentionPreset: '1w',
warmRetentionMs: RETENTION_PRESETS['1w'].ms,
warmMaxBytes: 0,
warmMaxPoints: 0,
autoPruneEnabled: true,
autoPruneIntervalMs: 60 * 60 * 1000,
lastPruneAt: null,
lastPruneDeleted: 0,
lastPruneBytesFreed: 0,
}
}
/**
* @param {Partial<RetentionConfig>|null|undefined} raw
* @returns {RetentionConfig}
*/
export function normalizeRetentionConfig(raw = {}) {
const d = defaultRetentionConfig()
const src = raw && typeof raw === 'object' ? raw : {}
const tier0Points = clampInt(src.tier0Points ?? d.tier0Points, 60, 604_800, d.tier0Points)
const tier1Points = clampInt(src.tier1Points ?? d.tier1Points, 10, 525_600, d.tier1Points)
const tier1Every = clampInt(src.tier1Every ?? d.tier1Every, 5, 3600, d.tier1Every)
let warmRetentionPreset = String(src.warmRetentionPreset || d.warmRetentionPreset)
if (!RETENTION_PRESETS[warmRetentionPreset]) warmRetentionPreset = 'custom'
let warmRetentionMs = Number(src.warmRetentionMs)
if (!Number.isFinite(warmRetentionMs) || warmRetentionMs < 0) warmRetentionMs = d.warmRetentionMs
if (warmRetentionPreset !== 'custom' && warmRetentionPreset !== 'forever') {
warmRetentionMs = RETENTION_PRESETS[warmRetentionPreset].ms
} else if (warmRetentionPreset === 'forever') {
warmRetentionMs = 0
}
const warmRetentionEnabled =
src.warmRetentionEnabled == null ? d.warmRetentionEnabled : Boolean(src.warmRetentionEnabled)
const warmMaxBytes = clampInt(src.warmMaxBytes ?? d.warmMaxBytes, 0, Number.MAX_SAFE_INTEGER, 0)
const warmMaxPoints = clampInt(src.warmMaxPoints ?? d.warmMaxPoints, 0, Number.MAX_SAFE_INTEGER, 0)
const autoPruneEnabled =
src.autoPruneEnabled == null ? d.autoPruneEnabled : Boolean(src.autoPruneEnabled)
const autoPruneIntervalMs = clampInt(
src.autoPruneIntervalMs ?? d.autoPruneIntervalMs,
60_000,
7 * 24 * 60 * 60 * 1000,
d.autoPruneIntervalMs
)
return {
tier0Points,
tier1Points,
tier1Every,
warmRetentionEnabled,
warmRetentionPreset,
warmRetentionMs,
warmMaxBytes,
warmMaxPoints,
autoPruneEnabled,
autoPruneIntervalMs,
lastPruneAt:
src.lastPruneAt != null && Number.isFinite(Number(src.lastPruneAt))
? Number(src.lastPruneAt)
: null,
lastPruneDeleted: clampInt(src.lastPruneDeleted ?? 0, 0, Number.MAX_SAFE_INTEGER, 0),
lastPruneBytesFreed: clampInt(src.lastPruneBytesFreed ?? 0, 0, Number.MAX_SAFE_INTEGER, 0),
}
}
/**
* Resolve cutoff timestamp for age-based prune (null = skip age prune).
* @param {RetentionConfig} cfg
* @param {number} [now]
*/
export function retentionCutoffMs(cfg, now = Date.now()) {
if (!cfg?.warmRetentionEnabled) return null
if (!cfg.warmRetentionMs || cfg.warmRetentionMs <= 0) return null
return now - cfg.warmRetentionMs
}
/**
* @param {number} points
* @param {Record<string, { points: number|null }>} presets
*/
export function matchPointsPreset(points, presets) {
for (const [id, p] of Object.entries(presets)) {
if (id === 'custom') continue
if (p.points === points) return id
}
return 'custom'
}
/**
* @param {number} bytes
*/
export function matchBytesPreset(bytes) {
for (const [id, p] of Object.entries(DISK_BUDGET_PRESETS)) {
if (id === 'custom') continue
if (p.bytes === bytes) return id
}
return 'custom'
}
function clampInt(v, min, max, fallback) {
const n = Math.floor(Number(v))
if (!Number.isFinite(n)) return fallback
return Math.min(max, Math.max(min, n))
}
+19
View File
@@ -164,8 +164,27 @@ export function validateMethodArgs(method, args = {}) {
case 'listPeerLinks': case 'listPeerLinks':
case 'getFleetHealth': case 'getFleetHealth':
case 'listChildPeers': case 'listChildPeers':
case 'getStorageInfo':
case 'getRetentionConfig':
return { ok: true, args } return { ok: true, args }
case 'setRetentionConfig': {
const patch = args && typeof args === 'object' ? { ...args } : {}
delete patch.success
return { ok: true, args: patch }
}
case 'pruneHistory': {
return {
ok: true,
args: {
dryRun: Boolean(args?.dryRun),
beforeMs: args?.beforeMs != null ? Number(args.beforeMs) : undefined,
force: Boolean(args?.force),
},
}
}
case 'linkPeer': { case 'linkPeer': {
const remotePublicKey = String(args.remotePublicKey || args.peerId || '').toLowerCase() const remotePublicKey = String(args.remotePublicKey || args.peerId || '').toLowerCase()
if (!/^[0-9a-f]{64}$/.test(remotePublicKey)) { if (!/^[0-9a-f]{64}$/.test(remotePublicKey)) {
+8 -1
View File
@@ -1,5 +1,5 @@
import test from 'brittle' import test from 'brittle'
import { formatMib, formatKilobitsPerSec } from '../shared/format.js' import { formatMib, formatKilobitsPerSec, formatBytes } from '../shared/format.js'
test('formatMib scales through binary units', async (t) => { test('formatMib scales through binary units', async (t) => {
t.is(formatMib(null), '—') t.is(formatMib(null), '—')
@@ -21,3 +21,10 @@ test('formatKilobitsPerSec scales through SI bit rates', async (t) => {
t.is(formatKilobitsPerSec(2_500_000), '2.5 Gb/s') t.is(formatKilobitsPerSec(2_500_000), '2.5 Gb/s')
t.is(formatKilobitsPerSec(3_000_000_000), '3 Tb/s') t.is(formatKilobitsPerSec(3_000_000_000), '3 Tb/s')
}) })
test('formatBytes scales binary units', async (t) => {
t.is(formatBytes(null), '—')
t.is(formatBytes(0), '0 B')
t.is(formatBytes(512), '512 B')
t.is(formatBytes(2048), '2 KiB')
})
+4
View File
@@ -41,6 +41,10 @@ test('method roles cover monitoring surface', (t) => {
'listChildPeers', 'listChildPeers',
'getWeights', 'getWeights',
'queryLogs', 'queryLogs',
'getStorageInfo',
'getRetentionConfig',
'setRetentionConfig',
'pruneHistory',
]) { ]) {
t.ok(MethodRoles[m], m) t.ok(MethodRoles[m], m)
t.is(Methods[m], m) t.is(Methods[m], m)
+144
View File
@@ -0,0 +1,144 @@
import test from 'brittle'
import fs from 'fs'
import path from 'path'
import { fileURLToPath } from 'url'
import {
defaultRetentionConfig,
normalizeRetentionConfig,
retentionCutoffMs,
RETENTION_PRESETS,
matchPointsPreset,
HOT_PRESETS,
} from '../shared/retention.js'
import { formatBytes } from '../shared/format.js'
import { MetricStore } from '../server/services/store.js'
import { PearDataModel } from '../server/db/model.js'
import Corestore from 'corestore'
import { validateMethodArgs } from '../shared/schema.js'
import { MethodRoles, Methods, Roles } from '../shared/protocol.js'
const __dirname = path.dirname(fileURLToPath(import.meta.url))
const TMP = path.join(__dirname, '..', 'tmp-retention-test')
test('normalizeRetentionConfig applies presets', (t) => {
const cfg = normalizeRetentionConfig({
warmRetentionPreset: '1m',
tier0Points: 900,
})
t.is(cfg.warmRetentionMs, RETENTION_PRESETS['1m'].ms)
t.is(cfg.tier0Points, 900)
t.is(matchPointsPreset(900, HOT_PRESETS), '15m')
})
test('retentionCutoffMs respects forever / disabled', (t) => {
const now = 1_700_000_000_000
t.is(retentionCutoffMs({ warmRetentionEnabled: false, warmRetentionMs: 1000 }, now), null)
t.is(retentionCutoffMs({ warmRetentionEnabled: true, warmRetentionMs: 0 }, now), null)
t.is(
retentionCutoffMs({ warmRetentionEnabled: true, warmRetentionMs: 1000 }, now),
now - 1000
)
})
test('formatBytes scales', (t) => {
t.is(formatBytes(0), '0 B')
t.is(formatBytes(1024), '1 KiB')
t.is(formatBytes(1536 * 1024 * 1024), '1.5 GiB')
})
test('MetricStore setRetention trims rings', (t) => {
const store = new MetricStore()
const ts = Date.now()
for (let i = 0; i < 100; i++) {
store.ingest([
{
chart: 'system.cpu',
context: 'system.cpu',
ts: ts + i * 1000,
values: { user: 1, idle: 99 },
},
])
}
t.ok(store.series.get('system.cpu').points.length >= 100)
store.setRetention({ tier0Max: 10, tier1Max: 5, tier1Every: 60 })
t.is(store.tier0Max, 10)
t.is(store.series.get('system.cpu').points.length, 10)
const mem = store.memoryStats()
t.is(mem.tier0Max, 10)
t.ok(mem.tier0Points <= 10)
})
test('protocol exposes data manager methods', (t) => {
for (const m of ['getStorageInfo', 'getRetentionConfig', 'setRetentionConfig', 'pruneHistory']) {
t.ok(MethodRoles[m], m)
t.is(Methods[m], m)
}
t.is(MethodRoles.setRetentionConfig, Roles.admin)
t.is(MethodRoles.pruneHistory, Roles.admin)
t.ok(validateMethodArgs('getStorageInfo', {}).ok)
t.ok(validateMethodArgs('setRetentionConfig', { tier0Points: 900 }).ok)
t.ok(validateMethodArgs('pruneHistory', { dryRun: true }).ok)
})
test('hyperdb deleteMetricPointsBefore removes old rows', async (t) => {
fs.rmSync(TMP, { recursive: true, force: true })
const store = new Corestore(path.join(TMP, 'corestore'))
await store.ready()
const core = store.get({ name: 'peardata-meta' })
const model = new PearDataModel(core, { autoUpdate: true })
await model.ready()
try {
const now = Date.now()
await model.putMetricPoints([
{
chart: 'system.cpu',
context: 'system.cpu',
ts: now - 10_000,
values: { user: 1 },
tier: 1,
},
{
chart: 'system.cpu',
context: 'system.cpu',
ts: now - 1000,
values: { user: 2 },
tier: 1,
},
{
chart: 'system.cpu',
context: 'system.cpu',
ts: now,
values: { user: 3 },
tier: 1,
},
])
const dry = await model.deleteMetricPointsBefore(now - 2000, {
charts: ['system.cpu'],
dryRun: true,
})
t.is(dry.deleted, 1)
const del = await model.deleteMetricPointsBefore(now - 2000, {
charts: ['system.cpu'],
dryRun: false,
})
t.is(del.deleted, 1)
const rows = await model.queryMetricPoints({
chart: 'system.cpu',
afterMs: now - 20_000,
beforeMs: now + 1000,
})
t.is(rows.length, 2)
t.is(rows[0].values.user, 2)
} finally {
await model.close().catch(() => {})
await store.close().catch(() => {})
fs.rmSync(TMP, { recursive: true, force: true })
}
})
test('defaultRetentionConfig has sane defaults', (t) => {
const d = defaultRetentionConfig()
t.is(d.warmRetentionPreset, '1w')
t.ok(d.autoPruneEnabled)
t.ok(d.tier0Points > 0)
})
+383
View File
@@ -0,0 +1,383 @@
/**
* Settings → Data Manager — agent retention, auto-prune, storage usage.
*/
import { Methods, Roles, roleAllows } from '../shared/protocol.js'
import { formatBytes } from '../shared/format.js'
import {
RETENTION_PRESETS,
HOT_PRESETS,
WARM_MEM_PRESETS,
DISK_BUDGET_PRESETS,
matchPointsPreset,
matchBytesPreset,
} from '../shared/retention.js'
/**
* @param {{
* $: (id: string) => HTMLElement|null,
* manager: { request: (m: string, a?: object) => Promise<any>, active: any },
* getRole: () => string,
* log: (msg: string) => void,
* }} deps
*/
export function createDataManager(deps) {
const { $, manager, getRole, log } = deps
/** @type {import('../shared/retention.js').RetentionConfig|null} */
let config = null
/** @type {any} */
let storage = null
let busy = false
const els = {
usageTotal: $('dm-usage-total'),
usageCorestore: $('dm-usage-corestore'),
usageMemory: $('dm-usage-memory'),
usageWarm: $('dm-usage-warm'),
usageMeta: $('dm-usage-meta'),
hotPreset: $('dm-hot-preset'),
hotPoints: $('dm-hot-points'),
warmMemPreset: $('dm-warm-mem-preset'),
warmMemPoints: $('dm-warm-mem-points'),
tier1Every: $('dm-tier1-every'),
warmEnabled: $('dm-warm-enabled'),
warmPreset: $('dm-warm-preset'),
warmCustomDays: $('dm-warm-custom-days'),
warmCustomWrap: $('dm-warm-custom-wrap'),
diskPreset: $('dm-disk-preset'),
diskCustomMb: $('dm-disk-custom-mb'),
diskCustomWrap: $('dm-disk-custom-wrap'),
warmMaxPoints: $('dm-warm-max-points'),
autoPrune: $('dm-auto-prune'),
pruneIntervalHours: $('dm-prune-interval-hours'),
lastPrune: $('dm-last-prune'),
status: $('dm-status'),
btnRefresh: $('dm-refresh'),
btnSave: $('dm-save'),
btnPruneDry: $('dm-prune-dry'),
btnPruneNow: $('dm-prune-now'),
gate: $('dm-admin-gate'),
}
function setStatus(msg, kind = '') {
if (!els.status) return
els.status.textContent = msg || ''
els.status.dataset.kind = kind
}
function canAdmin() {
return roleAllows(getRole() || Roles.viewer, Roles.admin)
}
function syncGate() {
const ok = canAdmin() && Boolean(manager.active)
if (els.gate) {
els.gate.classList.toggle('hidden', ok)
els.gate.textContent = !manager.active
? 'Connect to an agent to manage its data.'
: 'Admin role required to change retention or prune history.'
}
for (const el of [
els.hotPreset,
els.hotPoints,
els.warmMemPreset,
els.warmMemPoints,
els.tier1Every,
els.warmEnabled,
els.warmPreset,
els.warmCustomDays,
els.diskPreset,
els.diskCustomMb,
els.warmMaxPoints,
els.autoPrune,
els.pruneIntervalHours,
els.btnSave,
els.btnPruneDry,
els.btnPruneNow,
]) {
if (el) el.disabled = !ok || busy
}
if (els.btnRefresh) els.btnRefresh.disabled = !manager.active || busy
}
function fillSelect(select, presets, selected) {
if (!select) return
select.innerHTML = ''
for (const p of Object.values(presets)) {
const opt = document.createElement('option')
opt.value = p.id
opt.textContent = p.label
if (p.id === selected) opt.selected = true
select.appendChild(opt)
}
}
function renderUsage() {
if (!storage) {
if (els.usageTotal) els.usageTotal.textContent = '—'
if (els.usageCorestore) els.usageCorestore.textContent = '—'
if (els.usageMemory) els.usageMemory.textContent = '—'
if (els.usageWarm) els.usageWarm.textContent = '—'
if (els.usageMeta) els.usageMeta.textContent = ''
return
}
const u = storage.usage || {}
const mem = storage.memory || {}
const warm = storage.warm || {}
if (els.usageTotal) els.usageTotal.textContent = formatBytes(u.dataDirBytes)
if (els.usageCorestore) els.usageCorestore.textContent = formatBytes(u.corestoreBytes)
if (els.usageMemory) {
els.usageMemory.textContent = `${formatBytes(mem.approxBytes)} · ${mem.tier0Points ?? 0} hot / ${mem.tier1Points ?? 0} warm pts`
}
if (els.usageWarm) {
const age =
warm.oldestTs && warm.newestTs
? ` · ${formatAge(warm.newestTs - warm.oldestTs)} span`
: ''
els.usageWarm.textContent = warm.enabled
? `${warm.pointsSampled ?? 0} pts · ${warm.chartsSampled ?? 0} charts${age}${warm.sampleCapped ? ' (sampled)' : ''}`
: 'HyperDB off'
}
if (els.usageMeta) {
const dir = storage.dataDir || ''
els.usageMeta.textContent = dir ? `Data dir: ${dir}` : ''
}
}
function renderForm() {
if (!config) return
fillSelect(els.hotPreset, HOT_PRESETS, matchPointsPreset(config.tier0Points, HOT_PRESETS))
fillSelect(
els.warmMemPreset,
WARM_MEM_PRESETS,
matchPointsPreset(config.tier1Points, WARM_MEM_PRESETS)
)
fillSelect(els.warmPreset, RETENTION_PRESETS, config.warmRetentionPreset || '1w')
fillSelect(els.diskPreset, DISK_BUDGET_PRESETS, matchBytesPreset(config.warmMaxBytes || 0))
if (els.hotPoints) els.hotPoints.value = String(config.tier0Points)
if (els.warmMemPoints) els.warmMemPoints.value = String(config.tier1Points)
if (els.tier1Every) els.tier1Every.value = String(config.tier1Every)
if (els.warmEnabled) els.warmEnabled.checked = config.warmRetentionEnabled !== false
if (els.warmCustomDays) {
els.warmCustomDays.value = String(
Math.max(1, Math.round((config.warmRetentionMs || 0) / (24 * 60 * 60 * 1000)) || 7)
)
}
if (els.diskCustomMb) {
els.diskCustomMb.value = String(
config.warmMaxBytes > 0 ? Math.round(config.warmMaxBytes / (1024 * 1024)) : 1024
)
}
if (els.warmMaxPoints) els.warmMaxPoints.value = String(config.warmMaxPoints || 0)
if (els.autoPrune) els.autoPrune.checked = config.autoPruneEnabled !== false
if (els.pruneIntervalHours) {
els.pruneIntervalHours.value = String(
Math.max(1, Math.round((config.autoPruneIntervalMs || 3_600_000) / 3_600_000))
)
}
if (els.lastPrune) {
if (config.lastPruneAt) {
const when = new Date(config.lastPruneAt).toLocaleString()
els.lastPrune.textContent = `Last prune: ${when} · deleted ${config.lastPruneDeleted || 0} pts (~${formatBytes(config.lastPruneBytesFreed || 0)})`
} else {
els.lastPrune.textContent = 'Last prune: never'
}
}
toggleCustomRows()
}
function toggleCustomRows() {
const warmCustom = els.warmPreset?.value === 'custom'
const diskCustom = els.diskPreset?.value === 'custom'
els.warmCustomWrap?.classList.toggle('hidden', !warmCustom)
els.diskCustomWrap?.classList.toggle('hidden', !diskCustom)
if (els.hotPoints) {
els.hotPoints.disabled = !canAdmin() || busy || els.hotPreset?.value !== 'custom'
}
if (els.warmMemPoints) {
els.warmMemPoints.disabled = !canAdmin() || busy || els.warmMemPreset?.value !== 'custom'
}
}
function readForm() {
const hotId = els.hotPreset?.value || '1h'
const warmMemId = els.warmMemPreset?.value || '24h'
const warmId = els.warmPreset?.value || '1w'
const diskId = els.diskPreset?.value || 'unlimited'
let tier0Points =
hotId === 'custom'
? Number(els.hotPoints?.value)
: HOT_PRESETS[hotId]?.points ?? Number(els.hotPoints?.value)
let tier1Points =
warmMemId === 'custom'
? Number(els.warmMemPoints?.value)
: WARM_MEM_PRESETS[warmMemId]?.points ?? Number(els.warmMemPoints?.value)
let warmRetentionMs = 0
let warmRetentionPreset = warmId
if (warmId === 'forever') {
warmRetentionMs = 0
} else if (warmId === 'custom') {
const days = Math.max(1, Number(els.warmCustomDays?.value) || 7)
warmRetentionMs = days * 24 * 60 * 60 * 1000
} else {
warmRetentionMs = RETENTION_PRESETS[warmId]?.ms ?? 7 * 24 * 60 * 60 * 1000
}
let warmMaxBytes = 0
if (diskId === 'custom') {
warmMaxBytes = Math.max(0, Math.round(Number(els.diskCustomMb?.value) || 0) * 1024 * 1024)
} else {
warmMaxBytes = DISK_BUDGET_PRESETS[diskId]?.bytes ?? 0
}
const hours = Math.max(1, Number(els.pruneIntervalHours?.value) || 1)
return {
tier0Points,
tier1Points,
tier1Every: Number(els.tier1Every?.value) || 60,
warmRetentionEnabled: Boolean(els.warmEnabled?.checked),
warmRetentionPreset,
warmRetentionMs,
warmMaxBytes,
warmMaxPoints: Math.max(0, Math.floor(Number(els.warmMaxPoints?.value) || 0)),
autoPruneEnabled: Boolean(els.autoPrune?.checked),
autoPruneIntervalMs: hours * 60 * 60 * 1000,
}
}
async function refresh() {
if (!manager.active) {
config = null
storage = null
renderUsage()
setStatus('Not connected')
syncGate()
return
}
busy = true
syncGate()
setStatus('Loading…')
try {
const [info, ret] = await Promise.all([
manager.request(Methods.getStorageInfo, {}),
manager.request(Methods.getRetentionConfig, {}),
])
storage = info
config = ret?.config || info?.retention || null
renderUsage()
renderForm()
setStatus('Up to date')
} catch (err) {
setStatus(err.message || 'Failed to load', 'error')
log(`Data Manager: ${err.message || err}`)
} finally {
busy = false
syncGate()
toggleCustomRows()
}
}
async function save() {
if (!canAdmin()) return
busy = true
syncGate()
setStatus('Saving…')
try {
const patch = readForm()
const res = await manager.request(Methods.setRetentionConfig, patch)
config = res?.config || patch
setStatus('Saved on agent')
log('Retention config saved')
await refresh()
} catch (err) {
setStatus(err.message || 'Save failed', 'error')
log(`Data Manager save: ${err.message || err}`)
} finally {
busy = false
syncGate()
toggleCustomRows()
}
}
async function runPrune(dryRun) {
if (!canAdmin()) return
busy = true
syncGate()
setStatus(dryRun ? 'Dry-run…' : 'Pruning…')
try {
const res = await manager.request(Methods.pruneHistory, { dryRun, force: true })
if (res?.success === false) throw new Error(res.error || 'prune failed')
const verb = dryRun ? 'Would delete' : 'Deleted'
setStatus(`${verb} ${res.deleted ?? 0} points (~${formatBytes(res.approxBytesFreed || 0)})`)
log(`Prune${dryRun ? ' dry-run' : ''}: ${res.deleted ?? 0} deleted`)
if (!dryRun) await refresh()
} catch (err) {
setStatus(err.message || 'Prune failed', 'error')
log(`Data Manager prune: ${err.message || err}`)
} finally {
busy = false
syncGate()
toggleCustomRows()
}
}
function bind() {
els.btnRefresh?.addEventListener('click', () => refresh())
els.btnSave?.addEventListener('click', () => save())
els.btnPruneDry?.addEventListener('click', () => runPrune(true))
els.btnPruneNow?.addEventListener('click', () => {
if (
!confirm(
'Prune warm history on the connected agent now? Points older than the retention window will be deleted.'
)
) {
return
}
runPrune(false)
})
els.hotPreset?.addEventListener('change', () => {
const id = els.hotPreset.value
if (id !== 'custom' && HOT_PRESETS[id]?.points != null && els.hotPoints) {
els.hotPoints.value = String(HOT_PRESETS[id].points)
}
toggleCustomRows()
})
els.warmMemPreset?.addEventListener('change', () => {
const id = els.warmMemPreset.value
if (id !== 'custom' && WARM_MEM_PRESETS[id]?.points != null && els.warmMemPoints) {
els.warmMemPoints.value = String(WARM_MEM_PRESETS[id].points)
}
toggleCustomRows()
})
els.warmPreset?.addEventListener('change', toggleCustomRows)
els.diskPreset?.addEventListener('change', toggleCustomRows)
}
bind()
syncGate()
return {
refresh,
syncGate,
enter() {
syncGate()
refresh()
},
}
}
function formatAge(ms) {
if (!Number.isFinite(ms) || ms < 0) return '—'
const s = Math.floor(ms / 1000)
if (s < 60) return `${s}s`
const m = Math.floor(s / 60)
if (m < 60) return `${m}m`
const h = Math.floor(m / 60)
if (h < 48) return `${h}h`
const d = Math.floor(h / 24)
return `${d}d`
}
+51
View File
@@ -2282,6 +2282,57 @@ button.ghost.compact {
gap: 14px; gap: 14px;
} }
.settings-section-wide {
max-width: 720px;
}
.dm-usage {
display: grid;
grid-template-columns: repeat(2, minmax(0, 1fr));
gap: 10px;
}
.dm-usage-card {
display: grid;
gap: 4px;
padding: 10px 12px;
background: var(--bg-tertiary);
border: 1px solid var(--border-color);
border-radius: var(--border-radius-sm);
}
.dm-usage-card strong {
font-size: 1.05rem;
font-variant-numeric: tabular-nums;
color: var(--text-primary);
}
.dm-actions {
display: flex;
flex-wrap: wrap;
gap: 8px;
margin-top: 4px;
}
.dm-status[data-kind='error'] {
color: var(--accent-danger);
}
.dm-gate:not(.hidden) {
color: var(--accent-warning);
}
.danger-btn {
border-color: rgba(248, 113, 113, 0.45);
color: var(--accent-danger);
}
@media (max-width: 720px) {
.dm-usage {
grid-template-columns: 1fr;
}
}
.settings-hint { .settings-hint {
margin: 0; margin: 0;
font-size: 12px; font-size: 12px;
+22 -5
View File
@@ -1,13 +1,13 @@
# Settings # Settings
Settings persist under `~/.config/peardata` (or `PEARDATA_HOME`). Settings persist under `~/.config/peardata` (or `PEARDATA_HOME`) for desktop prefs. Agent storage policy is stored on the agent itself.
## Appearance ## Appearance
| Setting | Notes | | Setting | Notes |
|---------|-------| |---------|-------|
| **Theme** | Light / dark / system | | **Theme** | Light / dark |
| **Motion** | Reduce or allow UI motion | | **Accent / density / motion** | Visual prefs for the desktop shell |
## Charts / Overview ## Charts / Overview
@@ -15,7 +15,24 @@ Settings persist under `~/.config/peardata` (or `PEARDATA_HOME`).
|---------|-------| |---------|-------|
| **Overview spark depth** | Sample count when Charts has no overriding window | | **Overview spark depth** | Sample count when Charts has no overriding window |
| **Default spotlight chart** | Chart id for Overview explore (e.g. `system.io`) | | **Default spotlight chart** | Chart id for Overview explore (e.g. `system.io`) |
| Chart type defaults / wall prefs | Also stored with pins and card heights |
## Data Manager
Controls **agent-side** history so disks and RAM stay bounded. Requires an **admin** connection to the agent.
| Control | What it does |
|---------|----------------|
| **Usage cards** | Total data dir size, Corestore size, in-memory hot/warm points, warm HyperDB sample |
| **Hot memory** | 1s ring depth (15m → 24h presets or custom points) |
| **Warm memory** | Downsampled in-RAM ring + samples-per-bucket |
| **Disk history** | Auto-rotate by age (1 day → 1 year, forever, or custom days) |
| **Soft disk budget** | Optional Corestore size target that triggers extra prune |
| **Max warm points** | Hard cap on stored warm samples (`0` = unlimited) |
| **Auto prune** | Scheduled prune interval (hours) |
| **Dry-run / Prune now** | Preview or immediately delete expired warm points |
| **Save to agent** | Applies live + writes `retention.json` on the agent |
Env seeds (`PEARDATA_TIER*`, `PEARDATA_WARM_*`) only apply when no `retention.json` exists yet. See [CONFIGURATION.md](../docs/CONFIGURATION.md) and [STORAGE-HYPERDB.md](../docs/STORAGE-HYPERDB.md).
## Connections ## Connections
@@ -30,4 +47,4 @@ Desktop notification toggles for alerts (when the runtime supports them).
## Reset ## Reset
Clearing local prefs does not delete agent data on the server. Agent retention is controlled by agent env / warm store (see [docs/CONFIGURATION.md](../docs/CONFIGURATION.md)). Clearing local desktop prefs does not delete agent data. Use **Data → Prune now** or tighten retention to reclaim warm history on the agent.