diff --git a/CHANGELOG.md b/CHANGELOG.md index e0366aa..61e1639 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added +- **Phase N / N1 channel roster primitives** — `@murmurv2/core` now exposes `ChannelRosterStore` plus typed `ChannelRecord` / `ChannelMemberRecord` APIs. The roster keeps `channelId` distinct from legacy `conversationId`, stores `channels` / `channel_members` in a dedicated SQLite store, preserves existing message-history APIs, and reserves member-level `personaId`, `model`, `baseInstructionsHash`, and `eligibility` fields for N2 addressing and N3 personality binding. +- **Phase N / N2 addressing policy primitive** — `ChannelRosterStore.evaluateAddressing()` returns a shared reject/append/wake decision for `channelId` + explicit addressee flows: legacy no-channel remains broadcast, non-members are rejected, addressed members wake, and observers append history while staying muted. +- **Phase N / N3 personality binding** — `buildChannelThreadStartBinding()` projects a `ChannelMemberRecord` into Codex app-server `thread/start` overrides (`model`, `personality`, optional `baseInstructions`, and audit metadata). Daemon wiring is opt-in only (`channelRoster.enabled` or `MURMUR_CHANNEL_ROSTER=1`) and leaves legacy wake behavior unchanged by default. +- **Phase N / N6 MCP roster surface** — `@murmurv2/mcp-server` exposes `channel_create`, `channel_list`, `channel_members`, and `channel_evaluate_addressing`, backed by `MURMUR_CHANNEL_ROSTER_PATH` (default `DATA_DIR/channel-roster.db`) so agents and UI can manage rosters without direct SQLite access. + ### Pending - **Auth enforcement end-to-end** — the broker ingress hook + `authorizeInbound` exist; the daemon does not yet wire them (so `MURMUR_ENFORCE_AUTH` is not enforced end-to-end). Requires daemon roster/identity wiring + token provisioning. diff --git a/README.md b/README.md index 9b6a361..65364a3 100644 --- a/README.md +++ b/README.md @@ -52,6 +52,7 @@ A **murmuration** is one of nature's most extraordinary phenomena — thousands - **Scoped Channels & Session Affinity.** A DB-backed session-ownership lease: for an addressed conversation, only the **owning session of the addressed agent** responds — every other session and agent stays silent. Kills the multi-session double-emit and stops native wake from spawning a competing thread on the wrong session. Behind `MURMUR_SCOPED_CHANNELS` (default-OFF → fully backward-compatible). - **`SessionLeaseStore` — atomic single-owner claim.** Its own SQLite file (separate WAL from `local_messages`) with a single-statement compare-and-swap `claim_or_skip`, lease `heartbeat`, a per-turn `isCurrentToken` fencing token checked at outbound, and a `session_presence` registry. Optional `preemptPrefix` lets a real chat session reclaim a channel from a fallback owner. - **Native wake is now a lease-gated fallback.** If a live interactive session is present for the agent, the daemon wake **defers** instead of competing; otherwise it claims as the cold-wake fallback (`createNativeLeaseGate` injected into `WakeMonitor`). With the flag off, `WakeMonitor` behaves exactly as before. +- **Phase N channel roster, addressing, and personalities.** `ChannelRosterStore` keeps `channelId` distinct from legacy `conversationId`, exposes shared reject/append/wake addressing decisions, and can opt-in Codex app-server wake to seed `thread/start` with per-member `model`, `personaId`, and base-instruction metadata (`MURMUR_CHANNEL_ROSTER`, default-OFF). - **Verified: one claim across all delivery paths.** Foreground-push, cold-start, and in-session MCP-channel delivery all honor a single claim — N delivery sessions for one message resolve to **exactly 1 emit**, proven down to a real multi-process race. ## What's New in v2.3 @@ -478,6 +479,7 @@ See [protocol-v1.md](docs/protocol-v1.md) for the full specification. - [x] MCP server with 7 tools — full agent integration - [x] Native wake (live session) — Claude asyncRewake + Codex app-server UDS, with self-healing thread re-seed (`WakeMonitor`) - [x] **Scoped channels & session affinity** (v2.4) — DB-backed session-ownership lease: for an addressed conversation only the **owning session of the addressed agent** responds; native wake is demoted to a presence-deferring fallback (no competing thread). N delivery sessions → **exactly 1 emit** (live-verified). Behind `MURMUR_SCOPED_CHANNELS` (default-OFF). Lease ships in `@murmurv2/core`; delivery helpers and the cold-start spawn-on-inbound path are repo-shipped (`scripts/codex-murmur-*`) +- [x] **Phase N / N1-N3 + N6 channel roster, addressing, personalities, MCP** — typed `ChannelRosterStore` in `@murmurv2/core`: `channelId` is a routing/personality primitive distinct from legacy `conversationId`, with `channels` / `channel_members` in a dedicated SQLite store, shared `evaluateAddressing()` decisions for reject/append/wake gating, MCP roster tools, and opt-in Codex app-server `thread/start` binding for per-member `personaId`, `model`, and base-instruction metadata. - [x] Telegram/Discord/WhatsApp notification adapters - [x] Observability dashboard (real-time flow + 3D) + Prometheus metrics exporter (outbox depth, delivery latency, error rates) - [x] Reference deployment — Systemd + Docker, docker-compose (`deploy/docker-compose.messaging.yml`) + Kubernetes manifests (`deploy/kubernetes/`) diff --git a/package-lock.json b/package-lock.json index a7c043d..f427114 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "mur-mur-v2", - "version": "2.3.0", + "version": "2.4.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "mur-mur-v2", - "version": "2.3.0", + "version": "2.4.0", "workspaces": [ "packages/*" ], @@ -1787,6 +1787,15 @@ "@types/express": "^5.0.0" } }, + "packages/bridge-a2a/node_modules/@murmurv2/core": { + "version": "0.2.0", + "resolved": "https://registry.npmjs.org/@murmurv2/core/-/core-0.2.0.tgz", + "integrity": "sha512-ldczX36kucZCajlo4kYfR8MJAlixpITAA629iXbMSNRp5ylaTCdIXm+xoM5vAEBnV3+bE4Eu+sMCydOJWgo8kw==", + "license": "MIT", + "dependencies": { + "pg": "^8.16.3" + } + }, "packages/bridge-murmur": { "name": "@murmurv2/bridge-murmur", "version": "0.1.0", @@ -1801,6 +1810,15 @@ "nats": "^2.28.2" } }, + "packages/bridge-openclaw/node_modules/@murmurv2/core": { + "version": "0.2.0", + "resolved": "https://registry.npmjs.org/@murmurv2/core/-/core-0.2.0.tgz", + "integrity": "sha512-ldczX36kucZCajlo4kYfR8MJAlixpITAA629iXbMSNRp5ylaTCdIXm+xoM5vAEBnV3+bE4Eu+sMCydOJWgo8kw==", + "license": "MIT", + "dependencies": { + "pg": "^8.16.3" + } + }, "packages/bridge-telegram": { "name": "@murmurv2/bridge-telegram", "version": "0.1.0", @@ -1810,6 +1828,15 @@ "@murmurv2/security": "^0.1.0" } }, + "packages/bridge-telegram/node_modules/@murmurv2/core": { + "version": "0.2.0", + "resolved": "https://registry.npmjs.org/@murmurv2/core/-/core-0.2.0.tgz", + "integrity": "sha512-ldczX36kucZCajlo4kYfR8MJAlixpITAA629iXbMSNRp5ylaTCdIXm+xoM5vAEBnV3+bE4Eu+sMCydOJWgo8kw==", + "license": "MIT", + "dependencies": { + "pg": "^8.16.3" + } + }, "packages/broker-nats": { "name": "@murmurv2/broker-nats", "version": "0.2.0", @@ -1819,6 +1846,15 @@ "nats": "^2.28.2" } }, + "packages/broker-nats/node_modules/@murmurv2/core": { + "version": "0.2.0", + "resolved": "https://registry.npmjs.org/@murmurv2/core/-/core-0.2.0.tgz", + "integrity": "sha512-ldczX36kucZCajlo4kYfR8MJAlixpITAA629iXbMSNRp5ylaTCdIXm+xoM5vAEBnV3+bE4Eu+sMCydOJWgo8kw==", + "license": "MIT", + "dependencies": { + "pg": "^8.16.3" + } + }, "packages/broker-ws": { "name": "@murmurv2/broker-ws", "version": "0.1.0", @@ -1828,9 +1864,18 @@ "ws": "^8.21.0" } }, + "packages/broker-ws/node_modules/@murmurv2/core": { + "version": "0.2.0", + "resolved": "https://registry.npmjs.org/@murmurv2/core/-/core-0.2.0.tgz", + "integrity": "sha512-ldczX36kucZCajlo4kYfR8MJAlixpITAA629iXbMSNRp5ylaTCdIXm+xoM5vAEBnV3+bE4Eu+sMCydOJWgo8kw==", + "license": "MIT", + "dependencies": { + "pg": "^8.16.3" + } + }, "packages/core": { "name": "@murmurv2/core", - "version": "0.2.0", + "version": "0.3.0", "license": "MIT", "dependencies": { "pg": "^8.16.3" @@ -1857,13 +1902,22 @@ "nats": "^2.28.2" } }, + "packages/federation-nats/node_modules/@murmurv2/core": { + "version": "0.2.0", + "resolved": "https://registry.npmjs.org/@murmurv2/core/-/core-0.2.0.tgz", + "integrity": "sha512-ldczX36kucZCajlo4kYfR8MJAlixpITAA629iXbMSNRp5ylaTCdIXm+xoM5vAEBnV3+bE4Eu+sMCydOJWgo8kw==", + "license": "MIT", + "dependencies": { + "pg": "^8.16.3" + } + }, "packages/mcp-server": { "name": "@murmurv2/mcp-server", "version": "0.1.0", "license": "MIT", "dependencies": { "@murmurv2/broker-nats": "^0.2.0", - "@murmurv2/core": "^0.2.0", + "@murmurv2/core": "^0.3.0", "@murmurv2/security": "^0.1.0" } }, diff --git a/packages/core/src/channel.ts b/packages/core/src/channel.ts new file mode 100644 index 0000000..bd0e08d --- /dev/null +++ b/packages/core/src/channel.ts @@ -0,0 +1,566 @@ +// Murmur Phase N / N1 — typed channel roster. +// This is intentionally separate from local_messages: conversationId remains a history label, +// while channelId is the stable routing/personality primitive for N2/N3. +import { mkdirSync } from "node:fs"; +import { dirname } from "node:path"; +import { DatabaseSync } from "node:sqlite"; + +export type ChannelType = "dm" | "group" | "consult"; + +export interface ChannelRecord { + channelId: string; + conversationId: string; + type: ChannelType; + createdAt: string; + closedAt?: string; + metadata: Record; +} + +export interface ChannelMemberRecord { + channelId: string; + memberId: string; + memberSlot?: string; + agentId: string; + role?: string; + personaId?: string; + model?: string; + baseInstructionsHash?: string; + joinedAt: string; + leftAt?: string; + eligibility: Record; + metadata: Record; +} + +export interface ChannelMemberInput { + memberId: string; + memberSlot?: string; + agentId: string; + role?: string; + personaId?: string; + model?: string; + baseInstructionsHash?: string; + joinedAt?: string; + leftAt?: string; + eligibility?: Record; + metadata?: Record; +} + +export interface CreateChannelInput { + channelId: string; + conversationId: string; + type: ChannelType; + createdAt?: string; + closedAt?: string; + metadata?: Record; + members?: ChannelMemberInput[]; +} + +export type ChannelAddressingReason = + | "legacy-no-channel" + | "channel-not-found" + | "channel-closed" + | "self-not-member" + | "sender-not-member" + | "channel-broadcast" + | "addressee-not-member" + | "addressed-member" + | "observer-muted"; + +export interface ChannelAddressingInput { + channelId?: string; + selfAgentId: string; + senderAgentId?: string; + addresseeMemberId?: string; + addresseeAgentId?: string; +} + +export interface ChannelAddressingDecision { + allowAppend: boolean; + allowWake: boolean; + reject: boolean; + reason: ChannelAddressingReason; + channel?: ChannelRecord; + selfMember?: ChannelMemberRecord; + senderMember?: ChannelMemberRecord; + addresseeMember?: ChannelMemberRecord; +} + +export interface ChannelThreadStartBindingInput { + member: ChannelMemberRecord; + baseInstructions?: string | null; +} + +export interface ChannelThreadStartBinding { + model: string | null; + personality: string | null; + baseInstructions: string | null; + metadata: { + murmur_channel_id: string; + murmur_member_id: string; + murmur_agent_id: string; + murmur_member_slot?: string; + murmur_persona_id?: string; + murmur_model?: string; + murmur_base_instructions_hash?: string; + }; +} + +export const buildChannelThreadStartBinding = ({ member, baseInstructions = null }: ChannelThreadStartBindingInput): ChannelThreadStartBinding => ({ + model: member.model ?? null, + personality: member.personaId ?? null, + baseInstructions, + metadata: { + murmur_channel_id: member.channelId, + murmur_member_id: member.memberId, + murmur_agent_id: member.agentId, + ...(member.memberSlot ? { murmur_member_slot: member.memberSlot } : {}), + ...(member.personaId ? { murmur_persona_id: member.personaId } : {}), + ...(member.model ? { murmur_model: member.model } : {}), + ...(member.baseInstructionsHash ? { murmur_base_instructions_hash: member.baseInstructionsHash } : {}), + }, +}); + +const DDL = ` + PRAGMA journal_mode=WAL; + PRAGMA busy_timeout=10000; + PRAGMA foreign_keys=ON; + + CREATE TABLE IF NOT EXISTS channels ( + channel_id TEXT PRIMARY KEY, + conversation_id TEXT NOT NULL, + type TEXT NOT NULL CHECK (type IN ('dm', 'group', 'consult')), + created_at TEXT NOT NULL, + closed_at TEXT, + metadata_json TEXT NOT NULL DEFAULT '{}' + ); + CREATE INDEX IF NOT EXISTS idx_channels_conversation ON channels(conversation_id); + CREATE INDEX IF NOT EXISTS idx_channels_type ON channels(type); + + CREATE TABLE IF NOT EXISTS channel_members ( + channel_id TEXT NOT NULL, + member_id TEXT NOT NULL, + member_slot TEXT, + agent_id TEXT NOT NULL, + role TEXT, + persona_id TEXT, + model TEXT, + base_instructions_hash TEXT, + eligibility_json TEXT NOT NULL DEFAULT '{}', + joined_at TEXT NOT NULL, + left_at TEXT, + metadata_json TEXT NOT NULL DEFAULT '{}', + PRIMARY KEY (channel_id, member_id), + FOREIGN KEY (channel_id) REFERENCES channels(channel_id) ON DELETE CASCADE + ); + CREATE INDEX IF NOT EXISTS idx_channel_members_agent ON channel_members(agent_id); + CREATE INDEX IF NOT EXISTS idx_channel_members_slot ON channel_members(channel_id, member_slot); + CREATE INDEX IF NOT EXISTS idx_channel_members_persona ON channel_members(persona_id); +`; + +const jsonObject = (value?: Record): string => JSON.stringify(value ?? {}); + +const parseJsonObject = (raw: unknown): Record => { + if (typeof raw !== "string" || raw.length === 0) return {}; + try { + const parsed = JSON.parse(raw) as unknown; + return parsed && typeof parsed === "object" && !Array.isArray(parsed) ? parsed as Record : {}; + } catch { + return {}; + } +}; + +export class ChannelRosterStore { + private readonly db: DatabaseSync; + + constructor(dbPath = ".data/channel-roster.db") { + mkdirSync(dirname(dbPath), { recursive: true }); + this.db = new DatabaseSync(dbPath); + this.db.exec(DDL); + } + + createChannel(input: CreateChannelInput): ChannelRecord { + const createdAt = input.createdAt ?? new Date().toISOString(); + this.db.exec("BEGIN"); + try { + this.db + .prepare( + `INSERT INTO channels (channel_id, conversation_id, type, created_at, closed_at, metadata_json) + VALUES (?, ?, ?, ?, ?, ?)`, + ) + .run(input.channelId, input.conversationId, input.type, createdAt, input.closedAt ?? null, jsonObject(input.metadata)); + + for (const member of input.members ?? []) { + this.insertChannelMember(input.channelId, member, createdAt, false); + } + + this.db.exec("COMMIT"); + } catch (err) { + this.db.exec("ROLLBACK"); + throw err; + } + + const channel = this.getChannel(input.channelId); + if (!channel) throw new Error(`channel create failed: ${input.channelId}`); + return channel; + } + + getChannel(channelId: string): ChannelRecord | null { + const row = this.db + .prepare( + `SELECT + channel_id as channelId, + conversation_id as conversationId, + type, + created_at as createdAt, + closed_at as closedAt, + metadata_json as metadataJson + FROM channels + WHERE channel_id = ?`, + ) + .get(channelId) as + | { channelId: string; conversationId: string; type: ChannelType; createdAt: string; closedAt?: string | null; metadataJson: string } + | undefined; + return row ? this.toChannelRecord(row) : null; + } + + listChannelsForConversation(conversationId: string): ChannelRecord[] { + const rows = this.db + .prepare( + `SELECT + channel_id as channelId, + conversation_id as conversationId, + type, + created_at as createdAt, + closed_at as closedAt, + metadata_json as metadataJson + FROM channels + WHERE conversation_id = ? + ORDER BY created_at ASC`, + ) + .all(conversationId) as Array<{ channelId: string; conversationId: string; type: ChannelType; createdAt: string; closedAt?: string | null; metadataJson: string }>; + return rows.map((row) => this.toChannelRecord(row)); + } + + listChannelMembers(channelId: string): ChannelMemberRecord[] { + const rows = this.db + .prepare( + `SELECT + channel_id as channelId, + member_id as memberId, + member_slot as memberSlot, + agent_id as agentId, + role, + persona_id as personaId, + model, + base_instructions_hash as baseInstructionsHash, + eligibility_json as eligibilityJson, + joined_at as joinedAt, + left_at as leftAt, + metadata_json as metadataJson + FROM channel_members + WHERE channel_id = ? + ORDER BY member_id ASC`, + ) + .all(channelId) as Array<{ + channelId: string; + memberId: string; + memberSlot?: string | null; + agentId: string; + role?: string | null; + personaId?: string | null; + model?: string | null; + baseInstructionsHash?: string | null; + eligibilityJson: string; + joinedAt: string; + leftAt?: string | null; + metadataJson: string; + }>; + return rows.map((row) => this.toChannelMemberRecord(row)); + } + + getChannelMember(channelId: string, memberId: string): ChannelMemberRecord | null { + const row = this.db + .prepare( + `SELECT + channel_id as channelId, + member_id as memberId, + member_slot as memberSlot, + agent_id as agentId, + role, + persona_id as personaId, + model, + base_instructions_hash as baseInstructionsHash, + eligibility_json as eligibilityJson, + joined_at as joinedAt, + left_at as leftAt, + metadata_json as metadataJson + FROM channel_members + WHERE channel_id = ? AND member_id = ?`, + ) + .get(channelId, memberId) as + | { + channelId: string; + memberId: string; + memberSlot?: string | null; + agentId: string; + role?: string | null; + personaId?: string | null; + model?: string | null; + baseInstructionsHash?: string | null; + eligibilityJson: string; + joinedAt: string; + leftAt?: string | null; + metadataJson: string; + } + | undefined; + return row ? this.toChannelMemberRecord(row) : null; + } + + findActiveMembersForAgent(agentId: string): ChannelMemberRecord[] { + const rows = this.db + .prepare( + `SELECT + channel_id as channelId, + member_id as memberId, + member_slot as memberSlot, + agent_id as agentId, + role, + persona_id as personaId, + model, + base_instructions_hash as baseInstructionsHash, + eligibility_json as eligibilityJson, + joined_at as joinedAt, + left_at as leftAt, + metadata_json as metadataJson + FROM channel_members + WHERE agent_id = ? AND left_at IS NULL + ORDER BY joined_at ASC`, + ) + .all(agentId) as Array<{ + channelId: string; + memberId: string; + memberSlot?: string | null; + agentId: string; + role?: string | null; + personaId?: string | null; + model?: string | null; + baseInstructionsHash?: string | null; + eligibilityJson: string; + joinedAt: string; + leftAt?: string | null; + metadataJson: string; + }>; + return rows.map((row) => this.toChannelMemberRecord(row)); + } + + findActiveChannelMemberForAgent(channelId: string, agentId: string): ChannelMemberRecord | null { + const row = this.db + .prepare( + `SELECT + channel_id as channelId, + member_id as memberId, + member_slot as memberSlot, + agent_id as agentId, + role, + persona_id as personaId, + model, + base_instructions_hash as baseInstructionsHash, + eligibility_json as eligibilityJson, + joined_at as joinedAt, + left_at as leftAt, + metadata_json as metadataJson + FROM channel_members + WHERE channel_id = ? AND agent_id = ? AND left_at IS NULL + ORDER BY joined_at ASC + LIMIT 1`, + ) + .get(channelId, agentId) as + | { + channelId: string; + memberId: string; + memberSlot?: string | null; + agentId: string; + role?: string | null; + personaId?: string | null; + model?: string | null; + baseInstructionsHash?: string | null; + eligibilityJson: string; + joinedAt: string; + leftAt?: string | null; + metadataJson: string; + } + | undefined; + return row ? this.toChannelMemberRecord(row) : null; + } + + isChannelMember(channelId: string, agentId: string): boolean { + const row = this.db + .prepare(`SELECT 1 FROM channel_members WHERE channel_id = ? AND agent_id = ? AND left_at IS NULL LIMIT 1`) + .get(channelId, agentId) as { "1": number } | undefined; + return row !== undefined; + } + + evaluateAddressing(input: ChannelAddressingInput): ChannelAddressingDecision { + if (!input.channelId) { + return { allowAppend: true, allowWake: true, reject: false, reason: "legacy-no-channel" }; + } + + const channel = this.getChannel(input.channelId); + if (!channel) { + return { allowAppend: false, allowWake: false, reject: true, reason: "channel-not-found" }; + } + if (channel.closedAt) { + return { allowAppend: false, allowWake: false, reject: true, reason: "channel-closed", channel }; + } + + const selfMember = this.findActiveChannelMemberForAgent(input.channelId, input.selfAgentId); + if (!selfMember) { + return { allowAppend: false, allowWake: false, reject: true, reason: "self-not-member", channel }; + } + + const senderMember = input.senderAgentId ? this.findActiveChannelMemberForAgent(input.channelId, input.senderAgentId) : undefined; + if (input.senderAgentId && !senderMember) { + return { allowAppend: false, allowWake: false, reject: true, reason: "sender-not-member", channel, selfMember }; + } + + if (!input.addresseeMemberId && !input.addresseeAgentId) { + return { allowAppend: true, allowWake: true, reject: false, reason: "channel-broadcast", channel, selfMember, ...(senderMember ? { senderMember } : {}) }; + } + + const addresseeMember = input.addresseeMemberId + ? this.getChannelMember(input.channelId, input.addresseeMemberId) + : this.findActiveChannelMemberForAgent(input.channelId, input.addresseeAgentId ?? ""); + if (!addresseeMember || addresseeMember.leftAt) { + return { allowAppend: false, allowWake: false, reject: true, reason: "addressee-not-member", channel, selfMember, ...(senderMember ? { senderMember } : {}) }; + } + + const isAddressed = addresseeMember.agentId === input.selfAgentId; + return { + allowAppend: true, + allowWake: isAddressed, + reject: false, + reason: isAddressed ? "addressed-member" : "observer-muted", + channel, + selfMember, + ...(senderMember ? { senderMember } : {}), + addresseeMember, + }; + } + + upsertChannelMember(channelId: string, member: ChannelMemberInput): ChannelMemberRecord { + this.insertChannelMember(channelId, member, new Date().toISOString(), true); + const found = this.getChannelMember(channelId, member.memberId); + if (!found) throw new Error(`channel member upsert failed: ${channelId}/${member.memberId}`); + return found; + } + + closeChannel(channelId: string, closedAt = new Date().toISOString()): boolean { + const result = this.db + .prepare(`UPDATE channels SET closed_at = ? WHERE channel_id = ? AND closed_at IS NULL`) + .run(closedAt, channelId); + return result.changes > 0; + } + + close(): void { + this.db.close(); + } + + private insertChannelMember(channelId: string, member: ChannelMemberInput, defaultJoinedAt: string, upsert: boolean): void { + const joinedAt = member.joinedAt ?? defaultJoinedAt; + if (upsert) { + this.db + .prepare( + `INSERT INTO channel_members + (channel_id, member_id, member_slot, agent_id, role, persona_id, model, base_instructions_hash, eligibility_json, joined_at, left_at, metadata_json) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(channel_id, member_id) DO UPDATE SET + member_slot = excluded.member_slot, + agent_id = excluded.agent_id, + role = excluded.role, + persona_id = excluded.persona_id, + model = excluded.model, + base_instructions_hash = excluded.base_instructions_hash, + eligibility_json = excluded.eligibility_json, + left_at = excluded.left_at, + metadata_json = excluded.metadata_json`, + ) + .run( + channelId, + member.memberId, + member.memberSlot ?? null, + member.agentId, + member.role ?? null, + member.personaId ?? null, + member.model ?? null, + member.baseInstructionsHash ?? null, + jsonObject(member.eligibility), + joinedAt, + member.leftAt ?? null, + jsonObject(member.metadata), + ); + return; + } + + this.db + .prepare( + `INSERT INTO channel_members + (channel_id, member_id, member_slot, agent_id, role, persona_id, model, base_instructions_hash, eligibility_json, joined_at, left_at, metadata_json) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + ) + .run( + channelId, + member.memberId, + member.memberSlot ?? null, + member.agentId, + member.role ?? null, + member.personaId ?? null, + member.model ?? null, + member.baseInstructionsHash ?? null, + jsonObject(member.eligibility), + joinedAt, + member.leftAt ?? null, + jsonObject(member.metadata), + ); + } + + private toChannelRecord(row: { channelId: string; conversationId: string; type: ChannelType; createdAt: string; closedAt?: string | null; metadataJson: string }): ChannelRecord { + return { + channelId: row.channelId, + conversationId: row.conversationId, + type: row.type, + createdAt: row.createdAt, + ...(row.closedAt ? { closedAt: row.closedAt } : {}), + metadata: parseJsonObject(row.metadataJson), + }; + } + + private toChannelMemberRecord(row: { + channelId: string; + memberId: string; + memberSlot?: string | null; + agentId: string; + role?: string | null; + personaId?: string | null; + model?: string | null; + baseInstructionsHash?: string | null; + eligibilityJson: string; + joinedAt: string; + leftAt?: string | null; + metadataJson: string; + }): ChannelMemberRecord { + return { + channelId: row.channelId, + memberId: row.memberId, + ...(row.memberSlot ? { memberSlot: row.memberSlot } : {}), + agentId: row.agentId, + ...(row.role ? { role: row.role } : {}), + ...(row.personaId ? { personaId: row.personaId } : {}), + ...(row.model ? { model: row.model } : {}), + ...(row.baseInstructionsHash ? { baseInstructionsHash: row.baseInstructionsHash } : {}), + joinedAt: row.joinedAt, + ...(row.leftAt ? { leftAt: row.leftAt } : {}), + eligibility: parseJsonObject(row.eligibilityJson), + metadata: parseJsonObject(row.metadataJson), + }; + } +} diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index c639934..af54641 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -9,6 +9,9 @@ export * from "./discovery.js"; // Scoped channels — session-ownership lease (session affinity, fencing token, presence-deferring wake gate) export * from "./lease.js"; +// Phase N — typed channel roster (channelId distinct from legacy conversationId) +export * from "./channel.js"; + export type DeliveryMode = "at-least-once"; export interface EnvelopeV1 { diff --git a/packages/mcp-server/package.json b/packages/mcp-server/package.json index 7f322ec..bd76961 100644 --- a/packages/mcp-server/package.json +++ b/packages/mcp-server/package.json @@ -10,7 +10,7 @@ }, "dependencies": { "@murmurv2/broker-nats": "^0.2.0", - "@murmurv2/core": "^0.2.0", + "@murmurv2/core": "^0.3.0", "@murmurv2/security": "^0.1.0" }, "license": "MIT", @@ -26,5 +26,5 @@ "dist/src", "LICENSE" ], - "description": "Murmur V2 MCP server \u2014 7 tools for agent-to-agent messaging over the Model Context Protocol." + "description": "Murmur V2 MCP server \u2014 agent-to-agent messaging and scoped-channel roster tools over the Model Context Protocol." } diff --git a/packages/mcp-server/src/index.ts b/packages/mcp-server/src/index.ts index 68e774e..647c14c 100644 --- a/packages/mcp-server/src/index.ts +++ b/packages/mcp-server/src/index.ts @@ -3,6 +3,7 @@ import { readFileSync } from "node:fs"; import path from "node:path"; import { createInterface } from "node:readline"; import { + ChannelRosterStore, SQLiteDedupeOutboxStore, SQLiteMessageStore, stableEnvelopePayload, @@ -44,6 +45,7 @@ interface AgentConfig { const dataDir = process.env.DATA_DIR || ".data"; const configPath = path.join(dataDir, "agent-config.json"); const dbPath = process.env.MURMUR_STORE_PATH ?? path.join(dataDir, "murmur.db"); +const channelRosterPath = process.env.MURMUR_CHANNEL_ROSTER_PATH ?? path.join(dataDir, "channel-roster.db"); let agentConfig: AgentConfig | null = null; try { @@ -53,6 +55,7 @@ try { } const store = new SQLiteMessageStore(dbPath); +const channelRoster = new ChannelRosterStore(channelRosterPath); // Outbox store — shared with daemon, only created if agent config exists let outbox: SQLiteDedupeOutboxStore | null = null; @@ -152,6 +155,60 @@ const handleTool = async (name: string, args: Record): Promise< return { messages: messages.map(asMessage) }; } + if (name === "channel_create") { + const channelId = String(args.channelId ?? "").trim(); + if (!channelId) throw new Error("channelId is required"); + const conversationId = String(args.conversationId ?? "").trim(); + if (!conversationId) throw new Error("conversationId is required"); + const type = String(args.type ?? "").trim(); + if (!["dm", "group", "consult"].includes(type)) throw new Error("type must be one of: dm, group, consult"); + const members = Array.isArray(args.members) ? args.members as Array> : []; + const channel = channelRoster.createChannel({ + channelId, + conversationId, + type: type as "dm" | "group" | "consult", + metadata: typeof args.metadata === "object" && args.metadata && !Array.isArray(args.metadata) ? args.metadata as Record : undefined, + members: members.map((member) => ({ + memberId: String(member.memberId ?? "").trim(), + memberSlot: member.memberSlot === undefined ? undefined : String(member.memberSlot), + agentId: String(member.agentId ?? "").trim(), + role: member.role === undefined ? undefined : String(member.role), + personaId: member.personaId === undefined ? undefined : String(member.personaId), + model: member.model === undefined ? undefined : String(member.model), + baseInstructionsHash: member.baseInstructionsHash === undefined ? undefined : String(member.baseInstructionsHash), + eligibility: typeof member.eligibility === "object" && member.eligibility && !Array.isArray(member.eligibility) ? member.eligibility as Record : undefined, + metadata: typeof member.metadata === "object" && member.metadata && !Array.isArray(member.metadata) ? member.metadata as Record : undefined, + })), + }); + return { channel, members: channelRoster.listChannelMembers(channelId) }; + } + + if (name === "channel_list") { + const conversationId = String(args.conversationId ?? "").trim(); + if (!conversationId) throw new Error("conversationId is required"); + return { channels: channelRoster.listChannelsForConversation(conversationId) }; + } + + if (name === "channel_members") { + const channelId = String(args.channelId ?? "").trim(); + if (!channelId) throw new Error("channelId is required"); + return { members: channelRoster.listChannelMembers(channelId) }; + } + + if (name === "channel_evaluate_addressing") { + const selfAgentId = String(args.selfAgentId ?? "").trim(); + if (!selfAgentId) throw new Error("selfAgentId is required"); + return { + decision: channelRoster.evaluateAddressing({ + channelId: args.channelId === undefined ? undefined : String(args.channelId), + selfAgentId, + senderAgentId: args.senderAgentId === undefined ? undefined : String(args.senderAgentId), + addresseeMemberId: args.addresseeMemberId === undefined ? undefined : String(args.addresseeMemberId), + addresseeAgentId: args.addresseeAgentId === undefined ? undefined : String(args.addresseeAgentId), + }), + }; + } + // === New agent-to-agent tools (require agent config) === if (name === "murmur_send") { @@ -418,6 +475,75 @@ const tools = [ required: ["query"], }, }, + { + name: "channel_create", + description: "Create a typed Murmur channel roster entry with optional members.", + inputSchema: { + type: "object", + properties: { + channelId: { type: "string", description: "Stable channel/routing ID, distinct from conversationId" }, + conversationId: { type: "string", description: "Legacy history label associated with the channel" }, + type: { type: "string", enum: ["dm", "group", "consult"] }, + metadata: { type: "object" }, + members: { + type: "array", + items: { + type: "object", + properties: { + memberId: { type: "string" }, + memberSlot: { type: "string" }, + agentId: { type: "string" }, + role: { type: "string" }, + personaId: { type: "string" }, + model: { type: "string" }, + baseInstructionsHash: { type: "string" }, + eligibility: { type: "object" }, + metadata: { type: "object" }, + }, + required: ["memberId", "agentId"], + }, + }, + }, + required: ["channelId", "conversationId", "type"], + }, + }, + { + name: "channel_list", + description: "List typed channels associated with a legacy conversationId.", + inputSchema: { + type: "object", + properties: { + conversationId: { type: "string" }, + }, + required: ["conversationId"], + }, + }, + { + name: "channel_members", + description: "List active and historical members for a typed channel.", + inputSchema: { + type: "object", + properties: { + channelId: { type: "string" }, + }, + required: ["channelId"], + }, + }, + { + name: "channel_evaluate_addressing", + description: "Evaluate channel membership/addressing into reject/append/wake decisions.", + inputSchema: { + type: "object", + properties: { + channelId: { type: "string" }, + selfAgentId: { type: "string" }, + senderAgentId: { type: "string" }, + addresseeMemberId: { type: "string" }, + addresseeAgentId: { type: "string" }, + }, + required: ["selfAgentId"], + }, + }, { name: "murmur_send", description: diff --git a/scripts/codex-app-server-wake.mjs b/scripts/codex-app-server-wake.mjs index 41fb8b0..31d4d90 100644 --- a/scripts/codex-app-server-wake.mjs +++ b/scripts/codex-app-server-wake.mjs @@ -1,4 +1,5 @@ import WebSocket from "ws"; +import { buildChannelThreadStartBinding } from "@murmurv2/core"; const DEFAULT_TIMEOUT_MS = 10000; const INITIALIZE_TIMEOUT_MS = 10000; @@ -25,8 +26,8 @@ export const buildTurnStartRequest = ({ id = 1, threadId, text, metadata = {} }) }, }); -const buildThreadStartParams = () => ({ - model: null, +export const buildThreadStartParams = (binding = null) => ({ + model: binding?.model ?? null, modelProvider: null, cwd: null, runtimeWorkspaceRoots: null, @@ -36,9 +37,9 @@ const buildThreadStartParams = () => ({ permissions: null, config: null, serviceName: null, - baseInstructions: null, + baseInstructions: binding?.baseInstructions ?? null, developerInstructions: null, - personality: null, + personality: binding?.personality ?? null, ephemeral: false, sessionStartSource: null, threadSource: null, @@ -153,13 +154,52 @@ export class CodexAppServerClient { } } -export const createCodexAppServerInjector = ({ Client = CodexAppServerClient, log = () => {}, timeoutMs = DEFAULT_TIMEOUT_MS } = {}) => { +const pickSingleOpenChannel = (channels) => { + const open = (channels || []).filter((channel) => !channel.closedAt); + return open.length === 1 ? open[0] : null; +}; + +export const createChannelThreadStartBindingResolver = ({ rosterStore, agentId, baseInstructionsResolver = null, log = () => {} } = {}) => { + if (!rosterStore) return null; + return async (payload, peer = {}) => { + const channelId = payload?.channelId || peer.channelId || pickSingleOpenChannel(rosterStore.listChannelsForConversation(payload?.conversationId || ""))?.channelId; + if (!channelId) return null; + + const memberId = payload?.addresseeMemberId || peer.memberId; + const effectiveAgentId = agentId || peer.agentId; + const member = memberId + ? rosterStore.getChannelMember(channelId, memberId) + : effectiveAgentId + ? rosterStore.findActiveChannelMemberForAgent(channelId, effectiveAgentId) + : null; + if (!member || member.leftAt) return null; + + const baseInstructions = typeof baseInstructionsResolver === "function" + ? await baseInstructionsResolver(member, payload, peer) + : peer.baseInstructions ?? null; + const binding = buildChannelThreadStartBinding({ member, baseInstructions }); + log("info", "Codex app-server channel binding resolved", { + msgId: payload?.msgId, + channelId: member.channelId, + memberId: member.memberId, + agentId: member.agentId, + personaId: member.personaId ?? null, + model: member.model ?? null, + hasBaseInstructions: !!baseInstructions, + }); + return binding; + }; +}; + +export const createCodexAppServerInjector = ({ Client = CodexAppServerClient, log = () => {}, timeoutMs = DEFAULT_TIMEOUT_MS, resolveThreadStartBinding = null } = {}) => { return async (payload, peer) => { const socketPath = peer?.socketPath || peer?.target; if (!socketPath) throw new Error(`codex-app-server-socket-missing:${payload.from}`); const client = new Client({ socketPath, timeoutMs }); const text = buildCodexTurnText(payload); + const threadStartBinding = peer?.threadStartBinding ?? (resolveThreadStartBinding ? await resolveThreadStartBinding(payload, peer) : null); + const bindingMetadata = threadStartBinding?.metadata ?? {}; const startTurn = (threadId) => client.request("turn/start", { threadId, input: [{ type: "text", text, text_elements: [] }], @@ -167,12 +207,13 @@ export const createCodexAppServerInjector = ({ Client = CodexAppServerClient, lo murmur_msg_id: payload.msgId || "", murmur_conversation_id: payload.conversationId || "", murmur_from: payload.from || "", + ...bindingMetadata, }, }); let threadId = peer?.threadId; if (!threadId) { - const started = await client.request("thread/start", buildThreadStartParams()); + const started = await client.request("thread/start", buildThreadStartParams(threadStartBinding)); threadId = started?.thread?.id; if (!threadId) throw new Error(`codex-app-server-thread-start-missing:${payload.from}`); peer.threadId = threadId; @@ -185,7 +226,7 @@ export const createCodexAppServerInjector = ({ Client = CodexAppServerClient, lo } catch (err) { const e = err instanceof Error ? err : new Error(String(err)); if (!e.message.startsWith("codex-app-server-error:thread not found:")) throw e; - const started = await client.request("thread/start", buildThreadStartParams()); + const started = await client.request("thread/start", buildThreadStartParams(threadStartBinding)); threadId = started?.thread?.id; if (!threadId) throw new Error(`codex-app-server-thread-start-missing:${payload.from}`); peer.threadId = threadId; diff --git a/scripts/murmur-daemon.mjs b/scripts/murmur-daemon.mjs index 5cb9999..6b2c455 100644 --- a/scripts/murmur-daemon.mjs +++ b/scripts/murmur-daemon.mjs @@ -7,10 +7,10 @@ import { DatabaseSync } from "node:sqlite"; import path from "node:path"; import { setTimeout as sleep } from "node:timers/promises"; import { NatsBroker } from "@murmurv2/broker-nats"; -import { SQLiteDedupeOutboxStore, SQLiteMessageStore, stableEnvelopePayload } from "@murmurv2/core"; +import { ChannelRosterStore, SQLiteDedupeOutboxStore, SQLiteMessageStore, stableEnvelopePayload } from "@murmurv2/core"; import { decryptPayload, verifyEnvelopeSignature } from "@murmurv2/security"; import { NotifyQueue, flushNotifyQueue, normalizeNotifyTargets } from "./notify-router.mjs"; -import { createCodexAppServerInjector } from "./codex-app-server-wake.mjs"; +import { createChannelThreadStartBindingResolver, createCodexAppServerInjector } from "./codex-app-server-wake.mjs"; import { startJetStreamAdvisoryDlqIfEnabled } from "./murmur-jetstream-advisory.mjs"; import { WakeMonitor, createAuditShellHook, createShellHook, normalizeWakeConfig } from "./wake-monitor.mjs"; import { SessionLeaseStore, createNativeLeaseGate } from "./lease.mjs"; @@ -116,6 +116,16 @@ const nativeLeaseGate = leaseStore ? createNativeLeaseGate({ store: leaseStore, agentId, ttlMs: nativeLeaseTtlMs, log }) : null; if (scopedChannelsEnabled) log("info", "Scoped-channels native lease gate enabled", { ttlMs: nativeLeaseTtlMs }); + +const channelRosterConfig = config.channelRoster || {}; +const channelRosterEnabled = channelRosterConfig.enabled ?? process.env.MURMUR_CHANNEL_ROSTER === "1"; +const channelRosterPath = channelRosterConfig.path || process.env.MURMUR_CHANNEL_ROSTER_PATH || path.join(dataDir, "channel-roster.db"); +const channelRosterStore = channelRosterEnabled ? new ChannelRosterStore(channelRosterPath) : null; +const threadStartBindingResolver = channelRosterStore + ? createChannelThreadStartBindingResolver({ rosterStore: channelRosterStore, agentId, log }) + : null; +const codexAppServerInjector = createCodexAppServerInjector({ log, resolveThreadStartBinding: threadStartBindingResolver }); +if (channelRosterEnabled) log("info", "Channel roster thread-start binding enabled", { channelRosterPath }); const broker = new NatsBroker({ url: natsUrl, token: natsToken, @@ -182,7 +192,7 @@ const wakeMonitor = new WakeMonitor({ hook: createShellHook({ command: config.onReceive, log }), injector: async (payload, peer) => { if (peer.mode === "codex_app_server") { - return createCodexAppServerInjector({ log })(payload, peer); + return codexAppServerInjector(payload, peer); } throw new Error(`wake-native-mode-unsupported:${peer.mode}`); }, @@ -198,7 +208,7 @@ const proxyWakeMonitor = new WakeMonitor({ leaseGate: nativeLeaseGate, injector: async (payload, peer) => { if (peer.mode === "codex_app_server") { - return createCodexAppServerInjector({ log })(payload, peer); + return codexAppServerInjector(payload, peer); } throw new Error(`wake-native-mode-unsupported:${peer.mode}`); }, diff --git a/tests/codex-app-server-wake.test.mjs b/tests/codex-app-server-wake.test.mjs index edfb4c6..022213d 100644 --- a/tests/codex-app-server-wake.test.mjs +++ b/tests/codex-app-server-wake.test.mjs @@ -6,8 +6,16 @@ import test from "node:test"; import assert from "node:assert/strict"; import { once } from "node:events"; import { WebSocketServer } from "ws"; +import { ChannelRosterStore } from "../packages/core/dist/src/index.js"; import { WakeMonitor, normalizeWakeConfig } from "../scripts/wake-monitor.mjs"; -import { buildCodexTurnText, buildTurnStartRequest, CodexAppServerClient, createCodexAppServerInjector } from "../scripts/codex-app-server-wake.mjs"; +import { + buildCodexTurnText, + buildThreadStartParams, + buildTurnStartRequest, + CodexAppServerClient, + createChannelThreadStartBindingResolver, + createCodexAppServerInjector, +} from "../scripts/codex-app-server-wake.mjs"; const payload = { from: "agent-jarvis", @@ -52,6 +60,31 @@ test("buildTurnStartRequest builds Codex turn/start params", () => { assert.match(request.params.input[0].text, /msgId=msg-codex-1/); }); +test("buildThreadStartParams applies optional channel personality binding", () => { + const params = buildThreadStartParams({ + model: "gpt-5", + personality: "codex-writer", + baseInstructions: "Write concise engineering notes.", + metadata: { murmur_channel_id: "chan-1" }, + }); + + assert.equal(params.model, "gpt-5"); + assert.equal(params.personality, "codex-writer"); + assert.equal(params.baseInstructions, "Write concise engineering notes."); + assert.equal(params.modelProvider, null); + assert.equal(params.ephemeral, false); +}); + +test("buildThreadStartParams preserves legacy nulled defaults without binding", () => { + const params = buildThreadStartParams(); + + assert.equal(params.model, null); + assert.equal(params.personality, null); + assert.equal(params.baseInstructions, null); + assert.equal(params.modelProvider, null); + assert.equal(params.ephemeral, false); +}); + test("Codex app-server client initializes before turn/start over WS-over-UDS", async () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "murmur-codex-wake-")); const socketPath = path.join(dir, "codex.sock"); @@ -218,6 +251,95 @@ test("Codex app-server injector seeds missing app-server threads", async () => { assert.equal(calls[1].params.threadId, "fresh-thread"); }); +test("Codex app-server injector seeds thread with resolved channel member binding", async () => { + const calls = []; + class FakeClient { + async request(method, params) { + calls.push({ method, params }); + if (method === "thread/start") return { thread: { id: "fresh-thread" } }; + return { turn: { id: "turn-1" } }; + } + } + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "murmur-codex-roster-")); + const roster = new ChannelRosterStore(path.join(dir, "channel-roster.db")); + roster.createChannel({ + channelId: "chan:codex:writer", + conversationId: payload.conversationId, + type: "dm", + members: [{ + memberId: "codex-writer", + memberSlot: "agent-codex-volt:writer", + agentId: "agent-codex-volt", + personaId: "codex-writer", + model: "gpt-5", + baseInstructionsHash: "sha256:writer-v1", + }], + }); + const injector = createCodexAppServerInjector({ + Client: FakeClient, + resolveThreadStartBinding: createChannelThreadStartBindingResolver({ + rosterStore: roster, + agentId: "agent-codex-volt", + baseInstructionsResolver: () => "Write concise engineering notes.", + }), + }); + + const result = await injector(payload, { mode: "codex_app_server", socketPath: "/tmp/codex.sock" }); + + assert.deepEqual(result, { turn: { id: "turn-1" } }); + assert.equal(calls[0].method, "thread/start"); + assert.equal(calls[0].params.model, "gpt-5"); + assert.equal(calls[0].params.personality, "codex-writer"); + assert.equal(calls[0].params.baseInstructions, "Write concise engineering notes."); + assert.equal(calls[1].method, "turn/start"); + assert.equal(calls[1].params.responsesapiClientMetadata.murmur_channel_id, "chan:codex:writer"); + assert.equal(calls[1].params.responsesapiClientMetadata.murmur_member_id, "codex-writer"); + assert.equal(calls[1].params.responsesapiClientMetadata.murmur_base_instructions_hash, "sha256:writer-v1"); +}); + +test("channel thread-start binding resolver returns null without member or agent identity", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "murmur-codex-roster-no-id-")); + const roster = new ChannelRosterStore(path.join(dir, "channel-roster.db")); + roster.createChannel({ + channelId: "chan:codex:no-id", + conversationId: payload.conversationId, + type: "dm", + members: [{ memberId: "codex-writer", agentId: "agent-codex-volt" }], + }); + const resolveBinding = createChannelThreadStartBindingResolver({ rosterStore: roster }); + + const binding = await resolveBinding(payload, {}); + + assert.equal(binding, null); + roster.close(); +}); + +test("Codex app-server injector ignores remote payload thread-start binding", async () => { + const calls = []; + class FakeClient { + async request(method, params) { + calls.push({ method, params }); + if (method === "thread/start") return { thread: { id: "fresh-thread" } }; + return { turn: { id: "turn-1" } }; + } + } + const injector = createCodexAppServerInjector({ Client: FakeClient }); + + await injector({ + ...payload, + threadStartBinding: { + model: "remote-controlled-model", + personality: "remote-controlled-persona", + baseInstructions: "remote instructions", + }, + }, { mode: "codex_app_server", socketPath: "/tmp/codex.sock" }); + + assert.equal(calls[0].method, "thread/start"); + assert.equal(calls[0].params.model, null); + assert.equal(calls[0].params.personality, null); + assert.equal(calls[0].params.baseInstructions, null); +}); + test("Codex app-server injector fails loud without socket", async () => { const injector = createCodexAppServerInjector(); diff --git a/tests/core-channels-roster.test.mjs b/tests/core-channels-roster.test.mjs new file mode 100644 index 0000000..aac1a62 --- /dev/null +++ b/tests/core-channels-roster.test.mjs @@ -0,0 +1,318 @@ +import assert from "node:assert/strict"; +import { mkdtemp, rm } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import test from "node:test"; +import { buildChannelThreadStartBinding, ChannelRosterStore, SQLiteMessageStore } from "../packages/core/dist/src/index.js"; + +async function withRoster(fn) { + const dir = await mkdtemp(path.join(os.tmpdir(), "murmur-channels-")); + try { + const store = new ChannelRosterStore(path.join(dir, "channel-roster.db")); + return await fn(store); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +test("ChannelRosterStore creates typed channel roster with personality-ready member metadata", async () => { + await withRoster(async (store) => { + const channel = store.createChannel({ + channelId: "chan:research:alpha", + conversationId: "codex:task:alpha", + type: "group", + createdAt: "2026-06-24T10:00:00.000Z", + metadata: { topic: "phase-n" }, + members: [ + { + memberId: "owner", + memberSlot: "agent-codex-volt", + agentId: "agent-codex-volt", + role: "implementer", + personaId: "codex-senior-engineer", + model: "gpt-5", + baseInstructionsHash: "sha256:codex", + eligibility: { canAddress: true }, + metadata: { lane: "implementation" }, + }, + { + memberId: "reviewer", + memberSlot: "agent-jarvis", + agentId: "agent-jarvis", + role: "reviewer", + personaId: "jarvis-ops", + model: "claude", + }, + ], + }); + + assert.deepEqual(channel, { + channelId: "chan:research:alpha", + conversationId: "codex:task:alpha", + type: "group", + createdAt: "2026-06-24T10:00:00.000Z", + metadata: { topic: "phase-n" }, + }); + + const members = store.listChannelMembers("chan:research:alpha"); + assert.equal(members.length, 2); + assert.equal(members[0].memberId, "owner"); + assert.equal(members[0].memberSlot, "agent-codex-volt"); + assert.equal(members[0].personaId, "codex-senior-engineer"); + assert.equal(members[0].baseInstructionsHash, "sha256:codex"); + assert.deepEqual(members[0].eligibility, { canAddress: true }); + assert.deepEqual(members[0].metadata, { lane: "implementation" }); + assert.equal(store.isChannelMember("chan:research:alpha", "agent-codex-volt"), true); + assert.equal(store.isChannelMember("chan:research:alpha", "agent-unknown"), false); + }); +}); + +test("channelId is separate from conversationId and legacy conversation listing still works", async () => { + const dir = await mkdtemp(path.join(os.tmpdir(), "murmur-channels-legacy-")); + try { + const messages = new SQLiteMessageStore(path.join(dir, "murmur.db")); + const roster = new ChannelRosterStore(path.join(dir, "channel-roster.db")); + await messages.append({ + conversationId: "codex:task:legacy", + msgId: "msg-1", + direction: "inbound", + sender: "agent-jarvis", + text: "legacy message", + createdAt: "2026-06-24T10:01:00.000Z", + transport: "nats", + }); + roster.createChannel({ + channelId: "chan:dm:jarvis-codex", + conversationId: "codex:task:legacy", + type: "dm", + createdAt: "2026-06-24T10:02:00.000Z", + }); + + const conversations = await messages.listConversations(); + assert.equal(conversations.length, 1); + assert.equal(conversations[0].conversationId, "codex:task:legacy"); + assert.equal(conversations[0].messageCount, 1); + + const channels = roster.listChannelsForConversation("codex:task:legacy"); + assert.equal(channels.length, 1); + assert.equal(channels[0].channelId, "chan:dm:jarvis-codex"); + assert.equal(channels[0].conversationId, "codex:task:legacy"); + } finally { + await rm(dir, { recursive: true, force: true }); + } +}); + +test("channel members can be rebound by memberId without changing legacy messages", async () => { + await withRoster(async (store) => { + store.createChannel({ + channelId: "chan:consult:qa", + conversationId: "codex:task:qa", + type: "consult", + createdAt: "2026-06-24T10:03:00.000Z", + members: [{ memberId: "reviewer", memberSlot: "qa", agentId: "agent-stas", role: "qa" }], + }); + + const rebound = store.upsertChannelMember("chan:consult:qa", { + memberId: "reviewer", + memberSlot: "ops", + agentId: "agent-jarvis", + role: "ops-review", + personaId: "jarvis-readiness", + }); + + assert.equal(rebound.memberId, "reviewer"); + assert.equal(rebound.memberSlot, "ops"); + assert.equal(rebound.agentId, "agent-jarvis"); + assert.equal(rebound.personaId, "jarvis-readiness"); + assert.equal(store.isChannelMember("chan:consult:qa", "agent-stas"), false); + assert.equal(store.isChannelMember("chan:consult:qa", "agent-jarvis"), true); + }); +}); + +test("closing a channel is idempotent and preserves roster rows", async () => { + await withRoster(async (store) => { + store.createChannel({ + channelId: "chan:closed", + conversationId: "codex:task:closed", + type: "group", + createdAt: "2026-06-24T10:04:00.000Z", + members: [{ memberId: "owner", memberSlot: "agent-codex-volt", agentId: "agent-codex-volt" }], + }); + + assert.equal(store.closeChannel("chan:closed", "2026-06-24T10:05:00.000Z"), true); + assert.equal(store.closeChannel("chan:closed", "2026-06-24T10:06:00.000Z"), false); + + const channel = store.getChannel("chan:closed"); + assert.equal(channel.closedAt, "2026-06-24T10:05:00.000Z"); + const members = store.listChannelMembers("chan:closed"); + assert.equal(members.length, 1); + }); +}); + +test("addressing policy preserves legacy no-channel broadcast behavior", async () => { + await withRoster(async (store) => { + assert.deepEqual(store.evaluateAddressing({ selfAgentId: "agent-codex-volt" }), { + allowAppend: true, + allowWake: true, + reject: false, + reason: "legacy-no-channel", + }); + }); +}); + +test("addressing policy rejects channel traffic from or to non-members", async () => { + await withRoster(async (store) => { + store.createChannel({ + channelId: "chan:members-only", + conversationId: "codex:task:members-only", + type: "group", + members: [ + { memberId: "codex", agentId: "agent-codex-volt" }, + { memberId: "jarvis", agentId: "agent-jarvis" }, + ], + }); + + assert.equal(store.evaluateAddressing({ + channelId: "chan:members-only", + selfAgentId: "agent-outsider", + senderAgentId: "agent-jarvis", + addresseeAgentId: "agent-codex-volt", + }).reason, "self-not-member"); + + assert.equal(store.evaluateAddressing({ + channelId: "chan:members-only", + selfAgentId: "agent-codex-volt", + senderAgentId: "agent-outsider", + addresseeAgentId: "agent-codex-volt", + }).reason, "sender-not-member"); + + assert.equal(store.evaluateAddressing({ + channelId: "chan:members-only", + selfAgentId: "agent-codex-volt", + senderAgentId: "agent-jarvis", + addresseeAgentId: "agent-outsider", + }).reason, "addressee-not-member"); + }); +}); + +test("addressing policy wakes only the explicit addressee and mutes observers", async () => { + await withRoster(async (store) => { + store.createChannel({ + channelId: "chan:addressed", + conversationId: "codex:task:addressed", + type: "group", + members: [ + { memberId: "codex", agentId: "agent-codex-volt", role: "implementer" }, + { memberId: "jarvis", agentId: "agent-jarvis", role: "reviewer" }, + { memberId: "stas", agentId: "agent-stas", role: "observer" }, + ], + }); + + const addressed = store.evaluateAddressing({ + channelId: "chan:addressed", + selfAgentId: "agent-codex-volt", + senderAgentId: "agent-jarvis", + addresseeMemberId: "codex", + }); + assert.equal(addressed.reason, "addressed-member"); + assert.equal(addressed.allowAppend, true); + assert.equal(addressed.allowWake, true); + assert.equal(addressed.reject, false); + assert.equal(addressed.addresseeMember.agentId, "agent-codex-volt"); + + const observer = store.evaluateAddressing({ + channelId: "chan:addressed", + selfAgentId: "agent-stas", + senderAgentId: "agent-jarvis", + addresseeMemberId: "codex", + }); + assert.equal(observer.reason, "observer-muted"); + assert.equal(observer.allowAppend, true); + assert.equal(observer.allowWake, false); + assert.equal(observer.reject, false); + }); +}); + +test("addressing policy treats channel messages without addressee as channel broadcast", async () => { + await withRoster(async (store) => { + store.createChannel({ + channelId: "chan:broadcast", + conversationId: "codex:task:broadcast", + type: "group", + members: [{ memberId: "codex", agentId: "agent-codex-volt" }], + }); + + const decision = store.evaluateAddressing({ + channelId: "chan:broadcast", + selfAgentId: "agent-codex-volt", + senderAgentId: "agent-codex-volt", + }); + assert.equal(decision.reason, "channel-broadcast"); + assert.equal(decision.allowAppend, true); + assert.equal(decision.allowWake, true); + assert.equal(decision.reject, false); + }); +}); + +test("thread start binding projects channel member persona/model/instructions without wake wiring", async () => { + await withRoster(async (store) => { + store.createChannel({ + channelId: "chan:persona", + conversationId: "codex:task:persona", + type: "dm", + members: [{ + memberId: "codex-writer", + memberSlot: "agent-codex-volt:writer", + agentId: "agent-codex-volt", + personaId: "codex-writer", + model: "gpt-5", + baseInstructionsHash: "sha256:writer-v1", + }], + }); + + const member = store.getChannelMember("chan:persona", "codex-writer"); + const binding = buildChannelThreadStartBinding({ + member, + baseInstructions: "Write concise engineering notes.", + }); + + assert.deepEqual(binding, { + model: "gpt-5", + personality: "codex-writer", + baseInstructions: "Write concise engineering notes.", + metadata: { + murmur_channel_id: "chan:persona", + murmur_member_id: "codex-writer", + murmur_agent_id: "agent-codex-volt", + murmur_member_slot: "agent-codex-volt:writer", + murmur_persona_id: "codex-writer", + murmur_model: "gpt-5", + murmur_base_instructions_hash: "sha256:writer-v1", + }, + }); + }); +}); + +test("thread start binding remains null-safe for generic channel members", async () => { + const binding = buildChannelThreadStartBinding({ + member: { + channelId: "chan:generic", + memberId: "codex", + agentId: "agent-codex-volt", + joinedAt: "2026-06-24T10:00:00.000Z", + eligibility: {}, + metadata: {}, + }, + }); + + assert.deepEqual(binding, { + model: null, + personality: null, + baseInstructions: null, + metadata: { + murmur_channel_id: "chan:generic", + murmur_member_id: "codex", + murmur_agent_id: "agent-codex-volt", + }, + }); +}); diff --git a/tests/mcp-channel-roster.test.mjs b/tests/mcp-channel-roster.test.mjs new file mode 100644 index 0000000..fa50c92 --- /dev/null +++ b/tests/mcp-channel-roster.test.mjs @@ -0,0 +1,90 @@ +import assert from "node:assert/strict"; +import { mkdtemp, rm } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { spawn } from "node:child_process"; +import test from "node:test"; + +function callMcp(proc, id, name, args) { + proc.stdin.write(`${JSON.stringify({ + jsonrpc: "2.0", + id, + method: "tools/call", + params: { name, arguments: args }, + })}\n`); +} + +function nextResponse(proc) { + return new Promise((resolve, reject) => { + let buf = ""; + const onData = (chunk) => { + buf += String(chunk); + const idx = buf.indexOf("\n"); + if (idx < 0) return; + proc.stdout.off("data", onData); + const line = buf.slice(0, idx); + try { + resolve(JSON.parse(line)); + } catch (err) { + reject(err); + } + }; + proc.stdout.on("data", onData); + }); +} + +async function callTool(proc, id, name, args) { + const responseP = nextResponse(proc); + callMcp(proc, id, name, args); + const response = await responseP; + assert.equal(response.id, id); + assert.ok(!response.error, response.error?.message); + return JSON.parse(response.result.content[0].text); +} + +test("MCP channel roster tools create/list/evaluate addressing", async () => { + const dir = await mkdtemp(path.join(os.tmpdir(), "murmur-mcp-channel-")); + const proc = spawn(process.execPath, ["packages/mcp-server/dist/src/index.js"], { + cwd: process.cwd(), + env: { + ...process.env, + DATA_DIR: dir, + MURMUR_STORE_PATH: path.join(dir, "murmur.db"), + MURMUR_CHANNEL_ROSTER_PATH: path.join(dir, "channel-roster.db"), + }, + stdio: ["pipe", "pipe", "pipe"], + }); + try { + const created = await callTool(proc, 1, "channel_create", { + channelId: "chan:mcp:test", + conversationId: "codex:task:mcp-test", + type: "group", + members: [ + { memberId: "codex", agentId: "agent-codex-volt", role: "implementer" }, + { memberId: "jarvis", agentId: "agent-jarvis", role: "reviewer" }, + ], + }); + assert.equal(created.channel.channelId, "chan:mcp:test"); + assert.equal(created.members.length, 2); + + const listed = await callTool(proc, 2, "channel_list", { conversationId: "codex:task:mcp-test" }); + assert.equal(listed.channels.length, 1); + assert.equal(listed.channels[0].channelId, "chan:mcp:test"); + + const members = await callTool(proc, 3, "channel_members", { channelId: "chan:mcp:test" }); + assert.equal(members.members[0].memberId, "codex"); + + const decision = await callTool(proc, 4, "channel_evaluate_addressing", { + channelId: "chan:mcp:test", + selfAgentId: "agent-jarvis", + senderAgentId: "agent-codex-volt", + addresseeMemberId: "codex", + }); + assert.equal(decision.decision.reason, "observer-muted"); + assert.equal(decision.decision.allowAppend, true); + assert.equal(decision.decision.allowWake, false); + } finally { + proc.kill(); + await rm(dir, { recursive: true, force: true }); + } +});