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

  1. One owner. A channel names one service. That service is the only publisher. It passes its name and model: Channel::public(owner, model) is kalt.{owner} and group kalt-{owner}, started at 0.
  2. The model is the message type. The service chooses Model::Domain or Model::Source. That is the only type on its stream.
  3. 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) is kalt.private.{owner} and group kalt-private-{owner}, started at 0.
  4. Each reader acks alone. The owner group is kalt-{owner}, started at 0. A service calls subscribe_to(owner, model, start) to read another public stream as kalt-{service}. start is $ for new entries, or 0 when that service applies history.
  5. The publish is the stream entry. The channels process stores it. POST /api/channels/resend appends those rows again with replay=1. The archiver acks a replay and does not insert it.

Streams

VisibilityModelStreamGroupStart
Publicdomain or sourcekalt.{owner}kalt-{owner}0
Privatesame as the ownerkalt.private.{owner}kalt-private-{owner}0
Subscribercopied from the ownerkalt.{owner}kalt-{subscriber}$
Archivecopied from the ownerkalt.{owner}kalt-channels0

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.
  • type matches the channel model: domain or source.

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 uses file!(). 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. bindFetch sends it as X-Trace-Id. On Redis the field name is trace_id.
  • type — domain or source, 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 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_domain in Api calls ch.pub with 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 group kalt-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.