Mongodb aggregation changes not being persisted in go

Viewed 31

I'm running an aggregation to remove stale documents but the changes don't actually affect the database. The query ignores already expired documents, so the number of results should change after each query runs, but it doesn't.

func CheckShipmentExpiryDates(c *mongo.Client) (int, error) {
    numberOfExpiredShipments := 0
    coll := c.Database(os.Getenv("DATABASE")).Collection("shipments")
    update := bson.M{"$set": bson.M{"status": "EXPIRED", "updated_at": time.Now()}}
    pipeline := []bson.M{
        {"$lookup": bson.M{
            "from": "shipment_quotes",
            "let":  bson.M{"shipmentID": "$_id"},
            "pipeline": []bson.M{
                {"$match": bson.M{"$expr": bson.M{"$and": []bson.M{{"$eq": []string{"$shipment_id", "$$shipmentID"}}, {"$eq": []string{"$status", "WON"}}}}}},
            },
            "as": "quotes",
        }},
        {"$match": bson.M{"expiration_date": bson.M{"$exists": true}}},
        {"$match": bson.M{"$expr": bson.M{"$and": []bson.M{
            {"$ne": []string{"$status", "EXPIRED"}},
            {"$lt": []interface{}{"$expiration_date", time.Now()}},
            {"$eq": []interface{}{bson.M{"$size": "$quotes"}, 0}},
            {"expiration_date": bson.M{"$type": 9}},
        }}}},
        update,
    }

    err := c.UseSession(context.TODO(), func(sessionContext mongo.SessionContext) error {
        if err := sessionContext.StartTransaction(); err != nil {
            return err
        }
        cursor, err := coll.Aggregate(sessionContext, pipeline)
        if err != nil {
            _ = sessionContext.AbortTransaction(sessionContext)
            return err
        }

        var shipments []bson.M
        if err := cursor.All(sessionContext, &shipments); err != nil {
            _ = sessionContext.AbortTransaction(sessionContext)
            return err
        }

        fmt.Println("~First shipment's status", shipments[0]["shipment_unique_number"], shipments[0]["status"])

        numberOfExpiredShipments = len(shipments)

        fmt.Println(sessionContext.CommitTransaction(sessionContext))
        return nil
    })

    return numberOfExpiredShipments, err
}

As you can see, I'm logging the first result and checking it against the database in real time, using compass, but the changes aren't actually being persisted. The query runs over and over again, returning the same number of expired shipments.

mc, mongoErr := connection.MongoInit()
    if mongoErr != nil {
        panic(mongoErr)
    }
    utils.InitDB(mc)
    defer func() {
        if err := mc.Disconnect(context.TODO()); err != nil {
            panic(err)
        }
    }()

    n := connection.NewNotificationCenter()
    sseInit(mc, googleApi, n)
    graphSchema, err := schema.InjectSchema(mutationInit(mc, googleApi), queryInit(mc, googleApi))
    if err != nil {
        panic(err)
    }
    restApiUseCase := mutationsRestApiInit(mc, googleApi)
    connection.InjectGraphqlHandler(graphSchema, n, restApiUseCase)

    initIncrementStartdate(mc)
    initShipmentExpiredCron(mc)
func initShipmentExpiredCron(mg *mongo.Client) {
    c := cron.New()
    c.AddFunc("*/5 * * * *", func() {
        expiredShipments, err := utils.CheckShipmentExpiryDates(mg)
        if err != nil {
            log.Println("CRON ERROR: An error occured while trying to check the expiry date for each shipment")
            log.Println(err)
        } else {
            // Print how many shipments are expired
            log.Println("CRON SUCCESS: The following number of shipments have expired: ", expiredShipments)
        }
    })
    c.Start()
}

I really don't understand what's wrong with it.

1 Answers

rickhg12hs was right, I needed to also use $merge. Oddly, this makes the aggregation not return anything, so now I don't know the number of shipments that have expired, but that was not really necessary. This is the pipeline at the end

pipeline := []bson.M{{
  "$lookup": bson.M{
    "from": "shipment_quotes",
    "let": bson.M{
      "shipmentID": "$_id"
    },
    "pipeline": []bson.M{{
      "$match": bson.M{
        "$expr": bson.M{
          "$and": []bson.M{{
            "$eq": []string{
              "$shipment_id", "$$shipmentID"
            }}, {
              "$eq": []string{
                "$status", "WON"
              }
            }
          }}
        }
      },
    },
    "as": "quotes",
  }}, {
    "$match": bson.M{
      "expiration_date": bson.M{
        "$exists": true
      }
    }
  }, {
    "$match": bson.M{
      "$expr": bson.M{
        "$and": []bson.M{{
          "$ne": []string{
            "$status", "EXPIRED"
          }
        }, {
          "$lt": []interface{}{
            "$expiration_date", time.Now()
          }
        }, {
          "$eq": []interface{}{
            bson.M{
              "$size": "$quotes"
            }, 0
          }
        }, {
          "expiration_date": bson.M{
            "$type": 9
          }
        },
      }
    }
  }},
  update,
  {
    "$merge": bson.M{
      "into": "shipments",
      "on": "_id"
    }
  },
}
Related