Skip to content

@codesoul-co/hypha-core / modules/runtime/message-bus

模块用法

用于执行该边界的运行时行为。Message bus 模块公开 1 类、7 函数、6 接口。

从包入口导入

ts
import {
  InMemoryMessageBus,
  addMilliseconds,
  busError,
  createRuntimeMessageEnvelope,
  isAtOrBefore,
  nonEmpty,
  nonNegative,
  positive,
} from '@codesoul-co/hypha-core';

import type {
  InMemoryMessageBusOptions,
  MessageBus,
  MessageDelivery,
  MessagePublishRequest,
  MessagePublishResult,
  MessageSubscriptionRequest,
} from '@codesoul-co/hypha-core';

使用要点

  • 6 个类型/接口用于应用代码、Adapter 或测试中的静态契约;请使用 import type,运行时不应依赖它们。
  • 1 个类提供可实例化的运行时实现;构造参数与公开方法在各自条目中完整列出。
  • 7 个函数是该模块的直接操作入口;每个 overload 的必需/可选参数与返回类型均在下方列出。

公共导出

Symbol种类签名说明
InMemoryMessageBusnew InMemoryMessageBus(options?: InMemoryMessageBusOptions): InMemoryMessageBusIn Memory Message Bus 类,共公开 7 个构造函数或成员;精确签名见本条目的声明与成员表。
addMilliseconds函数addMilliseconds(timestamp: string, milliseconds: number): stringAdd Milliseconds 函数,提供 1 个公开调用签名;参数与返回类型见下表。
busError函数busError(code: string, message: string, context?: Record<string, unknown>): FrameworkErrorBus Error 函数,提供 1 个公开调用签名;参数与返回类型见下表。
createRuntimeMessageEnvelope函数createRuntimeMessageEnvelope<TPayload>(input: RuntimeMessageEnvelopeInput<TPayload>): RuntimeMessageEnvelope<TPayload>Create Runtime Message Envelope 函数,提供 1 个公开调用签名;参数与返回类型见下表。
isAtOrBefore函数isAtOrBefore(left: string, right: string): booleanIs At Or Before 函数,提供 1 个公开调用签名;参数与返回类型见下表。
nonEmpty函数nonEmpty(value: unknown, label: string): asserts value is stringNon Empty 函数,提供 1 个公开调用签名;参数与返回类型见下表。
nonNegative函数nonNegative(value: number, label?: string): numberNon Negative 函数,提供 1 个公开调用签名;参数与返回类型见下表。
positive函数positive(value: number, label: string): numberPositive 函数,提供 1 个公开调用签名;参数与返回类型见下表。
InMemoryMessageBusOptions接口interface InMemoryMessageBusOptionsIn Memory Message Bus Options 接口,共包含 6 个公开字段或方法。
MessageBus接口interface MessageBusMessage Bus 接口,共包含 5 个公开字段或方法。
MessageDelivery接口interface MessageDeliveryMessage Delivery 接口,共包含 9 个公开字段或方法。
MessagePublishRequest接口interface MessagePublishRequestMessage Publish Request 接口,共包含 1 个公开字段或方法。
MessagePublishResult接口interface MessagePublishResultMessage Publish Result 接口,共包含 6 个公开字段或方法。
MessageSubscriptionRequest接口interface MessageSubscriptionRequestMessage Subscription Request 接口,共包含 8 个公开字段或方法。

InMemoryMessageBus

In Memory Message Bus 类,共公开 7 个构造函数或成员;精确签名见本条目的声明与成员表。

声明

text
export declare class InMemoryMessageBus implements MessageBus {
    constructor(options?: InMemoryMessageBusOptions);
    publish<TPayload>(request: MessagePublishRequest<TPayload>): Promise<MessagePublishResult>;
    publishBatch<TPayload>(requests: MessagePublishRequest<TPayload>[]): Promise<MessagePublishResult[]>;
    subscribe<TPayload>(request: MessageSubscriptionRequest): AsyncIterable<MessageDelivery<TPayload>>;
    health(): Promise<ProviderHealth>;
    close(): Promise<void>;
    listDeadLetters(consumerGroup: string): RuntimeMessageEnvelope[];
}

