package hevctaskmod import ( "fmt" "strings" "time" "unicode/utf8" "91porn-server/common/db" "91porn-server/common/log" "91porn-server/models" "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" ) const ( table = models.HevcTask maxErrorMsgLen = 2048 ) var mdb *db.MongoDB func coll(t *db.MongoTool) *db.MongoTool { if t == nil { return mdb.Coll(table) } return t.Coll(table) } func Init() { mdb = db.Init(table) initIndex() } func initIndex() { indexes := []mongo.IndexModel{ { Keys: bson.D{{Key: "videoId", Value: 1}}, Options: options.Index().SetUnique(true), }, { Keys: bson.D{{Key: "status", Value: 1}, {Key: "updatedAt", Value: -1}}, }, } if _, err := coll(nil).CreateIndex(indexes); err != nil { panic(fmt.Sprintf("%s model set index err ==>[%+v]", table, err)) } } // Upsert 按视频 ID 写入或更新一条转码任务。 func Upsert(videoID primitive.ObjectID, videoName, submitURL, fileID, h265URL string, status H265TaskStatus) error { return UpsertWithError(videoID, videoName, submitURL, fileID, h265URL, status, "") } // UpsertWithError 更新转码任务并记录最近一次错误。 func UpsertWithError(videoID primitive.ObjectID, videoName, submitURL, fileID, h265URL string, status H265TaskStatus, errorMsg string) error { now := time.Now() set := bson.M{ "videoName": videoName, "submitUrl": submitURL, "status": status, "updatedAt": now, } unset := bson.M{} if fileID != "" { set["fileId"] = fileID } if h265URL != "" { set["h265Url"] = h265URL } errorMsg = truncateErrorMsg(errorMsg) if errorMsg != "" { set["errorMsg"] = errorMsg } else { unset["errorMsg"] = "" } update := bson.M{ "$set": set, "$setOnInsert": bson.M{"videoId": videoID, "createdAt": now}, } if len(unset) > 0 { update["$unset"] = unset } if _, err := coll(nil).UpsertOne(bson.M{"videoId": videoID}, update); err != nil { log.Warn("H265 task upsert failed", log.Any("videoId", videoID), log.E(err)) return err } return nil } func truncateErrorMsg(errorMsg string) string { errorMsg = strings.TrimSpace(errorMsg) if len(errorMsg) <= maxErrorMsgLen { return errorMsg } cut := maxErrorMsgLen for cut > 0 && !utf8.ValidString(errorMsg[:cut]) { cut-- } return errorMsg[:cut] } // UpdateStatus 只更新已有任务的状态。 func UpdateStatus(videoID primitive.ObjectID, status H265TaskStatus) error { if _, err := coll(nil).UpdateOne(bson.M{"videoId": videoID}, bson.M{"$set": bson.M{ "status": status, "updatedAt": time.Now(), }}); err != nil { log.Warn("H265 task status update failed", log.Any("videoId", videoID), log.E(err)) return err } return nil }