From 3d72133c47ce5b7c2da67a59ab155a6562dae421 Mon Sep 17 00:00:00 2001 From: Einar Date: Sat, 3 Oct 2026 09:14:44 +0200 Subject: [PATCH 1/3] Route commands and aggregates through Chronicle event source definitions Add the eventSourceDefinition decorator, per-event source and stream replacement as a unit, startup validation, and aggregate source/stream guarding. Pin @cratis/chronicle 6.49.0 for development and adjust the unsupported-projection specs for the in-memory arithmetic support added in that release. --- Samples/Library/package.json | 2 +- Source/Chronicle/AggregateRoot.ts | 10 ++- Source/Chronicle/ChronicleArtifacts.ts | 2 + Source/Chronicle/ChronicleCommand.ts | 13 +++- .../Chronicle/ChronicleCommandDefinition.ts | 7 ++ Source/Chronicle/ChronicleResponseHandler.ts | 28 +++++-- Source/Chronicle/EventSourceReference.ts | 11 +++ Source/Chronicle/EventSourceSelector.ts | 9 +++ Source/Chronicle/commandAggregate.ts | 31 +++++++- Source/Chronicle/commandEventRouting.ts | 47 ++++++++++++ Source/Chronicle/eventSourceDefinition.ts | 29 +++++++ Source/Chronicle/eventSourceRoute.ts | 76 +++++++++++++++++++ Source/Chronicle/eventSourceSdk.ts | 32 ++++++++ Source/Chronicle/index.ts | 3 + Source/Chronicle/package.json | 2 +- .../given/a_reduced_command.ts | 6 +- .../given/a_chronicle_query.ts | 4 +- Source/Chronicle/withChronicle.ts | 16 ++++ Source/Cratis/package.json | 2 +- yarn.lock | 24 +++--- 20 files changed, 320 insertions(+), 34 deletions(-) create mode 100644 Source/Chronicle/EventSourceReference.ts create mode 100644 Source/Chronicle/EventSourceSelector.ts create mode 100644 Source/Chronicle/commandEventRouting.ts create mode 100644 Source/Chronicle/eventSourceDefinition.ts create mode 100644 Source/Chronicle/eventSourceRoute.ts create mode 100644 Source/Chronicle/eventSourceSdk.ts diff --git a/Samples/Library/package.json b/Samples/Library/package.json index f33889bb..3a7a27ce 100644 --- a/Samples/Library/package.json +++ b/Samples/Library/package.json @@ -15,7 +15,7 @@ "@cratis/arc.core": "workspace:^", "@cratis/arc.express": "workspace:^", "@cratis/arc.testing": "workspace:^", - "@cratis/chronicle": "6.35.0", + "@cratis/chronicle": "6.49.0", "@cratis/fundamentals": "7.22.0", "express": "^5.1.0", "rxjs": "^7.8.2" diff --git a/Source/Chronicle/AggregateRoot.ts b/Source/Chronicle/AggregateRoot.ts index 3e417aee..dab3fd10 100644 --- a/Source/Chronicle/AggregateRoot.ts +++ b/Source/Chronicle/AggregateRoot.ts @@ -4,6 +4,7 @@ import { EventSequenceNumber } from '@cratis/chronicle/eventSequences'; import type { ConcurrencyScope, EventForEventSourceId } from '@cratis/chronicle/eventSequences'; import { AggregateRootCommitResult } from './AggregateRootCommitResult.js'; import type { EventContext } from '@cratis/chronicle/events'; +import type { EventSourceReference } from './EventSourceReference.js'; type EventClass = new (...args: never[]) => T; export const rehydrateAggregate = Symbol('rehydrate aggregate'); @@ -14,6 +15,7 @@ export class AggregateRoot { #pending: EventForEventSourceId[] = []; #staged = 0; #tail = EventSequenceNumber.beforeFirst.value; + #reference?: EventSourceReference; #route: Omit = {}; #handlers = new Map void>(); /** True when no recorded events were found for the selected event source. */ @@ -27,9 +29,10 @@ export class AggregateRoot { get eventTypes(): EventClass[] { return [...this.#handlers.keys()]; } /** @internal */ [rehydrateAggregate](sourceId: string, tail: bigint, route: Omit, - events: readonly { type: EventClass; content: object; context: EventContext }[]): void { + events: readonly { type: EventClass; content: object; context: EventContext }[], reference?: EventSourceReference): void { if (this.#sourceId) throw new Error('An aggregate can only be rehydrated once'); this.#sourceId = sourceId; + this.#reference = reference; this.#tail = tail; this.#route = route; this.isNew = tail === EventSequenceNumber.beforeFirst.value || tail === EventSequenceNumber.unset.value; @@ -39,7 +42,10 @@ export class AggregateRoot { apply(event: object): void { if (!this.#sourceId) throw new Error('The aggregate is not active'); this.#dispatch(event.constructor as EventClass, event); - this.#pending.push({ eventSourceId: this.#sourceId, event }); + // Appended events record the definition the aggregate was loaded through. + this.#pending.push({ eventSourceId: this.#sourceId, event, ...this.#reference ? { + eventSource: this.#reference.source as EventForEventSourceId['eventSource'], + ...this.#reference.stream === undefined ? {} : { eventStream: this.#reference.stream } } : {} }); } /** Return pending events; events not returned are also enrolled in the command unit of work. */ commit(): AggregateRootCommitResult { diff --git a/Source/Chronicle/ChronicleArtifacts.ts b/Source/Chronicle/ChronicleArtifacts.ts index fca622dc..fd411138 100644 --- a/Source/Chronicle/ChronicleArtifacts.ts +++ b/Source/Chronicle/ChronicleArtifacts.ts @@ -53,6 +53,8 @@ export class ChronicleArtifacts implements IClientArtifactsProvider { } return [...types]; } + /** Classes decorated with `@eventSource`; empty with an SDK that predates event source definitions. */ + get eventSources(): Constructor[] { return this.of(DecoratorType.EventSource); } get reactors(): Constructor[] { return this.of(DecoratorType.Reactor); } get reducers(): Constructor[] { return this.of(DecoratorType.Reducer); } /** Whether a model is populated by a registered declarative or model-bound projection. */ diff --git a/Source/Chronicle/ChronicleCommand.ts b/Source/Chronicle/ChronicleCommand.ts index 749fd802..0126091a 100644 --- a/Source/Chronicle/ChronicleCommand.ts +++ b/Source/Chronicle/ChronicleCommand.ts @@ -12,6 +12,8 @@ import type { ChronicleProduced } from './ChronicleProduced.js'; import { waitForProjectionCompletion } from './waitForProjectionCompletion.js'; import { isRoutedEvent } from './eventForEventSourceId.js'; import { AggregateRootCommitResult } from './AggregateRootCommitResult.js'; +import { routeEntry } from './commandEventRouting.js'; +import { assertEventSourcesSupported, validateEventSourceReference } from './eventSourceRoute.js'; import { EventsWithConcurrencyScopes } from './EventsWithConcurrencyScopes.js'; function memberName(propertyName: string): string { @@ -82,7 +84,8 @@ function containsAppendValue(value: unknown, store: IEventStore, seen = new Set< } export function defineChronicleCommand(definition: ChronicleCommandDefinition): CommandDefinition { - const { client, eventStore, namespaceForContext, produce, completionTimeoutMs, ...command } = definition; + const { client, eventStore, namespaceForContext, produce, completionTimeoutMs, eventSource, ...command } = definition; + if (eventSource) validateEventSourceReference(`Command ${definition.name}`, eventSource, {}); if (!eventStore) throw new Error('A Chronicle event store is required'); return defineCommand({ ...command, @@ -98,7 +101,8 @@ export function defineChronicleCommand(definition: Chron const namespace = namespaceForContext(context); if (!namespace) throw new Error('The Chronicle namespace resolver returned no namespace'); if (!Array.isArray(produced.events)) throw new Error('A Chronicle command must produce an event list'); - const events = produced.events.map(snapshotEvent); + const events = produced.events.map(snapshotEvent).map(entry => eventSource + ? { ...entry, ...routeEntry(entry, { legacy: {}, reference: eventSource }) } : entry); const commandResponse = produced.response; if (events.length && isOutcome(commandResponse)) throw new Error('A Chronicle command cannot persist events and return an Arc outcome'); if (!events.length) return commandResponse; @@ -108,6 +112,7 @@ export function defineChronicleCommand(definition: Chron throw new Error('The event type is not registered in the selected Chronicle event store'); } } + assertEventSourcesSupported(store, events); context.signal.throwIfAborted(); const results = events.length === 1 ? [await store.eventLog.append(events[0]!.eventSourceId, events[0]!.event, singleOptions(events[0]!, context))] @@ -139,6 +144,8 @@ function snapshotEvent(entry: EventForEventSourceId): EventForEventSourceId { ...(eventSourceType === undefined ? {} : { eventSourceType }), ...(eventStreamType === undefined ? {} : { eventStreamType }), ...(eventStreamId === undefined ? {} : { eventStreamId }), + ...(entry.eventSource === undefined ? {} : { eventSource: entry.eventSource }), + ...(entry.eventStream === undefined ? {} : { eventStream: entry.eventStream }), ...(subject === undefined ? {} : { subject }), ...(tags === undefined ? {} : { tags: [...tags] }), ...(occurred === undefined ? {} : { occurred: new Date(occurred.getTime()) }) }; @@ -146,5 +153,5 @@ function snapshotEvent(entry: EventForEventSourceId): EventForEventSourceId { function singleOptions(entry: EventForEventSourceId, context: ExecutionContext): AppendOptions { return { correlationId: context.correlationId, sourceType: entry.eventSourceType, streamType: entry.eventStreamType, - streamId: entry.eventStreamId, subject: entry.subject, occurred: entry.occurred, tags: entry.tags }; + streamId: entry.eventStreamId, eventSource: entry.eventSource, eventStream: entry.eventStream, subject: entry.subject, occurred: entry.occurred, tags: entry.tags }; } diff --git a/Source/Chronicle/ChronicleCommandDefinition.ts b/Source/Chronicle/ChronicleCommandDefinition.ts index bc6411cb..0a587169 100644 --- a/Source/Chronicle/ChronicleCommandDefinition.ts +++ b/Source/Chronicle/ChronicleCommandDefinition.ts @@ -3,6 +3,7 @@ import type { CommandDefinition, ExecutionContext } from '@cratis/arc.core'; import type { IChronicleClient } from '@cratis/chronicle'; import type { z } from 'zod'; +import type { EventSourceReference } from './EventSourceReference.js'; import type { ChronicleProduced } from './ChronicleProduced.js'; /** Application-controlled namespace selection, not inferred from untrusted request headers. */ @@ -11,6 +12,12 @@ export interface ChronicleCommandDefinition extends Omit readonly eventStore: string; /** Opt in to a bounded wait for kernel observer completion after a successful append (milliseconds). */ readonly completionTimeoutMs?: number; + /** + * Append through a Chronicle event source definition (and optionally one of its streams): each event records its event + * source, and Chronicle derives the concurrency scope from the definition. An event that names its own source or stream + * replaces this default. Requires `@cratis/chronicle` 6.49.0 or later. + */ + readonly eventSource?: EventSourceReference; readonly namespaceForContext: (context: ExecutionContext) => string; readonly produce: (input: z.output, context: ExecutionContext, provided: unknown) => ChronicleProduced | Promise>; } diff --git a/Source/Chronicle/ChronicleResponseHandler.ts b/Source/Chronicle/ChronicleResponseHandler.ts index d70d55a9..09d321e1 100644 --- a/Source/Chronicle/ChronicleResponseHandler.ts +++ b/Source/Chronicle/ChronicleResponseHandler.ts @@ -15,6 +15,11 @@ import { EventSourceIdResponse } from './eventSourceIdResponse.js'; import { eventForEventSourceId, isRoutedEvent } from './eventForEventSourceId.js'; import { AggregateRootCommitResult } from './AggregateRootCommitResult.js'; import { waitForProjectionCompletion } from './waitForProjectionCompletion.js'; +import { routeEntry } from './commandEventRouting.js'; +import type { CommandEventRouting } from './commandEventRouting.js'; +import { eventSourceReferenceFor } from './eventSourceDefinition.js'; +import { assertEventSourcesSupported, resolveEventSourceRoute, validateEventSourceReference } from './eventSourceRoute.js'; +import type { IEventStore } from '@cratis/chronicle'; import { chronicleIdentity } from './chronicleIdentity.js'; function eventLike(value: unknown): boolean { return typeof value === 'object' && value !== null && (isRoutedEvent(value) || hasEventType(value.constructor)); @@ -48,25 +53,27 @@ export class ChronicleResponseHandler implements CommandResponseValueHandler { const subjectField = getSubjectPropertyName((context.command as object).constructor); const subject = chronicleIdentity(command.getSubject?.(), 'subject') ?? (subjectField ? chronicleIdentity(Reflect.get(context.command as object, subjectField), 'subject') : undefined) ?? route.subject; + const reference = eventSourceReferenceFor((context.command as object).constructor); + if (reference) validateEventSourceReference(`Command ${(context.command as object).constructor.name}`, reference, route); + const routing: CommandEventRouting = { legacy: route, reference, streamId }; const entries: EventForEventSourceId[] = values.map(item => { const original: EventForEventSourceId = isRoutedEvent(item) ? item : { eventSourceId: selectedId, event: item as object }; if (typeof original.eventSourceId !== 'string' || !original.eventSourceId.trim()) throw new Error('Every appended event must have a nonempty event source id'); - return { eventSourceId: original.eventSourceId, event: original.event, - eventSourceType: original.eventSourceType ?? route.eventSourceType, - eventStreamType: original.eventStreamType ?? route.eventStreamType, - eventStreamId: original.eventStreamId ?? streamId, + return { eventSourceId: original.eventSourceId, event: original.event, ...routeEntry(original, routing), subject: original.subject ?? subject ?? original.eventSourceId, tags: original.tags, occurred: original.occurred }; }); + assertEventSourcesSupported(store, entries); const scopes: Record = { ...exact?.scopes, ...(value instanceof AggregateRootCommitResult ? value.scopes : {}) }; if (route.concurrentSource || route.concurrentStreamType || route.concurrentStreamId) { await Promise.all([...new Set(entries.map(entry => entry.eventSourceId))].map(async id => { if (scopes[id]) return; const sample = entries.find(entry => entry.eventSourceId === id)!; - const source = route.concurrentSource ? sample.eventSourceType : undefined; - const type = route.concurrentStreamType ? sample.eventStreamType : undefined; + const effective = effectiveRouting(store, sample, context); + const source = route.concurrentSource ? effective.eventSourceType : undefined; + const type = route.concurrentStreamType ? effective.eventStreamType : undefined; const stream = route.concurrentStreamId ? sample.eventStreamId : undefined; scopes[id] = { sequenceNumber: (await store.eventLog.getTailSequenceNumber(id, source, type, stream)).value, eventSourceId: true, eventSourceType: source, eventStreamType: type, eventStreamId: stream }; @@ -91,3 +98,12 @@ export class ChronicleResponseHandler implements CommandResponseValueHandler { return outcome; } } + +/** The source and stream type an entry is stored under, resolving a definition through the store. */ +function effectiveRouting(store: IEventStore, entry: EventForEventSourceId, context: CommandContext): + { eventSourceType?: string; eventStreamType?: string } { + if (entry.eventSource === undefined) return { eventSourceType: entry.eventSourceType, eventStreamType: entry.eventStreamType }; + const resolved = resolveEventSourceRoute(store, `Command ${(context.command as object).constructor.name}`, + { source: entry.eventSource, stream: entry.eventStream ?? entry.eventStreamType }); + return { eventSourceType: resolved.name, eventStreamType: resolved.stream }; +} diff --git a/Source/Chronicle/EventSourceReference.ts b/Source/Chronicle/EventSourceReference.ts new file mode 100644 index 00000000..88a76d00 --- /dev/null +++ b/Source/Chronicle/EventSourceReference.ts @@ -0,0 +1,11 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import type { EventSourceSelector } from './EventSourceSelector.js'; + +/** A command's or aggregate's reference to a Chronicle event source definition and, optionally, one of its streams. */ +export interface EventSourceReference { + /** The event source definition. */ + readonly source: EventSourceSelector; + /** The name of a stream declared by the definition; it becomes the event stream type of appended events. */ + readonly stream?: string; +} diff --git a/Source/Chronicle/EventSourceSelector.ts b/Source/Chronicle/EventSourceSelector.ts new file mode 100644 index 00000000..6dc7bc5c --- /dev/null +++ b/Source/Chronicle/EventSourceSelector.ts @@ -0,0 +1,9 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import type { Constructor } from '@cratis/fundamentals'; + +/** + * Names a Chronicle event source definition: the class decorated with `@eventSource`, its registered name, or a thunk + * returning the class. Use the thunk when the definition's module imports the declaring class (a circular import). + */ +export type EventSourceSelector = Constructor | string | (() => Constructor); diff --git a/Source/Chronicle/commandAggregate.ts b/Source/Chronicle/commandAggregate.ts index c6cccf95..53aad1ce 100644 --- a/Source/Chronicle/commandAggregate.ts +++ b/Source/Chronicle/commandAggregate.ts @@ -8,6 +8,23 @@ import { AggregateRoot, rehydrateAggregate } from './AggregateRoot.js'; import { ChronicleScopedStore } from './ChronicleStores.js'; import { ChronicleUnitOfWork } from './ChronicleUnitOfWork.js'; import { eventRoutingFor } from './eventRouting.js'; +import { eventSourceReferenceFor } from './eventSourceDefinition.js'; +import type { IEventStore } from '@cratis/chronicle'; +import type { EventSourceReference } from './EventSourceReference.js'; +import { resolveEventSourceRoute } from './eventSourceRoute.js'; + +/** The aggregate's own definition, or the command's; both must agree when both are declared. */ +function aggregateReference(store: IEventStore, aggregate: Function, command: Function): EventSourceReference | undefined { + const own = eventSourceReferenceFor(aggregate); + const inherited = eventSourceReferenceFor(command); + if (!own || !inherited) return own ?? inherited; + const label = `Aggregate ${aggregate.name} and command ${command.name}`; + const left = resolveEventSourceRoute(store, label, { source: own.source }); + const right = resolveEventSourceRoute(store, label, { source: inherited.source }); + if (left.name !== right.name || (own.stream !== undefined && inherited.stream !== undefined && own.stream !== inherited.stream)) + throw new Error(`${label} select different event source definitions or streams`); + return { source: own.source, stream: own.stream ?? inherited.stream }; +} const loaded = new WeakMap>>(); /** Inject a rehydrated aggregate for the command key into handle() or provide(). */ @@ -28,8 +45,16 @@ export function commandAggregate(type: new () => T): Se const route = eventRoutingFor((context.command as object).constructor); const command = context.command as { getEventStreamId?: () => string }; const streamId = command.getEventStreamId?.() ?? route.eventStreamId; - const source = route.eventSourceType; - const streamType = route.eventStreamType; + const commandType = (context.command as object).constructor; + const reference = aggregateReference(store, type, commandType); + let source = route.eventSourceType; + let streamType = route.eventStreamType; + if (reference) { + // The aggregate is guarded, and rehydrated, only from the declared source and stream. + const resolved = resolveEventSourceRoute(store, `Aggregate ${type.name}`, reference, route); + source = resolved.name; + streamType = resolved.stream ?? streamType; + } const tail = await store.eventLog.getTailSequenceNumber(context.key, source, streamType, streamId); const handlers = aggregate.eventTypes; const events = handlers.length ? await store.eventLog.getForEventSourceIdAndEventTypes( @@ -45,7 +70,7 @@ export function commandAggregate(type: new () => T): Se }); if (!eventType) throw new Error(`Unknown aggregate event type ${entry.eventType.toString()}`); return { type: eventType, content: entry.content, context: entry.context }; - })); + }), reference); ChronicleUnitOfWork.active()?.track(aggregate); return aggregate; } diff --git a/Source/Chronicle/commandEventRouting.ts b/Source/Chronicle/commandEventRouting.ts new file mode 100644 index 00000000..89ad5327 --- /dev/null +++ b/Source/Chronicle/commandEventRouting.ts @@ -0,0 +1,47 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import type { EventForEventSourceId } from '@cratis/chronicle/eventSequences'; +import type { EventRouting } from './eventRouting.js'; +import type { EventSourceReference } from './EventSourceReference.js'; +import { resolveEventSourceSelector } from './eventSourceDefinition.js'; + +/** What an event inherits from its command: legacy string routing, an optional definition, and the stream id. */ +export interface CommandEventRouting { + readonly legacy: EventRouting; + readonly reference?: EventSourceReference; + readonly streamId?: string; +} + +/** The routing fields of one appended event. */ +export type EntryRouting = Pick; + +function defined(fields: { [K in keyof T]: T[K] | undefined }): T { + return Object.fromEntries(Object.entries(fields).filter(([, value]) => value !== undefined)) as T; +} + +/** + * Route one event. The event's source and stream are one unit: an event that carries its own definition or its own raw + * source/stream type replaces everything the command would otherwise supply for them, so a command's definition is never + * silently combined with, or rewritten onto, an explicit per-event override. The stream id stays a separate default. + */ +export function routeEntry(original: EventForEventSourceId, command: CommandEventRouting): EntryRouting { + const eventStreamId = original.eventStreamId ?? command.streamId; + const ownDefinition = original.eventSource !== undefined || original.eventStream !== undefined; + const ownRaw = original.eventSourceType !== undefined || original.eventStreamType !== undefined; + if (ownDefinition) { + return defined({ eventSource: original.eventSource ?? (command.reference ? resolveEventSourceSelector(command.reference.source) : undefined), + eventStream: original.eventStream, eventSourceType: original.eventSourceType, + eventStreamType: original.eventStreamType, eventStreamId }); + } + if (ownRaw) { + const inherited = command.reference ? {} : command.legacy; + return defined({ eventSourceType: original.eventSourceType ?? inherited.eventSourceType, + eventStreamType: original.eventStreamType ?? inherited.eventStreamType, eventStreamId }); + } + if (command.reference) { + return defined({ eventSource: resolveEventSourceSelector(command.reference.source), eventStream: command.reference.stream, + // Legacy stream type without a declared stream names the stream; Chronicle validates it against the definition. + eventStreamType: command.reference.stream === undefined ? command.legacy.eventStreamType : undefined, eventStreamId }); + } + return defined({ eventSourceType: command.legacy.eventSourceType, eventStreamType: command.legacy.eventStreamType, eventStreamId }); +} diff --git a/Source/Chronicle/eventSourceDefinition.ts b/Source/Chronicle/eventSourceDefinition.ts new file mode 100644 index 00000000..b07c53cd --- /dev/null +++ b/Source/Chronicle/eventSourceDefinition.ts @@ -0,0 +1,29 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import type { Constructor } from '@cratis/fundamentals'; +import type { EventSourceReference } from './EventSourceReference.js'; +import type { EventSourceSelector } from './EventSourceSelector.js'; + +const references = new WeakMap(); + +/** + * Select the Chronicle event source definition (and optionally one of its streams) that a command or an aggregate + * root appends through. The reference is validated at startup and resolved by Chronicle when appending, so each + * event records its event source. Source and stream belong to routing; they are never declared on event types. + * Supports legacy and standard class decorators. + * @param source - The `@eventSource` class, its name, or a thunk returning the class for circular imports. + * @param stream - A stream declared by the definition. + */ +export function eventSourceDefinition(source: EventSourceSelector, stream?: string): + (target: object, context?: ClassDecoratorContext) => void { + return target => { references.set(target, { source, ...(stream === undefined ? {} : { stream }) }); }; +} + +/** The event source definition a command or aggregate class declared, if any. */ +export function eventSourceReferenceFor(type: object): EventSourceReference | undefined { return references.get(type); } + +/** Resolve a thunk to its class; a class or a name is returned unchanged. */ +export function resolveEventSourceSelector(source: EventSourceSelector): Constructor | string { + // Classes have a prototype; arrow-function thunks do not. + return typeof source === 'function' && source.prototype === undefined ? (source as () => Constructor)() : source as Constructor | string; +} diff --git a/Source/Chronicle/eventSourceRoute.ts b/Source/Chronicle/eventSourceRoute.ts new file mode 100644 index 00000000..1bd7c65c --- /dev/null +++ b/Source/Chronicle/eventSourceRoute.ts @@ -0,0 +1,76 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import type { Constructor } from '@cratis/fundamentals'; +import type { IEventStore } from '@cratis/chronicle'; +import type { EventRouting } from './eventRouting.js'; +import type { EventSourceReference } from './EventSourceReference.js'; +import { resolveEventSourceSelector } from './eventSourceDefinition.js'; +import { declaredEventSource, supportsEventSources, unsupported } from './eventSourceSdk.js'; + +/** A reference resolved against a registered definition. */ +export interface EventSourceRoute { + /** The registered name, which is the event source type of appended events. */ + readonly name: string; + /** The declared stream name, which is the event stream type of appended events. */ + readonly stream?: string; +} + +function describe(selector: Constructor | string): string { return typeof selector === 'string' ? selector : selector.name; } + +/** Reject the legacy string attributes that contradict, rather than repeat, the definition. */ +function checkAgainstLegacy(owner: string, name: string, streams: readonly string[], stream: string | undefined, legacy: EventRouting): void { + if (legacy.eventSourceType !== undefined && legacy.eventSourceType !== name) + throw new Error(`${owner} declares event source type '${legacy.eventSourceType}', which contradicts its event source definition '${name}'`); + if (stream !== undefined && legacy.eventStreamType !== undefined && legacy.eventStreamType !== stream) + throw new Error(`${owner} declares event stream type '${legacy.eventStreamType}', which contradicts its stream '${stream}' of event source '${name}'`); + if (stream === undefined && legacy.eventStreamType !== undefined && !streams.includes(legacy.eventStreamType)) + throw new Error(`${owner} declares event stream type '${legacy.eventStreamType}', which event source '${name}' does not declare (${streams.join(', ') || 'no streams'})`); +} + +function checkStream(owner: string, name: string, streams: readonly string[], stream: string | undefined): void { + if (stream !== undefined && !streams.includes(stream)) + throw new Error(`${owner} selects stream '${stream}', which event source '${name}' does not declare (${streams.join(', ') || 'no streams'})`); +} + +/** + * Validate a reference without a connection: the definition must be an `@eventSource` class (or, for a name, one of the + * known definitions), the stream must belong to it, and legacy string attributes must not contradict it. + * @param known - The definitions Arc registers with Chronicle; a name cannot be validated without them. + */ +export function validateEventSourceReference(owner: string, reference: EventSourceReference, legacy: EventRouting, + known?: readonly Constructor[]): void { + if (!supportsEventSources) throw new Error(unsupported(owner)); + const selector = resolveEventSourceSelector(reference.source); + let definition = typeof selector === 'string' ? undefined : declaredEventSource(selector); + if (typeof selector !== 'string' && !definition) + throw new Error(`${owner} refers to ${selector.name}, which is not an event source definition; decorate it with @eventSource()`); + if (typeof selector === 'string') { + if (!known) return; + const matches = known.map(declaredEventSource).filter(candidate => candidate?.name === selector); + if (!matches.length) + throw new Error(`${owner} refers to unknown event source definition '${selector}'; known definitions: ${ + known.map(type => declaredEventSource(type)?.name ?? type.name).join(', ') || 'none'}`); + definition = matches[0]; + } + checkStream(owner, definition!.name, definition!.streams, reference.stream); + checkAgainstLegacy(owner, definition!.name, definition!.streams, reference.stream, legacy); +} + +/** Resolve a reference through the store's registered definitions; fails clearly when the SDK predates them. */ +export function resolveEventSourceRoute(store: IEventStore, owner: string, reference: EventSourceReference, legacy: EventRouting = {}): EventSourceRoute { + const selector = resolveEventSourceSelector(reference.source); + if (!store.eventSources) throw new Error(unsupported(owner)); + let definition; + try { definition = store.eventSources.getFor(selector); } + catch (error) { throw new Error(`${owner} refers to event source definition '${describe(selector)}' that is not registered with Chronicle`, { cause: error }); } + const streams = definition.streams.map(stream => stream.name); + checkStream(owner, definition.name, streams, reference.stream); + checkAgainstLegacy(owner, definition.name, streams, reference.stream, legacy); + return { name: definition.name, ...(reference.stream === undefined ? {} : { stream: reference.stream }) }; +} + +/** Fail clearly rather than append without routing when entries select a definition the store cannot honor. */ +export function assertEventSourcesSupported(store: IEventStore, entries: readonly { readonly eventSource?: unknown; readonly eventStream?: string }[]): void { + if (!store.eventSources && entries.some(entry => entry.eventSource !== undefined || entry.eventStream !== undefined)) + throw new Error(unsupported('An appended event')); +} diff --git a/Source/Chronicle/eventSourceSdk.ts b/Source/Chronicle/eventSourceSdk.ts new file mode 100644 index 00000000..9c732952 --- /dev/null +++ b/Source/Chronicle/eventSourceSdk.ts @@ -0,0 +1,32 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import type { Constructor } from '@cratis/fundamentals'; + +interface EventSourceSdk { + getEventSourceMetadata?(target: Function): { readonly name: string } | undefined; + getEventStreamsFor?(target: Function): ReadonlyArray<{ readonly name: string }>; +} + +// A dynamic import cannot stop this package from linking against an older SDK; the helpers then report the missing support. +const sdk = await import('@cratis/chronicle') as unknown as EventSourceSdk; + +/** The name and declared streams of an `@eventSource` class. */ +export interface DeclaredEventSource { + readonly name: string; + readonly streams: readonly string[]; +} + +/** Whether the installed Chronicle SDK understands event source definitions. */ +export const supportsEventSources = typeof sdk.getEventSourceMetadata === 'function'; + +/** Read the definition a class declares, or undefined when it is not an `@eventSource` class. */ +export function declaredEventSource(type: Constructor): DeclaredEventSource | undefined { + if (!sdk.getEventSourceMetadata || !sdk.getEventStreamsFor) throw new Error(unsupported(type.name)); + const metadata = sdk.getEventSourceMetadata(type); + return metadata ? { name: metadata.name, streams: sdk.getEventStreamsFor(type).map(stream => stream.name) } : undefined; +} + +/** The message for a feature that needs a newer Chronicle SDK. */ +export function unsupported(subject: string): string { + return `${subject} uses a Chronicle event source definition, which requires @cratis/chronicle 6.49.0 or later`; +} diff --git a/Source/Chronicle/index.ts b/Source/Chronicle/index.ts index c61560e6..51ea1e06 100644 --- a/Source/Chronicle/index.ts +++ b/Source/Chronicle/index.ts @@ -15,6 +15,9 @@ export { ChronicleArtifacts } from './ChronicleArtifacts.js'; export { chronicleArtifactActivator } from './chronicleArtifactActivator.js'; export type { ChronicleArtifactActivator } from './chronicleArtifactActivator.js'; export { eventSourceType, eventStreamType, eventStreamId, eventSubject } from './eventRouting.js'; +export { eventSourceDefinition } from './eventSourceDefinition.js'; +export type { EventSourceReference } from './EventSourceReference.js'; +export type { EventSourceSelector } from './EventSourceSelector.js'; export { notAudited } from './notAudited.js'; export { EventsWithConcurrencyScopes, eventsWithConcurrencyScopes } from './EventsWithConcurrencyScopes.js'; export type { ChronicleRegistration } from './ChronicleOptions.js'; diff --git a/Source/Chronicle/package.json b/Source/Chronicle/package.json index 0090ecff..cf527896 100644 --- a/Source/Chronicle/package.json +++ b/Source/Chronicle/package.json @@ -50,7 +50,7 @@ "@cratis/arc.core": "workspace:^", "@cratis/arc.mongodb": "workspace:^", "@cratis/arc.testing": "workspace:^", - "@cratis/chronicle": "6.35.0", + "@cratis/chronicle": "6.49.0", "@cratis/fundamentals": "7.22.0", "@opentelemetry/api": "^1.9.0", "mongodb": "^6.21.0", diff --git a/Source/Chronicle/testing/for_ChronicleCommandScenario/given/a_reduced_command.ts b/Source/Chronicle/testing/for_ChronicleCommandScenario/given/a_reduced_command.ts index 723bf4dd..e2cb0f7e 100644 --- a/Source/Chronicle/testing/for_ChronicleCommandScenario/given/a_reduced_command.ts +++ b/Source/Chronicle/testing/for_ChronicleCommandScenario/given/a_reduced_command.ts @@ -3,7 +3,7 @@ import { field } from '@cratis/fundamentals'; import { eventType } from '@cratis/chronicle/events'; import { reducer } from '@cratis/chronicle/reducers'; -import { fromEvent, increment } from '@cratis/chronicle/projections'; +import { fromEvent } from '@cratis/chronicle/projections'; import { readModel } from '@cratis/chronicle/readModels'; import { command, commandReadModel, CommandValidator, inject, key, readModelForValidation, validator } from '@cratis/arc.core'; import { ChronicleCommandScenario } from '../../ChronicleCommandScenario.js'; @@ -17,7 +17,7 @@ class ItemReducer { itemAdded(event: ItemAdded, current?: ItemState): ItemState { return { count: (current?.count ?? 0) + event.amount }; } } @fromEvent(ItemAdded) class ProjectedState { @field(String) id = ''; @field(Number) amount = 0; } -@fromEvent(ItemAdded) class UnsupportedState { @field(String) id = ''; @field(Number) @increment(ItemAdded) count = 0; } +@fromEvent(ItemAdded) class UnsupportedState { @field(String) id = ''; @field(String) amount = ''; } @readModel() class UnreducedState { @field(Number) amount = 0; } @command() class CheckItem { @field(String) @key() id = ''; @@ -42,7 +42,7 @@ class ItemReducer { @command() class CheckUnsupportedItem { @field(String) @key() id = ''; @inject(commandReadModel(UnsupportedState)) - handle(state: UnsupportedState): number { return state.count; } + handle(state: UnsupportedState): string { return state.amount; } } @command() class CheckUnreducedItem { @field(String) @key() id = ''; diff --git a/Source/Chronicle/testing/for_ChronicleQueryScenario/given/a_chronicle_query.ts b/Source/Chronicle/testing/for_ChronicleQueryScenario/given/a_chronicle_query.ts index 3b11adea..a7c75d41 100644 --- a/Source/Chronicle/testing/for_ChronicleQueryScenario/given/a_chronicle_query.ts +++ b/Source/Chronicle/testing/for_ChronicleQueryScenario/given/a_chronicle_query.ts @@ -3,7 +3,7 @@ import { field } from '@cratis/fundamentals'; import { eventType } from '@cratis/chronicle/events'; import { reducer } from '@cratis/chronicle/reducers'; -import { fromEvent, increment, setFromContext } from '@cratis/chronicle/projections'; +import { fromEvent, setFromContext } from '@cratis/chronicle/projections'; import { readModel as chronicleReadModel } from '@cratis/chronicle/readModels'; import { argument, query, readModel, service } from '@cratis/arc.core'; import { ChronicleReadModels } from '../../../ChronicleReadModels.js'; @@ -22,7 +22,7 @@ export class BalanceReducer { } @fromEvent(BalanceChanged) export class UnsupportedBalance { @field(String) id = ''; - @field(Number) @increment(BalanceChanged) count = 0; + @field(String) amount = ''; } @reducer('query-projection-precedence', undefined, ProjectedBalance) export class ProjectionPrecedenceReducer { diff --git a/Source/Chronicle/withChronicle.ts b/Source/Chronicle/withChronicle.ts index ffcc58a3..8c2d1e4b 100644 --- a/Source/Chronicle/withChronicle.ts +++ b/Source/Chronicle/withChronicle.ts @@ -4,6 +4,9 @@ import { ArcApplicationBuilder, readModelCollectionNameResolver, serviceToken } import type { ArcServer, ReadModelInterceptor } from '@cratis/arc.core'; import { ArcApplicationBuilder as FetchArcApplicationBuilder } from '@cratis/arc.core/fetch'; import type { Constructor } from '@cratis/fundamentals'; +import { eventRoutingFor } from './eventRouting.js'; +import { eventSourceReferenceFor, resolveEventSourceSelector } from './eventSourceDefinition.js'; +import { validateEventSourceReference } from './eventSourceRoute.js'; import { ChronicleArtifacts } from './ChronicleArtifacts.js'; import { ChronicleReadModels } from './ChronicleReadModels.js'; import { ChronicleReadModelInterceptor } from './ChronicleReadModelInterceptor.js'; @@ -51,8 +54,21 @@ export function withChronicle(builder: ArcApplicationBuilder, options: Partial(); + const routedTypes = new Set(); + // Fail at startup, not on the first append, when a command's event source definition cannot work. Referencing a + // definition class makes it a registered artifact, so discovery needs no separate step and no import cycle. + builder.addBuiltObserver(() => { + for (const type of routedTypes) { + const reference = eventSourceReferenceFor(type)!; + const selector = resolveEventSourceSelector(reference.source); + if (!registration.client && typeof selector !== 'string') artifacts.register(selector); + validateEventSourceReference(`Command ${type.name}`, reference, eventRoutingFor(type), + registration.client ? undefined : artifacts.eventSources); + } + }); builder.addArtifactObserver(type => { const matched = artifacts.register(type as Constructor); + if (eventSourceReferenceFor(type)) routedTypes.add(type as Constructor); // Deferred scoped fallbacks: explicit, options.services and decorated lifetimes win at build time. for (const artifact of [...artifacts.reactors, ...artifacts.reducers]) builder.services.addScopedFallback(artifact); for (const model of artifacts.readModels) { diff --git a/Source/Cratis/package.json b/Source/Cratis/package.json index 3ef076aa..be9b162b 100644 --- a/Source/Cratis/package.json +++ b/Source/Cratis/package.json @@ -37,7 +37,7 @@ "zod": "^4.1.0" }, "devDependencies": { - "@cratis/chronicle": "6.35.0", + "@cratis/chronicle": "6.49.0", "@cratis/fundamentals": "7.22.0", "rxjs": "^7.8.2", "zod": "^4.1.0" diff --git a/yarn.lock b/yarn.lock index 9dbe5941..fff8f4b7 100644 --- a/yarn.lock +++ b/yarn.lock @@ -113,7 +113,7 @@ __metadata: "@cratis/arc.core": "workspace:^" "@cratis/arc.mongodb": "workspace:^" "@cratis/arc.testing": "workspace:^" - "@cratis/chronicle": "npm:6.35.0" + "@cratis/chronicle": "npm:6.49.0" "@cratis/fundamentals": "npm:7.22.0" "@opentelemetry/api": "npm:^1.9.0" mongodb: "npm:^6.21.0" @@ -354,7 +354,7 @@ __metadata: "@cratis/arc.express": "workspace:^" "@cratis/arc.react": "npm:22.45.0" "@cratis/arc.testing": "workspace:^" - "@cratis/chronicle": "npm:6.35.0" + "@cratis/chronicle": "npm:6.49.0" "@cratis/components": "npm:4.23.0" "@cratis/fundamentals": "npm:7.22.0" "@testing-library/dom": "npm:^10.4.2" @@ -392,23 +392,23 @@ __metadata: languageName: node linkType: hard -"@cratis/chronicle.contracts@npm:19.26.2": - version: 19.26.2 - resolution: "@cratis/chronicle.contracts@npm:19.26.2" +"@cratis/chronicle.contracts@npm:19.30.0": + version: 19.30.0 + resolution: "@cratis/chronicle.contracts@npm:19.30.0" dependencies: "@bufbuild/protobuf": "npm:^2.14.0" "@grpc/grpc-js": "npm:^1.14.4" nice-grpc: "npm:^2.1.17" - checksum: 10c0/0b06c998441ea1d13e9f4a8ffbca3ae7f1eddea74642a42b1feb6ffb8e8802047cddc9548a66cfb68287d902be4990d19c3941d0612fe13e8807fcc58da90c8e + checksum: 10c0/f5610eabdb40654f4cb45d9c9bd676c35e093b34d0221374db5a98bb5b8b86131278ae335082ef8bae9d5e10762b83aaf7f6f4e1bbe5f58dd42d1b3fdddb6ae7 languageName: node linkType: hard -"@cratis/chronicle@npm:6.35.0": - version: 6.35.0 - resolution: "@cratis/chronicle@npm:6.35.0" +"@cratis/chronicle@npm:6.49.0": + version: 6.49.0 + resolution: "@cratis/chronicle@npm:6.49.0" dependencies: "@bufbuild/protobuf": "npm:^2.16.0" - "@cratis/chronicle.contracts": "npm:19.26.2" + "@cratis/chronicle.contracts": "npm:19.30.0" "@grpc/grpc-js": "npm:^1.14.5" "@opentelemetry/api": "npm:^1.9.1" glob: "npm:^13.0.6" @@ -418,7 +418,7 @@ __metadata: undici: "npm:^8.11.2" peerDependencies: "@cratis/fundamentals": ^7.20.0 - checksum: 10c0/3a5d8179f020c64f62b3ae929d7d1781fc842220792b3430aad90e67df6f626bd8085f3af2b6d407846c75b7868145fa82293e5aed48ddfe0dce8448ff2e80bd + checksum: 10c0/928aad4ddeb0312c4121fa66fd739cbd8cb3e13baa027f0a2a7d23d6248ed64f405922cfc2137cae974965c6ab3f32f2c987bfc193e49d83c2e9bbea86711e5a languageName: node linkType: hard @@ -458,7 +458,7 @@ __metadata: "@cratis/arc.chronicle": "workspace:^" "@cratis/arc.core": "workspace:^" "@cratis/arc.testing": "workspace:^" - "@cratis/chronicle": "npm:6.35.0" + "@cratis/chronicle": "npm:6.49.0" "@cratis/fundamentals": "npm:7.22.0" rxjs: "npm:^7.8.2" zod: "npm:^4.1.0" From 3b8a4cd80cba53e5af0b31868a56b771ea426c7f Mon Sep 17 00:00:00 2001 From: Einar Date: Sat, 3 Oct 2026 09:24:48 +0200 Subject: [PATCH 2/3] Specify and document event source definition routing Add command, aggregate, nested transaction, mixed override, legacy and older-store specs, a live kernel check, and the Event source definitions page. --- CONTRIBUTING.md | 2 +- Documentation/chronicle/add-event-sourcing.md | 2 +- .../aggregates/injecting-into-commands.md | 2 +- .../chronicle/commands/event-metadata.md | 3 + .../commands/event-source-definitions.md | 89 +++++++++++++++ Documentation/chronicle/commands/toc.yml | 2 + Documentation/decorators.md | 1 + Documentation/reference/capabilities.md | 5 +- Documentation/reference/packages.md | 2 +- .../EventSourceRoutingArtifacts.ts | 58 ++++++++++ .../event-source-routing.live.test.mjs | 82 ++++++++++++++ Source/Chronicle/commandAggregate.ts | 7 +- Source/Chronicle/eventSourceSdk.ts | 4 +- ...command_that_selects_another_definition.ts | 17 +++ .../with_a_declared_definition.ts | 22 ++++ .../with_a_single_event.ts | 32 ++++++ .../with_an_unknown_stream.ts | 15 +++ .../with_a_declared_definition.ts | 19 ++++ .../with_a_definition_name.ts | 17 +++ .../with_a_lazy_definition.ts | 16 +++ .../with_a_legacy_declaration.ts | 19 ++++ .../with_a_per_event_definition.ts | 17 +++ .../with_a_raw_per_event_override.ts | 19 ++++ .../with_explicit_concurrency_flags.ts | 20 ++++ .../with_mixed_events.ts | 19 ++++ .../with_a_definition_routed_command.ts | 18 +++ .../with_a_legacy_command.ts | 17 +++ ...th_nested_commands_that_use_definitions.ts | 38 +++++++ .../with_a_class_that_is_not_a_definition.ts | 15 +++ ...ition_class_only_the_command_references.ts | 21 ++++ ...an_unknown_name_and_an_arc_owned_client.ts | 19 ++++ .../with_an_unknown_stream.ts | 15 +++ .../with_legacy_attributes_that_contradict.ts | 15 +++ ...y_attributes_that_repeat_the_definition.ts | 8 ++ .../Chronicle/given/event_source_routing.ts | 105 ++++++++++++++++++ Source/Chronicle/run-integration.sh | 4 +- 36 files changed, 754 insertions(+), 12 deletions(-) create mode 100644 Documentation/chronicle/commands/event-source-definitions.md create mode 100644 Source/Chronicle/Integration/EventSourceRoutingArtifacts.ts create mode 100644 Source/Chronicle/Integration/event-source-routing.live.test.mjs create mode 100644 Source/Chronicle/for_AggregateRoot/when_declaring_an_event_source_definition/with_a_command_that_selects_another_definition.ts create mode 100644 Source/Chronicle/for_AggregateRoot/when_declaring_an_event_source_definition/with_a_declared_definition.ts create mode 100644 Source/Chronicle/for_ChronicleCommand/when_routing_through_an_event_source_definition/with_a_single_event.ts create mode 100644 Source/Chronicle/for_ChronicleCommand/when_routing_through_an_event_source_definition/with_an_unknown_stream.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_routing_through_an_event_source_definition/with_a_declared_definition.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_routing_through_an_event_source_definition/with_a_definition_name.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_routing_through_an_event_source_definition/with_a_lazy_definition.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_routing_through_an_event_source_definition/with_a_legacy_declaration.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_routing_through_an_event_source_definition/with_a_per_event_definition.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_routing_through_an_event_source_definition/with_a_raw_per_event_override.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_routing_through_an_event_source_definition/with_explicit_concurrency_flags.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_routing_through_an_event_source_definition/with_mixed_events.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_the_store_predates_event_source_definitions/with_a_definition_routed_command.ts create mode 100644 Source/Chronicle/for_ChronicleResponseHandler/when_the_store_predates_event_source_definitions/with_a_legacy_command.ts create mode 100644 Source/Chronicle/for_ChronicleUnitOfWork/when_executing/with_nested_commands_that_use_definitions.ts create mode 100644 Source/Chronicle/for_withChronicle/when_validating_event_source_definitions/with_a_class_that_is_not_a_definition.ts create mode 100644 Source/Chronicle/for_withChronicle/when_validating_event_source_definitions/with_a_definition_class_only_the_command_references.ts create mode 100644 Source/Chronicle/for_withChronicle/when_validating_event_source_definitions/with_an_unknown_name_and_an_arc_owned_client.ts create mode 100644 Source/Chronicle/for_withChronicle/when_validating_event_source_definitions/with_an_unknown_stream.ts create mode 100644 Source/Chronicle/for_withChronicle/when_validating_event_source_definitions/with_legacy_attributes_that_contradict.ts create mode 100644 Source/Chronicle/for_withChronicle/when_validating_event_source_definitions/with_legacy_attributes_that_repeat_the_definition.ts create mode 100644 Source/Chronicle/given/event_source_routing.ts diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 62c9d465..b770095c 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -77,7 +77,7 @@ A hosted run does not replace local verification. The hosted CI workflow also su ### Upstream declaration errors -The core, host adapters, MongoDB, testing, proxy generator, and ESLint plugin consumer files use `skipLibCheck: false`. Chronicle and Drizzle consumer files use `skipLibCheck: true` **only** for third-party declaration errors; the script first runs their NodeNext compilation with `skipLibCheck: false`, prints the upstream diagnostics, and rejects errors in Arc declarations or consumer code. The pinned `@cratis/chronicle@6.35.0` and `@cratis/chronicle.contracts@19.26.2` declarations no longer require an exception; any Chronicle declaration error fails the guard. `drizzle-orm@0.45.3` has `gel-core/columns/date-duration.d.ts(1,35)` TS2307 (missing `gel`) and `pg-core/query-builders/query.d.ts(23,22)` TS2420 (`PgRelationalQuery` lacks `getSQL`), among other internal declaration errors. Fix these in their owning packages before removing the temporary integration exception. +The core, host adapters, MongoDB, testing, proxy generator, and ESLint plugin consumer files use `skipLibCheck: false`. Chronicle and Drizzle consumer files use `skipLibCheck: true` **only** for third-party declaration errors; the script first runs their NodeNext compilation with `skipLibCheck: false`, prints the upstream diagnostics, and rejects errors in Arc declarations or consumer code. The pinned `@cratis/chronicle@6.49.0` and `@cratis/chronicle.contracts@19.26.2` declarations no longer require an exception; any Chronicle declaration error fails the guard. `drizzle-orm@0.45.3` has `gel-core/columns/date-duration.d.ts(1,35)` TS2307 (missing `gel`) and `pg-core/query-builders/query.d.ts(23,22)` TS2420 (`PgRelationalQuery` lacks `getSQL`), among other internal declaration errors. Fix these in their owning packages before removing the temporary integration exception. ## Conventions diff --git a/Documentation/chronicle/add-event-sourcing.md b/Documentation/chronicle/add-event-sourcing.md index 52eabbec..057d8d78 100644 --- a/Documentation/chronicle/add-event-sourcing.md +++ b/Documentation/chronicle/add-event-sourcing.md @@ -57,7 +57,7 @@ cd ../my-arc-app Install it together with the Chronicle SDK and RxJS: ```bash -npm install ../arc-packages/arc.chronicle.tgz @cratis/chronicle@~6.35.0 rxjs@^7.8.2 +npm install ../arc-packages/arc.chronicle.tgz @cratis/chronicle@~6.49.0 rxjs@^7.8.2 ``` | Package | What it gives you | diff --git a/Documentation/chronicle/aggregates/injecting-into-commands.md b/Documentation/chronicle/aggregates/injecting-into-commands.md index c47e3b18..45e9fd9a 100644 --- a/Documentation/chronicle/aggregates/injecting-into-commands.md +++ b/Documentation/chronicle/aggregates/injecting-into-commands.md @@ -64,7 +64,7 @@ You may `return order.commit()` to make the commit visible. Do **not** also retu The aggregate belongs to the command's key: the `@key()` field, `getKey()`, or `getEventSourceId()`. A command without a key fails with the exception `A command key is required for Order` before `handle()` runs. -Loading uses the same route the command's returned events use: the current tenant's namespace, the command's `@eventSourceType`, `@eventStreamType`, and its stream ID from `getEventStreamId()` or `@eventStreamId`. Arc reads the events of the types the aggregate handles, in order, and replays them. See [Event metadata](../commands/event-metadata.md). +Loading uses the same route the command's returned events use: the current tenant's namespace, the command's `@eventSourceType`, `@eventStreamType`, and its stream ID from `getEventStreamId()` or `@eventStreamId`. Arc reads the events of the types the aggregate handles, in order, and replays them. See [Event metadata](../commands/event-metadata.md). An aggregate can instead declare an [event source definition](../commands/event-source-definitions.md), which guards and rehydrates only that source and stream. Arc loads each aggregate type once per command. A second parameter of the same type receives the same instance. For a second aggregate **type** on the same key, bind another `commandAggregate(Type)`. There is no way to load an aggregate for a different ID; the key decides. diff --git a/Documentation/chronicle/commands/event-metadata.md b/Documentation/chronicle/commands/event-metadata.md index f763156c..57e9b07c 100644 --- a/Documentation/chronicle/commands/event-metadata.md +++ b/Documentation/chronicle/commands/event-metadata.md @@ -20,6 +20,8 @@ An event records more than its payload. It also records which entity it belongs | Caused by | The signed-in principal, or Chronicle's system identity for an anonymous caller | None | | Causation | An `Arc.Command` entry with the command name and its values; see [Causation and auditing](causation.md) | None | +To select a registered Chronicle event source definition instead of free-form strings, see [Event source definitions](event-source-definitions.md). + Routing decorators come from `@cratis/arc.chronicle` and apply to every event the command returns. A value set on an `eventForEventSourceId` entry wins over the command's default for that entry only. ## Set command-wide defaults @@ -90,4 +92,5 @@ export class Onboarding { - [Returning events](index.md) - [Resolving the event source ID](../resolving-event-source-id.md) +- [Event source definitions](event-source-definitions.md) - [Concurrency](concurrency.md), where the same routing decorators opt into tail checks diff --git a/Documentation/chronicle/commands/event-source-definitions.md b/Documentation/chronicle/commands/event-source-definitions.md new file mode 100644 index 00000000..70a5f5b9 --- /dev/null +++ b/Documentation/chronicle/commands/event-source-definitions.md @@ -0,0 +1,89 @@ +--- +title: Event source definitions +description: Route a command's events, and an aggregate's, through a Chronicle event source definition and stream, and know what Arc validates at startup and what Chronicle enforces on append. +--- + +A string such as `@eventSourceType('Account')` names where an event goes, but nothing checks that the name exists or that the stream belongs to it. A Chronicle event source definition is that check: a registered class that declares a source and its streams. A command selects one with `@eventSourceDefinition`, and each event it returns records the definition it was appended through. + +Source and stream are routing. You declare them on the command or the aggregate, never on an event type. + +This feature needs `@cratis/chronicle` 6.49.0 or later. Older SDKs keep working for string routing, and a command that selects a definition fails with a message naming the required version. + +## Declare a definition and select it + +Declare the definition with the SDK's decorators, then select it on the command. + +```typescript +import { field } from '@cratis/fundamentals'; +import { ConcurrencyDimensions, eventSource, eventStream } from '@cratis/chronicle'; +import { command, key } from '@cratis/arc.core'; +import { eventSourceDefinition } from '@cratis/arc.chronicle'; + +@eventSource() +@eventStream('Transactions', { concurrency: ConcurrencyDimensions.eventStreamType | ConcurrencyDimensions.eventStreamId }) +export class Account {} + +@eventSourceDefinition(Account, 'Transactions') +@command() +export class Deposit { + @field(String) @key() id = ''; + handle(): FundsDeposited { return new FundsDeposited(); } +} +``` + +`FundsDeposited` is an `@eventType()` class. The event is appended with source `Account` and stream type `Transactions`. Chronicle records the definition on the event, and a reactor or projection reads it from `EventContext.eventSource`. + +The first argument is the class, its registered name, or a function returning the class. Use the function form, `@eventSourceDefinition(() => Account, 'Transactions')`, when the definition's module imports the command and the plain class would be read before it is defined. + +Referencing the class registers it with Chronicle, so you do not add it to discovery separately. A name can only be resolved against definitions that are registered, so a name Arc cannot find fails at startup. + +## What fails at startup + +Arc checks every command that selects a definition when the application is built, without a connection, and refuses to build when: + +- the class is not decorated with `@eventSource()`; +- the stream is not one the definition declares; +- a name matches no registered definition; or +- `@eventSourceType` or `@eventStreamType` on the same command contradicts the definition. Repeating the definition's own name is allowed. + +An aggregate, which Arc only meets when a command uses it, is checked the first time it loads. + +## Concurrency comes from the definition + +When a command selects a definition and sets no concurrency flags, Arc passes no scope and Chronicle derives one from the dimensions the definition or stream declares. Explicit flags win: `@eventStreamId('2026-05', { concurrency: true })` builds the same explicit scope as without a definition, using the definition's source name, and Chronicle then derives nothing. + +Chronicle's client refuses events for one event source ID that would need different automatic scopes within a single batch. Arc does not work around that; the command fails and appends nothing. Pass an explicit scope with `eventsWithConcurrencyScopes`, or return the events in separate commands. + +## Override one event + +An event entry that names its own `eventSource` and `eventStream`, or its own raw `eventSourceType` or `eventStreamType`, takes over source and stream for that event as one unit. The command's definition is neither merged into it nor used to rewrite it. + +```typescript +import { eventForEventSourceId } from '@cratis/arc.chronicle'; + +handle() { + return eventForEventSourceId({ eventSourceId: this.id, event: new FundsPosted(), + eventSource: Ledger, eventStream: 'Postings' }); +} +``` + +A raw `eventStreamType` on an entry wins over the command's definition the same way, and the entry is appended without a definition. The stream ID stays a separate default from `getEventStreamId()`. + +## Aggregates + +Declare the definition on the aggregate class to guard and rehydrate only that source and stream. + +```typescript +@eventSourceDefinition(Account, 'Transactions') +export class Wallet extends AggregateRoot { + constructor() { super(); this.on(FundsDeposited, () => {}); } +} +``` + +Loading reads the tail and the events for the declared source and stream only, the concurrency scope carries the same source and stream, and every event the aggregate applies records the definition. A command and its aggregate may each declare a definition when they agree; two that name different sources or streams fail instead of one silently winning. + +## Related + +- [Event metadata](event-metadata.md) +- [Concurrency](concurrency.md) +- [Defining an aggregate root](../aggregates/defining-an-aggregate-root.md) diff --git a/Documentation/chronicle/commands/toc.yml b/Documentation/chronicle/commands/toc.yml index 30c7e332..659b47d9 100644 --- a/Documentation/chronicle/commands/toc.yml +++ b/Documentation/chronicle/commands/toc.yml @@ -2,6 +2,8 @@ href: index.md - name: Event metadata href: event-metadata.md +- name: Event source definitions + href: event-source-definitions.md - name: Subject href: subject.md - name: Concurrency diff --git a/Documentation/decorators.md b/Documentation/decorators.md index 6c5099a4..05e27dca 100644 --- a/Documentation/decorators.md +++ b/Documentation/decorators.md @@ -87,6 +87,7 @@ From `@cratis/arc.chronicle` (experimental): | `@eventSourceType('Type', { concurrency? })` | command class | Default event source type for returned events | Command event metadata | | `@eventStreamType('Type', { concurrency? })` | command class | Default event stream type | Command event metadata | | `@eventStreamId('id', { concurrency? })` | command class | Default event stream ID | Command event metadata | +| `@eventSourceDefinition(Source, 'Stream'?)` | command or aggregate class | Route events through a Chronicle event source definition and stream | [Event source definitions](chronicle/commands/event-source-definitions.md) | | `@eventSubject('subject')` | command class | Default compliance subject | Command event metadata | | `@notAudited()` | command field | Keeps the value out of the causation chain | `[NotAudited]` | diff --git a/Documentation/reference/capabilities.md b/Documentation/reference/capabilities.md index e79033cd..92475983 100644 --- a/Documentation/reference/capabilities.md +++ b/Documentation/reference/capabilities.md @@ -123,10 +123,11 @@ Evidence paths are relative to the repository root. Spec folders follow `for_