公开成员

成员种类签名说明
close方法close(): Promise<void>公开方法;参数与返回类型以签名列为准。
constructor构造函数(options?: InMemoryMessageBusOptions): InMemoryMessageBus创建该类的实例。
health方法health(): Promise<ProviderHealth>公开方法;参数与返回类型以签名列为准。
listDeadLetters方法listDeadLetters(consumerGroup: string): RuntimeMessageEnvelope[]公开方法;参数与返回类型以签名列为准。
publish方法publish<TPayload>(request: MessagePublishRequest<TPayload>): Promise<MessagePublishResult>公开方法;参数与返回类型以签名列为准。
publishBatch方法publishBatch<TPayload>(requests: MessagePublishRequest<TPayload>[]): Promise<MessagePublishResult[]>公开方法;参数与返回类型以签名列为准。
subscribe方法subscribe<TPayload>(request: MessageSubscriptionRequest): AsyncIterable<MessageDelivery<TPayload>>公开方法;参数与返回类型以签名列为准。

addMilliseconds

Add Milliseconds 函数,提供 1 个公开调用签名;参数与返回类型见下表。

声明

text
export declare function addMilliseconds(timestamp: string, milliseconds: number): string;

调用签名

text
addMilliseconds(timestamp: string, milliseconds: number): string

参数

参数类型必需说明
timestampstring必需参数;接受的值由类型列定义。
millisecondsnumber必需参数;接受的值由类型列定义。

返回值

  • 类型: string
  • 说明: 返回值契约由上述类型定义。

busError

Bus Error 函数,提供 1 个公开调用签名;参数与返回类型见下表。

声明

text
export declare function busError(code: string, message: string, context?: Record<string, unknown>): FrameworkError;

调用签名

text
busError(code: string, message: string, context?: Record<string, unknown>): FrameworkError

参数

参数类型必需说明
codestring必需参数;接受的值由类型列定义。
messagestring必需参数;接受的值由类型列定义。
contextRecord<string, unknown>可选参数;接受的值由类型列定义。

返回值

  • 类型: FrameworkError
  • 说明: 返回值契约由上述类型定义。

createRuntimeMessageEnvelope

Create Runtime Message Envelope 函数,提供 1 个公开调用签名;参数与返回类型见下表。

声明

text
export declare function createRuntimeMessageEnvelope<TPayload>(input: RuntimeMessageEnvelopeInput<TPayload>): RuntimeMessageEnvelope<TPayload>;

调用签名

text
createRuntimeMessageEnvelope<TPayload>(input: RuntimeMessageEnvelopeInput<TPayload>): RuntimeMessageEnvelope<TPayload>

参数

参数类型必需说明
inputRuntimeMessageEnvelopeInput<TPayload>必需参数;接受的值由类型列定义。

返回值

  • 类型: RuntimeMessageEnvelope<TPayload>
  • 说明: 返回值契约由上述类型定义。

isAtOrBefore

Is At Or Before 函数,提供 1 个公开调用签名;参数与返回类型见下表。

声明

text
export declare function isAtOrBefore(left: string, right: string): boolean;

调用签名

text
isAtOrBefore(left: string, right: string): boolean

参数

参数类型必需说明
leftstring必需参数;接受的值由类型列定义。
rightstring必需参数;接受的值由类型列定义。

返回值

  • 类型: boolean
  • 说明: 返回值契约由上述类型定义。

nonEmpty

Non Empty 函数,提供 1 个公开调用签名;参数与返回类型见下表。

声明

text
export declare function nonEmpty(value: unknown, label: string): asserts value is string;

调用签名

text
nonEmpty(value: unknown, label: string): asserts value is string

参数

