Data Model
Status: Draft v0.1 | Date: 2026-10-10 | Owner: Founding team
This document defines every durable store in AgentBus: the PostgreSQL 19 control plane, the
ClickHouse analytics and audit store, S3-compatible object storage, the per-tenant encryption
scheme, and the sidecar's local SQLite store. Identifiers, prefixes, and vocabulary follow
00-CONVENTIONS.md.
Design rules that apply everywhere:
tenant_idis on every tenant-scoped table and every ClickHouse row. No exceptions.- Primary keys are UUIDv7 (
uuidtype). Time-ordered keys keep B-tree inserts append-only. - Payloads are never stored in cleartext. Postgres stores envelope metadata plus an encrypted payload reference; the bytes live in S3 under the tenant's data key.
- Nothing is hard-deleted from audit tables. Retention jobs delete; humans do not.
- Every write that matters to a customer also produces an event on
sys.audit.>.
1. PostgreSQL 19 control plane
1.1 Why Postgres for this, ClickHouse for that
Postgres holds everything that needs transactions, uniqueness, foreign keys, and row-level
security: identity, tokens, agents, task state, policies, grants. ClickHouse holds everything
that is append-only and queried by time range and aggregate: delivery events, usage, audit
chain. Messages sit in both: Postgres has the current state of each message for the API,
ClickHouse has the full timeline for analytics and agentbus trace.
1.2 Extensions and roles
CREATE EXTENSION IF NOT EXISTS pgcrypto; -- gen_random_bytes for token generation
CREATE EXTENSION IF NOT EXISTS btree_gist; -- exclusion constraints on time ranges
-- Application roles. The gateway connects as agentbus_app; migrations run as agentbus_owner.
-- agentbus_app has no BYPASSRLS, so row-level security is always enforced.
CREATE ROLE agentbus_owner LOGIN;
CREATE ROLE agentbus_app LOGIN NOBYPASSRLS;
CREATE ROLE agentbus_support LOGIN NOBYPASSRLS; -- support console, read-mostly
1.3 Tenancy and identity
CREATE TYPE plan_t AS ENUM ('starter', 'business', 'enterprise');
CREATE TYPE tenant_status_t AS ENUM ('active', 'suspended', 'deleting', 'deleted');
CREATE TABLE tenants (
id uuid PRIMARY KEY, -- ten_
slug text NOT NULL UNIQUE CHECK (slug ~ '^[a-z0-9-]{3,40}$'),
display_name text NOT NULL,
plan plan_t NOT NULL DEFAULT 'starter',
status tenant_status_t NOT NULL DEFAULT 'active',
region text NOT NULL DEFAULT 'eu-west', -- data residency, see §6
retention_days int NOT NULL DEFAULT 7, -- derived from plan, overridable on enterprise
nats_stream text NOT NULL, -- T_<id without hyphens>
nats_account text, -- dedicated NATS account (enterprise), else NULL
kek_ref text NOT NULL, -- key encryption key reference, see §5
custom_domain text UNIQUE, -- enterprise, e.g. bus.acme.com
settings jsonb NOT NULL DEFAULT '{}'::jsonb, -- support_access_consent, default_visibility, ...
created_at timestamptz NOT NULL DEFAULT now(),
updated_at timestamptz NOT NULL DEFAULT now(),
deleted_at timestamptz
);
CREATE TABLE workspaces (
id uuid PRIMARY KEY, -- ws_
tenant_id uuid NOT NULL REFERENCES tenants(id),
slug text NOT NULL CHECK (slug ~ '^[a-z0-9-]{3,40}$'),
display_name text NOT NULL,
description text,
settings jsonb NOT NULL DEFAULT '{}'::jsonb,
created_at timestamptz NOT NULL DEFAULT now(),
deleted_at timestamptz,
UNIQUE (tenant_id, slug)
);
-- Users are tenant-scoped. The same human in two tenants is two rows; the IdP subject links them.
CREATE TABLE users (
id uuid PRIMARY KEY, -- usr_
tenant_id uuid NOT NULL REFERENCES tenants(id),
idp text NOT NULL, -- 'clerk', 'oidc:<issuer>', 'saml:<issuer>'
idp_subject text NOT NULL, -- the IdP's stable subject claim
email text NOT NULL,
display_name text,
status text NOT NULL DEFAULT 'active' CHECK (status IN ('invited','active','disabled')),
created_at timestamptz NOT NULL DEFAULT now(),
last_login_at timestamptz,
UNIQUE (tenant_id, idp, idp_subject),
UNIQUE (tenant_id, email)
);
CREATE TYPE tenant_role_t AS ENUM ('owner', 'admin', 'member', 'auditor', 'billing');
CREATE TYPE ws_role_t AS ENUM ('admin', 'member', 'viewer');
-- Tenant-level role.
CREATE TABLE tenant_memberships (
tenant_id uuid NOT NULL REFERENCES tenants(id),
user_id uuid NOT NULL REFERENCES users(id),
role tenant_role_t NOT NULL DEFAULT 'member',
created_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (tenant_id, user_id)
);
-- Workspace-level role. A user must be a tenant member to be a workspace member.
CREATE TABLE workspace_memberships (
tenant_id uuid NOT NULL,
workspace_id uuid NOT NULL REFERENCES workspaces(id),
user_id uuid NOT NULL REFERENCES users(id),
role ws_role_t NOT NULL DEFAULT 'member',
invited_by uuid REFERENCES users(id),
created_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (workspace_id, user_id),
FOREIGN KEY (tenant_id, user_id) REFERENCES tenant_memberships(tenant_id, user_id)
);
-- Pending invitations (email-based). SSO-mapped joins do not use this table.
CREATE TABLE invitations (
id uuid PRIMARY KEY,
tenant_id uuid NOT NULL REFERENCES tenants(id),
workspace_id uuid REFERENCES workspaces(id),
email text NOT NULL,
role ws_role_t NOT NULL DEFAULT 'member',
token_hash bytea NOT NULL, -- SHA-256 of the one-time link token
invited_by uuid NOT NULL REFERENCES users(id),
expires_at timestamptz NOT NULL,
accepted_at timestamptz,
created_at timestamptz NOT NULL DEFAULT now()
);
1.4 Credentials
-- Long-lived integration tokens: one per harness/integration, created in the console or CLI.
-- The plaintext is shown exactly once. ab_live_<40> / ab_test_<40>.
CREATE TABLE integration_tokens (
id uuid PRIMARY KEY, -- tok_
tenant_id uuid NOT NULL REFERENCES tenants(id),
workspace_id uuid REFERENCES workspaces(id), -- NULL = tenant-wide (admin only)
owner_user_id uuid NOT NULL REFERENCES users(id),
name text NOT NULL, -- 'claude-code on laptop', 'ci-runner'
env text NOT NULL CHECK (env IN ('live','test')),
token_hash bytea NOT NULL UNIQUE, -- SHA-256(plaintext); high-entropy so no KDF needed
prefix text NOT NULL, -- 'ab_live_' + first 6 chars, for display and leak scanning
last4 text NOT NULL,
scopes text[] NOT NULL DEFAULT '{send,receive,register}', -- send, receive, register, read_audit, admin
expires_at timestamptz,
last_used_at timestamptz,
last_used_ip inet,
revoked_at timestamptz,
revoked_reason text,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX ON integration_tokens (tenant_id, owner_user_id) WHERE revoked_at IS NULL;
-- Short-lived agent credentials are JWTs and are not stored. Their signing keys are:
CREATE TABLE signing_keys (
kid text PRIMARY KEY,
alg text NOT NULL DEFAULT 'EdDSA',
public_key bytea NOT NULL,
private_key_enc bytea NOT NULL, -- wrapped by the platform KEK
active boolean NOT NULL DEFAULT true,
created_at timestamptz NOT NULL DEFAULT now(),
retired_at timestamptz
);
1.5 Agents
CREATE TYPE visibility_t AS ENUM ('private', 'workspace', 'tenant', 'public');
CREATE TYPE agent_status_t AS ENUM ('registered', 'online', 'idle', 'offline', 'disabled', 'deleted');
CREATE TYPE harness_t AS ENUM ('claude', 'codex', 'opencode', 'gemini', 'custom');
CREATE TABLE agents (
id uuid PRIMARY KEY, -- agt_
tenant_id uuid NOT NULL REFERENCES tenants(id),
workspace_id uuid NOT NULL REFERENCES workspaces(id),
owner_user_id uuid NOT NULL REFERENCES users(id),
name text NOT NULL CHECK (name ~ '^[a-z0-9-]{3,40}$'),
address text NOT NULL, -- agent://<tenant-slug>/<ws-slug>/<name>, denormalised
display_name text,
description text NOT NULL DEFAULT '', -- shown in the workspace directory
capabilities jsonb NOT NULL DEFAULT '[]'::jsonb, -- [{"name":"code.review","version":"1","description":"..."}]
agent_card jsonb NOT NULL DEFAULT '{}'::jsonb, -- A2A-compatible card, generated + user fields
harness harness_t NOT NULL DEFAULT 'custom',
harness_version text,
mode text NOT NULL DEFAULT 'interactive' CHECK (mode IN ('interactive','service')),
visibility visibility_t NOT NULL DEFAULT 'workspace',
status agent_status_t NOT NULL DEFAULT 'registered',
nats_consumer text NOT NULL, -- agt_<id>
inbox_subject text NOT NULL, -- t.<tenant>.ws.<ws>.agent.<id>.inbox
delivery jsonb NOT NULL DEFAULT '{"strict_order":false}'::jsonb,
last_seen_at timestamptz,
last_host text, -- hostname fingerprint from the sidecar, for doctor
created_at timestamptz NOT NULL DEFAULT now(),
updated_at timestamptz NOT NULL DEFAULT now(),
deleted_at timestamptz,
UNIQUE (workspace_id, name),
UNIQUE (address)
);
CREATE INDEX ON agents (tenant_id, workspace_id) WHERE deleted_at IS NULL;
CREATE INDEX ON agents USING gin (capabilities jsonb_path_ops);
-- Ed25519 keys the sidecar generates. Multiple rows allow rotation without a gap.
CREATE TABLE agent_keys (
id uuid PRIMARY KEY,
tenant_id uuid NOT NULL REFERENCES tenants(id),
agent_id uuid NOT NULL REFERENCES agents(id),
public_key bytea NOT NULL, -- 32 bytes
fingerprint text NOT NULL, -- SHA-256 base64url, shown in console
spiffe_id text NOT NULL, -- spiffe://<tenant-slug>.agentbus.exchange/ws/<ws>/agent/<id>
created_at timestamptz NOT NULL DEFAULT now(),
expires_at timestamptz NOT NULL, -- 90 days default; sidecar rotates at 60
revoked_at timestamptz,
UNIQUE (agent_id, public_key)
);
CREATE INDEX ON agent_keys (agent_id) WHERE revoked_at IS NULL;
1.6 Conversations, messages, tasks
CREATE TABLE conversations (
id uuid PRIMARY KEY, -- cnv_
tenant_id uuid NOT NULL REFERENCES tenants(id),
workspace_id uuid NOT NULL REFERENCES workspaces(id),
started_by uuid NOT NULL REFERENCES agents(id),
title text, -- optional, set by first message's data.title
participants uuid[] NOT NULL DEFAULT '{}', -- maintained by trigger on messages insert
message_count int NOT NULL DEFAULT 0,
last_message_at timestamptz,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX ON conversations (tenant_id, last_message_at DESC);
CREATE TYPE message_status_t AS ENUM (
'accepted', 'fetched', 'delivered', 'read', 'expired', 'dead_lettered', 'discarded', 'bridged'
);
-- Envelope metadata only. The payload (`data`) is encrypted and lives in S3; payload_ref points at it.
-- Partitioned by month on accepted_at for cheap retention drops.
CREATE TABLE messages (
id uuid NOT NULL, -- msg_ = CloudEvents id
tenant_id uuid NOT NULL,
workspace_id uuid NOT NULL,
conversation_id uuid NOT NULL,
type text NOT NULL, -- agentbus.task.request.v1 ...
schema_url text NOT NULL,
from_agent_id uuid NOT NULL,
to_agent_id uuid, -- NULL for topic:// (post-MVP)
to_address text NOT NULL,
correlation_id uuid,
reply_to text,
idempotency_key text,
capability text,
capability_ver text,
priority text NOT NULL DEFAULT 'normal' CHECK (priority IN ('low','normal','high')),
budget jsonb,
expires_at timestamptz,
traceparent text NOT NULL,
trace_id text GENERATED ALWAYS AS (split_part(traceparent, '-', 2)) STORED,
sig bytea NOT NULL,
signing_key_id uuid NOT NULL,
payload_ref text NOT NULL, -- s3 key, see §4
payload_bytes int NOT NULL,
payload_sha256 bytea NOT NULL, -- of the plaintext, for integrity after decrypt
dek_version int NOT NULL, -- which tenant DEK encrypted it
status message_status_t NOT NULL DEFAULT 'accepted',
stream_seq bigint, -- NULL until NATS PubAck
attempts int NOT NULL DEFAULT 0,
bridged_from uuid, -- source tenant for cross-tenant (post-MVP)
accepted_at timestamptz NOT NULL DEFAULT now(),
delivered_at timestamptz,
read_at timestamptz,
PRIMARY KEY (tenant_id, accepted_at, id)
) PARTITION BY RANGE (accepted_at);
-- Monthly partitions, created 3 months ahead by the retention job.
CREATE TABLE messages_2026_10 PARTITION OF messages FOR VALUES FROM ('2026-10-01') TO ('2026-11-01');
CREATE TABLE messages_2026_11 PARTITION OF messages FOR VALUES FROM ('2026-11-01') TO ('2026-12-01');
CREATE UNIQUE INDEX messages_id_uq ON messages (tenant_id, id, accepted_at); -- dedupe lookups
CREATE INDEX messages_to_status ON messages (to_agent_id, status, accepted_at) -- inbox queries
WHERE status IN ('accepted','fetched','delivered');
CREATE INDEX messages_conv ON messages (conversation_id, accepted_at);
CREATE INDEX messages_trace ON messages (trace_id);
CREATE INDEX messages_unpublished ON messages (accepted_at) WHERE stream_seq IS NULL; -- reconciler
CREATE TYPE task_state_t AS ENUM (
'submitted', 'delivered', 'accepted', 'in_progress', 'completed', 'failed', 'rejected',
'cancelled', 'expired'
);
CREATE TABLE tasks (
id uuid PRIMARY KEY, -- tsk_
tenant_id uuid NOT NULL REFERENCES tenants(id),
workspace_id uuid NOT NULL REFERENCES workspaces(id),
conversation_id uuid NOT NULL REFERENCES conversations(id),
request_msg_id uuid NOT NULL,
requester_agent_id uuid NOT NULL REFERENCES agents(id),
recipient_agent_id uuid NOT NULL REFERENCES agents(id),
idempotency_key text NOT NULL,
capability text NOT NULL,
capability_ver text NOT NULL DEFAULT '1',
state task_state_t NOT NULL DEFAULT 'submitted',
budget jsonb, -- {"currency":"USD","max":"2.00"}
spent numeric(12,6) NOT NULL DEFAULT 0,
grant_id uuid, -- post-MVP cross-tenant
result_msg_id uuid,
error_code text,
progress jsonb, -- last task.progress payload summary
submitted_at timestamptz NOT NULL DEFAULT now(),
delivered_at timestamptz,
accepted_at timestamptz,
started_at timestamptz,
finished_at timestamptz,
expires_at timestamptz,
duration_ms bigint GENERATED ALWAYS AS (
CASE WHEN finished_at IS NOT NULL AND accepted_at IS NOT NULL
THEN (EXTRACT(EPOCH FROM finished_at - accepted_at) * 1000)::bigint END) STORED,
UNIQUE (recipient_agent_id, idempotency_key) -- the task idempotency guarantee
);
CREATE INDEX ON tasks (tenant_id, state, submitted_at DESC);
CREATE INDEX ON tasks (recipient_agent_id, state) WHERE state IN ('submitted','delivered','accepted','in_progress');
-- Append-only task transitions, for the console timeline and for disputes.
CREATE TABLE task_transitions (
id uuid PRIMARY KEY,
tenant_id uuid NOT NULL,
task_id uuid NOT NULL REFERENCES tasks(id),
from_state task_state_t,
to_state task_state_t NOT NULL,
by_msg_id uuid,
reason text,
at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX ON task_transitions (task_id, at);
1.7 Delivery state and receipts
The authoritative timeline is in ClickHouse. Postgres keeps only the current delivery state needed by the API and the reconciler.
CREATE TABLE delivery_state (
tenant_id uuid NOT NULL,
msg_id uuid NOT NULL,
agent_id uuid NOT NULL,
attempt int NOT NULL DEFAULT 0,
last_fetched_at timestamptz,
transport_acked_at timestamptz,
application_acked_at timestamptz,
nak_until timestamptz,
dlq_reason text,
PRIMARY KEY (msg_id, agent_id)
);
CREATE TABLE receipts (
id uuid PRIMARY KEY, -- rcp_
tenant_id uuid NOT NULL,
msg_id uuid NOT NULL,
stream_seq bigint,
issued_to_token uuid NOT NULL,
issued_at timestamptz NOT NULL DEFAULT now(),
signature bytea NOT NULL -- gateway signs {id,msg_id,seq,issued_at} with signing_keys
);
CREATE INDEX ON receipts (msg_id);
1.8 Policies and grants
CREATE TABLE policies (
id uuid PRIMARY KEY, -- pol_
tenant_id uuid NOT NULL REFERENCES tenants(id),
workspace_id uuid REFERENCES workspaces(id), -- NULL = tenant-wide
name text NOT NULL,
version int NOT NULL DEFAULT 1,
cedar text NOT NULL, -- Cedar policy text, validated against the schema on save
enabled boolean NOT NULL DEFAULT true,
created_by uuid NOT NULL REFERENCES users(id),
created_at timestamptz NOT NULL DEFAULT now(),
UNIQUE (tenant_id, workspace_id, name, version)
);
-- Every policy version is kept; `enabled` on the newest version is what runs.
-- Policy decisions are logged to ClickHouse events as message.policy_denied / policy_allowed.
-- Post-MVP. Shaped now so the Cedar schema and FK targets are stable.
CREATE TYPE grant_status_t AS ENUM ('requested', 'active', 'suspended', 'revoked', 'expired');
CREATE TABLE grants (
id uuid PRIMARY KEY, -- grt_
tenant_id uuid NOT NULL REFERENCES tenants(id), -- the publishing (granting) tenant
publisher_agent_id uuid NOT NULL REFERENCES agents(id),
consumer_tenant_id uuid NOT NULL REFERENCES tenants(id),
consumer_ws_id uuid, -- NULL = any workspace in consumer tenant
consumer_agent_id uuid, -- NULL = any agent in scope
allowed_types text[] NOT NULL DEFAULT '{agentbus.task.request.v1,agentbus.message.v1}',
capabilities text[] NOT NULL,
rate_per_minute int NOT NULL DEFAULT 60,
pricing jsonb, -- {"model":"per_task","amount":"0.50","currency":"USD"}
approval text NOT NULL CHECK (approval IN ('auto','manual','paid')),
status grant_status_t NOT NULL DEFAULT 'requested',
requested_by uuid,
approved_by uuid,
starts_at timestamptz NOT NULL DEFAULT now(),
expires_at timestamptz,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX ON grants (consumer_tenant_id, status);
1.9 Blobs, audit actions, plans
CREATE TABLE blobs (
id uuid PRIMARY KEY, -- rendered as blb_...
tenant_id uuid NOT NULL REFERENCES tenants(id),
uploaded_by uuid NOT NULL REFERENCES agents(id),
s3_key text NOT NULL,
content_type text NOT NULL,
size_bytes bigint NOT NULL,
sha256 bytea NOT NULL,
dek_version int NOT NULL,
ref_count int NOT NULL DEFAULT 0, -- messages referencing it; GC when 0 past retention
expires_at timestamptz NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
);
-- Actions by humans with elevated access: console admins, support staff, platform operators.
-- Visible to the tenant in their own audit view. Never deleted by retention.
CREATE TABLE audit_actions (
id uuid PRIMARY KEY,
tenant_id uuid REFERENCES tenants(id), -- NULL for platform-level actions
actor_type text NOT NULL CHECK (actor_type IN ('user','support','platform','system')),
actor_id text NOT NULL, -- usr_ id or staff id
action text NOT NULL, -- 'token.create', 'support.view_payload', 'policy.update'
target_type text,
target_id text,
consent_ref uuid, -- support actions must reference a consent window
details jsonb NOT NULL DEFAULT '{}'::jsonb,
ip inet,
user_agent text,
at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX ON audit_actions (tenant_id, at DESC);
-- Time-boxed consent for support staff to view payloads. Created by a tenant admin in the console.
CREATE TABLE support_consents (
id uuid PRIMARY KEY,
tenant_id uuid NOT NULL REFERENCES tenants(id),
granted_by uuid NOT NULL REFERENCES users(id),
scope jsonb NOT NULL, -- {"agents":[...],"from":..,"to":..,"payloads":true}
starts_at timestamptz NOT NULL DEFAULT now(),
ends_at timestamptz NOT NULL,
revoked_at timestamptz,
ticket_ref text
);
CREATE TABLE plans (
code plan_t PRIMARY KEY,
limits jsonb NOT NULL -- {"workspaces":1,"agents":10,"retention_days":7,"pending_per_agent":1000,"stream_bytes":1073741824,"publish_per_minute":60,"blob_bytes":104857600}
);
CREATE TABLE subscriptions (
id uuid PRIMARY KEY,
tenant_id uuid NOT NULL REFERENCES tenants(id) UNIQUE,
plan plan_t NOT NULL,
billing_ref text, -- Stripe customer/subscription id
status text NOT NULL DEFAULT 'active',
current_period_start timestamptz,
current_period_end timestamptz,
overrides jsonb NOT NULL DEFAULT '{}'::jsonb, -- enterprise negotiated limits
created_at timestamptz NOT NULL DEFAULT now()
);
1.10 Row-level security
Every tenant-scoped table has RLS enabled and a policy keyed on a session setting. The gateway
sets the setting per request inside the transaction, immediately after BEGIN:
-- Gateway, per request:
BEGIN;
SELECT set_config('app.tenant_id', 'a1b2c3...', true); -- true = transaction-local
-- ... queries ...
COMMIT;
-- Helper used by every policy.
CREATE FUNCTION current_tenant_id() RETURNS uuid
LANGUAGE sql STABLE AS $$
SELECT NULLIF(current_setting('app.tenant_id', true), '')::uuid
$$;
-- Pattern applied to every table with a tenant_id column.
ALTER TABLE agents ENABLE ROW LEVEL SECURITY;
ALTER TABLE agents FORCE ROW LEVEL SECURITY;
CREATE POLICY tenant_isolation ON agents
USING (tenant_id = current_tenant_id())
WITH CHECK (tenant_id = current_tenant_id());
-- Support role: read-only, and only when a consent window is open for the tenant.
CREATE POLICY support_read ON agents FOR SELECT TO agentbus_support
USING (EXISTS (
SELECT 1 FROM support_consents c
WHERE c.tenant_id = agents.tenant_id
AND now() BETWEEN c.starts_at AND c.ends_at
AND c.revoked_at IS NULL
));
Rules:
FORCE ROW LEVEL SECURITYso even the table owner is subject to policies in tests.- The gateway role has no
BYPASSRLS. A request that forgets to setapp.tenant_idsees zero rows, which fails loudly in tests rather than leaking. - Cross-tenant operations (bridging, marketplace listing lookups) run through a dedicated
agentbus_bridgerole with explicit, narrow policies, never by unsetting the tenant. - Platform-level tables (
plans,signing_keys) have no RLS and are not readable byagentbus_appbeyond what views expose.
1.11 Retention and maintenance jobs
| Job | Schedule | Action |
|---|---|---|
partitions.create | daily | Create messages_YYYY_MM partitions 3 months ahead |
partitions.drop | daily | Detach and drop partitions entirely older than the longest tenant retention; per-tenant rows inside a live partition are deleted by retention.messages |
retention.messages | hourly | DELETE FROM messages WHERE tenant_id = $1 AND accepted_at < now() - retention_days in 10k batches; S3 payload objects deleted by lifecycle rule keyed on tag |
retention.blobs | hourly | Delete blobs with ref_count = 0 past expires_at |
retention.tokens | daily | Hard-delete revoked tokens after 90 days (audit keeps the event) |
reconcile.unpublished | every 10 s | Re-publish messages with null stream_seq older than 10 s |
reconcile.presence | every 30 s | Mirror NATS KV presence into agents.status / last_seen_at |
tasks.expire | every minute | Move tasks past expires_at in non-terminal state to expired |
keys.rotate | daily | Retire signing keys older than 180 days after a 30-day overlap |
tenant.delete | on demand | §5.5 crypto-shred, then queue hard deletes |
2. ClickHouse analytics and audit
2.1 Cluster and ingestion
- One ClickHouse cluster per region, replicated (
ReplicatedMergeTreevia ClickHouse Keeper), two shards is enough for the first years. - The ingester is a Go worker that consumes
SYS_AUDITfrom NATS (durable consumerclickhouse-ingest) and inserts in batches of 10k rows or 1 s, whichever first. Inserts are idempotent viaevent_idandReplacingMergeTreewhere needed. - ClickHouse never holds payload plaintext. Payload size, hash, and type are enough for every analytics question we want to answer.
2.2 Storage policy: hot NVMe, cold S3
<!-- config.d/storage.xml -->
<clickhouse>
<storage_configuration>
<disks>
<hot><path>/var/lib/clickhouse/hot/</path></hot>
<cold>
<type>s3</type>
<endpoint>https://s3.<region>.example/agentbus-ch-cold/</endpoint>
<metadata_path>/var/lib/clickhouse/disks/cold/</metadata_path>
<cache_enabled>true</cache_enabled>
<data_cache_max_size>107374182400</data_cache_max_size>
</cold>
</disks>
<policies>
<tiered>
<volumes>
<hot><disk>hot</disk><max_data_part_size_bytes>10737418240</max_data_part_size_bytes></hot>
<cold><disk>cold</disk></cold>
</volumes>
<move_factor>0.1</move_factor>
</tiered>
</policies>
</storage_configuration>
</clickhouse>
ClickHouse Cloud provides equivalent tiering automatically; the DDL below is the same either way.
2.3 events table
CREATE TABLE agentbus.events ON CLUSTER main
(
event_id UUID,
event_time DateTime64(3, 'UTC'),
event LowCardinality(String), -- message.accepted, message.delivered, task.state_changed ...
tenant_id UUID,
workspace_id UUID,
message_id UUID,
conversation_id UUID,
task_id Nullable(UUID),
trace_id String,
span_id String,
from_agent_id UUID,
to_agent_id Nullable(UUID),
type LowCardinality(String), -- envelope type
capability LowCardinality(Nullable(String)),
attempt UInt8,
actor LowCardinality(String), -- gateway, sidecar, dlq_worker, api
payload_bytes UInt32,
latency_ms Nullable(UInt32), -- time since previous event for this message
error_code LowCardinality(Nullable(String)),
policy_id Nullable(UUID),
bridged_from Nullable(UUID),
details String CODEC(ZSTD(3)), -- JSON, small, never payload
retention_days UInt16 MATERIALIZED dictGet('agentbus.tenant_retention', 'retention_days', tenant_id),
ingested_at DateTime DEFAULT now()
)
ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/events', '{replica}', ingested_at)
PARTITION BY toYYYYMM(event_time)
ORDER BY (tenant_id, event_time, message_id, event)
TTL toDateTime(event_time) + INTERVAL 30 DAY TO VOLUME 'cold',
toDateTime(event_time) + toIntervalDay(retention_days) DELETE
SETTINGS storage_policy = 'tiered', ttl_only_drop_parts = 0, index_granularity = 8192;
-- Dictionary backing retention_days, refreshed from Postgres every 5 minutes.
CREATE DICTIONARY agentbus.tenant_retention
(
tenant_id UUID,
retention_days UInt16
)
PRIMARY KEY tenant_id
SOURCE(POSTGRESQL(host 'pg-replica' port 5432 user 'agentbus_ro' password '' db 'agentbus' table 'tenants'))
LIFETIME(MIN 300 MAX 600)
LAYOUT(HASHED());
Notes:
retention_daysis materialised at insert, so later plan changes affect new rows only. A plan upgrade that lengthens retention triggers anALTER TABLE ... UPDATE retention_daysfor that tenant (mutations are fine at tenant granularity).- The 30-day hot window is a platform setting. Everything older is on S3 with a local cache,
which is where most
agentbus tracelookups for old receipts land. ORDER BY (tenant_id, event_time, ...)makes every tenant-scoped time-range query a contiguous read.
2.4 usage table
CREATE TABLE agentbus.usage ON CLUSTER main
(
usage_id UUID,
recorded_at DateTime64(3, 'UTC'),
tenant_id UUID,
workspace_id UUID,
agent_id UUID,
task_id Nullable(UUID),
message_id Nullable(UUID),
conversation_id Nullable(UUID),
harness LowCardinality(String),
provider LowCardinality(String), -- anthropic, openai, google, other
model LowCardinality(String),
input_tokens UInt64,
output_tokens UInt64,
cache_read_tokens UInt64,
cache_write_tokens UInt64,
reasoning_tokens UInt64,
cost_usd Decimal(14, 6), -- computed server-side from the pricing table unless provider-reported
cost_source LowCardinality(String), -- 'computed', 'harness_reported', 'estimate'
reported_by LowCardinality(String), -- 'harness', 'proxy', 'agent', 'estimate'
confidence LowCardinality(String), -- 'exact', 'approximate', 'unverified'
bus_bytes_in UInt64, -- bytes through the bus for this unit of work
bus_bytes_out UInt64,
bus_messages UInt32,
duration_ms Nullable(UInt64),
retention_days UInt16 MATERIALIZED dictGet('agentbus.tenant_retention', 'retention_days', tenant_id),
details String CODEC(ZSTD(3))
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/usage', '{replica}')
PARTITION BY toYYYYMM(recorded_at)
ORDER BY (tenant_id, recorded_at, agent_id)
TTL toDateTime(recorded_at) + INTERVAL 30 DAY TO VOLUME 'cold',
toDateTime(recorded_at) + toIntervalDay(greatest(retention_days, 400)) DELETE -- billing data kept >= 13 months
SETTINGS storage_policy = 'tiered';
-- Model pricing, maintained by platform. Joined at insert time by the ingester.
CREATE TABLE agentbus.model_pricing
(
provider LowCardinality(String),
model LowCardinality(String),
effective_from Date,
input_per_mtok Decimal(10, 4),
output_per_mtok Decimal(10, 4),
cache_read_per_mtok Decimal(10, 4),
cache_write_per_mtok Decimal(10, 4)
)
ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/model_pricing', '{replica}')
ORDER BY (provider, model, effective_from);
reported_by and confidence are first-class because the bus cannot observe LLM tokens
itself (08-HARNESS-ADAPTERS.md, "Metering"). Every dashboard shows them.
2.5 audit_chain table
Tamper evidence for the customer-visible audit log. Each tenant has one chain. Each row's hash commits to the previous row's hash, so deleting or editing any row breaks every later row.
CREATE TABLE agentbus.audit_chain ON CLUSTER main
(
tenant_id UUID,
seq UInt64, -- per-tenant monotonic, assigned by the chain writer
at DateTime64(3, 'UTC'),
event_id UUID, -- FK into events or audit_actions
event_kind LowCardinality(String), -- 'message_event' | 'audit_action' | 'checkpoint'
actor String,
summary String, -- short, no payload
prev_hash FixedString(32),
row_hash FixedString(32), -- SHA-256(tenant_id || seq || at || event_id || event_kind || actor || summary || prev_hash)
signature Nullable(String) -- Ed25519 by the platform signing key on every 1000th row (checkpoint)
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/audit_chain', '{replica}')
PARTITION BY toYYYYMM(at)
ORDER BY (tenant_id, seq)
SETTINGS storage_policy = 'tiered';
-- No DELETE TTL. Audit chain rows outlive message retention. Cold tier only.
Chain writer and verification:
- A single-writer-per-tenant Go worker (sharded by tenant id) consumes
sys.audit.>andaudit_actionsinserts, assignsseq, computesrow_hash, and inserts. The writer keeps the last(seq, row_hash)per tenant in NATS KVAUDIT_HEADso it survives restarts. - Every 1000 rows the writer inserts a
checkpointrow signed with the platform signing key and publishes the checkpoint hash to a public transparency log endpoint (GET /audit/checkpoints, 04-API-SPEC.md), so a tenant can later prove the platform did not rewrite history either. agentbus audit verify --from <seq> --to <seq>(and the console's "Verify" button) recomputes hashes in order and checks checkpoint signatures. Verification output names the first bad row.- Enterprise tenants can request the checkpoint hashes be mirrored to their own storage.
2.6 Materialized views
-- Daily cost per agent, feeds the console cost page.
CREATE MATERIALIZED VIEW agentbus.mv_cost_daily ON CLUSTER main
ENGINE = ReplicatedSummingMergeTree('/clickhouse/tables/{shard}/mv_cost_daily', '{replica}')
PARTITION BY toYYYYMM(day)
ORDER BY (tenant_id, workspace_id, agent_id, day, model, reported_by)
AS SELECT
tenant_id, workspace_id, agent_id,
toDate(recorded_at) AS day,
model, reported_by,
sum(input_tokens) AS input_tokens,
sum(output_tokens) AS output_tokens,
sum(cache_read_tokens) AS cache_read_tokens,
sum(cost_usd) AS cost_usd,
sum(bus_messages) AS bus_messages,
sum(bus_bytes_in + bus_bytes_out) AS bus_bytes
FROM agentbus.usage
GROUP BY tenant_id, workspace_id, agent_id, day, model, reported_by;
-- Who talked to whom: directed edges per day, for the graph view and for anomaly alerts.
CREATE MATERIALIZED VIEW agentbus.mv_edges_daily ON CLUSTER main
ENGINE = ReplicatedSummingMergeTree('/clickhouse/tables/{shard}/mv_edges_daily', '{replica}')
PARTITION BY toYYYYMM(day)
ORDER BY (tenant_id, day, from_agent_id, to_agent_id, type)
AS SELECT
tenant_id,
toDate(event_time) AS day,
from_agent_id, to_agent_id, type,
count() AS messages,
sum(payload_bytes) AS bytes,
countIf(event = 'message.dead_lettered') AS dead_lettered,
countIf(event = 'message.expired') AS expired
FROM agentbus.events
WHERE event IN ('message.accepted','message.dead_lettered','message.expired') AND to_agent_id IS NOT NULL
GROUP BY tenant_id, day, from_agent_id, to_agent_id, type;
-- Delivery latency percentiles per tenant per hour, for SLO dashboards.
CREATE MATERIALIZED VIEW agentbus.mv_delivery_latency_hourly ON CLUSTER main
ENGINE = ReplicatedAggregatingMergeTree('/clickhouse/tables/{shard}/mv_delivery_latency_hourly', '{replica}')
PARTITION BY toYYYYMM(hour)
ORDER BY (tenant_id, hour)
AS SELECT
tenant_id,
toStartOfHour(event_time) AS hour,
quantilesState(0.5, 0.95, 0.99)(latency_ms) AS accept_to_deliver_ms,
countState() AS deliveries
FROM agentbus.events
WHERE event = 'message.delivered'
GROUP BY tenant_id, hour;
2.7 Query examples the console runs
-- Trace timeline for a receipt
SELECT event, event_time, attempt, actor, latency_ms, error_code, details
FROM agentbus.events
WHERE tenant_id = {tenant:UUID} AND message_id = {msg:UUID}
ORDER BY event_time;
-- Cost by agent, last 30 days, with confidence
SELECT agent_id, reported_by, sum(cost_usd) AS usd, sum(input_tokens + output_tokens) AS tokens
FROM agentbus.mv_cost_daily
WHERE tenant_id = {tenant:UUID} AND day >= today() - 30
GROUP BY agent_id, reported_by ORDER BY usd DESC;
3. Object storage layout
One bucket per region. Tenants are prefixes; IAM never grants tenants direct bucket access, all access is via gateway-issued pre-signed URLs with 15-minute lifetimes.
s3://agentbus-<region>/
tenants/<tenant_id>/
payloads/YYYY/MM/DD/<msg_id>.bin # encrypted envelope `data`, tagged retention=<days>
blobs/<blob_id> # encrypted claim-check objects, tagged retention=<days>
archive/YYYY/MM/events-<shard>.parquet # monthly export of events for the tenant (Business+)
archive/YYYY/MM/usage-<shard>.parquet
exports/<export_id>/... # on-demand tenant exports, 7-day lifetime
platform/
audit-checkpoints/<tenant_id>/<seq>.json # signed checkpoints, write-once (object lock)
clickhouse-cold/ # ClickHouse cold volume, managed by ClickHouse
Lifecycle rules:
| Prefix | Rule |
|---|---|
payloads/, blobs/ | Expire by object tag retention=<days> (rules per distinct retention value: 7, 30, 90, 180, 365, custom) |
archive/ | Transition to infrequent access after 30 days, keep for 13 months, then expire unless enterprise override |
exports/ | Expire after 7 days |
audit-checkpoints/ | Object lock, compliance mode, 7 years |
Every object is written with server-side encryption on top of the application-level encryption in §5; SSE is defence in depth, not the tenant boundary.
4. Message payload path
sender sidecar ──(TLS)──▶ gateway
│ 1. validate envelope, extract `data`
│ 2. plaintext sha256 → messages.payload_sha256
│ 3. encrypt with tenant DEK (AES-256-GCM, AAD = msg_id || tenant_id)
│ 4. PUT s3://.../payloads/.../<msg_id>.bin
│ 5. insert messages row with payload_ref
▼
NATS message = envelope with `data` replaced by {"$ref": payload_ref}
│
recipient sidecar ◀─(TLS)── gateway fetches on behalf of sidecar, re-inlines `data` after
decrypt and policy check, delivers the full envelope
Payloads under 64 KiB are also carried inline in the NATS message (encrypted) to save an S3 round-trip on delivery; the S3 object is still written as the system of record. This is an optimisation flag in the gateway, off by default in the MVP.
5. Encryption at rest
5.1 Key hierarchy
Platform root key (KMS, HSM-backed, never leaves KMS)
└── Tenant KEK (one per tenant; in platform KMS by default, in the customer's KMS for BYOK)
└── Tenant DEK v1, v2, ... (AES-256 keys, generated by the gateway, stored wrapped by the KEK)
└── Per-object encryption: AES-256-GCM with a random 96-bit nonce per object
CREATE TABLE tenant_keys (
tenant_id uuid NOT NULL REFERENCES tenants(id),
version int NOT NULL,
dek_wrapped bytea NOT NULL, -- DEK encrypted by the KEK
kek_ref text NOT NULL, -- 'kms://platform/<key-id>' or 'kms://aws/arn:...' or 'vault://...'
kek_provider text NOT NULL CHECK (kek_provider IN ('platform','aws_kms','gcp_kms','vault')),
status text NOT NULL DEFAULT 'active' CHECK (status IN ('active','retired','destroyed')),
created_at timestamptz NOT NULL DEFAULT now(),
retired_at timestamptz,
PRIMARY KEY (tenant_id, version)
);
5.2 Runtime
- The gateway unwraps the active DEK at startup per tenant on first use and caches it in
memory for 10 minutes. A BYOK tenant whose KMS is unreachable gets
AB-9020 key_unavailableon every operation that touches payloads; metadata operations still work. - Encrypt: AES-256-GCM, nonce random per object, AAD binds the ciphertext to
msg_idandtenant_idso a ciphertext cannot be transplanted between messages or tenants. - The
dek_versiononmessagesandblobsselects the key to unwrap on read.
5.3 Rotation
- DEK rotation: create version n+1, mark n
retired. New writes use n+1; reads use whichever version the row names. Retired DEKs are never destroyed while rows reference them. Default cadence 90 days, or on demand from the console. - KEK rotation (platform): KMS rotates the KEK material; all wrapped DEKs are re-wrapped by a background job. No payload bytes are touched.
- BYOK KEK rotation: the customer rotates in their KMS; we re-wrap on the next unwrap failure or on their request.
5.4 Bring your own key (Enterprise)
- Tenant admin registers a KEK reference and grants our service principal
DecryptandEncrypt(wrap/unwrap) on it.agentbus doctor --tenantand the console verify the grant. - Every unwrap call is logged in the customer's KMS audit trail; they can see every time we touched their key.
- Revoking our access makes all tenant payloads unreadable within the 10-minute cache window. Metadata, envelopes, and audit remain readable because they never used the DEK.
5.5 Crypto-shredding on tenant deletion
- Tenant status set to
deleting; all tokens revoked; all agents disabled. - All
tenant_keysrows set todestroyed; wrapped DEKs overwritten with zeros; for platform KEKs the KMS key is scheduled for deletion. - From this moment every payload and blob is unrecoverable, even before S3 objects are physically removed.
- Hard deletes of Postgres rows, S3 objects, and ClickHouse partitions proceed over the next
30 days.
audit_chainandaudit_actionsare retained in anonymised form (tenant id kept, actor ids hashed) for 7 years as required for the platform's own compliance.
6. Data residency
tenants.regionselects the region cluster: a Postgres primary, NATS cluster, ClickHouse cluster, and S3 bucket per region. MVP launches one region; the schema and routing layer assume many.- The global control plane (
tenants,users,subscriptions, DNS routing table) lives in the home region and is replicated read-only to others sokomsary.agentbus.exchangecan route a request to the right region from the token prefix lookup. - Tenant data never crosses regions. Cross-tenant bridging between tenants in different
regions is refused in the MVP (
AB-7006 cross_region_not_supported) and is a post-MVP design item. - Enterprise custom domains are CNAMEs to the tenant's region gateway.
7. Sidecar local store (SQLite)
The sidecar keeps a per-agent SQLite database at
$XDG_STATE_HOME/agentbus/agents/<agent_id>/state.db (WAL mode, synchronous = FULL for
the inbox table). It is the durable half of the two-level ack (05-DELIVERY-SEMANTICS.md).
PRAGMA journal_mode = WAL;
PRAGMA synchronous = FULL;
CREATE TABLE meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL -- agent_id, address, key_fingerprint, server_time_offset_ms, schema_version
);
-- Inbound messages, written before transport ack.
CREATE TABLE inbox (
msg_id TEXT PRIMARY KEY,
stream_seq INTEGER NOT NULL,
conversation_id TEXT NOT NULL,
type TEXT NOT NULL,
from_address TEXT NOT NULL,
trace_id TEXT NOT NULL,
expires_at TEXT,
envelope BLOB NOT NULL, -- full envelope JSON, payload inline (decrypted by gateway on delivery)
received_at TEXT NOT NULL,
transport_acked INTEGER NOT NULL DEFAULT 0,
presented_at TEXT, -- when handed to the adapter
read_at TEXT, -- application ack time
redelivery INTEGER NOT NULL DEFAULT 0,
out_of_order INTEGER NOT NULL DEFAULT 0,
attempt INTEGER NOT NULL DEFAULT 1
);
CREATE INDEX inbox_unread ON inbox (read_at) WHERE read_at IS NULL;
CREATE INDEX inbox_seq ON inbox (stream_seq);
-- Outbound messages and acks waiting for the gateway.
CREATE TABLE outbox (
id TEXT PRIMARY KEY, -- msg_id for sends, 'ack:'||msg_id||':'||level for acks
kind TEXT NOT NULL CHECK (kind IN ('send','ack','nak','heartbeat','usage')),
body BLOB NOT NULL,
created_at TEXT NOT NULL,
attempts INTEGER NOT NULL DEFAULT 0,
next_attempt_at TEXT NOT NULL,
last_error TEXT
);
CREATE INDEX outbox_due ON outbox (next_attempt_at);
-- Receipts for everything we sent, so `agentbus trace` works offline for the local half.
CREATE TABLE receipts (
rcp_id TEXT PRIMARY KEY,
msg_id TEXT NOT NULL,
stream_seq INTEGER,
accepted_at TEXT NOT NULL,
to_address TEXT NOT NULL,
type TEXT NOT NULL,
local_events TEXT NOT NULL -- JSON array of {event, at}
);
-- Dedupe horizon, 7 days.
CREATE TABLE seen_msg_ids (
msg_id TEXT PRIMARY KEY,
seen_at TEXT NOT NULL,
app_acked INTEGER NOT NULL DEFAULT 0
);
-- Adapter-specific state: Codex thread id, OpenCode session id, Claude session id, cursor positions.
CREATE TABLE adapter_state (
adapter TEXT NOT NULL,
key TEXT NOT NULL,
value TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (adapter, key)
);
-- Usage records captured locally before upload.
CREATE TABLE usage_local (
id TEXT PRIMARY KEY,
recorded_at TEXT NOT NULL,
body BLOB NOT NULL,
uploaded INTEGER NOT NULL DEFAULT 0
);
Local retention: inbox rows are deleted 7 days after read_at; outbox rows after success;
receipts after 30 days; seen_msg_ids after 7 days. agentbus doctor reports the database
size and any rows stuck in the outbox.
The credential file ($XDG_CONFIG_HOME/agentbus/credentials.json) is separate from the
state database, mode 0600, and holds the integration token and the agent's Ed25519 private
key. On macOS and Windows the sidecar uses the OS keychain when available.
8. Retention matrix
| Data | Starter | Business | Enterprise | Notes |
|---|---|---|---|---|
| Message metadata (Postgres) | 7 d | 90 d | configurable, up to 7 y | drives agentbus trace and console |
| Payloads and blobs (S3) | 7 d | 90 d | configurable | encrypted, crypto-shreddable |
| Delivery events (ClickHouse) | 7 d | 90 d | configurable | 30 d hot, then cold |
| Usage and cost (ClickHouse) | 13 mo | 13 mo | configurable, min 13 mo | billing reconciliation |
| Audit chain | 7 y | 7 y | 7 y, mirrored on request | never shortened |
| Audit actions (support/admin) | 7 y | 7 y | 7 y | never shortened |
| Dead-letter queue | 7 d | 90 d | configurable | same as messages |
| NATS stream MaxAge | 7 d | 90 d | configurable | matches message retention |
| Sidecar local inbox | 7 d after read | same | same | local only |
| Revoked tokens | 90 d | 90 d | 90 d | event kept in audit |
| Tenant exports | 7 d | 7 d | 7 d |
Plan limits that are not retention (agents, pending depth, stream bytes, rate) are in
plans.limits and summarised in 01-PRD.md and 05-DELIVERY-SEMANTICS.md §11; those
three places must agree.
9. Migration and schema versioning
- Postgres migrations are plain SQL files applied by a Go migrator at gateway start, forward
only, with a
schema_migrationstable. Every migration is backward compatible with the previous gateway version so rolling deploys work. - ClickHouse DDL changes are applied by the same migrator against
ON CLUSTER; additive columns only, withALTER TABLE ... ADD COLUMN ... DEFAULT. - The sidecar SQLite schema has a
meta.schema_version; the sidecar migrates on startup and refuses to downgrade. - Envelope schema versions (
agentbus.*.v1) are independent of storage schema versions. A new envelope major adds rows, not columns, because the envelope is stored as metadata plus opaque payload.