Channels
platform/channels (kalt-channels) names a stream from the owner and model a service passes in. There is one bus on the mesh. lib/ch is part of this package: ch.pub appends the stream, ch.sub reads one consumer group, and subscribe_to registers a reader. A service creates its owner group when it starts, and its reader group with subscribe_to. The channels process reads group kalt-channels on every public stream and stores each entry in its own database.
Principles
- One owner. A channel names one service. That service is the only publisher. It passes its name and model:
Channel::public(owner, model)iskalt.{owner}and groupkalt-{owner}, started at0. - The model is the message type. The service chooses
Model::DomainorModel::Source. That is the onlytypeon its stream. - Public or private. On a public channel, other services read with their own groups. On a private channel, the publisher and the reader are the same service.
Channel::private(owner, model)iskalt.private.{owner}and groupkalt-private-{owner}, started at0. - Each reader acks alone. The owner group is
kalt-{owner}, started at0. A service callssubscribe_to(owner, model, start)to read another public stream askalt-{service}.startis$for new entries, or0when that service applies history. - The publish is the stream entry. The channels process stores it.
POST /api/channels/resendappends those rows again withreplay=1. The archiver acks a replay and does not insert it.
Streams
| Visibility | Model | Stream | Group | Start |
|---|---|---|---|---|
| Public | domain or source | kalt.{owner} | kalt-{owner} | 0 |
| Private | same as the owner | kalt.private.{owner} | kalt-private-{owner} | 0 |
| Subscriber | copied from the owner | kalt.{owner} | kalt-{subscriber} | $ |
| Archive | copied from the owner | kalt.{owner} | kalt-channels | 0 |
Channel::private(owner, model) is one handle. The service that owns it publishes with it and reads the same group. The stream is kalt.private.{owner}, the group is kalt-private-{owner}, started at 0.
subscribe_to(owner, model, start) is how a service reads someone else's public stream. The service names the owner and the model. The group is kalt-{service} on kalt.{owner}.
ch.pub
AppState.ch binds Channel::public(service, model). The service passes both. HTTP handlers publish there. The owner's worker reads group kalt-{service} and applies content.
A service calls subscribe_to(owner, model, start) to read another public stream. That creates group kalt-{service} on kalt.{owner} and returns a handle that reads with sub and acks. Publishing stays on the owner's handle.
ch.pub (r#pub in Rust, because pub is a keyword) accepts the write when three things hold, then appends the stream:
- The bound service is the stream's owner.
- The handle is the owner's, public or private.
typematches the channel model:domainorsource.
The channels ledger row uses the Redis id, the stream name, source, trace_id, type, and content. The primary key is (stream, id). pub_upsert(id) appends another entry. It does not replace an earlier one. pub_once(id) appends only when kalt.once:{stream}:{id} is absent, then sets that key. GET /api/channels/ledger lists the rows. POST /api/channels/resend delivers each row again. Pass stream to limit either call to one queue. Private streams are not stored.
The envelope is already the promoted layer. Every entry is those fields, written as Redis stream fields:
{
"service": "publications",
"source": "app/composables/usePublication.ts",
"traceId": ["nK7_abC1"],
"type": "domain",
"content": { "domain": "Publication", "body": { "kind": "title", "id": "AbC12~_x", "tz": 1710000000000, "value": "Title" } }
}
Envelope
- service — the publishing binary. It is the channel owner.
- source — the file that called the API. The browser sends
X-Msg-Source. If that header is missing, the publisher usesfile!(). This field is the call site. It is separate from a source-aligned channel. - traceId — string array, set in the browser and passed through. Observe uses it to stitch a user action across hops.
bindFetchsends it asX-Trace-Id. On Redis the field name istrace_id. - type —
domainorsource, the same word as the channel model. - content — JSON payload. A public domain message is a
StreamEvent.
service, source, trace_id, type, and content are reserved stream fields.
Promoted properties
Promoted properties are scalars copied out of the message body so a consumer can drop a message without parsing JSON. They come from content and content.body (domain, kind, id, …). Nested objects and the write value stay inside content.
On Redis they sit beside the envelope as sibling fields. They are a flat set of stream fields, and they never repeat envelope keys.
Types
The type is the channel model. Publish rejects any other type on that stream.
- domain — a fact on a domain channel (publication title, doctrine rels, blueprint ops).
emit_domainin Api callsch.pubwith this type. - source — a fact on a source channel. The integration that owns the stream publishes there, including on its private stream, with the same type. It binds
Channel::private(owner, model)for both the write and the ack groupkalt-private-{owner}.
Frontend
app/web/app/lib/msg.ts holds the session traceId array. bindFetch(import.meta.url) attaches both headers on every okFetch:
import { bindFetch } from '~/lib/ok-fetch'
const okFetch = bindFetch(import.meta.url)
await okFetch('/api/publications/id/title', {
method: 'PUT',
body: { tz: Date.now(), value: 'Title' },
})
pushTraceId() appends another id when observe needs a nested span. The array is the whole trace, not a single string.