| 130 | } |
| 131 | |
| 132 | enqueue( |
| 133 | queueName: string, |
| 134 | payload: Buffer, |
| 135 | contentType: string, |
| 136 | options?: EnqueueOptions |
| 137 | ): { messageId: string } { |
| 138 | const messageId = randomBytes(16).toString('hex'); |
| 139 | const retentionMs = |
| 140 | (options?.retentionSeconds ?? 0) > 0 |
| 141 | ? options!.retentionSeconds! * 1000 |
| 142 | : DEFAULT_RETENTION; |
| 143 | const idempotencyRecordKey = options?.idempotencyKey |
| 144 | ? `${queueName}:${options.idempotencyKey}` |
| 145 | : undefined; |
| 146 | |
| 147 | if (idempotencyRecordKey) { |
| 148 | const record = this.idempotencyRecords.get(idempotencyRecordKey); |
| 149 | if (record && record.expiresAt > Date.now()) { |
| 150 | this.duplicateMessages.set(messageId, { |
| 151 | queueName, |
| 152 | originalMessageId: record.messageId, |
| 153 | expiresAt: Date.now() + retentionMs, |
| 154 | }); |
| 155 | output.debug( |
| 156 | `queues: skipped duplicate message for queue "${queueName}"` |
| 157 | ); |
| 158 | return { messageId }; |
| 159 | } |
| 160 | if (record) { |
| 161 | this.idempotencyRecords.delete(idempotencyRecordKey); |
| 162 | } |
| 163 | } |
| 164 | |
| 165 | if (idempotencyRecordKey) { |
| 166 | this.idempotencyRecords.set(idempotencyRecordKey, { |
| 167 | messageId, |
| 168 | expiresAt: Date.now() + retentionMs, |
| 169 | }); |
| 170 | } |
| 171 | |
| 172 | const message: StoredMessage = { |
| 173 | messageId, |
| 174 | payload, |
| 175 | contentType, |
| 176 | queueName, |
| 177 | createdAt: new Date().toISOString(), |
| 178 | retentionMs, |
| 179 | }; |
| 180 | |
| 181 | this.messages.set(messageId, message); |
| 182 | output.debug( |
| 183 | `queues: stored message ${messageId} for queue "${queueName}"` |
| 184 | ); |
| 185 | |
| 186 | const delaySeconds = options?.delaySeconds ?? 0; |
| 187 | const matchingGroups = this.consumerGroups.filter(g => |
| 188 | g.topicRegex.test(queueName) |
| 189 | ); |