|
@@ -67,19 +67,7 @@ export class BufferService {
|
|
|
this.bufferMessage(message);
|
|
|
}
|
|
|
if (this.connectionState.getValue().status === 'DIRECT_PUBLISH') {
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
+
|
|
|
}
|
|
|
}
|
|
|
|
|
@@ -92,25 +80,6 @@ export class BufferService {
|
|
|
}
|
|
|
if (state.status === 'DIRECT_PUBLISH') {
|
|
|
this.releaseBufferedMessages(this.messageStream)
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- private async isBufferNotEmpty(): Promise<boolean> {
|
|
|
- if (this.messageModel) {
|
|
|
-
|
|
|
- const count = await this.messageModel.estimatedDocumentCount().exec();
|
|
|
- return count > 0;
|
|
|
- } else {
|
|
|
-
|
|
|
- return this.messageBuffer.length > 0;
|
|
|
}
|
|
|
}
|
|
|
|
|
@@ -119,7 +88,9 @@ export class BufferService {
|
|
|
try {
|
|
|
|
|
|
await this.messageModel.create(message);
|
|
|
- console.log(`Message${(message.message as MessageLog).appData.msgId} saved to MongoDB buffer`);
|
|
|
+ this.messageModel.countDocuments({}).then((count) => {
|
|
|
+ console.log(`Message${(message.message as MessageLog).appData.msgId} saved to MongoDB buffer. There is ${count} messages in datatbase at the moment.`);
|
|
|
+ })
|
|
|
} catch (error) {
|
|
|
console.error('Error saving message to MongoDB:', error);
|
|
|
|