Adaptador de canal de Claude para synagent (v1): JALAR por push vía bridge subproceso, ENVIAR por tool synagent_send al bus MQTT.

A4S convierte trabajo multi-repo recurrente en capacidades incrementales para outer harnesses.
Las piezas sólo entran en la narrativa cuando son usables; las integraciones no definen la categoría de A4S.
El monorepo reúne capacidades que pueden evolucionar a ritmos distintos. Cada una documenta su alcance y sus requisitos sin prometer una suite completa, cobertura universal ni autonomía.
.workspace/ contiene la configuración efectiva y el conocimiento gobernado de este workspace.profiles/ publica perfiles reutilizables de configuración.methods/ contiene métodos de trabajo opt-in.skills/ distribuye workflows portátiles y sus herramientas deterministas.agents/ conserva roles portátiles con provenance explícita.output-styles/ define contratos de interacción.packages/pi-context-expert/ contiene una extensión para Pion que gestiona contexto, compaction y propuestas de reglas review-only; consume TypeSafe/Jev mediante el runtime nativo del host según ADR 0065.packages/pi-tool-row-presentation/ contiene una extensión para Pion que controla la densidad de tool rows mediante Settings > Tool rows y el atajo Ctrl+Alt+O, sin comando propio.test/ verifica contratos transversales del repositorio.Pi, Herdr, TypeSafe, Jev, Rootline y Backscroll se integran mediante contratos explícitos. Una capacidad puede requerir una integración concreta; esas dependencias no definen la categoría de A4S ni convierten el monorepo en un runtime propio.
única autoridad normativa del WoW: .workspace/config.yaml
derivados no normativos: README + profiles/pablontiv/PROFILE.md
mecanismos sin reglas propias: skills + methods
referencia de procedencia: Engineering Handbook 1.4
producto: capacidades A4S + integraciones externas
La configuración, los artefactos portátiles y los paquetes conservan límites explícitos. Los providers externos conservan su propia autoridad; A4S no los reemplaza ni se define por ellos.
Estos enlaces conservan contexto, decisiones y diseño. Los ADRs y la spec son registros no normativos; las referencias tampoco gobiernan la forma de trabajo:
Los ADRs y diseños del antiguo repositorio Handbook se preservan como historia en .workspace/docs/history/handbook/; no forman un segundo decision log.
Productos TypeScript:
npm ci
npm test
npm run typecheck
Las extensiones se validan contra Pion 1.0.0-ports.1. El aggregate raíz incluye los tests y el typecheck de pi-tool-row-presentation; sus tipos de desarrollo están fijados al artefacto inmutable de esa release.
Contratos de artefactos portátiles:
python -m pip install --disable-pip-version-check --no-deps -r requirements-test.txt
python -m unittest discover -s test -p "test_*.py" -v
python -m unittest discover -s profiles/pablontiv/tests -t profiles/pablontiv -p "test_*.py" -v
Cada skill conserva sus dependencias, helpers, fixtures y tests dentro de su propio directorio. skills/remove-gentle-context/ requiere un Python 3.11+ executable disponible como python, python3 o un comando equivalente de la plataforma.
El código original de A4S no se distribuye bajo una licencia de reutilización (UNLICENSED); su lectura pública no concede permiso para copiarlo, modificarlo ni redistribuirlo. El alcance se aclara en NOTICE. Los artefactos trasladados desde pablontiv/handbook conservan su licencia MIT en LICENSES/handbook-MIT.txt, y algunos subárboles incluyen licencias propias que prevalecen para esos artefactos.
hooks/register.ts 353 lines1// Adaptador de canal de Claude para synagent (ADR 0069: identidad por sesión
2// nativa, sin variables de entorno). Sobre bus MQTT (aedes).
3//
4// JALAR — bridge suscriptor (subproceso de por vida); cada mensaje canónico
5// llega como una línea en stdout y se inyecta con $.prompt.submit.
6// ENVIAR — herramienta synagent_send (invocada por el modelo) publica al bus.
7//
8// IDENTIDAD (ADR 0069): la instancia sale de $.session.id() (única por sesión →
9// sin colisión aunque haya cientos de sesiones); el proyecto, de un setting del
10// host (userConfig `project`) o del remoto `origin` del repo ($.session.repo()).
11// NADA de env ni de launcher. Si no se resuelve, el adaptador queda inactivo.
12//
13// Topología v1: synagent/v1/<proyecto>/<instancia> (directo, durable),
14// /all (proyecto, transient), synagent/v1/all (global, transient por defecto).
15
16import type { EngineInterface, PluginOptions, Register } from 'claude-code'
17
18import {
19 acceptInbound,
20 directAddress,
21 instanceFromHostSession,
22 makeOutbound,
23 MESSAGE_KINDS,
24 newId,
25 parseCanonical,
26 renderForAgent,
27 resolveProject,
28 serialize,
29 subscriptions,
30 type Identity,
31 type MessageKind,
32} from './adapter'
33
34const DEFAULT_BROKER_URL = 'mqtt://127.0.0.1:1884'
35const LEGACY_CLIENT_ID = 'synagent-legacy-claude'
36const PROCESSED_KEY = 'processed'
37const MAX_SEEN = 1000
38const MAX_BUF = 1_000_000
39
40let started = false
41// Session id para el que arrancamos. ADR 0069: la identidad sale de
42// $.session.id(). `/clear` NO dispara session.start pero cambia el session id y el
43// módulo puede seguir cargado; detectamos el cambio para re-resolver y re-suscribir.
44// (La compactación conserva el session id, así que no reinicia.)
45let startedSessionId: string | undefined
46// Generación del arranque actual: invalida los loops de bridges previos para que
47// su cierre no pise el estado del arranque nuevo.
48let generation = 0
49// Detiene el bridge-sub vivo (mata el subproceso vía HookStream.return()).
50let stopBridge: (() => void) | undefined
51let identity: Identity | undefined
52let globalEnabled = true
53// resolveBrokerUrl corre en register() sin `$`; el aviso se emite al arrancar.
54let brokerWarning: string | undefined
55// Como máximo una migración one-shot en vuelo por activación del módulo. Tras
56// éxito queda hecha; si falta el broker puede reintentarse en un prompt futuro.
57let retirementInFlightOrDone = false
58
59const bridgeDir = ($: EngineInterface): string => `${$.plugin.root}/bridge`
60
61function optString(options: PluginOptions, key: string): string | undefined {
62 const v = options[key]
63 return typeof v === 'string' && v.length > 0 ? v : undefined
64}
65
66// El bus no tiene auth: solo aceptamos brokers de loopback sin credenciales
67// (paridad con el adaptador Pi). Una URL no conforme cae al default.
68function isLoopbackMqtt(value: string): boolean {
69 let u: URL
70 try {
71 u = new URL(value)
72 } catch {
73 return false
74 }
75 return (
76 u.protocol === 'mqtt:'
77 && ['127.0.0.1', 'localhost', '[::1]'].includes(u.hostname)
78 && !u.username
79 && !u.password
80 && (u.pathname === '' || u.pathname === '/')
81 && !u.search
82 && !u.hash
83 )
84}
85
86function resolveBrokerUrl(options: PluginOptions): string {
87 const candidate = optString(options, 'brokerUrl') ?? DEFAULT_BROKER_URL
88 if (!isLoopbackMqtt(candidate)) {
89 brokerWarning = `synagent: brokerUrl no loopback/sin credenciales rechazado (${candidate}); usando ${DEFAULT_BROKER_URL}`
90 return DEFAULT_BROKER_URL
91 }
92 return candidate
93}
94
95// Identidad SOLO de APIs nativas de sesión (ADR 0069). Lanza si no resuelve.
96// instancia = $.session.id() (única por sesión);
97// proyecto = setting del host > nombre del repo del remoto origin.
98async function resolveIdentity($: EngineInterface, options: PluginOptions): Promise<Identity> {
99 const instance = instanceFromHostSession(await $.session.id())
100 const repo = await $.session.repo()
101 const project = resolveProject({
102 ...(optString(options, 'project') !== undefined ? { setting: optString(options, 'project') } : {}),
103 ...(repo?.remote ? { origin: repo.remote } : {}),
104 })
105 return { project, instance }
106}
107
108// Idempotencia por id (at-most-once); dedupe compartido entre clientes.
109async function reserve($: EngineInterface, id: string): Promise<boolean> {
110 const list = ((await $.store.get(PROCESSED_KEY)) as string[] | undefined) ?? []
111 if (list.includes(id)) return false
112 const next = [...list, id]
113 if (next.length > MAX_SEEN) next.splice(0, next.length - MAX_SEEN)
114 await $.store.set(PROCESSED_KEY, next)
115 return true
116}
117
118function retireLegacySession($: EngineInterface, brokerUrl: string): void {
119 if (retirementInFlightOrDone) return
120 retirementInFlightOrDone = true
121 const argv = [
122 'node',
123 `${bridgeDir($)}/bridge-sub.cjs`,
124 '--url',
125 brokerUrl,
126 '--retire-id',
127 LEGACY_CLIENT_ID,
128 ]
129 // No bloquear el arranque: el bridge usa clean=true, sin retries, y termina.
130 void $.process.run(argv).then((result) => {
131 if (result.exitCode !== 0) {
132 retirementInFlightOrDone = false
133 void $.ui.status(`synagent: no se pudo retirar la sesión legacy (${result.stderr || `exit ${result.exitCode}`})`)
134 }
135 }).catch((err) => {
136 retirementInFlightOrDone = false
137 void $.ui.status(`synagent: no se pudo retirar la sesión legacy (${err instanceof Error ? err.message : String(err)})`)
138 })
139}
140
141// Loop JALAR: cada línea stdout es un mensaje canónico; se filtra por `accept`,
142// se deduplica y se inyecta como turno.
143function runJalar(gen: number, $: EngineInterface, argv: readonly string[], accept: (m: ReturnType<typeof parseCanonical>) => boolean): void {
144 void (async () => {
145 try {
146 const bridge = $.process.spawn({ argv: [...argv] })
147 // Permite detener ESTE bridge (p. ej. al cambiar de sesión con /clear).
148 if (gen === generation) stopBridge = () => void bridge.return(undefined as never)
149 let buf = ''
150 for await (const chunk of bridge) {
151 if (chunk.stream !== 'stdout' || typeof chunk.text !== 'string') continue
152 buf += chunk.text
153 if (buf.length > MAX_BUF) buf = buf.slice(-MAX_BUF)
154 let i: number
155 while ((i = buf.indexOf('\n')) >= 0) {
156 const line = buf.slice(0, i)
157 buf = buf.slice(i + 1)
158 if (!line.trim()) continue
159 let msg
160 try {
161 msg = parseCanonical(line)
162 } catch {
163 continue
164 }
165 if (!accept(msg)) continue
166 try {
167 if (!(await reserve($, msg.id))) continue
168 await $.prompt.submit({ text: renderForAgent(msg) })
169 } catch (err) {
170 void $.ui.status(`synagent: entrega falló (${msg.id}): ${err instanceof Error ? err.message : String(err)}`)
171 }
172 }
173 }
174 if (gen === generation) {
175 started = false
176 void $.ui.status('synagent: bridge detenido; re-suscribe en el próximo prompt')
177 }
178 } catch (err) {
179 if (gen === generation) {
180 started = false
181 void $.ui.status(`synagent: bridge error: ${err instanceof Error ? err.message : String(err)}`)
182 }
183 }
184 })()
185}
186
187async function ensureStarted($: EngineInterface, brokerUrl: string, options: PluginOptions): Promise<void> {
188 // La sesión histórica tenía un clientId fijo distinto de los clientId v1. Se
189 // retira incluso si luego no puede resolverse la identidad actual.
190 retireLegacySession($, brokerUrl)
191
192 // ADR 0069: la identidad depende del session id. `/clear` lo cambia SIN disparar
193 // session.start y el módulo puede seguir cargado; si ya arrancamos para otra
194 // sesión, detenemos el bridge viejo y re-resolvemos. Compactación: mismo id → no-op.
195 const sid = await $.session.id()
196 if (started && startedSessionId === sid) return
197 if (started) {
198 stopBridge?.()
199 stopBridge = undefined
200 }
201 const gen = ++generation
202 started = true
203 startedSessionId = sid
204 if (brokerWarning) void $.ui.status(brokerWarning)
205
206 const bridgeSub = `${bridgeDir($)}/bridge-sub.cjs`
207
208 try {
209 identity = await resolveIdentity($, options)
210 globalEnabled = options['global'] !== false && options['global'] !== 'false'
211 } catch (err) {
212 identity = undefined
213 globalEnabled = true
214 started = false
215 void $.ui.status(
216 `synagent: identidad no resuelta (${err instanceof Error ? err.message : String(err)}); adaptador inactivo`,
217 )
218 return
219 }
220
221 const self = identity
222 const addr = directAddress(self)
223 const plan = subscriptions({ identity: self, global: globalEnabled })
224 const durableId = `synagent-${addr.replace(/\//g, ':')}`
225 const transientId = `synagent-t-${addr.replace(/\//g, ':')}`
226
227 const argv = ['node', bridgeSub, '--url', brokerUrl, '--durable-id', durableId, '--durable', ...plan.durable]
228 if (plan.transient.length > 0) argv.push('--transient-id', transientId, '--transient', ...plan.transient)
229
230 runJalar(gen, $, argv, msg => acceptInbound(self, msg, { global: globalEnabled }))
231
232 try {
233 await $.tool.register({
234 name: 'synagent_send',
235 description:
236 'Envía un mensaje a otro agente por el bus synagent (v1). Úsalo solo cuando el usuario pida notificar/avisar/mandar algo a otro agente o proyecto. No para conversación con el usuario actual. El "body" debe ser un hecho autoexplicativo.',
237 inputSchema: {
238 type: 'object',
239 properties: {
240 to: {
241 type: 'string',
242 description: 'Dirección v1: "<proyecto>/<instancia>" (directo), "<proyecto>/all" (proyecto) o "all" (global).',
243 },
244 body: { type: 'string', description: 'Mensaje; un hecho autoexplicativo, no un comando suelto.' },
245 kind: { type: 'string', enum: [...MESSAGE_KINDS], description: 'Tipo; default "prompt". "steer" solo a destino directo.' },
246 reply_to: { type: 'string', description: 'Opcional: id del mensaje al que respondes.' },
247 },
248 required: ['to', 'body'],
249 },
250 })
251 } catch (err) {
252 void $.ui.status(`synagent_send no registrado: ${err instanceof Error ? err.message : String(err)}`)
253 }
254
255 try {
256 await $.command.register({ name: 'mq-send', description: 'Publica un mensaje al bus synagent (debug)', argumentHint: '<to>: texto' })
257 } catch (err) {
258 void $.ui.status(`mq-send no registrado: ${err instanceof Error ? err.message : String(err)}`)
259 }
260
261 void $.ui.status(`synagent: ${addr} (JALAR activo, tool synagent_send + /mq-send)`)
262}
263
264// Publica un mensaje canónico v1 con el bridge-pub one-shot. makeOutbound valida
265// dirección y rechaza steer broadcast (lanza). Devuelve el id o un error.
266async function publishV1(
267 $: EngineInterface,
268 brokerUrl: string,
269 self: Identity,
270 to: string,
271 body: string,
272 kind: MessageKind,
273 replyTo: string | undefined,
274): Promise<{ ok: true; id: string } | { ok: false; error: string }> {
275 const now = await $.clock.now()
276 let built
277 try {
278 built = makeOutbound(self, to, { body, kind, id: newId(directAddress(self), now), ts: now, ...(replyTo ? { replyTo } : {}) })
279 } catch (err) {
280 return { ok: false, error: err instanceof Error ? err.message : String(err) }
281 }
282 let r
283 try {
284 r = await $.process.run(['node', `${bridgeDir($)}/bridge-pub.cjs`, built.topic, serialize(built.message), brokerUrl])
285 } catch (err) {
286 return { ok: false, error: `publish lanzó: ${err instanceof Error ? err.message : String(err)}` }
287 }
288 if (r.exitCode !== 0) return { ok: false, error: `publish falló (exit ${r.exitCode}): ${r.stderr || r.stdout}` }
289 return { ok: true, id: built.message.id }
290}
291
292export const register: Register = (on, options) => {
293 const brokerUrl = resolveBrokerUrl(options)
294
295 on('session.start', async ($, e, next) => {
296 await ensureStarted($, brokerUrl, options)
297 return next(e)
298 })
299
300 on('prompt.submit', async ($, e, next) => {
301 await ensureStarted($, brokerUrl, options)
302 return next(e)
303 })
304
305 // ENVIAR por INTENCIÓN: synagent_send. Sin matcher; filtramos por e.tool y
306 // dejamos pasar los demás con next. Los args van ESPARCIDOS en e.
307 on('tool.call', async ($, e, next) => {
308 const a = e as unknown as Record<string, unknown>
309 if (a.tool !== 'mcp__synagent-adapter-mqtt__synagent_send') return next(e)
310 const result = (text: string, isError = false) => ({ result: text, ...(isError && { isError }) })
311 if (!identity) return result('error: identidad v1 no resuelta; configura el setting project o verifica el remoto origin del repo', true)
312 const to = typeof a.to === 'string' ? a.to : ''
313 const body = typeof a.body === 'string' ? a.body : ''
314 const kind = typeof a.kind === 'string' ? a.kind : 'prompt'
315 const replyTo = typeof a.reply_to === 'string' ? a.reply_to : undefined
316 if (!to || !body) return result('error: synagent_send requiere "to" y "body"', true)
317 if (!MESSAGE_KINDS.includes(kind as MessageKind)) return result(`error: kind inválido: ${kind}`, true)
318 const sent = await publishV1($, brokerUrl, identity, to, body, kind as MessageKind, replyTo)
319 return result(sent.ok ? `enviado ${sent.id} a ${to} (kind=${kind})` : `error: ${sent.error}`, !sent.ok)
320 })
321
322 on('command.run', { command: 'mq-send' }, async ($, e) => {
323 const raw = (e.args ?? '').trim()
324 if (!raw) return { text: 'mq-send: /mq-send <to>: texto (<to> = "<proyecto>/<instancia>", "<proyecto>/all" o "all")' }
325 if (!identity) return { text: 'mq-send: identidad v1 no resuelta' }
326 const m = raw.match(/^([^:]+):\s*([\s\S]*)$/)
327 if (!m) return { text: 'mq-send: formato /mq-send <to>: texto' }
328 const to = (m[1] ?? '').trim()
329 const body = (m[2] ?? '').trim()
330 if (!body) return { text: 'mq-send: falta el texto' }
331 const sent = await publishV1($, brokerUrl, identity, to, body, 'prompt', undefined)
332 return { text: sent.ok ? `mq-send → publicado ${sent.id} a ${to}` : `mq-send ERROR: ${sent.error}` }
333 })
334
335 on('prompt.compose', async ($, e, next) => {
336 const r = await next(e)
337 return {
338 ...r,
339 sections: [
340 ...r.sections,
341 {
342 id: 'synagent',
343 scope: 'session',
344 text:
345 'Bus synagent (v1): para mandar a otro agente usa la herramienta synagent_send (solo cuando el usuario pida notificar/mensajear a otro agente o proyecto, no en conversación normal). '
346 + 'Direcciones: "<proyecto>/<instancia>" (directo), "<proyecto>/all" (proyecto), "all" (global). Los mensajes deben ser hechos autoexplicativos. '
347 + `Tu dirección: ${identity ? directAddress(identity) : 'no configurada (adaptador inactivo)'}.`
348 },
349 ],
350 }
351 })
352}
353hooks/adapter.ts 348 lines1// GENERADO — copia de packages/synagent/core/protocol.ts. NO editar a mano.
2// Regenerar: node adapters/claude/scripts/sync-protocol.mjs (desde packages/synagent/core).
3// El plugin se instala copiándose solo; por eso lleva su propia copia del contrato.
4
5// Contrato canónico de synagent, direccionamiento v1 y CORE host-neutral de
6// adaptadores (ADR 0069, supera a 0068).
7//
8// Fuente ÚNICA del contrato. Autocontenido (no importa nada): el plugin de
9// Claude lleva una COPIA GENERADA (adapters/claude/hooks/adapter.ts) anclada por
10// un test de paridad, porque el marketplace lo instala copiándose solo.
11//
12// Topología v1 (regla: topic == `${TOPIC_ROOT}/${version}/${to}`):
13// synagent/v1/<proyecto>/<instancia> — buzón directo
14// synagent/v1/<proyecto>/all — broadcast de proyecto
15// synagent/v1/all — broadcast global (activo por defecto)
16//
17// IDENTIDAD (ADR 0069): cada binding la deriva de APIs NATIVAS de sesión del
18// host, NO de variables de entorno ni de un launcher. La instancia viene del id
19// de sesión nativo (único por sesión → sin colisión aunque haya cientos); el
20// proyecto, de un setting del host o de la identidad canónica del repo. El alias
21// humano/tab es solo presentación, nunca identidad.
22//
23// CONTRATO DE ADAPTADOR (extensible a hosts futuros): un adaptador provee un
24// HarnessBinding (identity/reserve/deliver) y un transporte Publish/Subscribe;
25// la lógica de direccionamiento vive en funciones puras (makeOutbound/
26// acceptInbound/resolveProject/instanceFromHostSession) que la suite de
27// conformidad ejercita por igual en todo adaptador.
28
29export const MESSAGE_KINDS = ['prompt', 'steer', 'result', 'notify', 'ack'] as const
30
31export type MessageKind = (typeof MESSAGE_KINDS)[number]
32
33export interface CanonicalMessage {
34 id: string
35 from: string
36 to: string
37 kind: MessageKind
38 body: string
39 reply_to?: string
40 ts: number
41}
42
43// --- versión y raíces de topic -------------------------------------------------
44
45export const PROTOCOL_VERSION = 'v1'
46export const TOPIC_ROOT = 'synagent'
47export const LEGACY_TOPIC_ROOT = 'a4s/inbox'
48export const GLOBAL_ADDRESS = 'all'
49
50// Token NO-LOSSY (ADR 0069): sensible a mayúsculas y admite punto, para
51// representar ids de host (UUIDs, nombres con mayúsculas o '.') sin saneado.
52// Seguro para un nivel de topic MQTT: sin '/', '+' ni '#'. 1..256 chars.
53const TOKEN_RE = /^[A-Za-z0-9][A-Za-z0-9._-]{0,255}$/
54
55export function isToken(value: string): boolean {
56 return TOKEN_RE.test(value)
57}
58
59// --- gramática de direcciones --------------------------------------------------
60
61export type Address =
62 | { readonly scope: 'direct'; readonly project: string; readonly instance: string }
63 | { readonly scope: 'project'; readonly project: string }
64 | { readonly scope: 'global' }
65
66// Parsea el campo canónico `to`. `all` reservado al global; `<proyecto>/all` es
67// broadcast de proyecto, por lo que una instancia NUNCA puede llamarse `all` ni
68// un proyecto puede llamarse `all`.
69export function parseAddress(to: string): Address | null {
70 if (to === GLOBAL_ADDRESS) return { scope: 'global' }
71 const parts = to.split('/')
72 if (parts.length !== 2) return null
73 const project = parts[0] as string
74 const leaf = parts[1] as string
75 if (!isToken(project) || project === GLOBAL_ADDRESS) return null
76 if (leaf === GLOBAL_ADDRESS) return { scope: 'project', project }
77 if (!isToken(leaf)) return null
78 return { scope: 'direct', project, instance: leaf }
79}
80
81export function isAddress(to: string): boolean {
82 return parseAddress(to) !== null
83}
84
85export function isBroadcast(to: string): boolean {
86 const address = parseAddress(to)
87 return address !== null && address.scope !== 'direct'
88}
89
90export function formatAddress(address: Address): string {
91 if (address.scope === 'global') return GLOBAL_ADDRESS
92 if (address.scope === 'project') return `${address.project}/${GLOBAL_ADDRESS}`
93 return `${address.project}/${address.instance}`
94}
95
96// --- identidad del agente ------------------------------------------------------
97
98export interface Identity {
99 readonly project: string
100 readonly instance: string
101}
102
103export function directAddress(identity: Identity): string {
104 return `${identity.project}/${identity.instance}`
105}
106
107export function projectAddress(project: string): string {
108 return `${project}/${GLOBAL_ADDRESS}`
109}
110
111// --- mapa dirección → topic ----------------------------------------------------
112
113export function toTopic(to: string, version: string = PROTOCOL_VERSION): string {
114 if (!isAddress(to)) throw new Error(`dirección v1 inválida: ${to}`)
115 return `${TOPIC_ROOT}/${version}/${to}`
116}
117
118export function legacyTopic(address: string): string {
119 return `${LEGACY_TOPIC_ROOT}/${address}`
120}
121
122// --- plan de suscripción -------------------------------------------------------
123// durable (clean=false): buzón directo (+ legacy solo si se pide explícitamente).
124// transient (clean=true): broadcasts de proyecto + global por defecto → online-only.
125
126export interface SubscriptionPlan {
127 readonly durable: readonly string[]
128 readonly transient: readonly string[]
129}
130
131export function subscriptions(options: {
132 identity: Identity
133 global?: boolean
134 legacyAddress?: string
135 version?: string
136}): SubscriptionPlan {
137 const version = options.version ?? PROTOCOL_VERSION
138 const durable: string[] = [toTopic(directAddress(options.identity), version)]
139 if (options.legacyAddress) durable.push(legacyTopic(options.legacyAddress))
140 const transient: string[] = [toTopic(projectAddress(options.identity.project), version)]
141 if (options.global !== false) transient.push(toTopic(GLOBAL_ADDRESS, version))
142 return { durable, transient }
143}
144
145// --- enrutado de recepción -----------------------------------------------------
146
147// Igualdad exacta de dirección (buzón directo o legacy plano).
148export function isFor(message: CanonicalMessage, address: string): boolean {
149 return message.to === address
150}
151
152// ¿Este mensaje es para mí? Directo a mi identidad, broadcast de mi proyecto,
153// global (si está habilitado) o legacy plano (compatibilidad explícita).
154export function isForIdentity(
155 message: CanonicalMessage,
156 options: { identity: Identity; global?: boolean; legacyAddress?: string },
157): boolean {
158 if (options.legacyAddress && message.to === options.legacyAddress) return true
159 const address = parseAddress(message.to)
160 if (!address) return false
161 if (address.scope === 'direct') {
162 return address.project === options.identity.project && address.instance === options.identity.instance
163 }
164 if (address.scope === 'project') return address.project === options.identity.project
165 return options.global !== false
166}
167
168// steer solo tiene sentido directo: se rechaza en ENVÍO y en RECEPCIÓN.
169export function isBroadcastSteer(message: CanonicalMessage): boolean {
170 return message.kind === 'steer' && isBroadcast(message.to)
171}
172
173// --- resolución de identidad por sesión nativa (ADR 0069) ----------------------
174
175// Nombre canónico del repo desde la URL del remoto `origin` (último segmento sin
176// `.git`). Determinista y compartido entre hosts: el mismo repo da el mismo
177// token en Pi y en Claude. NO usa el basename del cwd. Devuelve null si no parsea.
178export function repoNameFromOrigin(origin: string): string | null {
179 const cleaned = origin.trim().replace(/\.git$/, '').replace(/\/+$/, '')
180 const match = cleaned.match(/[/:]([^/:]+)$/)
181 const name = match?.[1]
182 return name && isToken(name) ? name : null
183}
184
185// Proyecto: setting del host > nombre canónico del repo (origin). FALLA explícito
186// si no hay ninguno (sin env, sin basename de cwd, sin saneado lossy).
187export function resolveProject(options: { setting?: string; origin?: string }): string {
188 const setting = options.setting?.trim()
189 if (setting) {
190 if (!isToken(setting)) throw new Error(`proyecto (setting) no es un token válido: ${setting}`)
191 return setting
192 }
193 const origin = options.origin?.trim()
194 if (origin) {
195 const name = repoNameFromOrigin(origin)
196 if (!name) throw new Error(`no se pudo derivar un proyecto-token del remoto origin: ${origin}`)
197 return name
198 }
199 throw new Error('no se pudo resolver <proyecto>: define un setting de proyecto o un remoto origin del repo')
200}
201
202// Instancia: derivada del id de sesión NATIVO del host (único por sesión). Un
203// resume conserva el id; new/fork/clear lo re-resuelven. FALLA explícito si el id
204// nativo no es un token válido (no se trunca ni se sanea: provocaría colisión).
205export function instanceFromHostSession(sessionId: string): string {
206 const id = sessionId.trim()
207 if (!isToken(id)) throw new Error(`id de sesión del host no es un token válido para instancia: ${sessionId}`)
208 return id
209}
210
211// --- contrato de adaptador host-neutral ----------------------------------------
212// Un adaptador implementa estas piezas; la lógica de direccionamiento es común.
213
214// Lo que un host expone al core para entregar/deduplicar bajo su identidad.
215export interface HarnessBinding {
216 readonly identity: Identity
217 // Reserva un id ANTES de entregar (at-most-once). true si es nuevo.
218 reserve(id: string): boolean | Promise<boolean>
219 // Entrega un mensaje aceptado al host (inyecta turno / sendUserMessage).
220 deliver(message: CanonicalMessage): void | Promise<void>
221}
222
223export type Publish = (topic: string, payload: string) => void | Promise<void>
224export type Subscribe = (
225 topics: readonly string[],
226 onMessage: (payload: string) => void,
227) => void | Promise<void>
228
229export interface Transport {
230 readonly publish: Publish
231 readonly subscribe: Subscribe
232}
233
234// Construye un envío v1 (single-write): valida la dirección, rechaza steer a
235// broadcast y arma {topic, message}. Lanza ante dirección inválida o steer broadcast.
236export function makeOutbound(
237 self: Identity,
238 to: string,
239 options: { body: string; kind?: MessageKind; replyTo?: string; id: string; ts: number },
240): { topic: string; message: CanonicalMessage } {
241 if (!isAddress(to)) throw new Error(`dirección v1 inválida: ${to}`)
242 const kind = options.kind ?? 'prompt'
243 if (kind === 'steer' && isBroadcast(to)) throw new Error('steer solo a destino directo, no broadcast')
244 const message = createCanonical(options.body, {
245 id: options.id,
246 from: directAddress(self),
247 to,
248 ts: options.ts,
249 kind,
250 ...(options.replyTo ? { reply_to: options.replyTo } : {}),
251 })
252 return { topic: toTopic(to), message }
253}
254
255// ¿Aceptar un mensaje entrante? Enrutado para mi identidad y NO steer-broadcast.
256export function acceptInbound(
257 self: Identity,
258 message: CanonicalMessage,
259 options: { global?: boolean; legacyAddress?: string } = {},
260): boolean {
261 // Self-echo: un broadcast propio vuelve por mi propia suscripción. No entregar
262 // lo que yo mismo emití (comparando la dirección de origen con la mía).
263 if (message.from === directAddress(self)) return false
264 if (isBroadcastSteer(message)) return false
265 return isForIdentity(message, { identity: self, ...options })
266}
267
268// --- contrato canónico (mensaje) -----------------------------------------------
269
270export function createCanonical(
271 body: string,
272 options: {
273 id: string
274 from: string
275 to: string
276 ts: number
277 kind?: MessageKind
278 reply_to?: string
279 },
280): CanonicalMessage {
281 return {
282 id: options.id,
283 from: options.from,
284 to: options.to,
285 kind: options.kind ?? 'prompt',
286 body,
287 ...(options.reply_to ? { reply_to: options.reply_to } : {}),
288 ts: options.ts,
289 }
290}
291
292export function serialize(message: CanonicalMessage): string {
293 return JSON.stringify(message)
294}
295
296export function parseCanonical(text: string): CanonicalMessage {
297 let value: unknown
298 try {
299 value = JSON.parse(text)
300 } catch {
301 throw new Error('mensaje canónico no es JSON válido')
302 }
303 if (!value || typeof value !== 'object' || Array.isArray(value)) {
304 throw new Error('mensaje canónico debe ser un objeto')
305 }
306
307 const message = value as Record<string, unknown>
308 for (const field of ['id', 'from', 'to', 'kind', 'body', 'ts'] as const) {
309 if (message[field] === undefined || message[field] === null) throw new Error(`falta campo canónico: ${field}`)
310 }
311 for (const field of ['id', 'from', 'to', 'kind'] as const) {
312 if (typeof message[field] !== 'string' || message[field].length === 0) {
313 throw new Error(`campo canónico inválido: ${field}`)
314 }
315 }
316 if (typeof message.body !== 'string') throw new Error('campo canónico inválido: body')
317 if (typeof message.ts !== 'number' || !Number.isFinite(message.ts)) {
318 throw new Error('campo canónico inválido: ts')
319 }
320 if (!MESSAGE_KINDS.includes(message.kind as MessageKind)) {
321 throw new Error(`kind inválido: ${String(message.kind)}`)
322 }
323 if (message.reply_to !== undefined && (typeof message.reply_to !== 'string' || message.reply_to.length === 0)) {
324 throw new Error('campo canónico inválido: reply_to')
325 }
326
327 return {
328 id: message.id as string,
329 from: message.from as string,
330 to: message.to as string,
331 kind: message.kind as MessageKind,
332 body: message.body,
333 ...(message.reply_to === undefined ? {} : { reply_to: message.reply_to as string }),
334 ts: message.ts,
335 }
336}
337
338export function renderForAgent(message: CanonicalMessage): string {
339 const head = `[bus:${message.kind}] de ${message.from} (id ${message.id})`
340 const tail = message.reply_to ? `\n(responder a: ${message.reply_to})` : ''
341 return `${head}\n${message.body}${tail}`
342}
343
344// id opaco y legible: incluye la dirección de origen y el ts.
345export function newId(address: string, ts: number): string {
346 return `${address}-${ts}-${Math.random().toString(36).slice(2, 8)}`
347}
348