Using mongodb driver for change streams over a mongoose connection

Viewed 367

How do I use the mongodb driver to watch for changes in my database, while using mongoose for database connection. Mongoose documentation provides very little examples on how to use change streams. In my case the change streams only work when there are no pipeline options provided. I assumed it was a problem with the syntax of the pipeline object, but I carefully followed the examples on the mongodb docs. The problem is either with the pipeline object, or the how the change stream is being implemented with mongoose. Any help is appreciated.

Here's my current approach (which doesn't work):

mongoose
.connect(db, {
    useNewUrlParser: true,
    useFindAndModify: false,
    useUnifiedTopology: true
    // useCreateIndex: true
})
.then(() => {
    console.log("Connected to MongoDB...");
})
.catch(err => {
    console.log(err);
});

const pipeline = { 
   $match: {
      $or: [{ operationType: 'insert' },{ operationType: 'update' }], 
      'fullDocument.institution': uniId 
   } 
};

const options = { fullDocument: 'updateLookup' }

changeStream.on("change", next => {
        switch(next.operationType) {
          case 'insert':
            console.log('an insert happened...', "uni_ID: ", next.fullDocument.institution);
            let rooms = Object.keys(socket.rooms);
            console.log("rooms: ", rooms);

            nmsps.emit('insert', {
              type: 'insert',
              msg: 'New question available',
              newPost: next.fullDocument
            });
            break;

          case 'update':
            console.log('an update happened...');

            nmsps.emit('update', {
              type: 'update',
              postId: next.documentKey._id,
              updateInfo: next.updateDescription.updatedFields,
              msg: "Question has been updated."
            });
            break;

          case 'delete':
            console.log('a delete happened...');

            nmsps.emit('delete', {
              type: 'delete',
              deletedId: next.documentKey._id,
              msg: 'Question has been deleted.'
            });
            break;

          default:
            break;
        }
      })
0 Answers
Related