参数类型必需说明
valueunknown必需参数;接受的值由类型列定义。
labelstring必需参数;接受的值由类型列定义。

返回值

  • 类型: asserts value is string
  • 说明: 返回值契约由上述类型定义。

nonNegative

Non Negative 函数,提供 1 个公开调用签名;参数与返回类型见下表。

声明

text
export declare function nonNegative(value: number, label?: string): number;

调用签名

text
nonNegative(value: number, label?: string): number

参数

参数类型必需说明
valuenumber必需参数;接受的值由类型列定义。
labelstring可选参数;接受的值由类型列定义。

返回值

  • 类型: number
  • 说明: 返回值契约由上述类型定义。

positive

Positive 函数,提供 1 个公开调用签名;参数与返回类型见下表。

声明

text
export declare function positive(value: number, label: string): number;

调用签名

text
positive(value: number, label: string): number

参数

参数类型必需说明
valuenumber必需参数;接受的值由类型列定义。
labelstring必需参数;接受的值由类型列定义。

返回值

  • 类型: number
  • 说明: 返回值契约由上述类型定义。

InMemoryMessageBusOptions

In Memory Message Bus Options 接口,共包含 6 个公开字段或方法。

  • 种类: 接口
  • 导入: import type { InMemoryMessageBusOptions } from '@codesoul-co/hypha-core';
  • 源码模块: modules/runtime/message-bus

声明

text
export interface InMemoryMessageBusOptions {
    now?: () => string;
    maxDeliveryAttempts?: number;
    defaultAckDeadlineMs?: number;
    maxMessageBytes?: number;
    maxQueueDepth?: number;
    pollIntervalMs?: number;
}

契约成员

成员种类签名说明
defaultAckDeadlineMs属性defaultAckDeadlineMs?: number公开属性;类型、只读和可选状态以签名列为准。
maxDeliveryAttempts属性maxDeliveryAttempts?: number公开属性;类型、只读和可选状态以签名列为准。
maxMessageBytes属性maxMessageBytes?: number公开属性;类型、只读和可选状态以签名列为准。
maxQueueDepth属性maxQueueDepth?: number公开属性;类型、只读和可选状态以签名列为准。
now方法now?(): string公开方法;参数与返回类型以签名列为准。
pollIntervalMs属性pollIntervalMs?: number公开属性;类型、只读和可选状态以签名列为准。

MessageBus

Message Bus 接口,共包含 5 个公开字段或方法。

声明

text
export interface MessageBus {
    publish<TPayload>(request: MessagePublishRequest<TPayload>): Promise<MessagePublishResult>;
    publishBatch<TPayload>(requests: MessagePublishRequest<TPayload>[]): Promise<MessagePublishResult[]>;
    subscribe<TPayload>(request: MessageSubscriptionRequest): AsyncIterable<MessageDelivery<TPayload>>;
    health(): Promise<ProviderHealth>;
    close(): Promise<void>;
}

契约成员

成员种类签名说明
close方法close(): Promise<void>公开方法;参数与返回类型以签名列为准。
health方法health(): Promise<ProviderHealth>公开方法;参数与返回类型以签名列为准。
publish方法publish<TPayload>(request: MessagePublishRequest<TPayload>): Promise<MessagePublishResult>公开方法;参数与返回类型以签名列为准。
publishBatch方法publishBatch<TPayload>(requests: MessagePublishRequest<TPayload>[]): Promise<MessagePublishResult[]>公开方法;参数与返回类型以签名列为准。
subscribe方法subscribe<TPayload>(request: MessageSubscriptionRequest): AsyncIterable<MessageDelivery<TPayload>>公开方法;参数与返回类型以签名列为准。

MessageDelivery

Message Delivery 接口,共包含 9 个公开字段或方法。

声明

