openvibe-sdk/events
Generated at from
openvibe-contracts v0.114.0 and
openvibe-sdk v0.35.1.
From types/events.d.ts (server only). Declarations are shown verbatim.
type EventInput
What a producer passes: event_id, timestamp, source, version, payload and trace_id are filled in.
export type EventInput = Omit<EventEnvelope, 'event_id' | 'timestamp' | 'source' | 'version' | 'payload'> & Partial<Pick<EventEnvelope, 'event_id' | 'timestamp' | 'source' | 'version' | 'payload'>>;interface PublishResult
export interface PublishResult { event_id: string; seq: number; duplicate: boolean; }interface StoredEvent
export interface StoredEvent { seq: number; cursor: string; event: EventEnvelope; }interface Gap
export interface Gap { from_seq: number; to_seq: number; reason?: string; latest_seq?: number; }interface EventsPage
export interface EventsPage { events: StoredEvent[]; next_cursor: string; next_after_seq: number; latest_seq: number; gap?: Gap; }interface Subscription
export interface Subscription {
id: string;
consumer: string;
topic_pattern: string;
endpoint: string;
enabled: boolean;
retry_policy: { max_attempts?: number; backoff_ms?: number[] } | null;
/** Only in the create response. */
secret?: string;
[field: string]: unknown;
}interface CreateSubscription
export interface CreateSubscription { topicPattern: string; endpoint: string; secret?: string; retryPolicy?: { max_attempts?: number; backoff_ms?: number[] }; }interface EventsClient
export interface EventsClient {
prepare(envelope: EventInput, opts?: { traceId?: string; now?: number }): EventEnvelope;
publish(envelope: EventInput, opts?: { traceparent?: string }): Promise<PublishResult>;
publish(envelopes: EventInput[], opts?: { traceparent?: string }): Promise<{ results: PublishResult[] }>;
pull(opts?: { topic?: string | string[]; after?: string; afterSeq?: number; limit?: number }): Promise<EventsPage>;
/** onPage runs after every item of that page was yielded and handled: save page.next_cursor there. */
iterate(opts?: { topic?: string | string[]; after?: string; afterSeq?: number; limit?: number; maxPages?: number; onGap?: (gap: Gap) => void | Promise<void>; onPage?: (page: EventsPage) => void | Promise<void> }): AsyncGenerator<StoredEvent, void, unknown>;
get(eventId: string): Promise<StoredEvent | null>;
getCheckpoint(topic: string): Promise<{ consumer: string; topic: string; cursor: number; updated_at: string | null }>;
setCheckpoint(topic: string, cursor: number): Promise<{ consumer: string; topic: string; cursor: number }>;
subscriptions: {
create(input: CreateSubscription): Promise<Subscription>;
list(): Promise<Subscription[]>;
get(id: string): Promise<Subscription | null>;
disable(id: string): Promise<Subscription>;
enable(id: string): Promise<Subscription>;
};
subscribe(input: CreateSubscription): Promise<Subscription>;
deliveries(opts?: { status?: 'pending' | 'delivered' | 'failed' | 'dead'; subscriptionId?: string; afterSeq?: number; limit?: number }): Promise<{ deliveries: Array<Record<string, unknown>>; counts: Record<string, number> }>;
replay(opts: { subscriptionId: string; fromSeq?: number; eventIds?: string[] }): Promise<{ subscription_id: string; queued: number }>;
}function createEventsClient
export declare function createEventsClient(client: OpenVibeClient, opts?: { source?: string; baseUrl?: string }): EventsClient;function projectKey
'prj_01JAB…' -> 'p01jab…' (the second segment of a project's event types); null if not prj_<ULID>.
export declare function projectKey(projectId: string): string | null;function appSource
'app_01JAB…' or 'app:app_01JAB…' -> 'app-01jab…' (the app's event source); null if not an app id.
export declare function appSource(appId: string): string | null;type AppEventInput
An app event: event_type may be project-relative ('order.shipped'); actor and subject default to the app.
export type AppEventInput = Omit<EventInput, 'actor' | 'subject'> & Partial<Pick<EventEnvelope, 'actor' | 'subject'>>;interface AppEventsClient
export interface AppEventsClient {
projectId: string;
/** app_<ULID> */
appId: string;
/** p<lowercased project ULID> */
projectKey: string;
/** app-<lowercased app ULID> */
source: string;
/** app.<projectKey>. */
prefix: string;
/** A project-relative type or pattern -> the full one ('order.*' -> 'app.<key>.order.*'; '*' -> 'app.<key>.*'). */
topic(name: string): string;
prepare(envelope: AppEventInput, opts?: { traceId?: string; now?: number }): EventEnvelope;
publish(envelope: AppEventInput, opts?: { traceparent?: string }): Promise<PublishResult>;
publish(envelopes: AppEventInput[], opts?: { traceparent?: string }): Promise<{ results: PublishResult[] }>;
/** topic is project-relative (default '*'); platformTopics are first-party patterns (public events only). */
pull(opts?: { topic?: string | string[]; platformTopics?: string[]; after?: string; afterSeq?: number; limit?: number }): Promise<EventsPage>;
iterate(opts?: { topic?: string | string[]; platformTopics?: string[]; after?: string; afterSeq?: number; limit?: number; maxPages?: number; onGap?: (gap: Gap) => void | Promise<void>; onPage?: (page: EventsPage) => void | Promise<void> }): AsyncGenerator<StoredEvent, void, unknown>;
get(eventId: string): Promise<StoredEvent | null>;
getCheckpoint(topic: string): Promise<{ consumer: string; topic: string; cursor: number; updated_at: string | null }>;
setCheckpoint(topic: string, cursor: number): Promise<{ consumer: string; topic: string; cursor: number }>;
subscriptions: {
/** topicPattern is project-relative (default '*'); the endpoint must be public https. */
create(input: Partial<CreateSubscription> & { endpoint: string }): Promise<Subscription>;
list(): Promise<Subscription[]>;
get(id: string): Promise<Subscription | null>;
disable(id: string): Promise<Subscription>;
enable(id: string): Promise<Subscription>;
};
subscribe(input: Partial<CreateSubscription> & { endpoint: string }): Promise<Subscription>;
/** The unscoped client underneath. */
events: EventsClient;
}function createAppEvents
Events for one developer app of one project (token: events.app.publish | read | subscribe). onBehalfOf (usr_…, the token's on_behalf_of) makes that person the default actor.
export declare function createAppEvents(client: OpenVibeClient, opts: { projectId: string; appId: string; onBehalfOf?: string; baseUrl?: string }): AppEventsClient;function signDelivery
export declare function signDelivery(rawBody: RawBody, secret: string): string;function verifyDelivery
export declare function verifyDelivery(rawBody: RawBody, signatureHeader: string | undefined | null, secret: string): boolean;function signDeliveryV2
`t=<ts>,v2=<hex HMAC-SHA256 of "<ts>.<raw body>">`, the value of X-OpenVibe-Signature-V2 (timestamp in unix seconds, default now).
export declare function signDeliveryV2(rawBody: RawBody, secret: string, timestamp?: number): string;function signDeliveryHeaders
X-OpenVibe-Signature, X-OpenVibe-Timestamp and X-OpenVibe-Signature-V2 for a delivery body, as Events sends them (now in ms).
export declare function signDeliveryHeaders(rawBody: RawBody, secret: string, opts?: { now?: number }): {
'X-OpenVibe-Signature': string;
'X-OpenVibe-Timestamp': string;
'X-OpenVibe-Signature-V2': string;
};interface DeliveryV2Options
export interface DeliveryV2Options {
/** Seconds the timestamp may be away from `now`, either way. Default 300. */
toleranceSec?: number;
/** The current time in ms. Default Date.now(). */
now?: number;
}function verifyDeliveryV2
Constant-time check of X-OpenVibe-Signature-V2 over the raw body, and of its timestamp (±toleranceSec).
export declare function verifyDeliveryV2(rawBody: RawBody, headers: Record<string, any> | Headers, secret: string, opts?: DeliveryV2Options): boolean;interface ParseDeliveryOptions
export interface ParseDeliveryOptions extends DeliveryV2Options {
/** Refuse deliveries without X-OpenVibe-Signature-V2 (v1-only). Default false. A present v2 header must always verify. */
requireV2?: boolean;
}function parseDelivery
export declare function parseDelivery(rawBody: RawBody, headers: Record<string, any> | Headers, secret: string, opts?: ParseDeliveryOptions): { event: EventEnvelope; seq: number; subscriptionId: string | null; attempt: number } | null;interface SqliteDatabase
The subset of a better-sqlite3 Database the outbox and inbox use.
export interface SqliteDatabase {
prepare(sql: string): any;
exec(sql: string): any;
transaction<F extends (...args: any[]) => any>(fn: F): F;
readonly inTransaction: boolean;
}interface FlushStats
export interface FlushStats { sent: number; failed: number; rejected: number; }interface Outbox
export interface Outbox {
ensureSchema(): void;
/** Inside the caller's transaction; returns the complete envelope. */
enqueue(envelope: EventInput, opts?: { traceparent?: string }): EventEnvelope;
flush(): Promise<FlushStats>;
start(): void;
stop(): Promise<unknown>;
kick(): void;
pending(): number;
rejected(): number;
prune(olderThanMs?: number): number;
}function createOutbox
export declare function createOutbox(db: SqliteDatabase, opts: {
events: Pick<EventsClient, 'publish' | 'prepare'>;
table?: string; batchSize?: number; intervalMs?: number; backoffMs?: number[];
now?: () => number; onError?: (err: unknown, row?: unknown) => void; allowOutsideTransaction?: boolean;
}): Outbox;interface Inbox
export interface Inbox {
ensureSchema(): void;
once<T>(consumer: string, eventId: string, fn: () => T): { duplicate: true } | { duplicate: false; result: T };
seen(consumer: string, eventId: string): boolean;
}function createInbox
export declare function createInbox(db: SqliteDatabase, opts?: { table?: string; now?: () => number }): Inbox;interface PgOutbox
The PostgreSQL outbox (ADR-035) on an openvibe-sdk/db handle; several relays may share one table.
export interface PgOutbox {
schema(): string;
ensureSchema(): Promise<unknown>;
/** Pass the transaction handle db.tx gives you: the event exists if and only if the change commits. */
enqueue(t: import('./db').Tx, envelope: EventInput, opts?: { traceparent?: string }): Promise<EventEnvelope>;
flush(): Promise<FlushStats>;
start(): void;
stop(): Promise<unknown>;
kick(): void;
pending(): Promise<number | null>;
rejected(): Promise<number | null>;
prune(olderThanMs?: number): Promise<number>;
}function createPgOutbox
export declare function createPgOutbox(db: import('./db').Db, opts: {
events: Pick<EventsClient, 'publish' | 'prepare'>;
table?: string; batchSize?: number; intervalMs?: number; leaseMs?: number; backoffMs?: number[];
now?: () => number; onError?: (err: unknown, row?: unknown) => void;
}): PgOutbox;interface PgInbox
export interface PgInbox {
schema(): string;
ensureSchema(): Promise<unknown>;
once<T>(consumer: string, eventId: string, fn: (t: import('./db').Tx) => Promise<T>): Promise<{ duplicate: true } | { duplicate: false; result: T }>;
seen(consumer: string, eventId: string): Promise<boolean>;
}function createPgInbox
export declare function createPgInbox(db: import('./db').Db, opts?: { table?: string; now?: () => number }): PgInbox;function outboxSchema
DDL for a service migration.
export declare function outboxSchema(table?: string): string;interface ServiceOutboxOptions
export interface ServiceOutboxOptions {
db: import('./db').Db;
/** The service id and event source, e.g. 'trade'. */
source: string;
/** The relay runs only with eventsUrl and clientSecret; otherwise rows wait and status().enabled is false. */
eventsUrl?: string | null;
networkInternalUrl?: string;
clientId?: string | null;
clientSecret?: string | null;
table?: string;
intervalMs?: number;
now?: () => number;
fetch?: import('./core').FetchLike;
log?: Pick<Console, 'warn' | 'log'>;
/** emit refuses a type not in this list. */
eventTypes?: string[] | null;
/** emit refuses an envelope this rejects (e.g. openvibe-contracts validate). */
validate?: ((envelope: EventEnvelope) => { valid: boolean; errors?: Array<{ path?: string; message?: string }> }) | null;
autoDiscover?: boolean;
/** false: <source>.moderation.action withholds target.owner_subject. */
moderationOwnerSubject?: boolean;
}interface ServiceOutbox
export interface ServiceOutbox {
/** Inside the change's ambient openvibe-sdk/db transaction. */
emit(envelope: EventInput, opts?: { traceparent?: string }): Promise<EventEnvelope>;
emitIn(t: import('./db').Tx, envelope: EventInput, opts?: { traceparent?: string }): Promise<EventEnvelope>;
/** <source>.moderation.action. */
moderationAction(input: { action: string; target: { type: string; id: string | number; owner_subject?: string | null; [k: string]: unknown }; actorSubject?: string | null; reason?: string | null; details?: Record<string, unknown> }, opts?: { traceparent?: string }): Promise<EventEnvelope>;
outbox: PgOutbox;
events: EventsClient;
enabled: boolean;
start(): void;
stop(): Promise<unknown>;
kick(): Promise<void>;
status(): Promise<{ enabled: boolean; pending: number | null; rejected: number | null; last_error: string | null }>;
}function createServiceOutbox
export declare function createServiceOutbox(opts: ServiceOutboxOptions): ServiceOutbox;function inboxSchema
export declare function inboxSchema(table?: string): string;