94 lines
2.7 KiB
Go
94 lines
2.7 KiB
Go
package dbAccess
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"time"
|
|
|
|
"git.fjla.uk/owlboard/go-types/pkg/database"
|
|
"git.fjla.uk/owlboard/timetable-mgr/log"
|
|
"go.mongodb.org/mongo-driver/bson"
|
|
"go.mongodb.org/mongo-driver/mongo"
|
|
"go.mongodb.org/mongo-driver/mongo/options"
|
|
"go.uber.org/zap"
|
|
)
|
|
|
|
const Doctype = "CifMetadata"
|
|
|
|
// The type describing the CifMetadata 'type' in the database.
|
|
// This type will be moved to owlboard/go-types
|
|
type CifMetadata struct {
|
|
Doctype string `json:"type"`
|
|
LastUpdate time.Time `json:"lastUpdate"`
|
|
LastTimestamp int64 `json:"lastTimestamp"`
|
|
LastSequence int64 `json:"lastSequence"`
|
|
}
|
|
|
|
// Fetches the CifMetadata from the database, returns nil if no metadata exists - before first initialisation for example.
|
|
func GetCifMetadata() (*CifMetadata, error) {
|
|
database := MongoClient.Database(databaseName)
|
|
collection := database.Collection(metaCollection)
|
|
filter := bson.M{"type": Doctype}
|
|
var result CifMetadata
|
|
err := collection.FindOne(context.Background(), filter).Decode(&result)
|
|
if err != nil {
|
|
if errors.Is(err, mongo.ErrNoDocuments) {
|
|
return nil, nil
|
|
}
|
|
log.Msg.Error("Error fetching CIF Metadata")
|
|
return nil, err
|
|
}
|
|
|
|
return &result, nil
|
|
}
|
|
|
|
// Uses upsert to Insert/Update the CifMetadata in the database
|
|
func PutCifMetadata(metadata CifMetadata) bool {
|
|
database := MongoClient.Database(databaseName)
|
|
collection := database.Collection(metaCollection)
|
|
options := options.Update().SetUpsert(true)
|
|
filter := bson.M{"type": Doctype}
|
|
update := bson.M{
|
|
"type": Doctype,
|
|
"LastUpdate": metadata.LastUpdate,
|
|
"LastTimestamp": metadata.LastTimestamp,
|
|
"LastSequence": metadata.LastSequence,
|
|
}
|
|
|
|
_, err := collection.UpdateOne(context.Background(), filter, update, options)
|
|
|
|
if err != nil {
|
|
log.Msg.Error("Error updating CIF Metadata", zap.Error(err))
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
func DeleteCifEntries(deletions []database.DeleteQuery) error {
|
|
// Prepare deletion tasks
|
|
bulkDeletions := make([]mongo.WriteModel, 0, len(deletions))
|
|
|
|
for _, deleteQuery := range deletions {
|
|
filter := bson.M{
|
|
"trainUid": deleteQuery.TrainUid,
|
|
"scheduleStartDate": deleteQuery.ScheduleStartDate,
|
|
"stpIndicator": deleteQuery.StpIndicator,
|
|
}
|
|
bulkDeletions = append(bulkDeletions, mongo.NewDeleteManyModel().SetFilter(filter))
|
|
}
|
|
|
|
log.Msg.Info("Running `Delete` tasks from CIF Update", zap.Int("Required deletions", len(deletions)))
|
|
for i := 0; i < len(bulkDeletions); i += batchsize {
|
|
end := i + batchsize
|
|
if end > len(bulkDeletions) {
|
|
end = len(bulkDeletions)
|
|
}
|
|
_, err := MongoClient.Database(databaseName).Collection(TimetableCollection).BulkWrite(context.Background(), bulkDeletions[i:end])
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|