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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 10 additions & 18 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -26,31 +24,25 @@ 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

lint:
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
12 changes: 5 additions & 7 deletions config.yaml
Original file line number Diff line number Diff line change
@@ -1,17 +1,15 @@
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

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.
47 changes: 34 additions & 13 deletions resources/IngestEvents.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -14,15 +14,15 @@
*
* 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 { Resource, databases } from "harperdb";
import { createPublicKey, verify } from "node:crypto";

const BATCH_LIMIT = 100;
Expand All @@ -44,8 +44,11 @@ interface AgentStatus {
agentId: string;
name?: string;
role?: string;
type?: string;
status?: string;
activity?: string;
model?: string;
currentTask?: string;
lastSeen?: string;
}

Expand Down Expand Up @@ -88,9 +91,23 @@ function verifyEd25519Signature(
}

export class IngestEvents extends Resource {
// 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) {
const request = (this as any).request;
const authHeader: string | undefined = request?.headers?.get?.("authorization") ?? request?.headers?.authorization;
// 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;
Expand All @@ -106,9 +123,9 @@ export class IngestEvents extends Resource {
}

// Look up the office
const office = await (tables as any).ObsOffice.get(officeId).catch(() => null);
const office = await (databases as any).roster.Office.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" } });
return new Response(JSON.stringify({ error: "office not registered — POST /Office first" }), { status: 403, headers: { "Content-Type": "application/json" } });
}

// Verify Ed25519 signature
Expand Down Expand Up @@ -136,14 +153,17 @@ export class IngestEvents extends Resource {
for (const agent of agents) {
if (!agent.agentId) continue;
const snapshotId = `${officeId}:${agent.agentId}`;
await (tables as any).ObsAgentSnapshot.put({
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,
Expand All @@ -155,16 +175,17 @@ export class IngestEvents extends Resource {
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);
const existing = await (databases as any).roster.Event.get(feedId).catch(() => null);
if (existing) continue;
await (tables as any).ObsEventFeed.put({
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,
Expand All @@ -173,7 +194,7 @@ export class IngestEvents extends Resource {
}

// Update office lastSeen + agentCount
await (tables as any).ObsOffice.put({
await (databases as any).roster.Office.put({
...office,
status: "online",
lastSeen: now,
Expand Down
49 changes: 49 additions & 0 deletions resources/RosterView.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
import { Resource, databases } from "harperdb";

/**
* GET /RosterView — public-safe snapshot of the roster directory for the
* Office Space (observatory) renderer.
*
* Auth: PUBLIC (allowRead). This is the single public 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<string, unknown> {
const out: Record<string, unknown> = {};
for (const k of keys) if (obj?.[k] !== undefined) out[k] = obj[k];
return out;
}

async function scan(table: any): Promise<any[]> {
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 {
allowRead() { return true; }

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 };
}
}
13 changes: 8 additions & 5 deletions schemas/schema.graphql
Original file line number Diff line number Diff line change
@@ -1,37 +1,40 @@
type ObsOffice @table @export {
type Office @table(database: "roster") @export {
id: ID @primaryKey # e.g. "rockit"
name: String!
publicKey: String! # Ed25519 hex (32 bytes)
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 {
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"
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 {
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
Expand Down
75 changes: 75 additions & 0 deletions scripts/push-roster.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
// 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}`);
}
Loading