diff --git a/generated/fila/v1/AckError.ts b/generated/fila/v1/AckError.ts new file mode 100644 index 0000000..bb2f56b --- /dev/null +++ b/generated/fila/v1/AckError.ts @@ -0,0 +1,13 @@ +// Original file: proto/fila/v1/service.proto + +import type { AckErrorCode as _fila_v1_AckErrorCode, AckErrorCode__Output as _fila_v1_AckErrorCode__Output } from '../../fila/v1/AckErrorCode'; + +export interface AckError { + 'code'?: (_fila_v1_AckErrorCode); + 'message'?: (string); +} + +export interface AckError__Output { + 'code': (_fila_v1_AckErrorCode__Output); + 'message': (string); +} diff --git a/generated/fila/v1/AckErrorCode.ts b/generated/fila/v1/AckErrorCode.ts new file mode 100644 index 0000000..04a2113 --- /dev/null +++ b/generated/fila/v1/AckErrorCode.ts @@ -0,0 +1,20 @@ +// Original file: proto/fila/v1/service.proto + +export const AckErrorCode = { + ACK_ERROR_CODE_UNSPECIFIED: 'ACK_ERROR_CODE_UNSPECIFIED', + ACK_ERROR_CODE_MESSAGE_NOT_FOUND: 'ACK_ERROR_CODE_MESSAGE_NOT_FOUND', + ACK_ERROR_CODE_STORAGE: 'ACK_ERROR_CODE_STORAGE', + ACK_ERROR_CODE_PERMISSION_DENIED: 'ACK_ERROR_CODE_PERMISSION_DENIED', +} as const; + +export type AckErrorCode = + | 'ACK_ERROR_CODE_UNSPECIFIED' + | 0 + | 'ACK_ERROR_CODE_MESSAGE_NOT_FOUND' + | 1 + | 'ACK_ERROR_CODE_STORAGE' + | 2 + | 'ACK_ERROR_CODE_PERMISSION_DENIED' + | 3 + +export type AckErrorCode__Output = typeof AckErrorCode[keyof typeof AckErrorCode] diff --git a/generated/fila/v1/AckMessage.ts b/generated/fila/v1/AckMessage.ts new file mode 100644 index 0000000..dac5ff4 --- /dev/null +++ b/generated/fila/v1/AckMessage.ts @@ -0,0 +1,12 @@ +// Original file: proto/fila/v1/service.proto + + +export interface AckMessage { + 'queue'?: (string); + 'messageId'?: (string); +} + +export interface AckMessage__Output { + 'queue': (string); + 'messageId': (string); +} diff --git a/generated/fila/v1/AckRequest.ts b/generated/fila/v1/AckRequest.ts index 53f7b92..3c63ce1 100644 --- a/generated/fila/v1/AckRequest.ts +++ b/generated/fila/v1/AckRequest.ts @@ -1,12 +1,11 @@ // Original file: proto/fila/v1/service.proto +import type { AckMessage as _fila_v1_AckMessage, AckMessage__Output as _fila_v1_AckMessage__Output } from '../../fila/v1/AckMessage'; export interface AckRequest { - 'queue'?: (string); - 'messageId'?: (string); + 'messages'?: (_fila_v1_AckMessage)[]; } export interface AckRequest__Output { - 'queue': (string); - 'messageId': (string); + 'messages': (_fila_v1_AckMessage__Output)[]; } diff --git a/generated/fila/v1/AckResponse.ts b/generated/fila/v1/AckResponse.ts index 9829f92..eae1240 100644 --- a/generated/fila/v1/AckResponse.ts +++ b/generated/fila/v1/AckResponse.ts @@ -1,8 +1,11 @@ // Original file: proto/fila/v1/service.proto +import type { AckResult as _fila_v1_AckResult, AckResult__Output as _fila_v1_AckResult__Output } from '../../fila/v1/AckResult'; export interface AckResponse { + 'results'?: (_fila_v1_AckResult)[]; } export interface AckResponse__Output { + 'results': (_fila_v1_AckResult__Output)[]; } diff --git a/generated/fila/v1/AckResult.ts b/generated/fila/v1/AckResult.ts new file mode 100644 index 0000000..d503b9a --- /dev/null +++ b/generated/fila/v1/AckResult.ts @@ -0,0 +1,16 @@ +// Original file: proto/fila/v1/service.proto + +import type { AckSuccess as _fila_v1_AckSuccess, AckSuccess__Output as _fila_v1_AckSuccess__Output } from '../../fila/v1/AckSuccess'; +import type { AckError as _fila_v1_AckError, AckError__Output as _fila_v1_AckError__Output } from '../../fila/v1/AckError'; + +export interface AckResult { + 'success'?: (_fila_v1_AckSuccess | null); + 'error'?: (_fila_v1_AckError | null); + 'result'?: "success"|"error"; +} + +export interface AckResult__Output { + 'success'?: (_fila_v1_AckSuccess__Output | null); + 'error'?: (_fila_v1_AckError__Output | null); + 'result'?: "success"|"error"; +} diff --git a/generated/fila/v1/AckSuccess.ts b/generated/fila/v1/AckSuccess.ts new file mode 100644 index 0000000..673d632 --- /dev/null +++ b/generated/fila/v1/AckSuccess.ts @@ -0,0 +1,8 @@ +// Original file: proto/fila/v1/service.proto + + +export interface AckSuccess { +} + +export interface AckSuccess__Output { +} diff --git a/generated/fila/v1/BatchEnqueueRequest.ts b/generated/fila/v1/BatchEnqueueRequest.ts deleted file mode 100644 index afd9d70..0000000 --- a/generated/fila/v1/BatchEnqueueRequest.ts +++ /dev/null @@ -1,11 +0,0 @@ -// Original file: proto/fila/v1/service.proto - -import type { EnqueueRequest as _fila_v1_EnqueueRequest, EnqueueRequest__Output as _fila_v1_EnqueueRequest__Output } from '../../fila/v1/EnqueueRequest'; - -export interface BatchEnqueueRequest { - 'messages'?: (_fila_v1_EnqueueRequest)[]; -} - -export interface BatchEnqueueRequest__Output { - 'messages': (_fila_v1_EnqueueRequest__Output)[]; -} diff --git a/generated/fila/v1/BatchEnqueueResponse.ts b/generated/fila/v1/BatchEnqueueResponse.ts deleted file mode 100644 index 0a876cb..0000000 --- a/generated/fila/v1/BatchEnqueueResponse.ts +++ /dev/null @@ -1,11 +0,0 @@ -// Original file: proto/fila/v1/service.proto - -import type { BatchEnqueueResult as _fila_v1_BatchEnqueueResult, BatchEnqueueResult__Output as _fila_v1_BatchEnqueueResult__Output } from '../../fila/v1/BatchEnqueueResult'; - -export interface BatchEnqueueResponse { - 'results'?: (_fila_v1_BatchEnqueueResult)[]; -} - -export interface BatchEnqueueResponse__Output { - 'results': (_fila_v1_BatchEnqueueResult__Output)[]; -} diff --git a/generated/fila/v1/BatchEnqueueResult.ts b/generated/fila/v1/BatchEnqueueResult.ts deleted file mode 100644 index 9b08473..0000000 --- a/generated/fila/v1/BatchEnqueueResult.ts +++ /dev/null @@ -1,15 +0,0 @@ -// Original file: proto/fila/v1/service.proto - -import type { EnqueueResponse as _fila_v1_EnqueueResponse, EnqueueResponse__Output as _fila_v1_EnqueueResponse__Output } from '../../fila/v1/EnqueueResponse'; - -export interface BatchEnqueueResult { - 'success'?: (_fila_v1_EnqueueResponse | null); - 'error'?: (string); - 'result'?: "success"|"error"; -} - -export interface BatchEnqueueResult__Output { - 'success'?: (_fila_v1_EnqueueResponse__Output | null); - 'error'?: (string); - 'result'?: "success"|"error"; -} diff --git a/generated/fila/v1/ConsumeResponse.ts b/generated/fila/v1/ConsumeResponse.ts index 9296e26..bc1f398 100644 --- a/generated/fila/v1/ConsumeResponse.ts +++ b/generated/fila/v1/ConsumeResponse.ts @@ -3,11 +3,9 @@ import type { Message as _fila_v1_Message, Message__Output as _fila_v1_Message__Output } from '../../fila/v1/Message'; export interface ConsumeResponse { - 'message'?: (_fila_v1_Message | null); 'messages'?: (_fila_v1_Message)[]; } export interface ConsumeResponse__Output { - 'message': (_fila_v1_Message__Output | null); 'messages': (_fila_v1_Message__Output)[]; } diff --git a/generated/fila/v1/EnqueueError.ts b/generated/fila/v1/EnqueueError.ts new file mode 100644 index 0000000..82fe40e --- /dev/null +++ b/generated/fila/v1/EnqueueError.ts @@ -0,0 +1,13 @@ +// Original file: proto/fila/v1/service.proto + +import type { EnqueueErrorCode as _fila_v1_EnqueueErrorCode, EnqueueErrorCode__Output as _fila_v1_EnqueueErrorCode__Output } from '../../fila/v1/EnqueueErrorCode'; + +export interface EnqueueError { + 'code'?: (_fila_v1_EnqueueErrorCode); + 'message'?: (string); +} + +export interface EnqueueError__Output { + 'code': (_fila_v1_EnqueueErrorCode__Output); + 'message': (string); +} diff --git a/generated/fila/v1/EnqueueErrorCode.ts b/generated/fila/v1/EnqueueErrorCode.ts new file mode 100644 index 0000000..cc0f38d --- /dev/null +++ b/generated/fila/v1/EnqueueErrorCode.ts @@ -0,0 +1,23 @@ +// Original file: proto/fila/v1/service.proto + +export const EnqueueErrorCode = { + ENQUEUE_ERROR_CODE_UNSPECIFIED: 'ENQUEUE_ERROR_CODE_UNSPECIFIED', + ENQUEUE_ERROR_CODE_QUEUE_NOT_FOUND: 'ENQUEUE_ERROR_CODE_QUEUE_NOT_FOUND', + ENQUEUE_ERROR_CODE_STORAGE: 'ENQUEUE_ERROR_CODE_STORAGE', + ENQUEUE_ERROR_CODE_LUA: 'ENQUEUE_ERROR_CODE_LUA', + ENQUEUE_ERROR_CODE_PERMISSION_DENIED: 'ENQUEUE_ERROR_CODE_PERMISSION_DENIED', +} as const; + +export type EnqueueErrorCode = + | 'ENQUEUE_ERROR_CODE_UNSPECIFIED' + | 0 + | 'ENQUEUE_ERROR_CODE_QUEUE_NOT_FOUND' + | 1 + | 'ENQUEUE_ERROR_CODE_STORAGE' + | 2 + | 'ENQUEUE_ERROR_CODE_LUA' + | 3 + | 'ENQUEUE_ERROR_CODE_PERMISSION_DENIED' + | 4 + +export type EnqueueErrorCode__Output = typeof EnqueueErrorCode[keyof typeof EnqueueErrorCode] diff --git a/generated/fila/v1/EnqueueMessage.ts b/generated/fila/v1/EnqueueMessage.ts new file mode 100644 index 0000000..04b917d --- /dev/null +++ b/generated/fila/v1/EnqueueMessage.ts @@ -0,0 +1,14 @@ +// Original file: proto/fila/v1/service.proto + + +export interface EnqueueMessage { + 'queue'?: (string); + 'headers'?: ({[key: string]: string}); + 'payload'?: (Buffer | Uint8Array | string); +} + +export interface EnqueueMessage__Output { + 'queue': (string); + 'headers': ({[key: string]: string}); + 'payload': (Buffer); +} diff --git a/generated/fila/v1/EnqueueRequest.ts b/generated/fila/v1/EnqueueRequest.ts index 2ccd490..e6f0c8e 100644 --- a/generated/fila/v1/EnqueueRequest.ts +++ b/generated/fila/v1/EnqueueRequest.ts @@ -1,14 +1,11 @@ // Original file: proto/fila/v1/service.proto +import type { EnqueueMessage as _fila_v1_EnqueueMessage, EnqueueMessage__Output as _fila_v1_EnqueueMessage__Output } from '../../fila/v1/EnqueueMessage'; export interface EnqueueRequest { - 'queue'?: (string); - 'headers'?: ({[key: string]: string}); - 'payload'?: (Buffer | Uint8Array | string); + 'messages'?: (_fila_v1_EnqueueMessage)[]; } export interface EnqueueRequest__Output { - 'queue': (string); - 'headers': ({[key: string]: string}); - 'payload': (Buffer); + 'messages': (_fila_v1_EnqueueMessage__Output)[]; } diff --git a/generated/fila/v1/EnqueueResponse.ts b/generated/fila/v1/EnqueueResponse.ts index a44e5ef..271c440 100644 --- a/generated/fila/v1/EnqueueResponse.ts +++ b/generated/fila/v1/EnqueueResponse.ts @@ -1,10 +1,11 @@ // Original file: proto/fila/v1/service.proto +import type { EnqueueResult as _fila_v1_EnqueueResult, EnqueueResult__Output as _fila_v1_EnqueueResult__Output } from '../../fila/v1/EnqueueResult'; export interface EnqueueResponse { - 'messageId'?: (string); + 'results'?: (_fila_v1_EnqueueResult)[]; } export interface EnqueueResponse__Output { - 'messageId': (string); + 'results': (_fila_v1_EnqueueResult__Output)[]; } diff --git a/generated/fila/v1/EnqueueResult.ts b/generated/fila/v1/EnqueueResult.ts new file mode 100644 index 0000000..3691008 --- /dev/null +++ b/generated/fila/v1/EnqueueResult.ts @@ -0,0 +1,15 @@ +// Original file: proto/fila/v1/service.proto + +import type { EnqueueError as _fila_v1_EnqueueError, EnqueueError__Output as _fila_v1_EnqueueError__Output } from '../../fila/v1/EnqueueError'; + +export interface EnqueueResult { + 'messageId'?: (string); + 'error'?: (_fila_v1_EnqueueError | null); + 'result'?: "messageId"|"error"; +} + +export interface EnqueueResult__Output { + 'messageId'?: (string); + 'error'?: (_fila_v1_EnqueueError__Output | null); + 'result'?: "messageId"|"error"; +} diff --git a/generated/fila/v1/FilaService.ts b/generated/fila/v1/FilaService.ts index 9b23bad..868ec1b 100644 --- a/generated/fila/v1/FilaService.ts +++ b/generated/fila/v1/FilaService.ts @@ -4,14 +4,14 @@ import type * as grpc from '@grpc/grpc-js' import type { MethodDefinition } from '@grpc/proto-loader' import type { AckRequest as _fila_v1_AckRequest, AckRequest__Output as _fila_v1_AckRequest__Output } from '../../fila/v1/AckRequest'; import type { AckResponse as _fila_v1_AckResponse, AckResponse__Output as _fila_v1_AckResponse__Output } from '../../fila/v1/AckResponse'; -import type { BatchEnqueueRequest as _fila_v1_BatchEnqueueRequest, BatchEnqueueRequest__Output as _fila_v1_BatchEnqueueRequest__Output } from '../../fila/v1/BatchEnqueueRequest'; -import type { BatchEnqueueResponse as _fila_v1_BatchEnqueueResponse, BatchEnqueueResponse__Output as _fila_v1_BatchEnqueueResponse__Output } from '../../fila/v1/BatchEnqueueResponse'; import type { ConsumeRequest as _fila_v1_ConsumeRequest, ConsumeRequest__Output as _fila_v1_ConsumeRequest__Output } from '../../fila/v1/ConsumeRequest'; import type { ConsumeResponse as _fila_v1_ConsumeResponse, ConsumeResponse__Output as _fila_v1_ConsumeResponse__Output } from '../../fila/v1/ConsumeResponse'; import type { EnqueueRequest as _fila_v1_EnqueueRequest, EnqueueRequest__Output as _fila_v1_EnqueueRequest__Output } from '../../fila/v1/EnqueueRequest'; import type { EnqueueResponse as _fila_v1_EnqueueResponse, EnqueueResponse__Output as _fila_v1_EnqueueResponse__Output } from '../../fila/v1/EnqueueResponse'; import type { NackRequest as _fila_v1_NackRequest, NackRequest__Output as _fila_v1_NackRequest__Output } from '../../fila/v1/NackRequest'; import type { NackResponse as _fila_v1_NackResponse, NackResponse__Output as _fila_v1_NackResponse__Output } from '../../fila/v1/NackResponse'; +import type { StreamEnqueueRequest as _fila_v1_StreamEnqueueRequest, StreamEnqueueRequest__Output as _fila_v1_StreamEnqueueRequest__Output } from '../../fila/v1/StreamEnqueueRequest'; +import type { StreamEnqueueResponse as _fila_v1_StreamEnqueueResponse, StreamEnqueueResponse__Output as _fila_v1_StreamEnqueueResponse__Output } from '../../fila/v1/StreamEnqueueResponse'; export interface FilaServiceClient extends grpc.Client { Ack(argument: _fila_v1_AckRequest, metadata: grpc.Metadata, options: grpc.CallOptions, callback: grpc.requestCallback<_fila_v1_AckResponse__Output>): grpc.ClientUnaryCall; @@ -23,15 +23,6 @@ export interface FilaServiceClient extends grpc.Client { ack(argument: _fila_v1_AckRequest, options: grpc.CallOptions, callback: grpc.requestCallback<_fila_v1_AckResponse__Output>): grpc.ClientUnaryCall; ack(argument: _fila_v1_AckRequest, callback: grpc.requestCallback<_fila_v1_AckResponse__Output>): grpc.ClientUnaryCall; - BatchEnqueue(argument: _fila_v1_BatchEnqueueRequest, metadata: grpc.Metadata, options: grpc.CallOptions, callback: grpc.requestCallback<_fila_v1_BatchEnqueueResponse__Output>): grpc.ClientUnaryCall; - BatchEnqueue(argument: _fila_v1_BatchEnqueueRequest, metadata: grpc.Metadata, callback: grpc.requestCallback<_fila_v1_BatchEnqueueResponse__Output>): grpc.ClientUnaryCall; - BatchEnqueue(argument: _fila_v1_BatchEnqueueRequest, options: grpc.CallOptions, callback: grpc.requestCallback<_fila_v1_BatchEnqueueResponse__Output>): grpc.ClientUnaryCall; - BatchEnqueue(argument: _fila_v1_BatchEnqueueRequest, callback: grpc.requestCallback<_fila_v1_BatchEnqueueResponse__Output>): grpc.ClientUnaryCall; - batchEnqueue(argument: _fila_v1_BatchEnqueueRequest, metadata: grpc.Metadata, options: grpc.CallOptions, callback: grpc.requestCallback<_fila_v1_BatchEnqueueResponse__Output>): grpc.ClientUnaryCall; - batchEnqueue(argument: _fila_v1_BatchEnqueueRequest, metadata: grpc.Metadata, callback: grpc.requestCallback<_fila_v1_BatchEnqueueResponse__Output>): grpc.ClientUnaryCall; - batchEnqueue(argument: _fila_v1_BatchEnqueueRequest, options: grpc.CallOptions, callback: grpc.requestCallback<_fila_v1_BatchEnqueueResponse__Output>): grpc.ClientUnaryCall; - batchEnqueue(argument: _fila_v1_BatchEnqueueRequest, callback: grpc.requestCallback<_fila_v1_BatchEnqueueResponse__Output>): grpc.ClientUnaryCall; - Consume(argument: _fila_v1_ConsumeRequest, metadata: grpc.Metadata, options?: grpc.CallOptions): grpc.ClientReadableStream<_fila_v1_ConsumeResponse__Output>; Consume(argument: _fila_v1_ConsumeRequest, options?: grpc.CallOptions): grpc.ClientReadableStream<_fila_v1_ConsumeResponse__Output>; consume(argument: _fila_v1_ConsumeRequest, metadata: grpc.Metadata, options?: grpc.CallOptions): grpc.ClientReadableStream<_fila_v1_ConsumeResponse__Output>; @@ -55,25 +46,30 @@ export interface FilaServiceClient extends grpc.Client { nack(argument: _fila_v1_NackRequest, options: grpc.CallOptions, callback: grpc.requestCallback<_fila_v1_NackResponse__Output>): grpc.ClientUnaryCall; nack(argument: _fila_v1_NackRequest, callback: grpc.requestCallback<_fila_v1_NackResponse__Output>): grpc.ClientUnaryCall; + StreamEnqueue(metadata: grpc.Metadata, options?: grpc.CallOptions): grpc.ClientDuplexStream<_fila_v1_StreamEnqueueRequest, _fila_v1_StreamEnqueueResponse__Output>; + StreamEnqueue(options?: grpc.CallOptions): grpc.ClientDuplexStream<_fila_v1_StreamEnqueueRequest, _fila_v1_StreamEnqueueResponse__Output>; + streamEnqueue(metadata: grpc.Metadata, options?: grpc.CallOptions): grpc.ClientDuplexStream<_fila_v1_StreamEnqueueRequest, _fila_v1_StreamEnqueueResponse__Output>; + streamEnqueue(options?: grpc.CallOptions): grpc.ClientDuplexStream<_fila_v1_StreamEnqueueRequest, _fila_v1_StreamEnqueueResponse__Output>; + } export interface FilaServiceHandlers extends grpc.UntypedServiceImplementation { Ack: grpc.handleUnaryCall<_fila_v1_AckRequest__Output, _fila_v1_AckResponse>; - BatchEnqueue: grpc.handleUnaryCall<_fila_v1_BatchEnqueueRequest__Output, _fila_v1_BatchEnqueueResponse>; - Consume: grpc.handleServerStreamingCall<_fila_v1_ConsumeRequest__Output, _fila_v1_ConsumeResponse>; Enqueue: grpc.handleUnaryCall<_fila_v1_EnqueueRequest__Output, _fila_v1_EnqueueResponse>; Nack: grpc.handleUnaryCall<_fila_v1_NackRequest__Output, _fila_v1_NackResponse>; + StreamEnqueue: grpc.handleBidiStreamingCall<_fila_v1_StreamEnqueueRequest__Output, _fila_v1_StreamEnqueueResponse>; + } export interface FilaServiceDefinition extends grpc.ServiceDefinition { Ack: MethodDefinition<_fila_v1_AckRequest, _fila_v1_AckResponse, _fila_v1_AckRequest__Output, _fila_v1_AckResponse__Output> - BatchEnqueue: MethodDefinition<_fila_v1_BatchEnqueueRequest, _fila_v1_BatchEnqueueResponse, _fila_v1_BatchEnqueueRequest__Output, _fila_v1_BatchEnqueueResponse__Output> Consume: MethodDefinition<_fila_v1_ConsumeRequest, _fila_v1_ConsumeResponse, _fila_v1_ConsumeRequest__Output, _fila_v1_ConsumeResponse__Output> Enqueue: MethodDefinition<_fila_v1_EnqueueRequest, _fila_v1_EnqueueResponse, _fila_v1_EnqueueRequest__Output, _fila_v1_EnqueueResponse__Output> Nack: MethodDefinition<_fila_v1_NackRequest, _fila_v1_NackResponse, _fila_v1_NackRequest__Output, _fila_v1_NackResponse__Output> + StreamEnqueue: MethodDefinition<_fila_v1_StreamEnqueueRequest, _fila_v1_StreamEnqueueResponse, _fila_v1_StreamEnqueueRequest__Output, _fila_v1_StreamEnqueueResponse__Output> } diff --git a/generated/fila/v1/NackError.ts b/generated/fila/v1/NackError.ts new file mode 100644 index 0000000..2cc2888 --- /dev/null +++ b/generated/fila/v1/NackError.ts @@ -0,0 +1,13 @@ +// Original file: proto/fila/v1/service.proto + +import type { NackErrorCode as _fila_v1_NackErrorCode, NackErrorCode__Output as _fila_v1_NackErrorCode__Output } from '../../fila/v1/NackErrorCode'; + +export interface NackError { + 'code'?: (_fila_v1_NackErrorCode); + 'message'?: (string); +} + +export interface NackError__Output { + 'code': (_fila_v1_NackErrorCode__Output); + 'message': (string); +} diff --git a/generated/fila/v1/NackErrorCode.ts b/generated/fila/v1/NackErrorCode.ts new file mode 100644 index 0000000..e7738f7 --- /dev/null +++ b/generated/fila/v1/NackErrorCode.ts @@ -0,0 +1,20 @@ +// Original file: proto/fila/v1/service.proto + +export const NackErrorCode = { + NACK_ERROR_CODE_UNSPECIFIED: 'NACK_ERROR_CODE_UNSPECIFIED', + NACK_ERROR_CODE_MESSAGE_NOT_FOUND: 'NACK_ERROR_CODE_MESSAGE_NOT_FOUND', + NACK_ERROR_CODE_STORAGE: 'NACK_ERROR_CODE_STORAGE', + NACK_ERROR_CODE_PERMISSION_DENIED: 'NACK_ERROR_CODE_PERMISSION_DENIED', +} as const; + +export type NackErrorCode = + | 'NACK_ERROR_CODE_UNSPECIFIED' + | 0 + | 'NACK_ERROR_CODE_MESSAGE_NOT_FOUND' + | 1 + | 'NACK_ERROR_CODE_STORAGE' + | 2 + | 'NACK_ERROR_CODE_PERMISSION_DENIED' + | 3 + +export type NackErrorCode__Output = typeof NackErrorCode[keyof typeof NackErrorCode] diff --git a/generated/fila/v1/NackMessage.ts b/generated/fila/v1/NackMessage.ts new file mode 100644 index 0000000..2ce0501 --- /dev/null +++ b/generated/fila/v1/NackMessage.ts @@ -0,0 +1,14 @@ +// Original file: proto/fila/v1/service.proto + + +export interface NackMessage { + 'queue'?: (string); + 'messageId'?: (string); + 'error'?: (string); +} + +export interface NackMessage__Output { + 'queue': (string); + 'messageId': (string); + 'error': (string); +} diff --git a/generated/fila/v1/NackRequest.ts b/generated/fila/v1/NackRequest.ts index b7b03f4..2f2450d 100644 --- a/generated/fila/v1/NackRequest.ts +++ b/generated/fila/v1/NackRequest.ts @@ -1,14 +1,11 @@ // Original file: proto/fila/v1/service.proto +import type { NackMessage as _fila_v1_NackMessage, NackMessage__Output as _fila_v1_NackMessage__Output } from '../../fila/v1/NackMessage'; export interface NackRequest { - 'queue'?: (string); - 'messageId'?: (string); - 'error'?: (string); + 'messages'?: (_fila_v1_NackMessage)[]; } export interface NackRequest__Output { - 'queue': (string); - 'messageId': (string); - 'error': (string); + 'messages': (_fila_v1_NackMessage__Output)[]; } diff --git a/generated/fila/v1/NackResponse.ts b/generated/fila/v1/NackResponse.ts index dbc6271..fd00fe1 100644 --- a/generated/fila/v1/NackResponse.ts +++ b/generated/fila/v1/NackResponse.ts @@ -1,8 +1,11 @@ // Original file: proto/fila/v1/service.proto +import type { NackResult as _fila_v1_NackResult, NackResult__Output as _fila_v1_NackResult__Output } from '../../fila/v1/NackResult'; export interface NackResponse { + 'results'?: (_fila_v1_NackResult)[]; } export interface NackResponse__Output { + 'results': (_fila_v1_NackResult__Output)[]; } diff --git a/generated/fila/v1/NackResult.ts b/generated/fila/v1/NackResult.ts new file mode 100644 index 0000000..9a205d9 --- /dev/null +++ b/generated/fila/v1/NackResult.ts @@ -0,0 +1,16 @@ +// Original file: proto/fila/v1/service.proto + +import type { NackSuccess as _fila_v1_NackSuccess, NackSuccess__Output as _fila_v1_NackSuccess__Output } from '../../fila/v1/NackSuccess'; +import type { NackError as _fila_v1_NackError, NackError__Output as _fila_v1_NackError__Output } from '../../fila/v1/NackError'; + +export interface NackResult { + 'success'?: (_fila_v1_NackSuccess | null); + 'error'?: (_fila_v1_NackError | null); + 'result'?: "success"|"error"; +} + +export interface NackResult__Output { + 'success'?: (_fila_v1_NackSuccess__Output | null); + 'error'?: (_fila_v1_NackError__Output | null); + 'result'?: "success"|"error"; +} diff --git a/generated/fila/v1/NackSuccess.ts b/generated/fila/v1/NackSuccess.ts new file mode 100644 index 0000000..cd61ea8 --- /dev/null +++ b/generated/fila/v1/NackSuccess.ts @@ -0,0 +1,8 @@ +// Original file: proto/fila/v1/service.proto + + +export interface NackSuccess { +} + +export interface NackSuccess__Output { +} diff --git a/generated/fila/v1/StreamEnqueueRequest.ts b/generated/fila/v1/StreamEnqueueRequest.ts new file mode 100644 index 0000000..03d0b32 --- /dev/null +++ b/generated/fila/v1/StreamEnqueueRequest.ts @@ -0,0 +1,14 @@ +// Original file: proto/fila/v1/service.proto + +import type { EnqueueMessage as _fila_v1_EnqueueMessage, EnqueueMessage__Output as _fila_v1_EnqueueMessage__Output } from '../../fila/v1/EnqueueMessage'; +import type { Long } from '@grpc/proto-loader'; + +export interface StreamEnqueueRequest { + 'messages'?: (_fila_v1_EnqueueMessage)[]; + 'sequenceNumber'?: (number | string | Long); +} + +export interface StreamEnqueueRequest__Output { + 'messages': (_fila_v1_EnqueueMessage__Output)[]; + 'sequenceNumber': (string); +} diff --git a/generated/fila/v1/StreamEnqueueResponse.ts b/generated/fila/v1/StreamEnqueueResponse.ts new file mode 100644 index 0000000..56c5586 --- /dev/null +++ b/generated/fila/v1/StreamEnqueueResponse.ts @@ -0,0 +1,14 @@ +// Original file: proto/fila/v1/service.proto + +import type { EnqueueResult as _fila_v1_EnqueueResult, EnqueueResult__Output as _fila_v1_EnqueueResult__Output } from '../../fila/v1/EnqueueResult'; +import type { Long } from '@grpc/proto-loader'; + +export interface StreamEnqueueResponse { + 'sequenceNumber'?: (number | string | Long); + 'results'?: (_fila_v1_EnqueueResult)[]; +} + +export interface StreamEnqueueResponse__Output { + 'sequenceNumber': (string); + 'results': (_fila_v1_EnqueueResult__Output)[]; +} diff --git a/generated/service.ts b/generated/service.ts index 726d422..eb61748 100644 --- a/generated/service.ts +++ b/generated/service.ts @@ -1,5 +1,5 @@ import type * as grpc from '@grpc/grpc-js'; -import type { MessageTypeDefinition } from '@grpc/proto-loader'; +import type { EnumTypeDefinition, MessageTypeDefinition } from '@grpc/proto-loader'; import type { FilaServiceClient as _fila_v1_FilaServiceClient, FilaServiceDefinition as _fila_v1_FilaServiceDefinition } from './fila/v1/FilaService'; @@ -10,21 +10,34 @@ type SubtypeConstructor any, Subtype> export interface ProtoGrpcType { fila: { v1: { + AckError: MessageTypeDefinition + AckErrorCode: EnumTypeDefinition + AckMessage: MessageTypeDefinition AckRequest: MessageTypeDefinition AckResponse: MessageTypeDefinition - BatchEnqueueRequest: MessageTypeDefinition - BatchEnqueueResponse: MessageTypeDefinition - BatchEnqueueResult: MessageTypeDefinition + AckResult: MessageTypeDefinition + AckSuccess: MessageTypeDefinition ConsumeRequest: MessageTypeDefinition ConsumeResponse: MessageTypeDefinition + EnqueueError: MessageTypeDefinition + EnqueueErrorCode: EnumTypeDefinition + EnqueueMessage: MessageTypeDefinition EnqueueRequest: MessageTypeDefinition EnqueueResponse: MessageTypeDefinition + EnqueueResult: MessageTypeDefinition FilaService: SubtypeConstructor & { service: _fila_v1_FilaServiceDefinition } Message: MessageTypeDefinition MessageMetadata: MessageTypeDefinition MessageTimestamps: MessageTypeDefinition + NackError: MessageTypeDefinition + NackErrorCode: EnumTypeDefinition + NackMessage: MessageTypeDefinition NackRequest: MessageTypeDefinition NackResponse: MessageTypeDefinition + NackResult: MessageTypeDefinition + NackSuccess: MessageTypeDefinition + StreamEnqueueRequest: MessageTypeDefinition + StreamEnqueueResponse: MessageTypeDefinition } } google: { diff --git a/proto/fila/v1/service.proto b/proto/fila/v1/service.proto index fc0f710..7d1db79 100644 --- a/proto/fila/v1/service.proto +++ b/proto/fila/v1/service.proto @@ -6,20 +6,49 @@ import "fila/v1/messages.proto"; // Hot-path RPCs for producers and consumers. service FilaService { rpc Enqueue(EnqueueRequest) returns (EnqueueResponse); - rpc BatchEnqueue(BatchEnqueueRequest) returns (BatchEnqueueResponse); + rpc StreamEnqueue(stream StreamEnqueueRequest) returns (stream StreamEnqueueResponse); rpc Consume(ConsumeRequest) returns (stream ConsumeResponse); rpc Ack(AckRequest) returns (AckResponse); rpc Nack(NackRequest) returns (NackResponse); } -message EnqueueRequest { +// Individual message to enqueue. +message EnqueueMessage { string queue = 1; map headers = 2; bytes payload = 3; } +// Enqueue one or more messages. +message EnqueueRequest { + repeated EnqueueMessage messages = 1; +} + +// Per-message enqueue result. +message EnqueueResult { + oneof result { + string message_id = 1; + EnqueueError error = 2; + } +} + +// Typed enqueue error with structured error code. +message EnqueueError { + EnqueueErrorCode code = 1; + string message = 2; +} + +enum EnqueueErrorCode { + ENQUEUE_ERROR_CODE_UNSPECIFIED = 0; + ENQUEUE_ERROR_CODE_QUEUE_NOT_FOUND = 1; + ENQUEUE_ERROR_CODE_STORAGE = 2; + ENQUEUE_ERROR_CODE_LUA = 3; + ENQUEUE_ERROR_CODE_PERMISSION_DENIED = 4; +} + +// One result per input message. message EnqueueResponse { - string message_id = 1; + repeated EnqueueResult results = 1; } message ConsumeRequest { @@ -27,36 +56,87 @@ message ConsumeRequest { } message ConsumeResponse { - Message message = 1; // Single message (backward compatible, used when batch size is 1) - repeated Message messages = 2; // Batched messages (populated when server sends multiple at once) + repeated Message messages = 1; } -message AckRequest { +// Individual ack item. +message AckMessage { string queue = 1; string message_id = 2; } -message AckResponse {} +message AckRequest { + repeated AckMessage messages = 1; +} + +message AckResult { + oneof result { + AckSuccess success = 1; + AckError error = 2; + } +} -message NackRequest { +message AckSuccess {} + +message AckError { + AckErrorCode code = 1; + string message = 2; +} + +enum AckErrorCode { + ACK_ERROR_CODE_UNSPECIFIED = 0; + ACK_ERROR_CODE_MESSAGE_NOT_FOUND = 1; + ACK_ERROR_CODE_STORAGE = 2; + ACK_ERROR_CODE_PERMISSION_DENIED = 3; +} + +message AckResponse { + repeated AckResult results = 1; +} + +// Individual nack item. +message NackMessage { string queue = 1; string message_id = 2; string error = 3; } -message NackResponse {} +message NackRequest { + repeated NackMessage messages = 1; +} + +message NackResult { + oneof result { + NackSuccess success = 1; + NackError error = 2; + } +} -message BatchEnqueueRequest { - repeated EnqueueRequest messages = 1; +message NackSuccess {} + +message NackError { + NackErrorCode code = 1; + string message = 2; } -message BatchEnqueueResponse { - repeated BatchEnqueueResult results = 1; +enum NackErrorCode { + NACK_ERROR_CODE_UNSPECIFIED = 0; + NACK_ERROR_CODE_MESSAGE_NOT_FOUND = 1; + NACK_ERROR_CODE_STORAGE = 2; + NACK_ERROR_CODE_PERMISSION_DENIED = 3; } -message BatchEnqueueResult { - oneof result { - EnqueueResponse success = 1; - string error = 2; - } +message NackResponse { + repeated NackResult results = 1; +} + +// Stream enqueue — per-write batch with sequence tracking. +message StreamEnqueueRequest { + repeated EnqueueMessage messages = 1; + uint64 sequence_number = 2; +} + +message StreamEnqueueResponse { + uint64 sequence_number = 1; + repeated EnqueueResult results = 2; } diff --git a/src/batcher.ts b/src/batcher.ts index c9b3586..214478d 100644 --- a/src/batcher.ts +++ b/src/batcher.ts @@ -4,7 +4,6 @@ import { QueueNotFoundError, RPCError } from "./errors"; import type { EnqueueMessage } from "./types"; import type { FilaServiceClient } from "../generated/fila/v1/FilaService"; import type { EnqueueResponse__Output } from "../generated/fila/v1/EnqueueResponse"; -import type { BatchEnqueueResponse__Output } from "../generated/fila/v1/BatchEnqueueResponse"; /** Controls how the SDK batches enqueue() calls. */ export type BatchMode = @@ -19,7 +18,18 @@ interface BatchItem { reject: (err: Error) => void; } -function mapEnqueueError(err: grpc.ServiceError): Error { +/** + * Map a per-message EnqueueResult error to an SDK error. + * The unified proto uses typed EnqueueError with an error code. + */ +function mapResultError(code: string, message: string): Error { + if (code === "ENQUEUE_ERROR_CODE_QUEUE_NOT_FOUND") { + return new QueueNotFoundError(`enqueue: ${message}`); + } + return new RPCError(grpc.status.INTERNAL, message); +} + +function mapTransportError(err: grpc.ServiceError): Error { if (err.code === grpc.status.NOT_FOUND) { return new QueueNotFoundError(`enqueue: ${err.details}`); } @@ -28,7 +38,8 @@ function mapEnqueueError(err: grpc.ServiceError): Error { /** * Background batcher that collects enqueue() calls and flushes them - * as batch RPCs. Supports auto (opportunistic) and linger (timer-based) modes. + * via the unified Enqueue RPC (which accepts repeated messages). + * Supports auto (opportunistic) and linger (timer-based) modes. */ export class Batcher { private readonly grpcClient: FilaServiceClient; @@ -41,6 +52,7 @@ export class Batcher { private closed = false; private drainResolvers: Array<() => void> = []; private lingerTimer: ReturnType | null = null; + private inFlightCount = 0; constructor( grpcClient: FilaServiceClient, @@ -155,17 +167,18 @@ export class Batcher { private flushAll(): void { while (this.pending.length > 0) { const items = this.pending.splice(0, this.maxBatchSize); - // Fire-and-forget: flush concurrently. + this.inFlightCount++; this.flushBatch(items).then(() => { + this.inFlightCount--; this.notifyDrainComplete(); }); } - // Also check drain in case pending was already empty. + // Also check drain in case pending was already empty and nothing in-flight. this.notifyDrainComplete(); } private notifyDrainComplete(): void { - if (this.pending.length === 0 && this.drainResolvers.length > 0) { + if (this.pending.length === 0 && this.inFlightCount === 0 && this.drainResolvers.length > 0) { const resolvers = this.drainResolvers.splice(0); for (const resolve of resolvers) { resolve(); @@ -174,43 +187,12 @@ export class Batcher { } /** - * Flush a batch of items. Single item uses Enqueue RPC (preserves error - * types like QueueNotFoundError). Multiple items use BatchEnqueue. + * Flush a batch of items via the unified Enqueue RPC (repeated messages). + * All items -- single or multiple -- use the same RPC. */ - private async flushBatch(items: BatchItem[]): Promise { - if (items.length === 0) return; - - if (items.length === 1) { - return this.flushSingle(items[0]); - } + private flushBatch(items: BatchItem[]): Promise { + if (items.length === 0) return Promise.resolve(); - return this.flushMultiple(items); - } - - /** Flush a single item via the regular Enqueue RPC. */ - private flushSingle(item: BatchItem): Promise { - return new Promise((resolve) => { - this.grpcClient.enqueue( - { - queue: item.message.queue, - headers: item.message.headers, - payload: item.message.payload, - }, - this.callMetadata(), - (err: grpc.ServiceError | null, resp?: EnqueueResponse__Output) => { - if (err) { - item.reject(mapEnqueueError(err)); - } else { - item.resolve(resp!.messageId); - } - resolve(); - } - ); - }); - } - - /** Flush multiple items via the BatchEnqueue RPC. */ - private flushMultiple(items: BatchItem[]): Promise { const messages = items.map((item) => ({ queue: item.message.queue, headers: item.message.headers, @@ -218,16 +200,13 @@ export class Batcher { })); return new Promise((resolve) => { - this.grpcClient.batchEnqueue( + this.grpcClient.enqueue( { messages }, this.callMetadata(), - ( - err: grpc.ServiceError | null, - resp?: BatchEnqueueResponse__Output - ) => { + (err: grpc.ServiceError | null, resp?: EnqueueResponse__Output) => { if (err) { // Transport-level failure: all items get the error. - const mapped = new RPCError(err.code, err.details); + const mapped = mapTransportError(err); for (const item of items) { item.reject(mapped); } @@ -244,11 +223,11 @@ export class Batcher { ); continue; } - if (result.result === "success" && result.success) { - items[i].resolve(result.success.messageId!); + if (result.result === "messageId" && result.messageId) { + items[i].resolve(result.messageId); } else if (result.result === "error" && result.error) { items[i].reject( - new RPCError(grpc.status.INTERNAL, result.error) + mapResultError(result.error.code, result.error.message) ); } else { items[i].reject( @@ -262,5 +241,4 @@ export class Batcher { ); }); } - } diff --git a/src/client.ts b/src/client.ts index ad4b4a6..2cc4a75 100644 --- a/src/client.ts +++ b/src/client.ts @@ -9,9 +9,11 @@ import { QueueNotFoundError, RPCError, } from "./errors"; -import type { ConsumeMessage, EnqueueMessage, BatchEnqueueResult } from "./types"; +import type { ConsumeMessage, EnqueueMessage, EnqueueResult } from "./types"; import type { FilaServiceClient } from "../generated/fila/v1/FilaService"; import type { EnqueueResponse__Output } from "../generated/fila/v1/EnqueueResponse"; +import type { AckResponse__Output } from "../generated/fila/v1/AckResponse"; +import type { NackResponse__Output } from "../generated/fila/v1/NackResponse"; import type { ConsumeResponse__Output } from "../generated/fila/v1/ConsumeResponse"; import { Batcher, type BatchMode } from "./batcher"; @@ -85,52 +87,49 @@ function mapConsumeError(err: grpc.ServiceError): FilaError { return new RPCError(err.code, err.details); } -function mapAckError(err: grpc.ServiceError): FilaError { - if (err.code === grpc.status.NOT_FOUND) { - return new MessageNotFoundError(`ack: ${err.details}`); +/** + * Map a per-message EnqueueResult error code to an SDK error type. + */ +function mapEnqueueResultError(code: string, message: string): FilaError { + if (code === "ENQUEUE_ERROR_CODE_QUEUE_NOT_FOUND") { + return new QueueNotFoundError(`enqueue: ${message}`); } - return new RPCError(err.code, err.details); + return new RPCError(grpc.status.INTERNAL, message); } -function mapNackError(err: grpc.ServiceError): FilaError { - if (err.code === grpc.status.NOT_FOUND) { - return new MessageNotFoundError(`nack: ${err.details}`); +/** + * Map a per-message AckResult error code to an SDK error type. + */ +function mapAckResultError(code: string, message: string): FilaError { + if (code === "ACK_ERROR_CODE_MESSAGE_NOT_FOUND") { + return new MessageNotFoundError(`ack: ${message}`); } - return new RPCError(err.code, err.details); + return new RPCError(grpc.status.INTERNAL, message); +} + +/** + * Map a per-message NackResult error code to an SDK error type. + */ +function mapNackResultError(code: string, message: string): FilaError { + if (code === "NACK_ERROR_CODE_MESSAGE_NOT_FOUND") { + return new MessageNotFoundError(`nack: ${message}`); + } + return new RPCError(grpc.status.INTERNAL, message); } /** Map a ConsumeResponse to ConsumeMessage(s), skipping keepalive frames. */ function mapConsumeResponse( resp: ConsumeResponse__Output ): ConsumeMessage[] { - // Prefer the batched `messages` field when non-empty. - if (resp.messages && resp.messages.length > 0) { - const results: ConsumeMessage[] = []; - for (const msg of resp.messages) { - if (!msg || !msg.id) continue; - const metadata = msg.metadata; - results.push({ - id: msg.id, - headers: msg.headers ?? {}, - payload: Buffer.isBuffer(msg.payload) - ? msg.payload - : Buffer.from(msg.payload ?? ""), - fairnessKey: metadata?.fairnessKey ?? "", - attemptCount: metadata?.attemptCount ?? 0, - queue: metadata?.queueId ?? "", - }); - } - return results; - } - - // Fall back to singular `message` field (backward compatible). - const msg = resp.message; - if (!msg || !msg.id) { + if (!resp.messages || resp.messages.length === 0) { return []; // keepalive frame } - const metadata = msg.metadata; - return [ - { + + const results: ConsumeMessage[] = []; + for (const msg of resp.messages) { + if (!msg || !msg.id) continue; + const metadata = msg.metadata; + results.push({ id: msg.id, headers: msg.headers ?? {}, payload: Buffer.isBuffer(msg.payload) @@ -139,8 +138,9 @@ function mapConsumeResponse( fairnessKey: metadata?.fairnessKey ?? "", attemptCount: metadata?.attemptCount ?? 0, queue: metadata?.queueId ?? "", - }, - ]; + }); + } + return results; } /** Connection options for TLS, authentication, and batching. */ @@ -292,7 +292,7 @@ export class Client { * * When batching is enabled (default), the message is routed through the * batcher. At low load, messages are sent individually. At high load, - * messages cluster naturally into BatchEnqueue RPCs. + * messages cluster naturally into larger Enqueue RPCs. * * @param queue - Target queue name. * @param headers - Optional message headers. @@ -315,16 +315,27 @@ export class Client { }); } - // No batching: direct RPC. + // No batching: direct RPC with single message in the repeated field. return new Promise((resolve, reject) => { this.grpcClient.enqueue( - { queue, headers: headers ?? {}, payload }, + { messages: [{ queue, headers: headers ?? {}, payload }] }, this.callMetadata(), (err: grpc.ServiceError | null, resp?: EnqueueResponse__Output) => { if (err) { reject(mapEnqueueError(err)); + return; + } + const result = resp!.results[0]; + if (!result) { + reject(new RPCError(grpc.status.INTERNAL, "no result from server")); + return; + } + if (result.result === "messageId" && result.messageId) { + resolve(result.messageId); + } else if (result.result === "error" && result.error) { + reject(mapEnqueueResultError(result.error.code, result.error.message)); } else { - resolve(resp!.messageId); + reject(new RPCError(grpc.status.INTERNAL, "no result from server")); } } ); @@ -332,60 +343,53 @@ export class Client { } /** - * Enqueue a batch of messages in a single RPC call. + * Enqueue multiple messages in a single RPC call. * * Each message is independently validated and processed. A failed message - * does not affect the others in the batch. Returns one result per input - * message, in the same order. + * does not affect the others. Returns one result per input message, + * in the same order. * - * This is more efficient than calling enqueue() in a loop because it - * amortizes the RPC overhead across all messages. + * This always bypasses the batcher and issues a direct Enqueue RPC. * * @param messages - Array of messages to enqueue. * @returns Per-message results (success with messageId, or error with description). - * @throws {RPCError} For transport-level failures affecting the entire batch. + * @throws {RPCError} For transport-level failures affecting the entire call. */ - batchEnqueue(messages: EnqueueMessage[]): Promise { - // batchEnqueue always bypasses the batcher and uses a direct RPC. - // Create a temporary batcher-like object to reuse the RPC logic, - // or just call the gRPC client directly. - return this.doBatchEnqueue(messages); - } - - private doBatchEnqueue(messages: EnqueueMessage[]): Promise { + enqueueMany(messages: EnqueueMessage[]): Promise { const protoMessages = messages.map((m) => ({ queue: m.queue, headers: m.headers, payload: m.payload, })); - return new Promise((resolve, reject) => { - this.grpcClient.batchEnqueue( + return new Promise((resolve, reject) => { + this.grpcClient.enqueue( { messages: protoMessages }, this.callMetadata(), - (err: grpc.ServiceError | null, resp?) => { + (err: grpc.ServiceError | null, resp?: EnqueueResponse__Output) => { if (err) { reject(new RPCError(err.code, err.details)); return; } - const results: BatchEnqueueResult[] = resp!.results.map( - (r: { result?: string; success?: { messageId?: string } | null; error?: string }) => { - if (r.result === "success" && r.success) { - return { - success: true as const, - messageId: r.success.messageId!, - }; - } else if (r.result === "error" && r.error) { - return { success: false as const, error: r.error }; - } else { - return { - success: false as const, - error: "no result from server", - }; - } + const results: EnqueueResult[] = resp!.results.map((r) => { + if (r.result === "messageId" && r.messageId) { + return { + success: true as const, + messageId: r.messageId, + }; + } else if (r.result === "error" && r.error) { + return { + success: false as const, + error: r.error.message, + }; + } else { + return { + success: false as const, + error: "no result from server", + }; } - ); + }); resolve(results); } @@ -397,9 +401,9 @@ export class Client { * Open a streaming consumer on the specified queue. * * Returns an async iterable that yields messages as they become available. - * Nil message frames (keepalive signals) are skipped automatically. - * Batched delivery frames (multiple messages per ConsumeResponse) are - * transparently unpacked into individual messages. + * Empty response frames (keepalive signals) are skipped automatically. + * Delivery frames containing multiple messages are transparently unpacked + * into individual messages. * * If the server returns UNAVAILABLE with an `x-fila-leader-addr` metadata * header, the client transparently reconnects to the leader node and retries @@ -489,13 +493,24 @@ export class Client { ack(queue: string, msgId: string): Promise { return new Promise((resolve, reject) => { this.grpcClient.ack( - { queue, messageId: msgId }, + { messages: [{ queue, messageId: msgId }] }, this.callMetadata(), - (err: grpc.ServiceError | null) => { + (err: grpc.ServiceError | null, resp?: AckResponse__Output) => { if (err) { - reject(mapAckError(err)); - } else { + reject(new RPCError(err.code, err.details)); + return; + } + const result = resp!.results[0]; + if (!result) { + reject(new RPCError(grpc.status.INTERNAL, "no result from server")); + return; + } + if (result.result === "success") { resolve(); + } else if (result.result === "error" && result.error) { + reject(mapAckResultError(result.error.code, result.error.message)); + } else { + reject(new RPCError(grpc.status.INTERNAL, "no result from server")); } } ); @@ -513,13 +528,24 @@ export class Client { nack(queue: string, msgId: string, error: string): Promise { return new Promise((resolve, reject) => { this.grpcClient.nack( - { queue, messageId: msgId, error }, + { messages: [{ queue, messageId: msgId, error }] }, this.callMetadata(), - (err: grpc.ServiceError | null) => { + (err: grpc.ServiceError | null, resp?: NackResponse__Output) => { if (err) { - reject(mapNackError(err)); - } else { + reject(new RPCError(err.code, err.details)); + return; + } + const result = resp!.results[0]; + if (!result) { + reject(new RPCError(grpc.status.INTERNAL, "no result from server")); + return; + } + if (result.result === "success") { resolve(); + } else if (result.result === "error" && result.error) { + reject(mapNackResultError(result.error.code, result.error.message)); + } else { + reject(new RPCError(grpc.status.INTERNAL, "no result from server")); } } ); diff --git a/src/index.ts b/src/index.ts index e7a3302..9f479a0 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,6 +1,6 @@ export { Client } from "./client"; export type { ClientOptions } from "./client"; -export type { ConsumeMessage, EnqueueMessage, BatchEnqueueResult } from "./types"; +export type { ConsumeMessage, EnqueueMessage, EnqueueResult } from "./types"; export { FilaError, QueueNotFoundError, diff --git a/src/types.ts b/src/types.ts index cb9fd4b..99706b8 100644 --- a/src/types.ts +++ b/src/types.ts @@ -24,7 +24,7 @@ export interface EnqueueMessage { payload: Buffer; } -/** The result of a single message within a batch enqueue call. */ -export type BatchEnqueueResult = +/** The result of a single message within an enqueue call. */ +export type EnqueueResult = | { success: true; messageId: string } | { success: false; error: string }; diff --git a/test/batch.test.ts b/test/batch.test.ts index ce47540..40ea0e2 100644 --- a/test/batch.test.ts +++ b/test/batch.test.ts @@ -7,7 +7,7 @@ import { type TestServer, } from "./helpers"; -describe.skipIf(!FILA_SERVER_AVAILABLE)("Batch operations", () => { +describe.skipIf(!FILA_SERVER_AVAILABLE)("Enqueue operations", () => { let server: TestServer; beforeAll(async () => { @@ -18,16 +18,16 @@ describe.skipIf(!FILA_SERVER_AVAILABLE)("Batch operations", () => { server?.stop(); }); - describe("batchEnqueue", () => { + describe("enqueueMany", () => { it("enqueues multiple messages in a single RPC", async () => { - await server.createQueue("batch-multi"); + await server.createQueue("multi-enqueue"); const client = new Client(server.addr, { batchMode: "disabled" }); try { - const results = await client.batchEnqueue([ - { queue: "batch-multi", headers: { idx: "0" }, payload: Buffer.from("msg-0") }, - { queue: "batch-multi", headers: { idx: "1" }, payload: Buffer.from("msg-1") }, - { queue: "batch-multi", headers: { idx: "2" }, payload: Buffer.from("msg-2") }, + const results = await client.enqueueMany([ + { queue: "multi-enqueue", headers: { idx: "0" }, payload: Buffer.from("msg-0") }, + { queue: "multi-enqueue", headers: { idx: "1" }, payload: Buffer.from("msg-1") }, + { queue: "multi-enqueue", headers: { idx: "2" }, payload: Buffer.from("msg-2") }, ]); expect(results).toHaveLength(3); @@ -41,9 +41,9 @@ describe.skipIf(!FILA_SERVER_AVAILABLE)("Batch operations", () => { // Verify all messages are consumable. const received: string[] = []; let count = 0; - for await (const msg of client.consume("batch-multi")) { + for await (const msg of client.consume("multi-enqueue")) { received.push(msg.payload.toString()); - await client.ack("batch-multi", msg.id); + await client.ack("multi-enqueue", msg.id); count++; if (count >= 3) break; } @@ -56,12 +56,12 @@ describe.skipIf(!FILA_SERVER_AVAILABLE)("Batch operations", () => { }); it("returns per-message errors for nonexistent queues", async () => { - await server.createQueue("batch-partial"); + await server.createQueue("multi-partial"); const client = new Client(server.addr, { batchMode: "disabled" }); try { - const results = await client.batchEnqueue([ - { queue: "batch-partial", headers: {}, payload: Buffer.from("ok") }, + const results = await client.enqueueMany([ + { queue: "multi-partial", headers: {}, payload: Buffer.from("ok") }, { queue: "no-such-queue", headers: {}, payload: Buffer.from("fail") }, ]); @@ -77,13 +77,13 @@ describe.skipIf(!FILA_SERVER_AVAILABLE)("Batch operations", () => { }); it("returns message IDs in same order as input", async () => { - await server.createQueue("batch-order"); + await server.createQueue("multi-order"); const client = new Client(server.addr, { batchMode: "disabled" }); try { - const results = await client.batchEnqueue([ - { queue: "batch-order", headers: {}, payload: Buffer.from("first") }, - { queue: "batch-order", headers: {}, payload: Buffer.from("second") }, + const results = await client.enqueueMany([ + { queue: "multi-order", headers: {}, payload: Buffer.from("first") }, + { queue: "multi-order", headers: {}, payload: Buffer.from("second") }, ]); expect(results).toHaveLength(2); @@ -157,7 +157,7 @@ describe.skipIf(!FILA_SERVER_AVAILABLE)("Batch operations", () => { const client = new Client(server.addr); try { // Single message to nonexistent queue: should get QueueNotFoundError - // because single-item batches use Enqueue RPC. + // because the per-result error code is mapped to QueueNotFoundError. await expect( client.enqueue("no-such-queue-auto", null, Buffer.from("fail")) ).rejects.toThrow(QueueNotFoundError);