oc-lib/dbs/mongo/mongo.go

198 lines
5.8 KiB
Go
Raw Normal View History

2024-07-17 18:02:30 +02:00
package mongo
import (
"context"
"encoding/json"
"errors"
2024-07-18 11:51:12 +02:00
lib "oc-lib"
"oc-lib/dbs"
2024-07-17 18:02:30 +02:00
"os"
"time"
"go.mongodb.org/mongo-driver/bson"
"go.mongodb.org/mongo-driver/bson/primitive"
"go.mongodb.org/mongo-driver/mongo"
"go.mongodb.org/mongo-driver/mongo/options"
)
var (
2024-07-18 11:51:12 +02:00
mngoClient *mongo.Client
mngoDB *mongo.Database
MngoCtx context.Context
cancel context.CancelFunc
2024-07-17 18:02:30 +02:00
2024-07-18 11:51:12 +02:00
existingCollections []string
2024-07-17 18:02:30 +02:00
2024-07-18 11:51:12 +02:00
ResourceMap map[string]interface{}
2024-07-17 18:02:30 +02:00
)
2024-07-18 12:05:32 +02:00
func Init(collections []string) {
2024-07-17 18:02:30 +02:00
// var baseConfig string
var err error
var conf map[string]string
var DBname string
ResourceMap = make(map[string]interface{})
2024-07-18 11:51:12 +02:00
lib.Logger = lib.CreateLogger("oclib", "")
2024-07-17 18:02:30 +02:00
db_conf, err := os.ReadFile("tests/oclib_conf.json")
if err != nil {
2024-07-18 11:51:12 +02:00
lib.Logger.Fatal().Msg("Could not find configuration file")
2024-07-17 18:02:30 +02:00
}
2024-07-18 11:51:12 +02:00
json.Unmarshal(db_conf, &conf)
2024-07-18 11:56:54 +02:00
DBname = lib.GetConfig().DCNAME + "-" + lib.GetConfig().DBPOINT
2024-07-17 18:02:30 +02:00
2024-07-18 11:56:54 +02:00
lib.Logger.Info().Msg("Connecting to" + lib.GetConfig().MongoURL)
2024-07-17 18:02:30 +02:00
MngoCtx, cancel = context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
2024-07-18 11:56:54 +02:00
createClient(lib.GetConfig().MongoURL)
2024-07-17 18:02:30 +02:00
2024-07-18 11:51:12 +02:00
lib.Logger.Info().Msg("Connecting mongo client to db " + DBname)
2024-07-18 12:05:32 +02:00
prepareDB(collections, lib.GetConfig().DCNAME, lib.GetConfig().DBPOINT)
2024-07-17 18:02:30 +02:00
2024-07-18 11:51:12 +02:00
lib.Logger.Info().Msg("Database is READY")
2024-07-17 18:02:30 +02:00
}
2024-07-18 11:51:12 +02:00
func createClient(MongoURL string) {
2024-07-17 18:02:30 +02:00
var err error
// Allows us to use marshal and unmarshall with results of FindOne() and others
2024-07-18 11:51:12 +02:00
bsonOpts := &options.BSONOptions{
2024-07-17 18:02:30 +02:00
UseJSONStructTags: true,
2024-07-18 11:51:12 +02:00
NilSliceAsEmpty: true,
2024-07-17 18:02:30 +02:00
}
clientOptions := options.Client().ApplyURI(MongoURL).SetBSONOptions(bsonOpts)
2024-07-18 11:51:12 +02:00
mngoClient, err = mongo.Connect(MngoCtx, clientOptions)
2024-07-17 18:02:30 +02:00
if err != nil {
2024-07-18 11:51:12 +02:00
lib.Logger.Fatal().Msg("Mongodb NewClient " + MongoURL + ":" + "err")
2024-07-17 18:02:30 +02:00
panic(err)
}
// Ping the primary
if mngoClient, err = mongo.Connect(MngoCtx, clientOptions); err != nil {
2024-07-18 11:51:12 +02:00
lib.Logger.Fatal().Msg("Mongodb connect " + MongoURL + ":" + "err")
2024-07-17 18:02:30 +02:00
panic(err)
}
if err = mngoClient.Ping(MngoCtx, nil); err != nil {
2024-07-18 11:51:12 +02:00
lib.Logger.Fatal().Msg("Mongodb ping " + MongoURL + ":" + "err")
2024-07-17 18:02:30 +02:00
panic(err)
}
}
2024-07-18 12:05:32 +02:00
func prepareDB(list_collection []string, dc_name string, db_point string) {
2024-07-17 18:02:30 +02:00
var err error
2024-07-18 11:51:12 +02:00
DBname := dc_name + "-" + db_point
2024-07-17 18:02:30 +02:00
mngoDB = mngoClient.Database(DBname)
2024-07-18 11:51:12 +02:00
existingCollections, err = mngoDB.ListCollectionNames(MngoCtx, bson.D{})
2024-07-17 18:02:30 +02:00
if err != nil {
2024-07-18 11:51:12 +02:00
lib.Logger.Fatal().Msg("Error contacting MongoDB\n" + err.Error())
2024-07-17 18:02:30 +02:00
}
collectionMap := make(map[string]bool)
2024-07-18 11:51:12 +02:00
for _, name := range existingCollections {
collectionMap[name] = true
}
2024-07-17 18:02:30 +02:00
// Only do the collection definition process if it doesn't already exists
// we add the collection to the collection map from mongo/mongo_utils to provide faster access to the collection
2024-07-18 11:51:12 +02:00
for _, collection_name := range list_collection {
2024-07-17 18:02:30 +02:00
new_collection := mngoDB.Collection(collection_name)
if _, exists := collectionMap[collection_name]; !exists {
createCollection(collection_name, new_collection)
2024-07-18 11:51:12 +02:00
} else {
2024-07-17 18:02:30 +02:00
CollectionMap[collection_name] = new_collection
}
2024-07-18 11:51:12 +02:00
}
2024-07-17 18:02:30 +02:00
}
// Creates the collection with index specified in mongo/mongo_collections
// or use the basic collection creation function
2024-07-18 11:51:12 +02:00
func createCollection(collection_name string, new_collection *mongo.Collection) {
var err error
2024-07-17 18:02:30 +02:00
CollectionMap[collection_name] = new_collection
2024-07-18 11:51:12 +02:00
_, exists := IndexesMap[collection_name]
if exists {
2024-07-17 18:02:30 +02:00
if _, err = new_collection.Indexes().CreateMany(MngoCtx, IndexesMap[collection_name]); err != nil {
var cmdErr mongo.CommandError
if errors.As(err, &cmdErr) && cmdErr.Code != 85 {
2024-07-18 11:51:12 +02:00
lib.Logger.Fatal().Msg("Error creating indexes for " + collection_name + " collection : \n" + err.Error())
2024-07-17 18:02:30 +02:00
panic(err)
} else if !errors.As(err, &cmdErr) {
2024-07-18 11:51:12 +02:00
lib.Logger.Fatal().Msg("Unexpected error: " + err.Error())
2024-07-17 18:02:30 +02:00
panic(err)
}
}
} else {
mngoDB.CreateCollection(MngoCtx, collection_name)
}
}
2024-07-18 11:51:12 +02:00
func DeleteOne(id string, collection_name string) (int64, error) {
filter := bson.M{"_id": GetObjIDFromString(id)}
targetDBCollection := CollectionMap[collection_name]
opts := options.Delete().SetHint(bson.D{{Key: "_id", Value: 1}})
MngoCtx, cancel = context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
result, err := targetDBCollection.DeleteOne(MngoCtx, filter, opts)
if err != nil {
lib.Logger.Error().Msg("Couldn't insert resource: " + err.Error())
return 0, err
}
return result.DeletedCount, nil
}
func UpdateOne(set map[string]interface{}, id string, collection_name string) (string, error) {
filter := bson.M{"_id": GetObjIDFromString(id)}
targetDBCollection := CollectionMap[collection_name]
MngoCtx, cancel = context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
result, err := targetDBCollection.UpdateOne(MngoCtx, filter, dbs.InputToBson(set, true))
if err != nil {
lib.Logger.Error().Msg("Couldn't insert resource: " + err.Error())
return "", err
}
return result.UpsertedID.(primitive.ObjectID).Hex(), nil
}
2024-07-17 18:02:30 +02:00
func StoreOne(obj interface{}, collection_name string) (string, error) {
targetDBCollection := CollectionMap[collection_name]
MngoCtx, cancel = context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
2024-07-18 11:51:12 +02:00
result, err := targetDBCollection.InsertOne(MngoCtx, obj)
2024-07-17 18:02:30 +02:00
if err != nil {
2024-07-18 11:51:12 +02:00
lib.Logger.Error().Msg("Couldn't insert resource: " + err.Error())
2024-07-17 18:02:30 +02:00
return "", err
}
return result.InsertedID.(primitive.ObjectID).Hex(), nil
}
2024-07-18 11:51:12 +02:00
func LoadOne(id string, collection_name string) (res *mongo.SingleResult, err error) {
2024-07-17 18:02:30 +02:00
filter := bson.M{"_id": GetObjIDFromString(id)}
targetDBCollection := CollectionMap[collection_name]
MngoCtx, cancel = context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
res = targetDBCollection.FindOne(MngoCtx, filter)
if res.Err() != nil {
2024-07-18 11:51:12 +02:00
lib.Logger.Error().Msg("Couldn't find resource " + id + ". Error : " + res.Err().Error())
2024-07-17 18:02:30 +02:00
err = res.Err()
return nil, err
}
return res, nil
}