text
export interface MessageDelivery<TPayload = unknown> {
    envelope: RuntimeMessageEnvelope<TPayload>;
    deliveryId: string;
    attempt: number;
    receivedAt: string;
    ackDeadlineAt: string;
    ack(): Promise<void>;
    nack(options?: {
        delayMs?: number;
        reason?: string;
    }): Promise<void>;
    deadLetter(reason: string): Promise<void>;
    extendAckDeadline(extensionMs: number): Promise<void>;
}

契约成员

成员种类签名说明
ack方法ack(): Promise<void>公开方法;参数与返回类型以签名列为准。
ackDeadlineAt属性ackDeadlineAt: string公开属性;类型、只读和可选状态以签名列为准。
attempt属性attempt: number公开属性;类型、只读和可选状态以签名列为准。
deadLetter方法deadLetter(reason: string): Promise<void>公开方法;参数与返回类型以签名列为准。
deliveryId属性deliveryId: string公开属性;类型、只读和可选状态以签名列为准。
envelope属性envelope: RuntimeMessageEnvelope<TPayload>公开属性;类型、只读和可选状态以签名列为准。
extendAckDeadline方法extendAckDeadline(extensionMs: number): Promise<void>公开方法;参数与返回类型以签名列为准。
nack方法nack(options?: { delayMs?: number; reason?: string; }): Promise<void>公开方法;参数与返回类型以签名列为准。
receivedAt属性receivedAt: string公开属性;类型、只读和可选状态以签名列为准。

MessagePublishRequest

Message Publish Request 接口,共包含 1 个公开字段或方法。

声明

text
export interface MessagePublishRequest<TPayload = unknown> {
    envelope: RuntimeMessageEnvelopeInput<TPayload>;
}

契约成员

成员种类签名说明
envelope属性envelope: RuntimeMessageEnvelopeInput<TPayload>公开属性;类型、只读和可选状态以签名列为准。

MessagePublishResult

Message Publish Result 接口,共包含 6 个公开字段或方法。

声明

text
export interface MessagePublishResult {
    messageId: string;
    topic: string;
    partitionKey: string;
    sequence: number;
    publishedAt: string;
    reused: boolean;
}

契约成员

成员种类签名说明
messageId属性messageId: string公开属性;类型、只读和可选状态以签名列为准。
partitionKey属性partitionKey: string公开属性;类型、只读和可选状态以签名列为准。
publishedAt属性publishedAt: string公开属性;类型、只读和可选状态以签名列为准。
reused属性reused: boolean公开属性;类型、只读和可选状态以签名列为准。
sequence属性sequence: number公开属性;类型、只读和可选状态以签名列为准。
topic属性topic: string公开属性;类型、只读和可选状态以签名列为准。

MessageSubscriptionRequest

Message Subscription Request 接口,共包含 8 个公开字段或方法。

  • 种类: 接口
  • 导入: import type { MessageSubscriptionRequest } from '@codesoul-co/hypha-core';
  • 源码模块: modules/runtime/message-bus

声明

text
export interface MessageSubscriptionRequest {
    consumerId: string;
    consumerGroup?: string;
    topic: string;
    partitionKey?: string;
    maxMessages?: number;
    idleTimeoutMs?: number;
    ackDeadlineMs?: number;
    signal?: AbortSignal;
}

契约成员

成员种类签名说明
ackDeadlineMs属性ackDeadlineMs?: number公开属性;类型、只读和可选状态以签名列为准。
consumerGroup属性consumerGroup?: string公开属性;类型、只读和可选状态以签名列为准。
consumerId属性consumerId: string公开属性;类型、只读和可选状态以签名列为准。
idleTimeoutMs属性idleTimeoutMs?: number公开属性;类型、只读和可选状态以签名列为准。
maxMessages属性maxMessages?: number公开属性;类型、只读和可选状态以签名列为准。
partitionKey属性partitionKey?: string公开属性;类型、只读和可选状态以签名列为准。
signal属性signal?: AbortSignal公开属性;类型、只读和可选状态以签名列为准。
topic属性topic: string公开属性;类型、只读和可选状态以签名列为准。