@ -52,51 +52,48 @@ events into app events. Consumers subscribe **only** to app events.
The model is **Domain Events + Integration Events** with an
**Anti-Corruption Layer**:
| Active term | DDD term |
| ------------------------------ | ----------------------------------- |
| Module events (`SESSION_EVENT_*`, …) | Domain events (private, iterable) |
| App events (`APP_EVENT_*`) | Integration events (public, stable) |
| Active term | DDD term |
| ----------------------------------------- | - ----------------------------------- |
| Module events (`SESSION_EVENT_*`, …) | Domain events (private, iterable) |
| App events (`APP_EVENT_*`) | Integration events (public, stable) |
| `active-app/presets/*` or service bridges | Anti-corruption layer + event mapper |
| Per-consumer auto-reactions | Stateless process managers |
| Per-consumer auto-reactions | Stateless process managers |
If you have a DDD background the model maps 1:1.
## Layer split
```
arts/bus/ ← engine, framework-agnostic
types.ts EngineBus, BusEnvelope, BusListener,
BusPublishOptions, BusPublishResult,
EventPublisher, BusAnyListener
consts.ts BUS_*, listener error modes, diagnostics
engine-bus.ts createEngineBus()
no Svelte imports
errors.ts / matching.ts / diagnostics.ts
test/
libs/bus/ ← pure contracts (framework-agnostic, zero-dep)
types.ts BusEnvelope, BusListener, BusListenerErrorMode,
BusPublishOptions, BusPublishResult, BusListenerFailure,
EventPublisher, EventSubscriber, BusAnyListener, BusClock
consts.ts BUS_* literals, listener error modes, diagnostic events
errors.ts BusDisposedError + 4 siblings + guards
matching.ts event-name matching helpers
diagnostics.ts createBusDiagnostics()
arts/bus/ ← engine (framework-agnostic runtime)
types.ts EngineBus, EngineBusOptions (compose the libs contracts)
engine-bus.ts createEngineBus() — no Svelte imports
active-bus.svelte.ts reactive wrapper: createBusRecent()
{ lastEvent, count, clear, dispose }
silent-bus.ts no-op EngineBus (SSR / bus disabled)
arts/bus/svelte/ ← Svelte adapter
index.ts createSvelteEngineBus()
wraps listener invocation in untrack
arts/bus/active-bus.svelte.ts ← reactive wrappers (minimal surface)
createBusRecent() { lastEvent, count, clear, dispose }
index.ts createSvelteEngineBus() — wraps listeners in untrack
context.svelte.ts getBus / setBus via createContext
arts/active-app/events.ts ← APP_EVENT_* constants + payloads
+ APP_EVENT_RUNTIMES metadata
arts/bus/svelte/context.svelte.ts ← getBus / setBus via createContext
arts/active-app/presets/ ← orca reaction presets (wired by applyStandardOrca)
cache-clear-on-identity-change / cache-clear-on-revoke
connections-close-on-revoke / connections-reauth-on-identity-change
perm-invalidate-on-identity-change / session-auto-refresh / standard
arts/< module > /consts.ts ← module event constants
arts/< module > /types.ts ← module event payload types
arts/< module > /bus-helpers.ts ← typed publish/subscribe helpers
(the only path to call-site type safety)
arts/active-app/presets/ ← translators/reactions (anti-corruption layer)
identity-translator.ts module events → app events
permissions-translator.ts
tenant-translator.ts
connectivity-translator.ts
dispose-translator.ts
arts/< module > /bus-helpers.ts ← typed publish/subscribe helpers (call-site type safety)
```
**Hard rule**: `arts/bus/engine-bus.ts` does not import from `svelte` .
@ -108,64 +105,64 @@ engine must be testable without DOM or Svelte runtime.
```ts
export interface EngineBus {
publish< TType extends string , TPayload > (
type: TType,
payload: TPayload,
options?: BusPublishOptions
): BusPublishResult< TType , TPayload > ;
publishAsync< TType extends string , TPayload > (
type: TType,
payload: TPayload,
options?: BusPublishOptions
): Promise< BusPublishResult < TType , TPayload > >;
/**
* Publish with `causationId` automatically set to `parent.id` . Use
* inside translators to keep the causation chain populated.
*/
publishCausedBy< TType extends string , TPayload > (
parent: BusEnvelope,
type: TType,
payload: TPayload,
options?: BusPublishOptions
): BusPublishResult< TType , TPayload > ;
on< TType extends string , TPayload > (
type: TType,
listener: BusListener< TPayload > ,
options?: BusListenOptions
): BusSubscription;
/**
* Sugar for `$effect` : returns the unsubscribe function directly so a
* Svelte component writes `$effect(() => Bus.subscribe(type, fn))` .
*/
subscribe< TType extends string , TPayload > (
type: TType,
listener: BusListener< TPayload > ,
options?: BusListenOptions
): () => void;
onAny(listener: BusAnyListener< TEvents > , options?: BusListenOptions): BusSubscription;
once< TType extends string , TPayload > (
type: TType,
listener: BusListener< TPayload > ,
options?: BusListenOptions
): BusSubscription;
listenerCount(type?: string): number;
_clearForTesting(type?: string): void;
dispose(): void;
publish< TType extends string , TPayload > (
type: TType,
payload: TPayload,
options?: BusPublishOptions
): BusPublishResult< TType , TPayload > ;
publishAsync< TType extends string , TPayload > (
type: TType,
payload: TPayload,
options?: BusPublishOptions
): Promise< BusPublishResult < TType , TPayload > >;
/**
* Publish with `causationId` automatically set to `parent.id` . Use
* inside translators to keep the causation chain populated.
*/
publishCausedBy< TType extends string , TPayload > (
parent: BusEnvelope,
type: TType,
payload: TPayload,
options?: BusPublishOptions
): BusPublishResult< TType , TPayload > ;
on< TType extends string , TPayload > (
type: TType,
listener: BusListener< TPayload > ,
options?: BusListenOptions
): BusSubscription;
/**
* Sugar for `$effect` : returns the unsubscribe function directly so a
* Svelte component writes `$effect(() => Bus.subscribe(type, fn))` .
*/
subscribe< TType extends string , TPayload > (
type: TType,
listener: BusListener< TPayload > ,
options?: BusListenOptions
): () => void;
onAny(listener: BusAnyListener< TEvents > , options?: BusListenOptions): BusSubscription;
once< TType extends string , TPayload > (
type: TType,
listener: BusListener< TPayload > ,
options?: BusListenOptions
): BusSubscription;
listenerCount(type?: string): number;
_clearForTesting(type?: string): void;
dispose(): void;
}
export interface EventPublisher {
publish< TType extends string , TPayload > (
type: TType,
payload: TPayload,
options?: BusPublishOptions
): BusPublishResult< TType , TPayload > ;
publish< TType extends string , TPayload > (
type: TType,
payload: TPayload,
options?: BusPublishOptions
): BusPublishResult< TType , TPayload > ;
}
```
@ -173,15 +170,15 @@ export interface EventPublisher {
```ts
export interface BusEnvelope< TType extends string = string, TPayload = unknown > {
readonly id: string;
readonly type: TType;
readonly payload: TPayload;
readonly at: number; // epoch ms
readonly source: string;
readonly correlationId?: string;
readonly causationId?: string;
readonly context?: Readonly< Record < string , unknown > >;
readonly tags?: readonly string[];
readonly id: string;
readonly type: TType;
readonly payload: TPayload;
readonly at: number; // epoch ms
readonly source: string;
readonly correlationId?: string;
readonly causationId?: string;
readonly context?: Readonly< Record < string , unknown > >;
readonly tags?: readonly string[];
}
```
@ -197,18 +194,18 @@ carry those fields on every publish.
```ts
export interface EngineBusOptions {
readonly logger?: Logger;
readonly clock?: BusClock;
readonly idFactory?: () => string;
readonly maxListenersPerEvent?: number; // default 32
readonly maxReentrancyDepth?: number; // default 32
readonly listenerErrorMode?: BusListenerErrorMode;
/**
* Called for every listener invocation. The Svelte adapter wraps with
* `untrack` . Default: identity.
*/
readonly invokeListener?: (fn: () => void | Promise< void > ) => void | Promise< void > ;
readonly logger?: Logger;
readonly clock?: BusClock;
readonly idFactory?: () => string;
readonly maxListenersPerEvent?: number; // default 32
readonly maxReentrancyDepth?: number; // default 32
readonly listenerErrorMode?: BusListenerErrorMode;
/**
* Called for every listener invocation. The Svelte adapter wraps with
* `untrack` . Default: identity.
*/
readonly invokeListener?: (fn: () => void | Promise< void > ) => void | Promise< void > ;
}
```
@ -216,18 +213,34 @@ The engine never imports from `svelte`. The `invokeListener` hook is the
seam the Svelte adapter uses to wrap calls in `untrack` (see
[Svelte/SvelteKit Runtime Contract ](#sveltesveltekit-runtime-contract )).
### Listener error modes
`listenerErrorMode` — the per-engine default (`EngineBusOptions`) or overridden
per `publish` (`BusPublishOptions`) — decides what happens when a listener throws.
The three constants live in `$libs/bus` :
| Mode | Behaviour |
| -------------------------------- | -------------------------------------------------------------------------------------- |
| `log-and-continue` (**default**) | emit a `bus.listener.failed` diagnostic per failure, run the remaining listeners |
| `throw` | run every listener, then throw an aggregated `BusAggregateListenerError` if any failed |
| `collect` | silent — no diagnostic, no throw; the caller reads the failures off the result |
Every mode runs **all** listeners and records the `BusListenerFailure[]` on
`BusPublishResult.errors` ; they differ only in the side effect (diagnostic / throw
/ neither).
### Listener failure shape
`BusListenerFailure` carries enough context for distributed debugging:
```ts
export interface BusListenerFailure {
readonly listenerId?: string;
readonly type: string;
readonly envelopeId: string;
readonly correlationId?: string;
readonly causationId?: string;
readonly error: unknown;
readonly listenerId?: string;
readonly type: string;
readonly envelopeId: string;
readonly correlationId?: string;
readonly causationId?: string;
readonly error: unknown;
}
```
@ -237,8 +250,8 @@ export interface BusListenerFailure {
```ts
export type BusAnyListener< TEvents extends BusEventMap > = (
event: BusEventUnion< TEvents > ,
context: BusListenerContext
event: BusEventUnion< TEvents > ,
context: BusListenerContext
) => void | Promise< void > ;
```
@ -283,12 +296,12 @@ export const SESSION_EVENT_REFRESHED = 'session.refreshed';
// arts/session/types.ts
export interface SessLifecyclePayload {
readonly event: SessionEvent;
readonly generation: number;
readonly identity: {
readonly from: SessionIdentityState;
readonly to: SessionIdentityState;
};
readonly event: SessionEvent;
readonly generation: number;
readonly identity: {
readonly from: SessionIdentityState;
readonly to: SessionIdentityState;
};
}
// arts/session/bus-helpers.ts
@ -299,36 +312,36 @@ import type { BusPublishOptions, EventPublisher } from '$libs/bus';
// identity actually transitions; `session.revoked` / `session.expired` /
// `session.refreshed` only for the matching lifecycle event.
export function publishSessLifecycleEvent(
bus: SessEventPublisher,
payload: SessLifecyclePayload,
options: BusPublishOptions = {}
bus: SessEventPublisher,
payload: SessLifecyclePayload,
options: BusPublishOptions = {}
): void {
if (payload.event === SESSION_EVENT_LIFECYCLE_INITIAL) return;
const opts = { source: SESSION_MODULE, ...options };
bus.publish(SESSION_EVENT_CHANGED, payload, opts);
if (payload.identity.from !== payload.identity.to) {
bus.publish(SESSION_EVENT_IDENTITY_CHANGED, payload, opts);
}
if (payload.event === SESSION_EVENT_LIFECYCLE_REVOKED) {
bus.publish(SESSION_EVENT_REVOKED, payload, opts);
}
if (payload.event === SESSION_EVENT_LIFECYCLE_EXPIRED) {
bus.publish(SESSION_EVENT_EXPIRED, payload, opts);
}
if (payload.event === SESSION_EVENT_LIFECYCLE_REFRESHED) {
bus.publish(SESSION_EVENT_REFRESHED, payload, opts);
}
if (payload.event === SESSION_EVENT_LIFECYCLE_INITIAL) return;
const opts = { source: SESSION_MODULE, ...options };
bus.publish(SESSION_EVENT_CHANGED, payload, opts);
if (payload.identity.from !== payload.identity.to) {
bus.publish(SESSION_EVENT_IDENTITY_CHANGED, payload, opts);
}
if (payload.event === SESSION_EVENT_LIFECYCLE_REVOKED) {
bus.publish(SESSION_EVENT_REVOKED, payload, opts);
}
if (payload.event === SESSION_EVENT_LIFECYCLE_EXPIRED) {
bus.publish(SESSION_EVENT_EXPIRED, payload, opts);
}
if (payload.event === SESSION_EVENT_LIFECYCLE_REFRESHED) {
bus.publish(SESSION_EVENT_REFRESHED, payload, opts);
}
}
// Subscriber: one event, one listener. `onSessChanged` covers every
// lifecycle transition; subscribe to `SESSION_EVENT_IDENTITY_CHANGED` /
// `SESSION_EVENT_REVOKED` / etc directly when you only care about a slice.
export function onSessChanged(
bus: EngineBus< SessEventMap > ,
listener: (event: SessChangedEnvelope) => void | Promise< void >
bus: EngineBus< SessEventMap > ,
listener: (event: SessChangedEnvelope) => void | Promise< void >
): BusSubscription {
return bus.on(SESSION_EVENT_CHANGED, (event) => listener(event as SessChangedEnvelope));
return bus.on(SESSION_EVENT_CHANGED, (event) => listener(event as SessChangedEnvelope));
}
```
@ -359,8 +372,8 @@ never a module singleton** (see Rule 1 below).
import { createSvelteEngineBus } from '$bus/svelte';
const Bus = createSvelteEngineBus({
logger: Logger,
clock: Timers.clock
logger: Logger,
clock: Timers.clock
});
const App = { Logger, Lang, Format, Dom, Storage, Http, Timers, Bus, Cache };
@ -414,11 +427,14 @@ import type { EngineBus } from '../types';
const BUS_CONTEXT = Symbol('arts.bus.context');
export function setBus(bus) { setContext(BUS_CONTEXT, bus); return bus; }
export function setBus(bus) {
setContext(BUS_CONTEXT, bus);
return bus;
}
export function getBus() {
const bus = getContext(BUS_CONTEXT);
if (!bus) throw new BusNoContextError();
return bus;
const bus = getContext(BUS_CONTEXT);
if (!bus) throw new BusNoContextError();
return bus;
}
```
@ -427,8 +443,8 @@ Root component sets it; descendants read it:
```svelte
<!-- src/routes/+layout.svelte -->
< script lang = "ts" >
import { setBus } from '$bus';
setBus(App.bus);
import { setBus } from '$bus';
setBus(App.bus);
< / script >
```
@ -443,14 +459,14 @@ function directly) inside `$effect`:
```svelte
< script lang = "ts" >
import { SESSION_EVENT_IDENTITY_CHANGED } from '$session';
const Bus = getBus();
$effect(() =>
Bus.subscribe(SESSION_EVENT_IDENTITY_CHANGED, (event) => {
// react
})
);
import { SESSION_EVENT_IDENTITY_CHANGED } from '$session';
const Bus = getBus();
$effect(() =>
Bus.subscribe(SESSION_EVENT_IDENTITY_CHANGED, (event) => {
// react
})
);
< / script >
```
@ -476,10 +492,10 @@ import { untrack } from 'svelte';
import { createEngineBus, type EngineBusOptions } from '$bus';
export function createSvelteEngineBus(options: EngineBusOptions = {}) {
return createEngineBus({
...options,
invokeListener: (fn) => untrack(fn)
});
return createEngineBus({
...options,
invokeListener: (fn) => untrack(fn)
});
}
```
@ -523,18 +539,16 @@ Each `APP_EVENT_*` declares where it is allowed to fire:
```ts
// arts/active-app/events.ts
export const APP_EVENT_RUNTIMES: Readonly< Record < string , AppEventRuntime > > = {
[APP_EVENT_DISPOSE_STARTING]: APP_EVENT_RUNTIME_BOTH
[APP_EVENT_DISPOSE_STARTING]: APP_EVENT_RUNTIME_BOTH
};
export type AppEventRuntime =
| typeof APP_EVENT_RUNTIME_BOTH
| typeof APP_EVENT_RUNTIME_CLIENT;
export type AppEventRuntime = typeof APP_EVENT_RUNTIME_BOTH | typeof APP_EVENT_RUNTIME_CLIENT;
export function assertEventCanFire(type: string, where: 'server' | 'client'): void {
const runtime = APP_EVENT_RUNTIMES[type];
if (runtime & & runtime !== 'both' & & runtime !== where) {
throw new AappInvalidEventRuntimeError(type, runtime, where);
}
const runtime = APP_EVENT_RUNTIMES[type];
if (runtime & & runtime !== 'both' & & runtime !== where) {
throw new AappInvalidEventRuntimeError(type, runtime, where);
}
}
```
@ -549,16 +563,16 @@ strings; metadata lives in a parallel table.
Eight tests are required for the v0.1 gate:
| Test | What it verifies |
| -------------------------- | -- ------------------------------------------------------------------------------- |
| **SSR isolation** | A listener registered against request A's bus never sees request B's events. |
| **No singleton** | Importing the `bus` module twice yields no shared mutable state. |
| **Context isolation** | Two rendered app roots have different bus instances via `getBus()` . |
| **Svelte cleanup** | `$effect` -registered listener is unsubscribed when the component unmounts. |
| **Untrack** | Publishing during a `$effect` run does not create accidental reactive deps. |
| **Payload safety** | DEV-mode `structuredClone` check rejects non-cloneable payloads. |
| **Runtime guard** | Publishing a `client` -only event on the server throws in DEV. |
| **Flush test** | `flushSync(() => Bus.publish(...))` makes DOM updates observable synchronously. |
| Test | What it verifies |
| --------------------- | ------------------------------------------------------------------------------- |
| **SSR isolation** | A listener registered against request A's bus never sees request B's events. |
| **No singleton** | Importing the `bus` module twice yields no shared mutable state. |
| **Context isolation** | Two rendered app roots have different bus instances via `getBus()` . |
| **Svelte cleanup** | `$effect` -registered listener is unsubscribed when the component unmounts. |
| **Untrack** | Publishing during a `$effect` run does not create accidental reactive deps. |
| **Payload safety** | DEV-mode `structuredClone` check rejects non-cloneable payloads. |
| **Runtime guard** | Publishing a `client` -only event on the server throws in DEV. |
| **Flush test** | `flushSync(() => Bus.publish(...))` makes DOM updates observable synchronously. |
## Cross-module reactions live in orca, not on the bus
@ -601,10 +615,7 @@ applyStandardOrca(App);
Cherry-pick when the standard set is too aggressive:
```ts
import {
applyCacheClearOnIdentityChange,
applyPermInvalidateOnIdentityChange
} from '$active-app';
import { applyCacheClearOnIdentityChange, applyPermInvalidateOnIdentityChange } from '$active-app';
applyCacheClearOnIdentityChange(App);
applyPermInvalidateOnIdentityChange(App);
@ -627,20 +638,27 @@ type. v0.1 ships a minimal active wrapper:
```ts
// arts/bus/active-bus.svelte.ts
export function createBusRecent< TPayload > (bus: EngineBus, type: string) {
let lastEvent = $state< BusEnvelope < string , TPayload > | undefined>();
let count = $state(0);
const unsubscribe = bus.subscribe(type, (event) => {
lastEvent = event as BusEnvelope< string , TPayload > ;
count += 1;
});
return {
get lastEvent() { return lastEvent; },
get count() { return count; },
clear() { lastEvent = undefined; count = 0; },
dispose: unsubscribe
};
let lastEvent = $state< BusEnvelope < string , TPayload > | undefined>();
let count = $state(0);
const unsubscribe = bus.subscribe(type, (event) => {
lastEvent = event as BusEnvelope< string , TPayload > ;
count += 1;
});
return {
get lastEvent() {
return lastEvent;
},
get count() {
return count;
},
clear() {
lastEvent = undefined;
count = 0;
},
dispose: unsubscribe
};
}
```
@ -651,8 +669,8 @@ Surface stays minimal in v0.1: `lastEvent`, `count`, `clear()`,
Only one event survived the big-bang refactor:
| Constant | Meaning |
| ---------------------------- | ------------------------------------------------------ |
| Constant | Meaning |
| ---------------------------- | ------------------------------------------------------- |
| `APP_EVENT_DISPOSE_STARTING` | `App.dispose()` entered, last chance to flush listeners |
Everything else (identity changes, tenant switches, connectivity, cache
@ -706,7 +724,7 @@ import { flushSync } from 'svelte';
import { publishAppDisposeStarting } from '$active-app';
flushSync(() => {
publishAppDisposeStarting(Bus, { cause: 'test' });
publishAppDisposeStarting(Bus, { cause: 'test' });
});
expect(screen.getByText('disposing')).toBeVisible();
```
@ -718,8 +736,8 @@ For strict tests, throw on listener errors:
```ts
const Bus = createEngineBus({
listenerErrorMode: 'throw',
logger: testLogger
listenerErrorMode: 'throw',
logger: testLogger
});
```