Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,12 @@ services:
volumes:
- rabbit_data:/var/lib/rabbitmq
restart: on-failure
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
interval: 5s
timeout: 5s
retries: 10
start_period: 10s

redis:
image: redis:6.2.7-alpine
Expand Down
4 changes: 2 additions & 2 deletions packages/amqp/lib/AbstractAmqpConsumer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -306,7 +306,7 @@ export abstract class AbstractAmqpConsumer<
}

// Empty content for whatever reason
if (!resolveMessageResult.result || !resolveMessageResult.result.body) {
if (!resolveMessageResult.result?.body) {
return ABORT_EARLY_EITHER
}

Expand Down Expand Up @@ -347,7 +347,7 @@ export abstract class AbstractAmqpConsumer<
const resolvedMessage = resolveMessageResult.result

// Empty content for whatever reason
if (!resolvedMessage || !resolvedMessage.body) return ABORT_EARLY_EITHER
if (!resolvedMessage?.body) return ABORT_EARLY_EITHER

// @ts-expect-error
if (this.messageIdField in resolvedMessage.body) {
Expand Down
4 changes: 2 additions & 2 deletions packages/amqp/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
"test:coverage": "npm run test -- --coverage",
"lint": "biome check . && tsc",
"lint:fix": "biome check --write .",
"docker:start": "docker compose up -d rabbitmq",
"docker:start": "docker compose up -d --wait rabbitmq",
"docker:stop": "docker compose down",
"prepublishOnly": "npm run lint && npm run build"
},
Expand All @@ -45,7 +45,7 @@
"@types/node": "^25.0.2",
"@vitest/coverage-v8": "^4.0.15",
"amqplib": "^0.10.8",
"awilix": "^12.0.5",
"awilix": "^13.0.3",
"awilix-manager": "^6.1.0",
"rimraf": "^6.0.1",
"typescript": "^5.9.3",
Expand Down
2 changes: 1 addition & 1 deletion packages/amqp/test/utils/testContext.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ export async function registerDependencies(
dependencyOverrides: DependencyOverrides = {},
queuesEnabled = true,
) {
const diContainer = createContainer({
const diContainer = createContainer<Dependencies>({
injectionMode: 'PROXY',
})
const awilixManager = new AwilixManager({
Expand Down
2 changes: 1 addition & 1 deletion packages/core/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@
"@types/node": "^25.0.2",
"@types/tmp": "^0.2.6",
"@vitest/coverage-v8": "^4.0.15",
"awilix": "^12.0.5",
"awilix": "^13.0.3",
"awilix-manager": "^6.1.0",
"rimraf": "^6.0.1",
"typescript": "^5.9.2",
Expand Down
2 changes: 1 addition & 1 deletion packages/core/test/testContext.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ export const TestEvents = {
export type TestEventsType = (typeof TestEvents)[keyof typeof TestEvents][]

export async function registerDependencies(dependencyOverrides: DependencyOverrides = {}) {
const diContainer = createContainer({
const diContainer = createContainer<Dependencies>({
injectionMode: 'PROXY',
})
const awilixManager = new AwilixManager({
Expand Down
2 changes: 1 addition & 1 deletion packages/gcp-pubsub/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@
"@message-queue-toolkit/schemas": "*",
"@types/node": "^25.0.2",
"@vitest/coverage-v8": "^4.0.15",
"awilix": "^12.0.5",
"awilix": "^13.0.3",
"awilix-manager": "^6.1.0",
"ioredis": "^5.7.0",
"rimraf": "^6.0.1",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -230,32 +230,30 @@ describe('PubSubPermissionConsumer - Deduplication', () => {
expect(consumer.addCounter).toBe(1)
})

it(
'processes message when lock acquisition has non-timeout error',
{ timeout: 15000 },
async () => {
const messageId = randomUUID()

// Mock acquireLock to simulate a non-timeout error (e.g., Redis connection error)
// Non-timeout errors are swallowed and message is processed normally
vi.spyOn(messageDeduplicationStore, 'acquireLock').mockResolvedValue({
error: new Error('Redis connection error'),
})

const message: PERMISSIONS_ADD_MESSAGE_TYPE = {
id: messageId,
messageType: 'add',
timestamp: new Date().toISOString(),
userIds: ['user1'],
}

await publisher.publish(message)

// Message should be processed even though lock acquisition failed
const result = await consumer.handlerSpy.waitForMessageWithId(messageId, 'consumed')
expect(result.processingResult.status).toBe('consumed')
expect(consumer.addCounter).toBe(1)
},
)
it('processes message when lock acquisition has non-timeout error', {
timeout: 15000,
}, async () => {
const messageId = randomUUID()

// Mock acquireLock to simulate a non-timeout error (e.g., Redis connection error)
// Non-timeout errors are swallowed and message is processed normally
vi.spyOn(messageDeduplicationStore, 'acquireLock').mockResolvedValue({
error: new Error('Redis connection error'),
})

const message: PERMISSIONS_ADD_MESSAGE_TYPE = {
id: messageId,
messageType: 'add',
timestamp: new Date().toISOString(),
userIds: ['user1'],
}

await publisher.publish(message)

// Message should be processed even though lock acquisition failed
const result = await consumer.handlerSpy.waitForMessageWithId(messageId, 'consumed')
expect(result.processingResult.status).toBe('consumed')
expect(consumer.addCounter).toBe(1)
})
})
})
Original file line number Diff line number Diff line change
Expand Up @@ -125,89 +125,85 @@ describe('PubSubPermissionConsumer - Payload Offloading', () => {
expect(consumer.addCounter).toBe(1)
})

it(
'consumes offloaded message with array field and validates schema correctly',
{ timeout: 10000 },
async () => {
// Create a large array of userIds to trigger offloading (need > 10MB)
// Each userId needs to be ~1000 chars to make 10,500 items exceed 10MB
const largeUserIdArray = Array.from(
{ length: 10500 },
(_, i) => `user-${i}-${'x'.repeat(1000)}`,
)
it('consumes offloaded message with array field and validates schema correctly', {
timeout: 10000,
}, async () => {
// Create a large array of userIds to trigger offloading (need > 10MB)
// Each userId needs to be ~1000 chars to make 10,500 items exceed 10MB
const largeUserIdArray = Array.from(
{ length: 10500 },
(_, i) => `user-${i}-${'x'.repeat(1000)}`,
)

const message = {
id: 'large-array-message-1',
messageType: 'add',
timestamp: new Date().toISOString(),
userIds: largeUserIdArray,
} satisfies PERMISSIONS_ADD_MESSAGE_TYPE
const message = {
id: 'large-array-message-1',
messageType: 'add',
timestamp: new Date().toISOString(),
userIds: largeUserIdArray,
} satisfies PERMISSIONS_ADD_MESSAGE_TYPE

// Verify the message is large enough to trigger offloading
expect(JSON.stringify(message).length).toBeGreaterThan(largeMessageSizeThreshold)
// Verify the message is large enough to trigger offloading
expect(JSON.stringify(message).length).toBeGreaterThan(largeMessageSizeThreshold)

await publisher.publish(message)
await publisher.publish(message)

// Wait for the message to be consumed
const consumptionResult = await consumer.handlerSpy.waitForMessageWithId(
message.id,
'consumed',
)
// Wait for the message to be consumed
const consumptionResult = await consumer.handlerSpy.waitForMessageWithId(
message.id,
'consumed',
)

// Verify the full payload was received including the large array
expect(consumptionResult.message).toMatchObject({
id: message.id,
messageType: message.messageType,
userIds: largeUserIdArray,
})
expect(consumptionResult.message.userIds).toHaveLength(largeUserIdArray.length)
expect(consumer.addCounter).toBe(1)
},
)

it(
'validates schema correctly after retrieving offloaded payload',
{ timeout: 10000 },
async () => {
// Create a message with metadata that will be validated against the schema
const message = {
id: 'schema-validation-1',
messageType: 'add',
timestamp: new Date().toISOString(),
metadata: {
largeField: 'x'.repeat(largeMessageSizeThreshold + 1000),
},
userIds: ['test-user'],
} satisfies PERMISSIONS_ADD_MESSAGE_TYPE

expect(JSON.stringify(message).length).toBeGreaterThan(largeMessageSizeThreshold)

await publisher.publish(message)

const consumptionResult = await consumer.handlerSpy.waitForMessageWithId(
message.id,
'consumed',
)
// Verify the full payload was received including the large array
expect(consumptionResult.message).toMatchObject({
id: message.id,
messageType: message.messageType,
userIds: largeUserIdArray,
})
expect(consumptionResult.message.userIds).toHaveLength(largeUserIdArray.length)
expect(consumer.addCounter).toBe(1)
})

// Verify all fields were properly deserialized and validated
expect(consumptionResult.message).toMatchObject({
id: message.id,
messageType: message.messageType,
userIds: message.userIds,
metadata: {
largeField: message.metadata.largeField,
},
})

// Type guard to access metadata property
if (consumptionResult.message.messageType === 'add') {
expect(consumptionResult.message.metadata?.largeField).toHaveLength(
message.metadata.largeField.length,
)
}
expect(consumer.addCounter).toBe(1)
},
)
it('validates schema correctly after retrieving offloaded payload', {
timeout: 10000,
}, async () => {
// Create a message with metadata that will be validated against the schema
const message = {
id: 'schema-validation-1',
messageType: 'add',
timestamp: new Date().toISOString(),
metadata: {
largeField: 'x'.repeat(largeMessageSizeThreshold + 1000),
},
userIds: ['test-user'],
} satisfies PERMISSIONS_ADD_MESSAGE_TYPE

expect(JSON.stringify(message).length).toBeGreaterThan(largeMessageSizeThreshold)

await publisher.publish(message)

const consumptionResult = await consumer.handlerSpy.waitForMessageWithId(
message.id,
'consumed',
)

// Verify all fields were properly deserialized and validated
expect(consumptionResult.message).toMatchObject({
id: message.id,
messageType: message.messageType,
userIds: message.userIds,
metadata: {
largeField: message.metadata.largeField,
},
})

// Type guard to access metadata property
if (consumptionResult.message.messageType === 'add') {
expect(consumptionResult.message.metadata?.largeField).toHaveLength(
message.metadata.largeField.length,
)
}
expect(consumer.addCounter).toBe(1)
})
})

describe('payload retrieval errors', () => {
Expand Down
Loading
Loading