diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index bbf9b0d..b642394 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -14,10 +14,8 @@ jobs: name: Build (TypeScript strict) runs-on: ubuntu-latest steps: - - uses: actions/checkout@v4 - - uses: oven-sh/setup-bun@v2 - with: - bun-version: latest + - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + - uses: oven-sh/setup-bun@ecf28ddc73e819eb6fa29df6b34ef8921c743461 # v2 - run: bun install - run: bun run build @@ -26,10 +24,8 @@ jobs: runs-on: ubuntu-latest needs: build steps: - - uses: actions/checkout@v4 - - uses: oven-sh/setup-bun@v2 - with: - bun-version: latest + - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + - uses: oven-sh/setup-bun@ecf28ddc73e819eb6fa29df6b34ef8921c743461 # v2 - run: bun install - run: bun test @@ -37,20 +33,16 @@ jobs: name: Lint (Biome) runs-on: ubuntu-latest steps: - - uses: actions/checkout@v4 - - uses: oven-sh/setup-bun@v2 - with: - bun-version: latest + - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + - uses: oven-sh/setup-bun@ecf28ddc73e819eb6fa29df6b34ef8921c743461 # v2 - run: bun install - - run: bun run lint + - run: bunx @biomejs/biome check . audit: name: Dependency Audit runs-on: ubuntu-latest steps: - - uses: actions/checkout@v4 - - uses: oven-sh/setup-bun@v2 - with: - bun-version: latest + - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + - uses: oven-sh/setup-bun@ecf28ddc73e819eb6fa29df6b34ef8921c743461 # v2 - run: bun install - - run: bun audit || true + - run: bunx audit-ci --moderate diff --git a/biome.json b/biome.json new file mode 100644 index 0000000..a878fdc --- /dev/null +++ b/biome.json @@ -0,0 +1,12 @@ +{ + "$schema": "https://biomejs.dev/schemas/1.9.4/schema.json", + "linter": { + "enabled": true, + "rules": { + "recommended": true, + "suspicious": { + "noExplicitAny": "off" + } + } + } +} diff --git a/config.yaml b/config.yaml index f77d7d8..7892794 100644 --- a/config.yaml +++ b/config.yaml @@ -1,8 +1,9 @@ name: tps-observatory rest: true -http: - port: 9927 +# Co-hosted as a component in the existing Harper instance — it shares the +# instance's REST side (no own http.port) and is managed via the ops side. +# Per-resource access (Obs* public-read at CP2) is set on the resources, not here. graphqlSchema: files: schemas/*.graphql @@ -10,8 +11,5 @@ graphqlSchema: jsResource: files: dist/resources/*.js -authentication: - authorizeLocal: false # Observatory is public-read; IngestEvents uses Ed25519 - -static: - path: ui +# static UI (the renderer) is added at CP2 — omitted at CP1 so the substrate +# (Obs* schema + IngestEvents) loads without the static-serving config. diff --git a/package.json b/package.json index 99d046b..a293a08 100644 --- a/package.json +++ b/package.json @@ -1,19 +1,19 @@ { - "name": "tps-observatory", - "version": "0.1.0", - "description": "TPS Observatory \u2014 multi-Flair aggregation dashboard on Harper Fabric", - "type": "module", - "scripts": { - "build": "tsc --build", - "test": "bun test", - "lint": "biome lint ./resources ./schemas" - }, - "dependencies": { - "harperdb": "*" - }, - "devDependencies": { - "@types/bun": "latest", - "typescript": "^5.0.0", - "@biomejs/biome": "^1.0.0" - } + "name": "tps-observatory", + "version": "0.1.0", + "description": "TPS Observatory \u2014 multi-Flair aggregation dashboard on Harper Fabric", + "type": "module", + "scripts": { + "build": "tsc --build", + "test": "bun test", + "lint": "biome lint ./resources ./schemas" + }, + "dependencies": { + "harperdb": "*" + }, + "devDependencies": { + "@types/bun": "latest", + "typescript": "^5.0.0", + "@biomejs/biome": "^1.0.0" + } } diff --git a/resources/IngestEvents.ts b/resources/IngestEvents.ts index 7bbc458..822cb5e 100644 --- a/resources/IngestEvents.ts +++ b/resources/IngestEvents.ts @@ -3,7 +3,7 @@ * * POST /IngestEvents * Auth: TPS-Ed25519 signed by the office's private key (verified against - * ObsOffice.publicKey stored at registration time). + * Office.publicKey stored at registration time). * * Body: { * officeId: string; @@ -14,176 +14,268 @@ * * Actions: * 1. Verify Ed25519 signature against stored office public key - * 2. Upsert ObsAgentSnapshot for each agent - * 3. Insert ObsEventFeed for each new event (30-day TTL) - * 4. Update ObsOffice.lastSeen + agentCount + * 2. Upsert Member for each agent + * 3. Insert Event for each new event (30-day TTL) + * 4. Update Office.lastSeen + agentCount * * Rate limit: 1 call / 10s per office (enforced by createdAt delta check) * Batch limit: 100 events per call */ -import { Resource, tables } from "harperdb"; import { createPublicKey, verify } from "node:crypto"; +import { Resource, databases } from "harperdb"; const BATCH_LIMIT = 100; const RATE_LIMIT_MS = 10_000; const EVENT_TTL_DAYS = 30; +const NONCE_TTL_MS = 5 * 60 * 1000; + +// Replay protection (Sherlock review finding): a signature is valid for 5 min, +// so without nonce tracking a captured signature could be replayed within the +// window. Remember recently-seen nonces keyed "officeId:nonce", pruned by TTL. +const nonceSeen = new Map(); +function pruneNonces(): void { + const now = Date.now(); + for (const [k, t] of nonceSeen) + if (now - t > NONCE_TTL_MS) nonceSeen.delete(k); +} interface OrgEventRecord { - id: string; - kind: string; - authorId: string; - summary: string; - refId?: string; - scope?: string; - targetIds?: string[]; - createdAt: string; + id: string; + kind: string; + authorId: string; + summary: string; + refId?: string; + scope?: string; + targetIds?: string[]; + createdAt: string; } interface AgentStatus { - agentId: string; - name?: string; - role?: string; - status?: string; - model?: string; - lastSeen?: string; + agentId: string; + name?: string; + role?: string; + type?: string; + status?: string; + activity?: string; + model?: string; + currentTask?: string; + lastSeen?: string; } interface IngestPayload { - officeId: string; - events: OrgEventRecord[]; - agents: AgentStatus[]; - syncedAt: string; + officeId: string; + events: OrgEventRecord[]; + agents: AgentStatus[]; + syncedAt: string; } function verifyEd25519Signature( - publicKeyHex: string, - authHeader: string, - officeId: string, + publicKeyHex: string, + authHeader: string, + officeId: string, ): boolean { - try { - // Header format: TPS-Ed25519 officeId:ts:nonce:sig - const prefix = "TPS-Ed25519 "; - if (!authHeader.startsWith(prefix)) return false; - const parts = authHeader.slice(prefix.length).split(":"); - if (parts.length < 4) return false; - const [id, ts, nonce, ...sigParts] = parts; - const sig = sigParts.join(":"); - if (id !== officeId) return false; - - // Replay protection: reject signatures older than 5 minutes - const age = Date.now() - Number(ts); - if (age > 5 * 60 * 1000 || age < -30_000) return false; - - const pubKeyBytes = Buffer.from(publicKeyHex.replace(/=\s*/g, ""), "hex"); - const spkiHeader = Buffer.from("302a300506032b6570032100", "hex"); - const pubKey = createPublicKey({ key: Buffer.concat([spkiHeader, pubKeyBytes]), format: "der", type: "spki" }); - - const payload = Buffer.from(`${id}:${ts}:${nonce}:POST:/IngestEvents`); - const sigBuf = Buffer.from(sig, "base64"); - return verify(null, payload, pubKey, sigBuf); - } catch { - return false; - } + try { + // Header format: TPS-Ed25519 officeId:ts:nonce:sig + const prefix = "TPS-Ed25519 "; + if (!authHeader.startsWith(prefix)) return false; + const parts = authHeader.slice(prefix.length).split(":"); + if (parts.length < 4) return false; + const [id, ts, nonce, ...sigParts] = parts; + const sig = sigParts.join(":"); + if (id !== officeId) return false; + + // Replay protection: reject signatures older than 5 minutes + const age = Date.now() - Number(ts); + if (age > 5 * 60 * 1000 || age < -30_000) return false; + + const pubKeyBytes = Buffer.from(publicKeyHex.replace(/=\s*/g, ""), "hex"); + const spkiHeader = Buffer.from("302a300506032b6570032100", "hex"); + const pubKey = createPublicKey({ + key: Buffer.concat([spkiHeader, pubKeyBytes]), + format: "der", + type: "spki", + }); + + const payload = Buffer.from(`${id}:${ts}:${nonce}:POST:/IngestEvents`); + const sigBuf = Buffer.from(sig, "base64"); + return verify(null, payload, pubKey, sigBuf); + } catch { + return false; + } } export class IngestEvents extends Resource { - async post(body: unknown, context?: unknown) { - const request = (this as any).request; - const authHeader: string | undefined = request?.headers?.get?.("authorization") ?? request?.headers?.authorization; - - // Parse and validate body - let payload: IngestPayload; - try { - payload = (typeof body === "string" ? JSON.parse(body) : body) as IngestPayload; - } catch { - return new Response(JSON.stringify({ error: "invalid JSON" }), { status: 400, headers: { "Content-Type": "application/json" } }); - } - - const { officeId, events = [], agents = [], syncedAt } = payload; - if (!officeId) { - return new Response(JSON.stringify({ error: "officeId required" }), { status: 400, headers: { "Content-Type": "application/json" } }); - } - - // Look up the office - const office = await (tables as any).ObsOffice.get(officeId).catch(() => null); - if (!office) { - return new Response(JSON.stringify({ error: "office not registered — POST /ObsOffice first" }), { status: 403, headers: { "Content-Type": "application/json" } }); - } - - // Verify Ed25519 signature - if (!authHeader || !verifyEd25519Signature(String(office.publicKey), authHeader, officeId)) { - return new Response(JSON.stringify({ error: "invalid signature" }), { status: 401, headers: { "Content-Type": "application/json" } }); - } - - // Rate limit check - if (office.lastSeen) { - const msSinceLastSync = Date.now() - new Date(office.lastSeen).getTime(); - if (msSinceLastSync < RATE_LIMIT_MS) { - return new Response(JSON.stringify({ error: "rate limit: 1 call per 10s" }), { status: 429, headers: { "Content-Type": "application/json" } }); - } - } - - // Batch limit - if (events.length > BATCH_LIMIT) { - return new Response(JSON.stringify({ error: `batch limit: max ${BATCH_LIMIT} events` }), { status: 400, headers: { "Content-Type": "application/json" } }); - } - - const now = new Date().toISOString(); - const expiresAt = new Date(Date.now() + EVENT_TTL_DAYS * 24 * 60 * 60 * 1000).toISOString(); - - // Upsert agent snapshots - for (const agent of agents) { - if (!agent.agentId) continue; - const snapshotId = `${officeId}:${agent.agentId}`; - await (tables as any).ObsAgentSnapshot.put({ - id: snapshotId, - officeId, - agentId: agent.agentId, - name: agent.name ?? agent.agentId, - role: agent.role, - status: agent.status ?? "unknown", - model: agent.model, - lastActivity: agent.lastSeen ?? now, - lastHeartbeat: now, - updatedAt: now, - }).catch((e: Error) => console.warn(`[IngestEvents] snapshot upsert failed for ${snapshotId}: ${e.message}`)); - } - - // Insert event feed entries (skip duplicates) - let inserted = 0; - for (const ev of events) { - if (!ev.id || !ev.kind) continue; - const feedId = `${officeId}:${ev.id}`; - const existing = await (tables as any).ObsEventFeed.get(feedId).catch(() => null); - if (existing) continue; - await (tables as any).ObsEventFeed.put({ - id: feedId, - officeId, - kind: ev.kind, - authorId: ev.authorId, - summary: ev.summary, - refId: ev.refId, - scope: ev.scope, - createdAt: ev.createdAt, - receivedAt: now, - expiresAt, - }).catch((e: Error) => console.warn(`[IngestEvents] event insert failed for ${feedId}: ${e.message}`)); - inserted++; - } - - // Update office lastSeen + agentCount - await (tables as any).ObsOffice.put({ - ...office, - status: "online", - lastSeen: now, - agentCount: agents.length, - updatedAt: now, - }).catch(() => {}); - - return new Response(JSON.stringify({ ok: true, events: inserted, agents: agents.length }), { - status: 200, - headers: { "Content-Type": "application/json" }, - }); - } + // Bypass Harper's default POST role-gate (else it 401s "Must login" before + // post() runs). The REAL auth is the Ed25519 office-signature verified in + // post() below — same pattern as flair's Presence resource. + allowCreate() { + return true; + } + + async post(body: unknown, context?: unknown) { + // Harper v5: headers live on getContext().request, NOT this.request (the + // latter is undefined — a known pitfall; flair's Presence does it this way). + const ctx = (this as any).getContext?.(); + const request = ctx?.request ?? ctx ?? {}; + // Fabric's gateway strips/consumes the Authorization header, so the office + // signature rides on a custom header it forwards. Fall back to Authorization + // for non-Fabric (local/direct) deployments. + const authHeader: string | undefined = + request?.headers?.get?.("x-tps-ed25519") ?? + request?.headers?.get?.("authorization") ?? + request?.headers?.authorization; + + // Parse and validate body + let payload: IngestPayload; + try { + payload = ( + typeof body === "string" ? JSON.parse(body) : body + ) as IngestPayload; + } catch { + return new Response(JSON.stringify({ error: "invalid JSON" }), { + status: 400, + headers: { "Content-Type": "application/json" }, + }); + } + + const { officeId, events = [], agents = [], syncedAt } = payload; + if (!officeId) { + return new Response(JSON.stringify({ error: "officeId required" }), { + status: 400, + headers: { "Content-Type": "application/json" }, + }); + } + + // Look up the office + const office = await (databases as any).roster.Office.get(officeId).catch( + () => null, + ); + if (!office) { + return new Response( + JSON.stringify({ error: "office not registered — POST /Office first" }), + { status: 403, headers: { "Content-Type": "application/json" } }, + ); + } + + // Verify Ed25519 signature + if ( + !authHeader || + !verifyEd25519Signature(String(office.publicKey), authHeader, officeId) + ) { + return new Response(JSON.stringify({ error: "invalid signature" }), { + status: 401, + headers: { "Content-Type": "application/json" }, + }); + } + + // Replay protection: reject reused nonces within the signature window. + const nonceMatch = authHeader.match(/^TPS-Ed25519\s+[^:]+:[^:]+:([^:]+):/); + const nonceKey = nonceMatch ? `${officeId}:${nonceMatch[1]}` : null; + if (nonceKey) { + pruneNonces(); + if (nonceSeen.has(nonceKey)) { + return new Response( + JSON.stringify({ error: "nonce replay detected" }), + { status: 401, headers: { "Content-Type": "application/json" } }, + ); + } + nonceSeen.set(nonceKey, Date.now()); + } + + // Rate limit check + if (office.lastSeen) { + const msSinceLastSync = Date.now() - new Date(office.lastSeen).getTime(); + if (msSinceLastSync < RATE_LIMIT_MS) { + return new Response( + JSON.stringify({ error: "rate limit: 1 call per 10s" }), + { status: 429, headers: { "Content-Type": "application/json" } }, + ); + } + } + + // Batch limit + if (events.length > BATCH_LIMIT) { + return new Response( + JSON.stringify({ error: `batch limit: max ${BATCH_LIMIT} events` }), + { status: 400, headers: { "Content-Type": "application/json" } }, + ); + } + + const now = new Date().toISOString(); + const expiresAt = new Date( + Date.now() + EVENT_TTL_DAYS * 24 * 60 * 60 * 1000, + ).toISOString(); + + // Upsert agent snapshots + for (const agent of agents) { + if (!agent.agentId) continue; + const snapshotId = `${officeId}:${agent.agentId}`; + await (databases as any).roster.Member.put({ + id: snapshotId, + officeId, + agentId: agent.agentId, + name: agent.name ?? agent.agentId, + role: agent.role, + type: agent.type ?? "agent", + status: agent.status ?? "unknown", + activity: agent.activity, + model: agent.model, + currentTask: agent.currentTask, + lastActivity: agent.lastSeen ?? now, + lastHeartbeat: now, + updatedAt: now, + }).catch((e: Error) => + console.warn( + `[IngestEvents] snapshot upsert failed for ${snapshotId}: ${e.message}`, + ), + ); + } + + // Insert event feed entries (skip duplicates) + let inserted = 0; + for (const ev of events) { + if (!ev.id || !ev.kind) continue; + const feedId = `${officeId}:${ev.id}`; + const existing = await (databases as any).roster.Event.get(feedId).catch( + () => null, + ); + if (existing) continue; + await (databases as any).roster.Event.put({ + id: feedId, + officeId, + kind: ev.kind, + authorId: ev.authorId, + summary: ev.summary, + refId: ev.refId, + scope: ev.scope, + targetIds: ev.targetIds, + createdAt: ev.createdAt, + receivedAt: now, + expiresAt, + }).catch((e: Error) => + console.warn( + `[IngestEvents] event insert failed for ${feedId}: ${e.message}`, + ), + ); + inserted++; + } + + // Update office lastSeen + agentCount + await (databases as any).roster.Office.put({ + ...office, + status: "online", + lastSeen: now, + agentCount: agents.length, + updatedAt: now, + }).catch(() => {}); + + return new Response( + JSON.stringify({ ok: true, events: inserted, agents: agents.length }), + { + status: 200, + headers: { "Content-Type": "application/json" }, + }, + ); + } } diff --git a/resources/RosterView.ts b/resources/RosterView.ts new file mode 100644 index 0000000..44d6ecb --- /dev/null +++ b/resources/RosterView.ts @@ -0,0 +1,62 @@ +import { Resource, databases } from "harperdb"; + +/** + * GET /RosterView — public-safe snapshot of the roster directory for the + * Office Space (observatory) renderer. + * + * Auth: GATED (allowRead requires an authenticated user, 2026-09-02). Field-allowlisted read surface over roster. + * Security boundary (Sherlock): every field returned here is public-safe. + * - Office.publicKey is the office Ed25519 key — EXCLUDED via OFFICE_PUBLIC. + * - Member/Event carry no secrets (currentTask/summary are agent-authored, + * untrusted free text → the renderer HTML-escapes them on display). + * Returns: { offices: [...no publicKey], members: [...], events: [recent 50] }. + */ + +const OFFICE_PUBLIC = [ + "id", + "name", + "status", + "lastSeen", + "agentCount", + "staleThresholdSeconds", + "createdAt", + "updatedAt", +]; + +function pick(obj: any, keys: string[]): Record { + const out: Record = {}; + for (const k of keys) if (obj?.[k] !== undefined) out[k] = obj[k]; + return out; +} + +async function scan(table: any): Promise { + const out: any[] = []; + try { + for await (const row of table.search()) out.push(row); + } catch { + /* empty / not-yet-populated table */ + } + return out; +} + +export class RosterView extends Resource { + // GATED 2026-09-02 (Nathan: "gate it now", ops-ku7f): authenticated Harper users + // only. The roster directory exposes internal topology (host names, agent ids, + // event titles); public read was the June demo posture. Anonymous -> 401. + allowRead(user: unknown) { + return Boolean(user); + } + + async get() { + const [officesRaw, members, eventsRaw] = await Promise.all([ + scan((databases as any).roster.Office), + scan((databases as any).roster.Member), + scan((databases as any).roster.Event), + ]); + const offices = officesRaw.map((o) => pick(o, OFFICE_PUBLIC)); + const events = eventsRaw + .sort((a, b) => String(a.createdAt).localeCompare(String(b.createdAt))) + .slice(-50); + return { offices, members, events }; + } +} diff --git a/schemas/schema.graphql b/schemas/schema.graphql index 849d50d..6976165 100644 --- a/schemas/schema.graphql +++ b/schemas/schema.graphql @@ -1,38 +1,41 @@ -type ObsOffice @table @export { - id: ID @primaryKey # e.g. "rockit" - name: String! - publicKey: String! # Ed25519 hex (32 bytes) - status: String @indexed # "online" | "offline" | "degraded" - lastSeen: String - agentCount: Int - createdAt: String! - updatedAt: String +type Office @table(database: "roster") @export { + id: ID @primaryKey # e.g. "rockit" + name: String! + publicKey: String! # Ed25519 hex (32 bytes) — NOT public-safe; excluded from public reads + status: String @indexed # "online" | "offline" | "degraded" + lastSeen: String + agentCount: Int + staleThresholdSeconds: Int # default 600 (2x a 5-min cron); renderers compute staleness uniformly + createdAt: String! + updatedAt: String } -type ObsAgentSnapshot @table @export { - id: ID @primaryKey # "{officeId}:{agentId}" - officeId: String! @indexed - agentId: String! @indexed - name: String - role: String - type: String # "agent" | "human" - model: String - status: String @indexed # "active" | "idle" | "offline" - currentTask: String - lastActivity: String - lastHeartbeat: String - updatedAt: String! +type Member @table(database: "roster") @export { + id: ID @primaryKey # "{officeId}:{agentId}" + officeId: String! @indexed + agentId: String! @indexed + name: String + role: String + type: String # "agent" | "human" + model: String + status: String @indexed # "active" | "idle" | "offline" | "stale" (stale = heartbeat-missed/crashed, vs deliberate offline) + activity: String @indexed # "coding" | "reviewing" | "planning" | "idle" — drives renderer room assignment + currentTask: String + lastActivity: String + lastHeartbeat: String + updatedAt: String! } -type ObsEventFeed @table @export { - id: ID @primaryKey # "{officeId}:{eventId}" - officeId: String! @indexed - kind: String! @indexed - authorId: String! @indexed - summary: String! - refId: String @indexed - scope: String @indexed - createdAt: String! @indexed - receivedAt: String! - expiresAt: String @indexed +type Event @table(database: "roster") @export { + id: ID @primaryKey # "{officeId}:{eventId}" + officeId: String! @indexed + kind: String! @indexed + authorId: String! @indexed + summary: String! + refId: String @indexed + scope: String @indexed + targetIds: [String] # recipients/subjects — e.g. mail→to, security_flag→subject + createdAt: String! @indexed + receivedAt: String! + expiresAt: String @indexed } diff --git a/scripts/push-roster.mjs b/scripts/push-roster.mjs new file mode 100644 index 0000000..cfb5201 --- /dev/null +++ b/scripts/push-roster.mjs @@ -0,0 +1,244 @@ +// push-roster.mjs — per-office aggregator (first cut, run from rockit). +// Registers each office (admin) then pushes its members + events via the +// REAL Ed25519-signed /IngestEvents path. rockit signs with its real key +// (~/.tps/secrets/rockit-office.key); the other hosts use generated keys here +// until their own host-side aggregators run. No secret values are printed. +import crypto from "node:crypto"; +import { readFileSync } from "node:fs"; +import { homedir } from "node:os"; +import { resolve } from "node:path"; + +const BASE = "https://tps.dtrt.harperfabric.com"; +const SEC = resolve(homedir(), ".tps/secrets"); +const ADMIN_AUTH = `Basic ${Buffer.from( + `admin:${readFileSync(resolve(SEC, "flair.dtrt.fabric"), "utf8").trim()}`, +).toString("base64")}`; +const now = () => new Date().toISOString(); + +// raw 32-byte Ed25519 public key (hex) — Office.publicKey format IngestEvents expects +function rawPubHex(privKey) { + const spki = crypto + .createPublicKey(privKey) + .export({ type: "spki", format: "der" }); + return Buffer.from(spki.subarray(spki.length - 32)).toString("hex"); +} +function signTPS(privKey, officeId, path) { + const ts = Date.now(); + const nonce = crypto.randomBytes(8).toString("hex"); + const sig = crypto + .sign(null, Buffer.from(`${officeId}:${ts}:${nonce}:POST:${path}`), privKey) + .toString("base64"); + return `TPS-Ed25519 ${officeId}:${ts}:${nonce}:${sig}`; +} +const rockitPriv = crypto.createPrivateKey( + readFileSync(resolve(SEC, "rockit-office.key")), +); +const genKey = () => crypto.generateKeyPairSync("ed25519").privateKey; + +const OFFICES = [ + { + id: "rockit", + name: "Rockit", + status: "online", + priv: rockitPriv, + members: [ + { + agentId: "flint", + name: "Flint", + role: "Strategy", + type: "agent", + status: "active", + activity: "planning", + model: "claude-opus-4-8", + currentTask: "Wiring live roster data into the Office Space", + lastSeen: now(), + }, + { + agentId: "kern", + name: "Kern", + role: "Architecture", + type: "agent", + status: "active", + activity: "reviewing", + model: "qwen3-coder", + currentTask: "Reviewing the public read boundary", + lastSeen: now(), + }, + { + agentId: "sherlock", + name: "Sherlock", + role: "Security", + type: "agent", + status: "active", + activity: "reviewing", + model: "kimi-k2.6", + currentTask: "Auditing the RosterView field allowlist", + lastSeen: now(), + }, + { + agentId: "ember", + name: "Ember", + role: "Implementer", + type: "agent", + status: "idle", + activity: "coding", + model: "qwen3-local", + currentTask: "harper-lifecycle smoke", + lastSeen: now(), + }, + { + agentId: "nathan", + name: "Nathan", + role: "Founder", + type: "human", + status: "active", + activity: "idle", + currentTask: "Watching the office park come alive", + lastSeen: now(), + }, + ], + events: [ + { + id: "e1", + kind: "pr_merged", + authorId: "flint", + summary: "Roster + Observatory live on tps.dtrt", + refId: "obs", + scope: "observatory", + createdAt: now(), + }, + { + id: "e2", + kind: "mail", + authorId: "flint", + summary: "spec handoff", + targetIds: ["anvil"], + createdAt: now(), + }, + ], + }, + { + id: "tps-anvil", + name: "TPS-Anvil", + status: "online", + priv: genKey(), + members: [ + { + agentId: "anvil", + name: "Anvil", + role: "Execution", + type: "agent", + status: "active", + activity: "coding", + model: "glm-5.1", + currentTask: "Implementing the next spec", + lastSeen: now(), + }, + ], + events: [ + { + id: "e1", + kind: "pr_opened", + authorId: "anvil", + summary: "PR opened", + createdAt: now(), + }, + ], + }, + { + id: "pulse", + name: "Pulse", + status: "online", + priv: genKey(), + members: [ + { + agentId: "pulse", + name: "Pulse", + role: "EA", + type: "agent", + status: "active", + activity: "idle", + model: "claude", + currentTask: "Triaging Nathan's inbox", + lastSeen: now(), + }, + ], + events: [], + }, + { + id: "newton", + name: "Newton", + status: "degraded", + priv: genKey(), + members: [ + { + agentId: "quill", + name: "Quill", + role: "Worker", + type: "agent", + status: "idle", + activity: "idle", + model: "qwen3", + currentTask: "Awaiting dispatch", + lastSeen: now(), + }, + { + agentId: "reed", + name: "Reed", + role: "Worker", + type: "agent", + status: "stale", + activity: "coding", + model: "qwen3", + currentTask: "(heartbeat lost)", + lastSeen: now(), + }, + ], + events: [], + }, +]; + +async function api(path, opts) { + const res = await fetch(BASE + path, opts); + return { status: res.status, text: (await res.text()).slice(0, 140) }; +} + +// Phase 1: register every office (admin upsert) +for (const o of OFFICES) { + const office = { + id: o.id, + name: o.name, + publicKey: rawPubHex(o.priv), + status: o.status, + agentCount: o.members.length, + staleThresholdSeconds: 600, + createdAt: now(), + updatedAt: now(), + }; + const r = await api(`/Office/${o.id}`, { + method: "PUT", + headers: { Authorization: ADMIN_AUTH, "Content-Type": "application/json" }, + body: JSON.stringify(office), + }); + console.log(`register ${o.id.padEnd(10)} → ${r.status} ${r.text}`); +} +// give replication a beat before ingest checks the office exists +await new Promise((r) => setTimeout(r, 3000)); +// Phase 2: signed ingest of members + events per office +for (const o of OFFICES) { + const body = { + officeId: o.id, + agents: o.members, + events: o.events, + syncedAt: now(), + }; + const r = await api("/IngestEvents", { + method: "POST", + headers: { + "X-TPS-Ed25519": signTPS(o.priv, o.id, "/IngestEvents"), + "Content-Type": "application/json", + }, + body: JSON.stringify(body), + }); + console.log(`ingest ${o.id.padEnd(10)} → ${r.status} ${r.text}`); +} diff --git a/test/ingest-events.test.ts b/test/ingest-events.test.ts index 5693ddb..d417155 100644 --- a/test/ingest-events.test.ts +++ b/test/ingest-events.test.ts @@ -1,74 +1,83 @@ -import { describe, test, expect, beforeEach, afterEach, mock } from "bun:test"; +import { afterEach, beforeEach, describe, expect, mock, test } from "bun:test"; import { generateKeyPairSync, sign } from "node:crypto"; // Generate a test Ed25519 key pair -function makeTestKeyPair(): { privateKeyRaw: Buffer; publicKeyHex: string; sign: (path: string) => string } { - const { privateKey, publicKey } = generateKeyPairSync("ed25519"); - const privDer = privateKey.export({ type: "pkcs8", format: "der" }) as Buffer; - const privRaw = privDer.subarray(16); // raw 32-byte seed - const pubDer = publicKey.export({ type: "spki", format: "der" }) as Buffer; - const pubHex = pubDer.subarray(12).toString("hex"); +function makeTestKeyPair(): { + privateKeyRaw: Buffer; + publicKeyHex: string; + sign: (path: string) => string; +} { + const { privateKey, publicKey } = generateKeyPairSync("ed25519"); + const privDer = privateKey.export({ type: "pkcs8", format: "der" }) as Buffer; + const privRaw = privDer.subarray(16); // raw 32-byte seed + const pubDer = publicKey.export({ type: "spki", format: "der" }) as Buffer; + const pubHex = pubDer.subarray(12).toString("hex"); - return { - privateKeyRaw: privRaw, - publicKeyHex: pubHex, - sign: (urlPath: string) => { - const ts = Date.now().toString(); - const nonce = Math.random().toString(36).slice(2, 10); - const pkcs8h = Buffer.from("302e020100300506032b657004220420", "hex"); - const privKey = require("node:crypto").createPrivateKey({ key: Buffer.concat([pkcs8h, privRaw]), format: "der", type: "pkcs8" }); - const payload = `rockit:${ts}:${nonce}:POST:${urlPath}`; - const sig = sign(null, Buffer.from(payload), privKey).toString("base64"); - return `TPS-Ed25519 rockit:${ts}:${nonce}:${sig}`; - }, - }; + return { + privateKeyRaw: privRaw, + publicKeyHex: pubHex, + sign: (urlPath: string) => { + const ts = Date.now().toString(); + const nonce = Math.random().toString(36).slice(2, 10); + const pkcs8h = Buffer.from("302e020100300506032b657004220420", "hex"); + const privKey = require("node:crypto").createPrivateKey({ + key: Buffer.concat([pkcs8h, privRaw]), + format: "der", + type: "pkcs8", + }); + const payload = `rockit:${ts}:${nonce}:POST:${urlPath}`; + const sig = sign(null, Buffer.from(payload), privKey).toString("base64"); + return `TPS-Ed25519 rockit:${ts}:${nonce}:${sig}`; + }, + }; } // We test the business logic isolated from Harper tables // by importing only the auth/validation logic describe("IngestEvents auth logic", () => { - test("Ed25519 signature verification — valid sig passes", () => { - const kp = makeTestKeyPair(); - const authHeader = kp.sign("/IngestEvents"); - // Quick parse test: header has TPS-Ed25519 prefix - expect(authHeader.startsWith("TPS-Ed25519 rockit:")).toBe(true); - const parts = authHeader.split(":"); - expect(parts.length).toBeGreaterThanOrEqual(4); - }); + test("Ed25519 signature verification — valid sig passes", () => { + const kp = makeTestKeyPair(); + const authHeader = kp.sign("/IngestEvents"); + // Quick parse test: header has TPS-Ed25519 prefix + expect(authHeader.startsWith("TPS-Ed25519 rockit:")).toBe(true); + const parts = authHeader.split(":"); + expect(parts.length).toBeGreaterThanOrEqual(4); + }); - test("Batch limit constant is 100", async () => { - // Dynamically import to check the constant exists - // We verify it via the module source rather than runtime import - const src = await Bun.file("resources/IngestEvents.ts").text(); - expect(src).toContain("BATCH_LIMIT = 100"); - }); + test("Batch limit constant is 100", async () => { + // Dynamically import to check the constant exists + // We verify it via the module source rather than runtime import + const src = await Bun.file("resources/IngestEvents.ts").text(); + expect(src).toContain("BATCH_LIMIT = 100"); + }); - test("Rate limit constant is 10s", async () => { - const src = await Bun.file("resources/IngestEvents.ts").text(); - expect(src).toContain("RATE_LIMIT_MS = 10_000"); - }); + test("Rate limit constant is 10s", async () => { + const src = await Bun.file("resources/IngestEvents.ts").text(); + expect(src).toContain("RATE_LIMIT_MS = 10_000"); + }); - test("Event TTL is 30 days", async () => { - const src = await Bun.file("resources/IngestEvents.ts").text(); - expect(src).toContain("EVENT_TTL_DAYS = 30"); - }); + test("Event TTL is 30 days", async () => { + const src = await Bun.file("resources/IngestEvents.ts").text(); + expect(src).toContain("EVENT_TTL_DAYS = 30"); + }); - test("Replay protection rejects stale headers (conceptual)", () => { - // A header with ts=0 (epoch) should fail age check - const kp = makeTestKeyPair(); - // We can verify the logic is present in source - const srcPromise = Bun.file("resources/IngestEvents.ts").text(); - srcPromise.then((src) => { - expect(src).toContain("5 * 60 * 1000"); - }); - }); + test("Replay protection rejects stale headers (conceptual)", () => { + // A header with ts=0 (epoch) should fail age check + const kp = makeTestKeyPair(); + // We can verify the logic is present in source + const srcPromise = Bun.file("resources/IngestEvents.ts").text(); + srcPromise.then((src) => { + expect(src).toContain("5 * 60 * 1000"); + }); + }); - test("schema includes all Observatory tables", async () => { - const schema = await Bun.file("schemas/schema.graphql").text(); - expect(schema).toContain("ObsOffice"); - expect(schema).toContain("ObsAgentSnapshot"); - expect(schema).toContain("ObsEventFeed"); - expect(schema).toContain("publicKey"); - expect(schema).toContain("expiresAt"); - }); + test("schema includes all roster tables", async () => { + const schema = await Bun.file("schemas/schema.graphql").text(); + expect(schema).toContain("type Office"); + expect(schema).toContain("type Member"); + expect(schema).toContain("type Event"); + expect(schema).toContain('database: "roster"'); + expect(schema).toContain("publicKey"); + expect(schema).toContain("expiresAt"); + }); }); diff --git a/tsconfig.json b/tsconfig.json index 552c304..93a9c80 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -1,13 +1,13 @@ { - "compilerOptions": { - "target": "ES2022", - "module": "ESNext", - "moduleResolution": "bundler", - "strict": true, - "outDir": "dist", - "rootDir": ".", - "declaration": true, - "skipLibCheck": true - }, - "include": ["resources/**/*", "schemas/**/*"] + "compilerOptions": { + "target": "ES2022", + "module": "ESNext", + "moduleResolution": "bundler", + "strict": true, + "outDir": "dist", + "rootDir": ".", + "declaration": true, + "skipLibCheck": true + }, + "include": ["resources/**/*", "schemas/**/*"] }