import { CliArgument, CliCommand, CliOption, type ParentCliCommandDefinition, } from '@causa/cli'; import { WorkspaceFunction } from '@causa/workspace'; import { AllowMissing } from '@causa/workspace/validation'; import { Transform } from 'class-transformer'; import { IsBoolean, IsDate, IsInstance, IsInt, IsPositive, IsString, Validate, type ValidationArguments, ValidatorConstraint, type ValidatorConstraintInterface, } from 'class-validator'; /** * The definition for an event topic. * Definitions are found by looking for files matching the configured globs and regular expression. * A definition file contains the schema for a topic. * Its ID is constructed from parts of the file path to the definition. */ export type EventTopicDefinition = { /** * The ID of the event topic. * This is formatted using the {@link EventDefinition.formatParts} extracted from the * {@link EventDefinition.schemaFilePath} using the regular expression in the configuration. */ readonly id: string; /** * The file path to the schema definition for the topic. */ readonly schemaFilePath: string; /** * The parts extracted from the {@link EventDefinition.schemaFilePath} using the configured regular expression. */ readonly formatParts: Record; }; /** * Describes the event topics consumed and produced by a project. */ export type ReferencedEventTopics = { /** * The IDs of the topics that are consumed by the project. */ readonly consumed: EventTopicDefinition[]; /** * The IDs of the topics to which the project publishes events. */ readonly produced: EventTopicDefinition[]; }; /** * The `events` parent command, grouping all commands related to managing events, their topics, and the corresponding * schemas. */ export const eventsCommandDefinition: ParentCliCommandDefinition = { name: 'events', description: 'Manages events and topics.', }; /** * The base error for event topic-related errors. */ export class EventTopicError extends Error {} /** * An error thrown when two topic definition files lead to the same topic ID being rendered using the format string. */ export class DuplicateEventTopicError extends EventTopicError { constructor(readonly topicId: string) { super(`Found duplicate topic '${topicId}'.`); } } /** * An error thrown when the schema definition for the given event topics cannot be found. */ export class MissingEventTopicDefinitionsError extends Error { constructor(readonly topicIds: string[]) { super( `Missing definitions for topics ${topicIds .map((id) => `'${id}'`) .join(', ')}.`, ); } } /** * Lists all the event topics in the workspace. */ export abstract class EventTopicList extends WorkspaceFunction< Promise > {} /** * Returns all the event topics that are either consumed or produced by the current project. */ export abstract class EventTopicListReferencedInProject extends WorkspaceFunction< Promise > { /** * Maps the ID of the consumed and produced topics to their definitions, and returns the corresponding * {@link ReferencedEventTopics}. * This is a utility method that can be used by implementations that only retrieve the IDs of the topics from the * project's configuration. * * @param consumedIds The IDs of the topics that are consumed by the project. * @param producedIds The IDs of the topics that are produced by the project. * @returns The {@link ReferencedEventTopics} for the project. */ protected async mapToDefinitions( consumedIds: string[], producedIds: string[], ): Promise { const allDefinitions = await this._context.call(EventTopicList, {}); const missingDefinitions: string[] = []; const mapTopics = (topics: string[]) => topics.flatMap((topic) => { const definition = allDefinitions.find((d) => d.id === topic); if (!definition) { missingDefinitions.push(topic); return []; } return definition; }); const consumed = mapTopics(consumedIds); const produced = mapTopics(producedIds); if (missingDefinitions.length > 0) { throw new MissingEventTopicDefinitionsError(missingDefinitions); } return { consumed, produced }; } } /** * Temporary data stored in the backfill file, listing the resources that should be cleaned up when the backfill is * complete. */ export type BackfillTemporaryData = { /** * The ID of the temporary topic. * This is `null` if the main topic was used for publishing events. */ readonly temporaryTopicId: string | null; /** * A list of resource IDs that were created as part of the temporary triggers creation. * Those should be cleaned up when the backfill is complete. */ readonly temporaryTriggerResourceIds: string[]; }; /** * Backfills events for an event topic. * This supports several use cases, where existing event triggers for the topic may or may not process the backfilled * events. Temporary triggers can be created for the backfill, and deleted afterwards. The source of events to backfill * can be the default storage for the broker, or a custom source. * Returns the path to a file that contains the resources that should be deleted after the backfill has completed. */ @CliCommand({ parent: eventsCommandDefinition, name: 'backfill', description: `Backfills events for an event topic. Events can be backfilled using the existing topic (in which case all existing triggers receive the events), or through a temporary topic (in which case only the triggers passed to the command receive the events). Temporary triggers can be created for the backfill, and deleted afterwards using the 'cleanBackfill' command. If no source is specified, the events are retrieved from the default storage for the broker. Optionally, a custom source might be used. Finally, some event sources might support filtering of the events to backfill.`, summary: 'Backfills events for an event topic.', outputFn: (result) => { if (result) { console.log(result); } }, }) export abstract class EventTopicBackfill extends WorkspaceFunction< Promise > { /** * The full event topic name (e.g. `my-domain.my-event.v1`) for which events should be backfilled. */ @CliArgument({ name: 'eventTopic', position: 0, description: 'The full event topic name (e.g. `my-domain.my-event.v1`).', }) @IsString() readonly eventTopic!: string; /** * Whether a temporary topic should be created for the backfill. * When using a temporary topic, existing triggers do not receive the backfilled events, only the `triggers` passed to * the function. * If this is `false`, backfilled events are published to the existing topic, and all existing triggers will receive * the events. Temporary triggers can still be created by passing them in the `triggers`. */ @CliOption({ flags: '-c, --createTemporaryTopic', description: 'Whether a temporary topic should be created for the backfill instead of using the existing one.', }) @IsBoolean() @AllowMissing() readonly createTemporaryTopic?: boolean; /** * The source for events to publish. * By default, events will be fetched from the configured data storage using the `eventTopic` name. */ @CliOption({ flags: '-s, --source ', description: 'The source for events to publish. If not set, events are fetched from the default storage for the broker.', }) @IsString() @AllowMissing() readonly source?: string; /** * A filter for source events. * The format depends on the source type. */ @CliOption({ flags: '-f, --filter ', description: 'A filter for source events. The format and whether this is supported depends on the source type.', }) @IsString() @AllowMissing() readonly filter?: string; /** * A list of triggers to create for the backfill. * If a temporary topic is created, those will be the only triggers on the topic. */ @CliOption({ flags: '-t, --trigger ', description: 'A list of temporary triggers to create for the backfill. The format depends on the type of trigger.', }) @IsString({ each: true }) @AllowMissing() readonly triggers?: string[]; /** * The path to the file that will be written with the temporary resources to delete. */ @CliOption({ flags: '-o, --output ', description: `The path to the file that will be written with the temporary resources to delete. This is used by the 'cleanBackfill' command.`, }) @IsString() @AllowMissing() readonly output?: string; /** * Whether to wait for events to be processed after publishing, then clean up temporary resources inline. * When set, after a successful publish the implementation calls {@link EventTopicBrokerWaitForProcessing}, then * {@link EventTopicCleanBackfill}. On success no backfill file is written and the function returns an empty string. * On any failure during wait or clean, the backfill file is written so that `cleanBackfill` can be run manually. */ @CliOption({ flags: '--autoClean', description: 'Wait for events to be processed after publishing and clean up temporary resources in the same command.', }) @IsBoolean() @AllowMissing() readonly autoClean?: boolean; } /** * Cleans up resources created for a backfill based on the file outputted by the {@link EventTopicBackfill} function. */ @CliCommand({ parent: eventsCommandDefinition, name: 'cleanBackfill', description: `Cleans up temporary resources created for a backfill. This includes temporary triggers and topic.`, summary: 'Cleans up temporary resources created for a backfill.', }) export abstract class EventTopicCleanBackfill extends WorkspaceFunction< Promise > { /** * The path to the file that was written by the {@link EventTopicBackfill} function. */ @CliArgument({ name: 'file', position: 0, description: 'The path returned by the `backfill` command.', }) @IsString() readonly file!: string; } /** * Creates a topic for the configured broker. * Returns the (broker-specific) topic ID. */ export abstract class EventTopicBrokerCreateTopic extends WorkspaceFunction< Promise > { /** * A name for the topic. */ @IsString() readonly name!: string; } /** * Returns the broker-specific topic ID for an event topic. */ export abstract class EventTopicBrokerGetTopicId extends WorkspaceFunction< Promise > { /** * The full event topic name (e.g. `my-domain.my-event.v1`). */ @IsString() readonly eventTopic!: string; } /** * The structured form of a trigger passed to {@link EventTopicBrokerCreateTrigger}. * When the function is called with a value of this type, the context is guaranteed to be scoped to a project. */ export type EventTopicBrokerTrigger = { /** * The name of the trigger within the project. */ readonly name: string; /** * A free-form bag of string options for the trigger to create. */ readonly options: Record; }; /** * Validates that a value is either a string or a valid {@link EventTopicBrokerTrigger} object. */ @ValidatorConstraint({ name: 'isEventTopicBrokerTrigger' }) class IsEventTopicBrokerTriggerConstraint implements ValidatorConstraintInterface { validate(value: unknown): boolean { if (typeof value === 'string') { return true; } if (value === null || typeof value !== 'object' || Array.isArray(value)) { return false; } const { name, options } = value as Record; if (typeof name !== 'string') { return false; } if ( options === null || typeof options !== 'object' || Array.isArray(options) ) { return false; } return Object.values(options).every((v) => typeof v === 'string'); } defaultMessage(args: ValidationArguments): string { return `${args.property} must be a string or a valid EventTopicBrokerTrigger.`; } } /** * Creates a trigger on the given topic for the specified service. * Returns IDs of resources that should be deleted after the backfill has completed. */ export abstract class EventTopicBrokerCreateTrigger extends WorkspaceFunction< Promise > { /** * The unique ID for the backfilling operation. */ @IsString() readonly backfillId!: string; /** * The broker-specific topic ID used as trigger. */ @IsString() readonly topicId!: string; /** * Describes the trigger to create. * * A raw `string` is a free-form URI interpreted by the broker implementation. * * An {@link EventTopicBrokerTrigger} object is used when the function is called on a project-scoped context. */ @Validate(IsEventTopicBrokerTriggerConstraint) readonly trigger!: string | EventTopicBrokerTrigger; } /** * An error thrown when a trigger cannot be created, but some resources for the trigger might have already been created. * This should be thrown by implementations of the {@link EventTopicBrokerCreateTrigger} function. */ export class EventTopicTriggerCreationError extends Error { /** * Creates a new {@link EventTopicTriggerCreationError}. * * @param parent The parent error. * @param resourceIds The IDs of the resources that should be deleted. */ constructor( readonly parent: any, readonly resourceIds: string[], ) { super(parent.message); } } /** * A single event to publish as part of a backfill. */ export type BackfillEvent = { /** * The data to publish. */ readonly data: Buffer; /** * Optional attributes for the message. */ readonly attributes?: Record; }; /** * Creates the async iterable of {@link BackfillEvent}s to publish during a backfill. * Implementations are selected based on the {@link EventTopicCreateBackfillSource.source} string (or its absence, for * the broker's default storage). Filtering, when supported, is also applied by the implementation: the returned * iterable yields only the events that should actually be published. */ export abstract class EventTopicCreateBackfillSource extends WorkspaceFunction< Promise> > { /** * The full event topic name (e.g. `my-domain.my-event.v1`) for which events should be backfilled. * Implementations may use this to resolve schemas or locate default storage. */ @IsString() readonly eventTopic!: string; /** * An optional source descriptor (e.g. `json://path/*.jsonl`). * When omitted, implementations should fall back to the broker's default storage for the topic. */ @IsString() @AllowMissing() readonly source?: string; /** * An optional filter applied to source events. * The format and support depend on the selected implementation. */ @IsString() @AllowMissing() readonly filter?: string; } /** * Publishes events from an async iterable to the given topic. */ export abstract class EventTopicBrokerPublishEvents extends WorkspaceFunction< Promise > { /** * The broker-specific ID for the topic to which events should be published. */ @IsString() readonly topicId!: string; /** * The full event topic name. * Might be used as a hint for the source, to get the schema for publishing, etc. */ @IsString() readonly eventTopic!: string; /** * A factory returning the async iterable of events to publish. Built by {@link EventTopicCreateBackfillSource}. */ // Exposed as a thunk rather than the iterable directly: the function registry deep-clones args via // `class-transformer`'s `plainToInstance`, which tries to `new`-up the value's constructor for non-`Object` objects // and throws on async-generator instances. A function-typed value is passed through untouched. @IsInstance(Function) readonly source!: () => AsyncIterable; } /** * Waits for events published during a backfill to be processed by the triggers that received them. * Called by {@link EventTopicBackfill} when `autoClean` is set, before invoking {@link EventTopicCleanBackfill}. * The exact "processed" semantics, polling cadence, and timeout are decided by the broker-specific implementation. * Should throw if processing cannot be confirmed (e.g. on timeout), so the backfill file is preserved for manual * cleanup. */ export abstract class EventTopicBrokerWaitForProcessing extends WorkspaceFunction< Promise > { /** * The full event topic name (e.g. `my-domain.my-event.v1`) the events were published for. */ @IsString() readonly eventTopic!: string; /** * The temporary resources created during the backfill, as written to the backfill file. * Implementations typically use {@link BackfillTemporaryData.temporaryTriggerResourceIds} to know which subscriptions * to monitor, and {@link BackfillTemporaryData.temporaryTopicId} when set. */ readonly temporaryData!: BackfillTemporaryData; } /** * Deletes a resource that was created for a temporary trigger. */ export abstract class EventTopicBrokerDeleteTriggerResource extends WorkspaceFunction< Promise > { /** * An ID that describes the resource. */ @IsString() readonly id!: string; } /** * Deletes a topic using its broker-specific ID. */ export abstract class EventTopicBrokerDeleteTopic extends WorkspaceFunction< Promise > { /** * The broker-specific ID. */ @IsString() readonly id!: string; } /** * A single event returned by {@link EventTopicQueryEvents}. */ export type QueriedEvent = { /** * A unique ID for the event, stable across queries. * This may not be provided by all implementations. */ readonly id?: string; /** * The time at which the event was published to the topic. */ readonly timestamp: Date; /** * The event attributes. */ readonly attributes: Record; /** * The event payload. */ readonly data: any; }; /** * Queries an event topic for events that were published within a time range, optionally restricted by an * implementation-defined `filter`. */ export abstract class EventTopicQueryEvents extends WorkspaceFunction< Promise > { /** * The full event topic name to query. */ @IsString() readonly topic!: string; /** * The inclusive lower bound of the time range over which to look for events. */ @AllowMissing() @Transform(({ value }) => typeof value === 'string' ? new Date(value) : value, ) @IsDate() readonly from?: Date; /** * The exclusive upper bound of the time range over which to look for events. */ @AllowMissing() @Transform(({ value }) => typeof value === 'string' ? new Date(value) : value, ) @IsDate() readonly to?: Date; /** * An implementation-specific filter expression. */ @AllowMissing() @IsString() readonly filter?: string; /** * The maximum number of events to return. */ @AllowMissing() @IsInt() @IsPositive() readonly limit?: number; }