mirror of
https://github.com/logos-messaging/logos-messaging-js.git
synced 2026-08-10 03:53:21 +00:00
feat: add v3 protocol support tests and exports
- Export LightPushCodecV3 and LightPushCodecs from core package - Add v3 protocol support tests with status code validation - Add mock functions for v3 responses and RLN errors - Test mixed v2/v3 peer scenarios - Validate protocol error mapping and status code handling - Fix linting errors by adding explicit return types
This commit is contained in:
parent
539c92f1bb
commit
41301ea64a
@ -10,7 +10,12 @@ export * as waku_filter from "./lib/filter/index.js";
|
|||||||
export { FilterCore, FilterCodecs } from "./lib/filter/index.js";
|
export { FilterCore, FilterCodecs } from "./lib/filter/index.js";
|
||||||
|
|
||||||
export * as waku_light_push from "./lib/light_push/index.js";
|
export * as waku_light_push from "./lib/light_push/index.js";
|
||||||
export { LightPushCodec, LightPushCore } from "./lib/light_push/index.js";
|
export {
|
||||||
|
LightPushCodec,
|
||||||
|
LightPushCodecV3,
|
||||||
|
LightPushCodecs,
|
||||||
|
LightPushCore
|
||||||
|
} from "./lib/light_push/index.js";
|
||||||
|
|
||||||
export * as waku_store from "./lib/store/index.js";
|
export * as waku_store from "./lib/store/index.js";
|
||||||
export { StoreCore, StoreCodec } from "./lib/store/index.js";
|
export { StoreCore, StoreCodec } from "./lib/store/index.js";
|
||||||
|
|||||||
@ -2,9 +2,9 @@ import type { PeerId, Stream } from "@libp2p/interface";
|
|||||||
import type { IncomingStreamData } from "@libp2p/interface-internal";
|
import type { IncomingStreamData } from "@libp2p/interface-internal";
|
||||||
import {
|
import {
|
||||||
type ContentTopic,
|
type ContentTopic,
|
||||||
type CoreProtocolResult,
|
type FilterCoreResult,
|
||||||
|
FilterError,
|
||||||
type Libp2p,
|
type Libp2p,
|
||||||
ProtocolError,
|
|
||||||
type PubsubTopic
|
type PubsubTopic
|
||||||
} from "@waku/interfaces";
|
} from "@waku/interfaces";
|
||||||
import { WakuMessage } from "@waku/proto";
|
import { WakuMessage } from "@waku/proto";
|
||||||
@ -62,7 +62,7 @@ export class FilterCore {
|
|||||||
pubsubTopic: PubsubTopic,
|
pubsubTopic: PubsubTopic,
|
||||||
peerId: PeerId,
|
peerId: PeerId,
|
||||||
contentTopics: ContentTopic[]
|
contentTopics: ContentTopic[]
|
||||||
): Promise<CoreProtocolResult> {
|
): Promise<FilterCoreResult> {
|
||||||
const stream = await this.streamManager.getStream(peerId);
|
const stream = await this.streamManager.getStream(peerId);
|
||||||
|
|
||||||
const request = FilterSubscribeRpc.createSubscribeRequest(
|
const request = FilterSubscribeRpc.createSubscribeRequest(
|
||||||
@ -88,7 +88,7 @@ export class FilterCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.GENERIC_FAIL,
|
error: FilterError.SUBSCRIPTION_FAILED,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -103,8 +103,11 @@ export class FilterCore {
|
|||||||
);
|
);
|
||||||
return {
|
return {
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.REMOTE_PEER_REJECTED,
|
error: FilterError.REMOTE_PEER_REJECTED,
|
||||||
peerId: peerId
|
peerId: peerId,
|
||||||
|
statusCode: statusCode,
|
||||||
|
statusDesc: statusDesc,
|
||||||
|
requestId: requestId
|
||||||
},
|
},
|
||||||
success: null
|
success: null
|
||||||
};
|
};
|
||||||
@ -120,7 +123,7 @@ export class FilterCore {
|
|||||||
pubsubTopic: PubsubTopic,
|
pubsubTopic: PubsubTopic,
|
||||||
peerId: PeerId,
|
peerId: PeerId,
|
||||||
contentTopics: ContentTopic[]
|
contentTopics: ContentTopic[]
|
||||||
): Promise<CoreProtocolResult> {
|
): Promise<FilterCoreResult> {
|
||||||
let stream: Stream | undefined;
|
let stream: Stream | undefined;
|
||||||
try {
|
try {
|
||||||
stream = await this.streamManager.getStream(peerId);
|
stream = await this.streamManager.getStream(peerId);
|
||||||
@ -132,7 +135,7 @@ export class FilterCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.NO_STREAM_AVAILABLE,
|
error: FilterError.NO_STREAM_AVAILABLE,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -150,7 +153,7 @@ export class FilterCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.GENERIC_FAIL,
|
error: FilterError.GENERIC_FAIL,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -165,7 +168,7 @@ export class FilterCore {
|
|||||||
public async unsubscribeAll(
|
public async unsubscribeAll(
|
||||||
pubsubTopic: PubsubTopic,
|
pubsubTopic: PubsubTopic,
|
||||||
peerId: PeerId
|
peerId: PeerId
|
||||||
): Promise<CoreProtocolResult> {
|
): Promise<FilterCoreResult> {
|
||||||
const stream = await this.streamManager.getStream(peerId);
|
const stream = await this.streamManager.getStream(peerId);
|
||||||
|
|
||||||
const request = FilterSubscribeRpc.createUnsubscribeAllRequest(pubsubTopic);
|
const request = FilterSubscribeRpc.createUnsubscribeAllRequest(pubsubTopic);
|
||||||
@ -181,7 +184,7 @@ export class FilterCore {
|
|||||||
if (!res || !res.length) {
|
if (!res || !res.length) {
|
||||||
return {
|
return {
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.NO_RESPONSE,
|
error: FilterError.NO_RESPONSE,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
},
|
},
|
||||||
success: null
|
success: null
|
||||||
@ -197,7 +200,7 @@ export class FilterCore {
|
|||||||
);
|
);
|
||||||
return {
|
return {
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.REMOTE_PEER_REJECTED,
|
error: FilterError.REMOTE_PEER_REJECTED,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
},
|
},
|
||||||
success: null
|
success: null
|
||||||
@ -210,7 +213,7 @@ export class FilterCore {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
public async ping(peerId: PeerId): Promise<CoreProtocolResult> {
|
public async ping(peerId: PeerId): Promise<FilterCoreResult> {
|
||||||
let stream: Stream | undefined;
|
let stream: Stream | undefined;
|
||||||
try {
|
try {
|
||||||
stream = await this.streamManager.getStream(peerId);
|
stream = await this.streamManager.getStream(peerId);
|
||||||
@ -222,7 +225,7 @@ export class FilterCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.NO_STREAM_AVAILABLE,
|
error: FilterError.NO_STREAM_AVAILABLE,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -244,7 +247,7 @@ export class FilterCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.GENERIC_FAIL,
|
error: FilterError.GENERIC_FAIL,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -254,7 +257,7 @@ export class FilterCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.NO_RESPONSE,
|
error: FilterError.NO_RESPONSE,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -270,7 +273,7 @@ export class FilterCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.REMOTE_PEER_REJECTED,
|
error: FilterError.REMOTE_PEER_REJECTED,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|||||||
@ -1 +1,7 @@
|
|||||||
export { LightPushCore, LightPushCodec, PushResponse } from "./light_push.js";
|
export {
|
||||||
|
LightPushCore,
|
||||||
|
LightPushCodec,
|
||||||
|
LightPushCodecV3,
|
||||||
|
LightPushCodecs,
|
||||||
|
PushResponse
|
||||||
|
} from "./light_push.js";
|
||||||
|
|||||||
@ -1,13 +1,13 @@
|
|||||||
import type { PeerId, Stream } from "@libp2p/interface";
|
import type { PeerId, Stream } from "@libp2p/interface";
|
||||||
import {
|
import {
|
||||||
type CoreProtocolResult,
|
|
||||||
type IEncoder,
|
type IEncoder,
|
||||||
type IMessage,
|
type IMessage,
|
||||||
isSuccess as isV3Success,
|
isSuccess as isV3Success,
|
||||||
type Libp2p,
|
type Libp2p,
|
||||||
ProtocolError,
|
type LightPushCoreResult,
|
||||||
|
LightPushError,
|
||||||
type ThisOrThat,
|
type ThisOrThat,
|
||||||
toProtocolError
|
toLightPushError
|
||||||
} from "@waku/interfaces";
|
} from "@waku/interfaces";
|
||||||
import { PushResponse } from "@waku/proto";
|
import { PushResponse } from "@waku/proto";
|
||||||
import { isMessageSizeUnderCap } from "@waku/utils";
|
import { isMessageSizeUnderCap } from "@waku/utils";
|
||||||
@ -30,7 +30,12 @@ export const LightPushCodecV3 = "/vac/waku/lightpush/3.0.0";
|
|||||||
export const LightPushCodecs = [LightPushCodecV3, LightPushCodec];
|
export const LightPushCodecs = [LightPushCodecV3, LightPushCodec];
|
||||||
export { PushResponse };
|
export { PushResponse };
|
||||||
|
|
||||||
type PreparePushMessageResult = ThisOrThat<"query", PushRpc>;
|
type PreparePushMessageResult = ThisOrThat<
|
||||||
|
"query",
|
||||||
|
PushRpc,
|
||||||
|
"error",
|
||||||
|
LightPushError
|
||||||
|
>;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Implements the [Waku v2 Light Push protocol](https://rfc.vac.dev/spec/19/).
|
* Implements the [Waku v2 Light Push protocol](https://rfc.vac.dev/spec/19/).
|
||||||
@ -89,12 +94,12 @@ export class LightPushCore {
|
|||||||
try {
|
try {
|
||||||
if (!message.payload || message.payload.length === 0) {
|
if (!message.payload || message.payload.length === 0) {
|
||||||
log.error("Failed to send waku light push: payload is empty");
|
log.error("Failed to send waku light push: payload is empty");
|
||||||
return { query: null, error: ProtocolError.EMPTY_PAYLOAD };
|
return { query: null, error: LightPushError.EMPTY_PAYLOAD };
|
||||||
}
|
}
|
||||||
|
|
||||||
if (!(await isMessageSizeUnderCap(encoder, message))) {
|
if (!(await isMessageSizeUnderCap(encoder, message))) {
|
||||||
log.error("Failed to send waku light push: message is bigger than 1MB");
|
log.error("Failed to send waku light push: message is bigger than 1MB");
|
||||||
return { query: null, error: ProtocolError.SIZE_TOO_BIG };
|
return { query: null, error: LightPushError.SIZE_TOO_BIG };
|
||||||
}
|
}
|
||||||
|
|
||||||
const protoMessage = await encoder.toProtoObj(message);
|
const protoMessage = await encoder.toProtoObj(message);
|
||||||
@ -102,7 +107,7 @@ export class LightPushCore {
|
|||||||
log.error("Failed to encode to protoMessage, aborting push");
|
log.error("Failed to encode to protoMessage, aborting push");
|
||||||
return {
|
return {
|
||||||
query: null,
|
query: null,
|
||||||
error: ProtocolError.ENCODE_FAILED
|
error: LightPushError.ENCODE_FAILED
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -113,7 +118,7 @@ export class LightPushCore {
|
|||||||
|
|
||||||
return {
|
return {
|
||||||
query: null,
|
query: null,
|
||||||
error: ProtocolError.GENERIC_FAIL
|
error: LightPushError.GENERIC_FAIL
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -122,7 +127,7 @@ export class LightPushCore {
|
|||||||
encoder: IEncoder,
|
encoder: IEncoder,
|
||||||
message: IMessage,
|
message: IMessage,
|
||||||
peerId: PeerId
|
peerId: PeerId
|
||||||
): Promise<CoreProtocolResult> {
|
): Promise<LightPushCoreResult> {
|
||||||
const { query, error: preparationError } = await this.preparePushMessage(
|
const { query, error: preparationError } = await this.preparePushMessage(
|
||||||
encoder,
|
encoder,
|
||||||
message
|
message
|
||||||
@ -154,7 +159,7 @@ export class LightPushCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.NO_STREAM_AVAILABLE,
|
error: LightPushError.NO_STREAM_AVAILABLE,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -176,7 +181,7 @@ export class LightPushCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.STREAM_ABORTED,
|
error: LightPushError.STREAM_ABORTED,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -195,7 +200,7 @@ export class LightPushCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.DECODE_FAILED,
|
error: LightPushError.DECODE_FAILED,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -206,7 +211,7 @@ export class LightPushCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.NO_RESPONSE,
|
error: LightPushError.NO_RESPONSE,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -214,7 +219,7 @@ export class LightPushCore {
|
|||||||
|
|
||||||
if (protocol === LightPushCodecV3 && response.statusCode !== undefined) {
|
if (protocol === LightPushCodecV3 && response.statusCode !== undefined) {
|
||||||
if (!isV3Success(response.statusCode)) {
|
if (!isV3Success(response.statusCode)) {
|
||||||
const error = toProtocolError(response.statusCode);
|
const error = toLightPushError(response.statusCode);
|
||||||
log.error(
|
log.error(
|
||||||
`Remote peer rejected with v3 status code ${response.statusCode}: ${response.statusDesc || response.info}`
|
`Remote peer rejected with v3 status code ${response.statusCode}: ${response.statusDesc || response.info}`
|
||||||
);
|
);
|
||||||
@ -222,7 +227,9 @@ export class LightPushCore {
|
|||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error,
|
error,
|
||||||
peerId: peerId
|
peerId: peerId,
|
||||||
|
statusCode: response.statusCode,
|
||||||
|
statusDesc: response.statusDesc || response.info
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
@ -239,7 +246,7 @@ export class LightPushCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.RLN_PROOF_GENERATION,
|
error: LightPushError.RLN_PROOF_GENERATION,
|
||||||
peerId: peerId
|
peerId: peerId
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@ -250,8 +257,9 @@ export class LightPushCore {
|
|||||||
return {
|
return {
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.REMOTE_PEER_REJECTED,
|
error: LightPushError.REMOTE_PEER_REJECTED,
|
||||||
peerId: peerId
|
peerId: peerId,
|
||||||
|
statusDesc: response.info
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@ -1,5 +1,5 @@
|
|||||||
import { ProtocolError } from "./protocols.js";
|
import { LightPushError } from "./protocols.js";
|
||||||
import type { ISender, ISendOptions } from "./sender.js";
|
import type { ILightPushSender, ISendOptions } from "./sender.js";
|
||||||
|
|
||||||
export type LightPushProtocolOptions = ISendOptions & {
|
export type LightPushProtocolOptions = ISendOptions & {
|
||||||
/**
|
/**
|
||||||
@ -16,7 +16,7 @@ export type LightPushProtocolOptions = ISendOptions & {
|
|||||||
numPeersToUse?: number;
|
numPeersToUse?: number;
|
||||||
};
|
};
|
||||||
|
|
||||||
export type ILightPush = ISender & {
|
export type ILightPush = ILightPushSender & {
|
||||||
readonly multicodec: string;
|
readonly multicodec: string;
|
||||||
start: () => void;
|
start: () => void;
|
||||||
stop: () => void;
|
stop: () => void;
|
||||||
@ -53,35 +53,48 @@ export function isSuccess(statusCode: number | undefined): boolean {
|
|||||||
return statusCode === LightPushStatusCode.SUCCESS;
|
return statusCode === LightPushStatusCode.SUCCESS;
|
||||||
}
|
}
|
||||||
|
|
||||||
export function toProtocolError(
|
export function toLightPushError(
|
||||||
statusCode: LightPushStatusCode | number | undefined
|
statusCode: LightPushStatusCode | number | undefined
|
||||||
): ProtocolError {
|
): LightPushError {
|
||||||
if (!statusCode) {
|
if (!statusCode) {
|
||||||
return ProtocolError.GENERIC_FAIL;
|
return LightPushError.GENERIC_FAIL;
|
||||||
}
|
}
|
||||||
|
|
||||||
switch (statusCode) {
|
switch (statusCode) {
|
||||||
case LightPushStatusCode.SUCCESS:
|
case LightPushStatusCode.SUCCESS:
|
||||||
return ProtocolError.GENERIC_FAIL;
|
return LightPushError.GENERIC_FAIL;
|
||||||
case LightPushStatusCode.BAD_REQUEST:
|
case LightPushStatusCode.BAD_REQUEST:
|
||||||
|
return LightPushError.BAD_REQUEST;
|
||||||
case LightPushStatusCode.INVALID_MESSAGE:
|
case LightPushStatusCode.INVALID_MESSAGE:
|
||||||
|
return LightPushError.INVALID_MESSAGE;
|
||||||
case LightPushStatusCode.TOO_MANY_REQUESTS:
|
case LightPushStatusCode.TOO_MANY_REQUESTS:
|
||||||
return ProtocolError.REMOTE_PEER_REJECTED;
|
return LightPushError.TOO_MANY_REQUESTS;
|
||||||
case LightPushStatusCode.PAYLOAD_TOO_LARGE:
|
case LightPushStatusCode.PAYLOAD_TOO_LARGE:
|
||||||
return ProtocolError.SIZE_TOO_BIG;
|
return LightPushError.PAYLOAD_TOO_LARGE;
|
||||||
case LightPushStatusCode.UNSUPPORTED_TOPIC:
|
case LightPushStatusCode.UNSUPPORTED_TOPIC:
|
||||||
return ProtocolError.TOPIC_NOT_CONFIGURED;
|
return LightPushError.UNSUPPORTED_TOPIC;
|
||||||
case LightPushStatusCode.UNAVAILABLE:
|
case LightPushStatusCode.UNAVAILABLE:
|
||||||
|
return LightPushError.UNAVAILABLE;
|
||||||
case LightPushStatusCode.NO_PEERS:
|
case LightPushStatusCode.NO_PEERS:
|
||||||
return ProtocolError.NO_PEER_AVAILABLE;
|
return LightPushError.NO_PEERS;
|
||||||
case LightPushStatusCode.NO_RLN_PROOF:
|
case LightPushStatusCode.NO_RLN_PROOF:
|
||||||
return ProtocolError.RLN_PROOF_GENERATION;
|
return LightPushError.NO_RLN_PROOF;
|
||||||
case LightPushStatusCode.INTERNAL_ERROR:
|
case LightPushStatusCode.INTERNAL_ERROR:
|
||||||
default:
|
default:
|
||||||
return ProtocolError.GENERIC_FAIL;
|
return LightPushError.INTERNAL_ERROR;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Legacy function for backward compatibility
|
||||||
|
/**
|
||||||
|
* @deprecated Use toLightPushError instead
|
||||||
|
*/
|
||||||
|
export function toProtocolError(
|
||||||
|
statusCode: LightPushStatusCode | number | undefined
|
||||||
|
): LightPushError {
|
||||||
|
return toLightPushError(statusCode);
|
||||||
|
}
|
||||||
|
|
||||||
export function getStatusDescription(
|
export function getStatusDescription(
|
||||||
statusCode: number | undefined,
|
statusCode: number | undefined,
|
||||||
customDesc?: string
|
customDesc?: string
|
||||||
|
|||||||
@ -127,107 +127,138 @@ export type Callback<T extends IDecodedMessage> = (
|
|||||||
msg: T
|
msg: T
|
||||||
) => void | Promise<void>;
|
) => void | Promise<void>;
|
||||||
|
|
||||||
export enum ProtocolError {
|
// LightPush specific errors
|
||||||
//
|
export enum LightPushError {
|
||||||
// GENERAL ERRORS SECTION
|
// General errors
|
||||||
//
|
|
||||||
/**
|
|
||||||
* Could not determine the origin of the fault. Best to check connectivity and try again
|
|
||||||
* */
|
|
||||||
GENERIC_FAIL = "Generic error",
|
GENERIC_FAIL = "Generic error",
|
||||||
|
|
||||||
/**
|
|
||||||
* The remote peer rejected the message. Information provided by the remote peer
|
|
||||||
* is logged. Review message validity, or mitigation for `NO_PEER_AVAILABLE`
|
|
||||||
* or `DECODE_FAILED` can be used.
|
|
||||||
*/
|
|
||||||
REMOTE_PEER_REJECTED = "Remote peer rejected",
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Failure to protobuf decode the message. May be due to a remote peer issue,
|
|
||||||
* ensuring that messages are sent via several peer enable mitigation of this error.
|
|
||||||
*/
|
|
||||||
DECODE_FAILED = "Failed to decode",
|
DECODE_FAILED = "Failed to decode",
|
||||||
|
|
||||||
/**
|
|
||||||
* Failure to find a peer with suitable protocols. This may due to a connection issue.
|
|
||||||
* Mitigation can be: retrying after a given time period, display connectivity issue
|
|
||||||
* to user or listening for `peer:connected:bootstrap` or `peer:connected:peer-exchange`
|
|
||||||
* on the connection manager before retrying.
|
|
||||||
*/
|
|
||||||
NO_PEER_AVAILABLE = "No peer available",
|
NO_PEER_AVAILABLE = "No peer available",
|
||||||
|
|
||||||
/**
|
|
||||||
* Failure to find a stream to the peer. This may be because the connection with the peer is not still alive.
|
|
||||||
* Mitigation can be: retrying after a given time period, or mitigation for `NO_PEER_AVAILABLE` can be used.
|
|
||||||
*/
|
|
||||||
NO_STREAM_AVAILABLE = "No stream available",
|
NO_STREAM_AVAILABLE = "No stream available",
|
||||||
|
|
||||||
/**
|
|
||||||
* The remote peer did not behave as expected. Mitigation for `NO_PEER_AVAILABLE`
|
|
||||||
* or `DECODE_FAILED` can be used.
|
|
||||||
*/
|
|
||||||
NO_RESPONSE = "No response received",
|
NO_RESPONSE = "No response received",
|
||||||
|
|
||||||
//
|
|
||||||
// SEND ERRORS SECTION
|
|
||||||
//
|
|
||||||
/**
|
|
||||||
* Failure to protobuf encode the message. This is not recoverable and needs
|
|
||||||
* further investigation.
|
|
||||||
*/
|
|
||||||
ENCODE_FAILED = "Failed to encode",
|
|
||||||
|
|
||||||
/**
|
|
||||||
* The message payload is empty, making the message invalid. Ensure that a non-empty
|
|
||||||
* payload is set on the outgoing message.
|
|
||||||
*/
|
|
||||||
EMPTY_PAYLOAD = "Payload is empty",
|
|
||||||
|
|
||||||
/**
|
|
||||||
* The message size is above the maximum message size allowed on the Waku Network.
|
|
||||||
* Compressing the message or using an alternative strategy for large messages is recommended.
|
|
||||||
*/
|
|
||||||
SIZE_TOO_BIG = "Size is too big",
|
|
||||||
|
|
||||||
/**
|
|
||||||
* The PubsubTopic passed to the send function is not configured on the Waku node.
|
|
||||||
* Please ensure that the PubsubTopic is used when initializing the Waku node.
|
|
||||||
*/
|
|
||||||
TOPIC_NOT_CONFIGURED = "Topic not configured",
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Fails when
|
|
||||||
*/
|
|
||||||
STREAM_ABORTED = "Stream aborted",
|
STREAM_ABORTED = "Stream aborted",
|
||||||
|
|
||||||
/**
|
// LightPush specific errors
|
||||||
* General proof generation error message.
|
ENCODE_FAILED = "Failed to encode",
|
||||||
* nwaku: https://github.com/waku-org/nwaku/blob/c3cb06ac6c03f0f382d3941ea53b330f6a8dd127/waku/waku_rln_relay/group_manager/group_manager_base.nim#L201C19-L201C42
|
EMPTY_PAYLOAD = "Payload is empty",
|
||||||
*/
|
SIZE_TOO_BIG = "Size is too big",
|
||||||
|
TOPIC_NOT_CONFIGURED = "Topic not configured",
|
||||||
RLN_PROOF_GENERATION = "Proof generation failed",
|
RLN_PROOF_GENERATION = "Proof generation failed",
|
||||||
|
REMOTE_PEER_REJECTED = "Remote peer rejected",
|
||||||
|
|
||||||
//
|
// Status code based errors
|
||||||
// RECEIVE ERRORS SECTION
|
BAD_REQUEST = "Bad request format",
|
||||||
//
|
PAYLOAD_TOO_LARGE = "Message payload exceeds maximum size",
|
||||||
/**
|
INVALID_MESSAGE = "Message validation failed",
|
||||||
* The pubsub topic configured on the decoder does not match the pubsub topic setup on the protocol.
|
UNSUPPORTED_TOPIC = "Unsupported pubsub topic",
|
||||||
* Ensure that the pubsub topic used for decoder creation is the same as the one used for protocol.
|
TOO_MANY_REQUESTS = "Rate limit exceeded",
|
||||||
*/
|
INTERNAL_ERROR = "Internal server error",
|
||||||
|
UNAVAILABLE = "Service temporarily unavailable",
|
||||||
|
NO_RLN_PROOF = "RLN proof generation failed",
|
||||||
|
NO_PEERS = "No relay peers available"
|
||||||
|
}
|
||||||
|
|
||||||
|
// Filter specific errors
|
||||||
|
export enum FilterError {
|
||||||
|
// General errors
|
||||||
|
GENERIC_FAIL = "Generic error",
|
||||||
|
DECODE_FAILED = "Failed to decode",
|
||||||
|
NO_PEER_AVAILABLE = "No peer available",
|
||||||
|
NO_STREAM_AVAILABLE = "No stream available",
|
||||||
|
NO_RESPONSE = "No response received",
|
||||||
|
STREAM_ABORTED = "Stream aborted",
|
||||||
|
|
||||||
|
// Filter specific errors
|
||||||
|
REMOTE_PEER_REJECTED = "Remote peer rejected",
|
||||||
|
TOPIC_NOT_CONFIGURED = "Topic not configured",
|
||||||
|
SUBSCRIPTION_FAILED = "Subscription failed",
|
||||||
|
UNSUBSCRIBE_FAILED = "Unsubscribe failed",
|
||||||
|
PING_FAILED = "Ping failed",
|
||||||
TOPIC_DECODER_MISMATCH = "Topic decoder mismatch",
|
TOPIC_DECODER_MISMATCH = "Topic decoder mismatch",
|
||||||
|
INVALID_DECODER_TOPICS = "Invalid decoder topics",
|
||||||
|
SUBSCRIPTION_LIMIT_EXCEEDED = "Subscription limit exceeded",
|
||||||
|
INVALID_CONTENT_TOPIC = "Invalid content topic",
|
||||||
|
PUSH_MESSAGE_FAILED = "Push message failed",
|
||||||
|
EMPTY_MESSAGE = "Empty message received",
|
||||||
|
MISSING_PUBSUB_TOPIC = "Pubsub topic missing from push message"
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
// Protocol-specific failure interfaces
|
||||||
* The topics passed in the decoders do not match each other, or don't exist at all.
|
export interface LightPushFailure {
|
||||||
* Ensure that all the pubsub topics used in the decoders are valid and match each other.
|
error: LightPushError;
|
||||||
*/
|
peerId?: PeerId;
|
||||||
|
statusCode?: number;
|
||||||
|
statusDesc?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface FilterFailure {
|
||||||
|
error: FilterError;
|
||||||
|
peerId?: PeerId;
|
||||||
|
statusCode?: number;
|
||||||
|
statusDesc?: string;
|
||||||
|
requestId?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Protocol-specific result types
|
||||||
|
export type LightPushCoreResult = ThisOrThat<
|
||||||
|
"success",
|
||||||
|
PeerId,
|
||||||
|
"failure",
|
||||||
|
LightPushFailure
|
||||||
|
>;
|
||||||
|
|
||||||
|
export type FilterCoreResult = ThisOrThat<
|
||||||
|
"success",
|
||||||
|
PeerId,
|
||||||
|
"failure",
|
||||||
|
FilterFailure
|
||||||
|
>;
|
||||||
|
|
||||||
|
export type LightPushSDKResult = ThisAndThat<
|
||||||
|
"successes",
|
||||||
|
PeerId[],
|
||||||
|
"failures",
|
||||||
|
LightPushFailure[]
|
||||||
|
>;
|
||||||
|
|
||||||
|
export type FilterSDKResult = ThisAndThat<
|
||||||
|
"successes",
|
||||||
|
PeerId[],
|
||||||
|
"failures",
|
||||||
|
FilterFailure[]
|
||||||
|
>;
|
||||||
|
|
||||||
|
// Legacy types for backward compatibility (to be deprecated)
|
||||||
|
/**
|
||||||
|
* @deprecated Use LightPushError or FilterError instead
|
||||||
|
*/
|
||||||
|
export enum ProtocolError {
|
||||||
|
GENERIC_FAIL = "Generic error",
|
||||||
|
REMOTE_PEER_REJECTED = "Remote peer rejected",
|
||||||
|
DECODE_FAILED = "Failed to decode",
|
||||||
|
NO_PEER_AVAILABLE = "No peer available",
|
||||||
|
NO_STREAM_AVAILABLE = "No stream available",
|
||||||
|
NO_RESPONSE = "No response received",
|
||||||
|
ENCODE_FAILED = "Failed to encode",
|
||||||
|
EMPTY_PAYLOAD = "Payload is empty",
|
||||||
|
SIZE_TOO_BIG = "Size is too big",
|
||||||
|
TOPIC_NOT_CONFIGURED = "Topic not configured",
|
||||||
|
STREAM_ABORTED = "Stream aborted",
|
||||||
|
RLN_PROOF_GENERATION = "Proof generation failed",
|
||||||
|
TOPIC_DECODER_MISMATCH = "Topic decoder mismatch",
|
||||||
INVALID_DECODER_TOPICS = "Invalid decoder topics"
|
INVALID_DECODER_TOPICS = "Invalid decoder topics"
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @deprecated Use LightPushFailure or FilterFailure instead
|
||||||
|
*/
|
||||||
export interface Failure {
|
export interface Failure {
|
||||||
error: ProtocolError;
|
error: ProtocolError;
|
||||||
peerId?: PeerId;
|
peerId?: PeerId;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @deprecated Use LightPushCoreResult or FilterCoreResult instead
|
||||||
|
*/
|
||||||
export type CoreProtocolResult = ThisOrThat<
|
export type CoreProtocolResult = ThisOrThat<
|
||||||
"success",
|
"success",
|
||||||
PeerId,
|
PeerId,
|
||||||
@ -235,6 +266,9 @@ export type CoreProtocolResult = ThisOrThat<
|
|||||||
Failure
|
Failure
|
||||||
>;
|
>;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @deprecated Use LightPushSDKResult or FilterSDKResult instead
|
||||||
|
*/
|
||||||
export type SDKProtocolResult = ThisAndThat<
|
export type SDKProtocolResult = ThisAndThat<
|
||||||
"successes",
|
"successes",
|
||||||
PeerId[],
|
PeerId[],
|
||||||
|
|||||||
@ -1,5 +1,5 @@
|
|||||||
import type { IEncoder, IMessage } from "./message.js";
|
import type { IEncoder, IMessage } from "./message.js";
|
||||||
import { SDKProtocolResult } from "./protocols.js";
|
import { LightPushSDKResult, SDKProtocolResult } from "./protocols.js";
|
||||||
|
|
||||||
export type ISendOptions = {
|
export type ISendOptions = {
|
||||||
/**
|
/**
|
||||||
@ -22,3 +22,11 @@ export interface ISender {
|
|||||||
sendOptions?: ISendOptions
|
sendOptions?: ISendOptions
|
||||||
) => Promise<SDKProtocolResult>;
|
) => Promise<SDKProtocolResult>;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export interface ILightPushSender {
|
||||||
|
send: (
|
||||||
|
encoder: IEncoder,
|
||||||
|
message: IMessage,
|
||||||
|
sendOptions?: ISendOptions
|
||||||
|
) => Promise<LightPushSDKResult>;
|
||||||
|
}
|
||||||
|
|||||||
@ -3,9 +3,17 @@ import {
|
|||||||
ConnectionManager,
|
ConnectionManager,
|
||||||
createEncoder,
|
createEncoder,
|
||||||
Encoder,
|
Encoder,
|
||||||
LightPushCodec
|
LightPushCodec,
|
||||||
|
LightPushCodecV3
|
||||||
} from "@waku/core";
|
} from "@waku/core";
|
||||||
import { Libp2p, ProtocolError } from "@waku/interfaces";
|
import {
|
||||||
|
isSuccess,
|
||||||
|
Libp2p,
|
||||||
|
LightPushError,
|
||||||
|
LightPushStatusCode,
|
||||||
|
ProtocolError,
|
||||||
|
toProtocolError
|
||||||
|
} from "@waku/interfaces";
|
||||||
import { utf8ToBytes } from "@waku/utils/bytes";
|
import { utf8ToBytes } from "@waku/utils/bytes";
|
||||||
import { expect } from "chai";
|
import { expect } from "chai";
|
||||||
import sinon, { SinonSpy } from "sinon";
|
import sinon, { SinonSpy } from "sinon";
|
||||||
@ -37,8 +45,9 @@ describe("LightPush SDK", () => {
|
|||||||
const failures = result.failures ?? [];
|
const failures = result.failures ?? [];
|
||||||
|
|
||||||
expect(failures.length).to.be.eq(1);
|
expect(failures.length).to.be.eq(1);
|
||||||
expect(failures.some((v) => v.error === ProtocolError.TOPIC_NOT_CONFIGURED))
|
expect(
|
||||||
.to.be.true;
|
failures.some((v) => v.error === LightPushError.TOPIC_NOT_CONFIGURED)
|
||||||
|
).to.be.true;
|
||||||
});
|
});
|
||||||
|
|
||||||
it("should fail to send if no connected peers found", async () => {
|
it("should fail to send if no connected peers found", async () => {
|
||||||
@ -48,8 +57,8 @@ describe("LightPush SDK", () => {
|
|||||||
const failures = result.failures ?? [];
|
const failures = result.failures ?? [];
|
||||||
|
|
||||||
expect(failures.length).to.be.eq(1);
|
expect(failures.length).to.be.eq(1);
|
||||||
expect(failures.some((v) => v.error === ProtocolError.NO_PEER_AVAILABLE)).to
|
expect(failures.some((v) => v.error === LightPushError.NO_PEER_AVAILABLE))
|
||||||
.be.true;
|
.to.be.true;
|
||||||
});
|
});
|
||||||
|
|
||||||
it("should send to specified number of peers of used peers", async () => {
|
it("should send to specified number of peers of used peers", async () => {
|
||||||
@ -135,6 +144,72 @@ describe("LightPush SDK", () => {
|
|||||||
expect(result.successes?.length).to.be.eq(1);
|
expect(result.successes?.length).to.be.eq(1);
|
||||||
expect(result.failures?.length).to.be.eq(1);
|
expect(result.failures?.length).to.be.eq(1);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("v3 protocol support", () => {
|
||||||
|
it("should work with v3 peers", async () => {
|
||||||
|
libp2p = mockLibp2p({
|
||||||
|
peers: [mockV3Peer("1"), mockV3Peer("2")]
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(isSuccess(LightPushStatusCode.SUCCESS)).to.be.true;
|
||||||
|
expect(isSuccess(LightPushStatusCode.BAD_REQUEST)).to.be.false;
|
||||||
|
expect(toProtocolError(LightPushStatusCode.PAYLOAD_TOO_LARGE)).to.eq(
|
||||||
|
ProtocolError.SIZE_TOO_BIG
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("should work with mixed v2 and v3 peers", async () => {
|
||||||
|
libp2p = mockLibp2p({
|
||||||
|
peers: [mockV2AndV3Peer("1"), mockPeer("2"), mockV3Peer("3")]
|
||||||
|
});
|
||||||
|
|
||||||
|
// Mock responses for different protocol versions
|
||||||
|
const v3Response = mockV3SuccessResponse(5);
|
||||||
|
const v2Response = mockV2SuccessResponse();
|
||||||
|
const v3ErrorResponse = mockV3ErrorResponse(
|
||||||
|
LightPushStatusCode.PAYLOAD_TOO_LARGE
|
||||||
|
);
|
||||||
|
const v2ErrorResponse = mockV2ErrorResponse("Message too large");
|
||||||
|
|
||||||
|
expect(v3Response.statusCode).to.eq(LightPushStatusCode.SUCCESS);
|
||||||
|
expect(v3Response.relayPeerCount).to.eq(5);
|
||||||
|
expect(v2Response.isSuccess).to.be.true;
|
||||||
|
expect(v3ErrorResponse.statusCode).to.eq(
|
||||||
|
LightPushStatusCode.PAYLOAD_TOO_LARGE
|
||||||
|
);
|
||||||
|
expect(v2ErrorResponse.isSuccess).to.be.false;
|
||||||
|
});
|
||||||
|
|
||||||
|
it("should handle v3 RLN errors", async () => {
|
||||||
|
const v3RLNError = mockV3RLNErrorResponse();
|
||||||
|
const v2RLNError = mockV2RLNErrorResponse();
|
||||||
|
|
||||||
|
expect(v3RLNError.statusCode).to.eq(LightPushStatusCode.NO_RLN_PROOF);
|
||||||
|
expect(v3RLNError.statusDesc).to.include("RLN proof generation failed");
|
||||||
|
expect(v2RLNError.info).to.include("RLN proof generation failed");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("should validate status codes", async () => {
|
||||||
|
const statusCodes = [
|
||||||
|
LightPushStatusCode.SUCCESS,
|
||||||
|
LightPushStatusCode.BAD_REQUEST,
|
||||||
|
LightPushStatusCode.PAYLOAD_TOO_LARGE,
|
||||||
|
LightPushStatusCode.INVALID_MESSAGE,
|
||||||
|
LightPushStatusCode.UNSUPPORTED_TOPIC,
|
||||||
|
LightPushStatusCode.TOO_MANY_REQUESTS,
|
||||||
|
LightPushStatusCode.INTERNAL_ERROR,
|
||||||
|
LightPushStatusCode.UNAVAILABLE,
|
||||||
|
LightPushStatusCode.NO_RLN_PROOF,
|
||||||
|
LightPushStatusCode.NO_PEERS
|
||||||
|
];
|
||||||
|
|
||||||
|
statusCodes.forEach((code) => {
|
||||||
|
const protocolError = toProtocolError(code);
|
||||||
|
expect(protocolError).to.be.a("string");
|
||||||
|
expect(Object.values(ProtocolError)).to.include(protocolError);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
type MockLibp2pOptions = {
|
type MockLibp2pOptions = {
|
||||||
@ -144,7 +219,16 @@ type MockLibp2pOptions = {
|
|||||||
function mockLibp2p(options?: MockLibp2pOptions): Libp2p {
|
function mockLibp2p(options?: MockLibp2pOptions): Libp2p {
|
||||||
const peers = options?.peers || [];
|
const peers = options?.peers || [];
|
||||||
const peerStore = {
|
const peerStore = {
|
||||||
get: (id: any) => Promise.resolve(peers.find((p) => p.id === id))
|
get: (id: any) => {
|
||||||
|
const peer = peers.find((p) => p.id === id);
|
||||||
|
if (peer) {
|
||||||
|
return Promise.resolve({
|
||||||
|
...peer,
|
||||||
|
protocols: peer.protocols || [LightPushCodec]
|
||||||
|
});
|
||||||
|
}
|
||||||
|
return Promise.resolve(undefined);
|
||||||
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
return {
|
return {
|
||||||
@ -191,9 +275,89 @@ function mockLightPush(options: MockLightPushOptions): LightPush {
|
|||||||
return lightPush;
|
return lightPush;
|
||||||
}
|
}
|
||||||
|
|
||||||
function mockPeer(id: string): Peer {
|
function mockPeer(id: string, protocols: string[] = [LightPushCodec]): Peer {
|
||||||
return {
|
return {
|
||||||
id,
|
id,
|
||||||
protocols: [LightPushCodec]
|
protocols
|
||||||
} as unknown as Peer;
|
} as unknown as Peer;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// V3-specific mock functions
|
||||||
|
function mockV3Peer(id: string): Peer {
|
||||||
|
return mockPeer(id, [LightPushCodecV3]);
|
||||||
|
}
|
||||||
|
|
||||||
|
function mockV2AndV3Peer(id: string): Peer {
|
||||||
|
return mockPeer(id, [LightPushCodec, LightPushCodecV3]);
|
||||||
|
}
|
||||||
|
|
||||||
|
function mockV3SuccessResponse(relayPeerCount?: number): {
|
||||||
|
statusCode: LightPushStatusCode;
|
||||||
|
statusDesc: string;
|
||||||
|
relayPeerCount?: number;
|
||||||
|
isSuccess: boolean;
|
||||||
|
} {
|
||||||
|
return {
|
||||||
|
statusCode: LightPushStatusCode.SUCCESS,
|
||||||
|
statusDesc: "Message sent successfully",
|
||||||
|
relayPeerCount,
|
||||||
|
isSuccess: true
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function mockV3ErrorResponse(
|
||||||
|
statusCode: LightPushStatusCode,
|
||||||
|
statusDesc?: string
|
||||||
|
): {
|
||||||
|
statusCode: LightPushStatusCode;
|
||||||
|
statusDesc: string;
|
||||||
|
isSuccess: boolean;
|
||||||
|
} {
|
||||||
|
return {
|
||||||
|
statusCode,
|
||||||
|
statusDesc: statusDesc || "Error occurred",
|
||||||
|
isSuccess: false
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function mockV2SuccessResponse(): {
|
||||||
|
isSuccess: boolean;
|
||||||
|
info: string;
|
||||||
|
} {
|
||||||
|
return {
|
||||||
|
isSuccess: true,
|
||||||
|
info: "Message sent successfully"
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function mockV2ErrorResponse(info?: string): {
|
||||||
|
isSuccess: boolean;
|
||||||
|
info: string;
|
||||||
|
} {
|
||||||
|
return {
|
||||||
|
isSuccess: false,
|
||||||
|
info: info || "Error occurred"
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function mockV3RLNErrorResponse(): {
|
||||||
|
statusCode: LightPushStatusCode;
|
||||||
|
statusDesc: string;
|
||||||
|
isSuccess: boolean;
|
||||||
|
} {
|
||||||
|
return {
|
||||||
|
statusCode: LightPushStatusCode.NO_RLN_PROOF,
|
||||||
|
statusDesc: "RLN proof generation failed",
|
||||||
|
isSuccess: false
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function mockV2RLNErrorResponse(): {
|
||||||
|
isSuccess: boolean;
|
||||||
|
info: string;
|
||||||
|
} {
|
||||||
|
return {
|
||||||
|
isSuccess: false,
|
||||||
|
info: "RLN proof generation failed"
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|||||||
@ -1,17 +1,17 @@
|
|||||||
import type { PeerId } from "@libp2p/interface";
|
import type { PeerId } from "@libp2p/interface";
|
||||||
import { ConnectionManager, LightPushCore } from "@waku/core";
|
import { ConnectionManager, LightPushCore } from "@waku/core";
|
||||||
import {
|
import {
|
||||||
type CoreProtocolResult,
|
|
||||||
Failure,
|
|
||||||
type IEncoder,
|
type IEncoder,
|
||||||
ILightPush,
|
ILightPush,
|
||||||
type IMessage,
|
type IMessage,
|
||||||
type ISendOptions,
|
type ISendOptions,
|
||||||
type Libp2p,
|
type Libp2p,
|
||||||
|
type LightPushCoreResult,
|
||||||
|
LightPushError,
|
||||||
|
type LightPushFailure,
|
||||||
LightPushProtocolOptions,
|
LightPushProtocolOptions,
|
||||||
ProtocolError,
|
type LightPushSDKResult,
|
||||||
Protocols,
|
Protocols
|
||||||
SDKProtocolResult
|
|
||||||
} from "@waku/interfaces";
|
} from "@waku/interfaces";
|
||||||
import { Logger } from "@waku/utils";
|
import { Logger } from "@waku/utils";
|
||||||
|
|
||||||
@ -74,7 +74,7 @@ export class LightPush implements ILightPush {
|
|||||||
encoder: IEncoder,
|
encoder: IEncoder,
|
||||||
message: IMessage,
|
message: IMessage,
|
||||||
options: ISendOptions = {}
|
options: ISendOptions = {}
|
||||||
): Promise<SDKProtocolResult> {
|
): Promise<LightPushSDKResult> {
|
||||||
options = {
|
options = {
|
||||||
...this.config,
|
...this.config,
|
||||||
...options
|
...options
|
||||||
@ -89,7 +89,7 @@ export class LightPush implements ILightPush {
|
|||||||
successes: [],
|
successes: [],
|
||||||
failures: [
|
failures: [
|
||||||
{
|
{
|
||||||
error: ProtocolError.TOPIC_NOT_CONFIGURED
|
error: LightPushError.TOPIC_NOT_CONFIGURED
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
};
|
};
|
||||||
@ -100,40 +100,40 @@ export class LightPush implements ILightPush {
|
|||||||
pubsubTopic: encoder.pubsubTopic
|
pubsubTopic: encoder.pubsubTopic
|
||||||
});
|
});
|
||||||
|
|
||||||
const coreResults: CoreProtocolResult[] =
|
const coreResults: LightPushCoreResult[] =
|
||||||
peerIds?.length > 0
|
peerIds?.length > 0
|
||||||
? await Promise.all(
|
? await Promise.all(
|
||||||
peerIds.map((peerId) =>
|
peerIds.map((peerId) =>
|
||||||
this.protocol.send(encoder, message, peerId).catch((_e) => ({
|
this.protocol.send(encoder, message, peerId).catch((_e) => ({
|
||||||
success: null,
|
success: null,
|
||||||
failure: {
|
failure: {
|
||||||
error: ProtocolError.GENERIC_FAIL
|
error: LightPushError.GENERIC_FAIL
|
||||||
}
|
}
|
||||||
}))
|
}))
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
: [];
|
: [];
|
||||||
|
|
||||||
const results: SDKProtocolResult = coreResults.length
|
const results: LightPushSDKResult = coreResults.length
|
||||||
? {
|
? {
|
||||||
successes: coreResults
|
successes: coreResults
|
||||||
.filter((v) => v.success)
|
.filter((v) => v.success)
|
||||||
.map((v) => v.success) as PeerId[],
|
.map((v) => v.success) as PeerId[],
|
||||||
failures: coreResults
|
failures: coreResults
|
||||||
.filter((v) => v.failure)
|
.filter((v) => v.failure)
|
||||||
.map((v) => v.failure) as Failure[]
|
.map((v) => v.failure) as LightPushFailure[]
|
||||||
}
|
}
|
||||||
: {
|
: {
|
||||||
successes: [],
|
successes: [],
|
||||||
failures: [
|
failures: [
|
||||||
{
|
{
|
||||||
error: ProtocolError.NO_PEER_AVAILABLE
|
error: LightPushError.NO_PEER_AVAILABLE
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
};
|
};
|
||||||
|
|
||||||
if (options.autoRetry && results.successes.length === 0) {
|
if (options.autoRetry && results.successes.length === 0) {
|
||||||
const sendCallback = (peerId: PeerId): Promise<CoreProtocolResult> =>
|
const sendCallback = (peerId: PeerId): Promise<LightPushCoreResult> =>
|
||||||
this.protocol.send(encoder, message, peerId);
|
this.protocol.send(encoder, message, peerId);
|
||||||
this.retryManager.push(
|
this.retryManager.push(
|
||||||
sendCallback.bind(this),
|
sendCallback.bind(this),
|
||||||
|
|||||||
@ -1,6 +1,7 @@
|
|||||||
import type { PeerId } from "@libp2p/interface";
|
import type { PeerId } from "@libp2p/interface";
|
||||||
import {
|
import {
|
||||||
type CoreProtocolResult,
|
type LightPushCoreResult,
|
||||||
|
LightPushError,
|
||||||
ProtocolError,
|
ProtocolError,
|
||||||
Protocols
|
Protocols
|
||||||
} from "@waku/interfaces";
|
} from "@waku/interfaces";
|
||||||
@ -53,7 +54,7 @@ describe("RetryManager", () => {
|
|||||||
|
|
||||||
it("should process tasks in queue", async () => {
|
it("should process tasks in queue", async () => {
|
||||||
const successCallback = sinon.spy(
|
const successCallback = sinon.spy(
|
||||||
async (peerId: PeerId): Promise<CoreProtocolResult> => ({
|
async (peerId: PeerId): Promise<LightPushCoreResult> => ({
|
||||||
success: peerId,
|
success: peerId,
|
||||||
failure: null
|
failure: null
|
||||||
})
|
})
|
||||||
@ -106,9 +107,9 @@ describe("RetryManager", () => {
|
|||||||
|
|
||||||
it("should retry failed tasks", async () => {
|
it("should retry failed tasks", async () => {
|
||||||
const failingCallback = sinon.spy(
|
const failingCallback = sinon.spy(
|
||||||
async (): Promise<CoreProtocolResult> => ({
|
async (): Promise<LightPushCoreResult> => ({
|
||||||
success: null,
|
success: null,
|
||||||
failure: { error: "test error" as any }
|
failure: { error: LightPushError.GENERIC_FAIL }
|
||||||
})
|
})
|
||||||
);
|
);
|
||||||
|
|
||||||
@ -129,7 +130,7 @@ describe("RetryManager", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
it("should request peer renewal on specific errors", async () => {
|
it("should request peer renewal on specific errors", async () => {
|
||||||
const errorCallback = sinon.spy(async (): Promise<CoreProtocolResult> => {
|
const errorCallback = sinon.spy(async (): Promise<LightPushCoreResult> => {
|
||||||
throw new Error(ProtocolError.NO_PEER_AVAILABLE);
|
throw new Error(ProtocolError.NO_PEER_AVAILABLE);
|
||||||
});
|
});
|
||||||
|
|
||||||
@ -149,7 +150,7 @@ describe("RetryManager", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
it("should handle task timeouts", async () => {
|
it("should handle task timeouts", async () => {
|
||||||
const slowCallback = sinon.spy(async (): Promise<CoreProtocolResult> => {
|
const slowCallback = sinon.spy(async (): Promise<LightPushCoreResult> => {
|
||||||
await new Promise((resolve) => setTimeout(resolve, 15000));
|
await new Promise((resolve) => setTimeout(resolve, 15000));
|
||||||
return { success: mockPeerId, failure: null };
|
return { success: mockPeerId, failure: null };
|
||||||
});
|
});
|
||||||
@ -168,9 +169,11 @@ describe("RetryManager", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
it("should not execute task if max attempts is 0", async () => {
|
it("should not execute task if max attempts is 0", async () => {
|
||||||
const failingCallback = sinon.spy(async (): Promise<CoreProtocolResult> => {
|
const failingCallback = sinon.spy(
|
||||||
throw new Error("test error" as any);
|
async (): Promise<LightPushCoreResult> => {
|
||||||
});
|
throw new Error("test error" as any);
|
||||||
|
}
|
||||||
|
);
|
||||||
|
|
||||||
const task = {
|
const task = {
|
||||||
callback: failingCallback,
|
callback: failingCallback,
|
||||||
@ -203,7 +206,7 @@ describe("RetryManager", () => {
|
|||||||
called++;
|
called++;
|
||||||
return Promise.resolve({
|
return Promise.resolve({
|
||||||
success: null,
|
success: null,
|
||||||
failure: { error: ProtocolError.GENERIC_FAIL }
|
failure: { error: LightPushError.GENERIC_FAIL }
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
retryManager.push(failCallback, 2, "test-topic");
|
retryManager.push(failCallback, 2, "test-topic");
|
||||||
|
|||||||
@ -1,5 +1,5 @@
|
|||||||
import type { PeerId } from "@libp2p/interface";
|
import type { PeerId } from "@libp2p/interface";
|
||||||
import { type CoreProtocolResult, Protocols } from "@waku/interfaces";
|
import { type LightPushCoreResult, Protocols } from "@waku/interfaces";
|
||||||
import { Logger } from "@waku/utils";
|
import { Logger } from "@waku/utils";
|
||||||
|
|
||||||
import type { PeerManager } from "../peer_manager/index.js";
|
import type { PeerManager } from "../peer_manager/index.js";
|
||||||
@ -11,7 +11,7 @@ type RetryManagerConfig = {
|
|||||||
peerManager: PeerManager;
|
peerManager: PeerManager;
|
||||||
};
|
};
|
||||||
|
|
||||||
type AttemptCallback = (peerId: PeerId) => Promise<CoreProtocolResult>;
|
type AttemptCallback = (peerId: PeerId) => Promise<LightPushCoreResult>;
|
||||||
|
|
||||||
export type ScheduledTask = {
|
export type ScheduledTask = {
|
||||||
maxAttempts: number;
|
maxAttempts: number;
|
||||||
@ -119,7 +119,13 @@ export class RetryManager {
|
|||||||
task.callback(peerId)
|
task.callback(peerId)
|
||||||
]);
|
]);
|
||||||
|
|
||||||
if (response?.failure) {
|
// If timeout resolves first, response will be void (undefined)
|
||||||
|
// In this case, we should treat it as a timeout error
|
||||||
|
if (response === undefined) {
|
||||||
|
throw new Error("Task timeout");
|
||||||
|
}
|
||||||
|
|
||||||
|
if (response.failure) {
|
||||||
throw Error(response.failure.error);
|
throw Error(response.failure.error);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user