fix: issues
This commit is contained in:
@@ -245,6 +245,7 @@ func (svc *Service) Create(ctx context.Context, tenant *model.Tenants, user *mod
|
|||||||
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
||||||
|
|
||||||
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -268,6 +269,35 @@ func (svc *Service) GetPostByHash(ctx context.Context, tenantID int64, hash stri
|
|||||||
return svc.GetPostByID(ctx, postId)
|
return svc.GetPostByID(ctx, postId)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ForceDelete
|
||||||
|
func (svc *Service) ForceDelete(ctx context.Context, tenantID, userID, postID int64) error {
|
||||||
|
_, span := otel.Start(ctx, "users.service.ForceDelete")
|
||||||
|
defer span.End()
|
||||||
|
span.SetAttributes(
|
||||||
|
attribute.Int64("tenant.id", tenantID),
|
||||||
|
attribute.Int64("user.id", userID),
|
||||||
|
attribute.Int64("post.id", postID),
|
||||||
|
)
|
||||||
|
tbl := table.Posts
|
||||||
|
|
||||||
|
stmt := tbl.
|
||||||
|
DELETE().
|
||||||
|
WHERE(
|
||||||
|
tbl.ID.EQ(Int64(postID)).AND(
|
||||||
|
tbl.TenantID.EQ(Int64(tenantID)).AND(
|
||||||
|
tbl.UserID.EQ(Int64(userID)),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
||||||
|
|
||||||
|
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
// Delete
|
// Delete
|
||||||
func (svc *Service) Delete(ctx context.Context, tenantID, userID, postID int64) error {
|
func (svc *Service) Delete(ctx context.Context, tenantID, userID, postID int64) error {
|
||||||
_, span := otel.Start(ctx, "users.service.Delete")
|
_, span := otel.Start(ctx, "users.service.Delete")
|
||||||
@@ -294,6 +324,7 @@ func (svc *Service) Delete(ctx context.Context, tenantID, userID, postID int64)
|
|||||||
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
||||||
|
|
||||||
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -330,6 +361,7 @@ func (svc *Service) Update(ctx context.Context, tenantID, userID, postID int64,
|
|||||||
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
||||||
|
|
||||||
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -348,6 +380,7 @@ func (svc *Service) AttachAssets(ctx context.Context, tenantID, userID, postID i
|
|||||||
|
|
||||||
post, err := svc.ForceGetPostByID(ctx, postID)
|
post, err := svc.ForceGetPostByID(ctx, postID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -371,6 +404,7 @@ func (svc *Service) AttachAssets(ctx context.Context, tenantID, userID, postID i
|
|||||||
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
||||||
|
|
||||||
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -389,6 +423,7 @@ func (svc *Service) UpdateMeta(ctx context.Context, tenantID, userID, postID int
|
|||||||
|
|
||||||
post, err := svc.ForceGetPostByID(ctx, postID)
|
post, err := svc.ForceGetPostByID(ctx, postID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -411,6 +446,91 @@ func (svc *Service) UpdateMeta(ctx context.Context, tenantID, userID, postID int
|
|||||||
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
||||||
|
|
||||||
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return svc.Update(ctx, tenantID, userID, postID, post)
|
||||||
|
}
|
||||||
|
|
||||||
|
// UpdateStatus
|
||||||
|
func (svc *Service) UpdateStatus(ctx context.Context, tenantID, userID, postID int64, status fields.PostStatus) error {
|
||||||
|
_, span := otel.Start(ctx, "users.service.UpdateStatus")
|
||||||
|
defer span.End()
|
||||||
|
span.SetAttributes(
|
||||||
|
attribute.Int64("tenant.id", tenantID),
|
||||||
|
attribute.Int64("user.id", userID),
|
||||||
|
attribute.Int64("post.id", postID),
|
||||||
|
)
|
||||||
|
|
||||||
|
post, err := svc.ForceGetPostByID(ctx, postID)
|
||||||
|
if err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
post.Status = status
|
||||||
|
|
||||||
|
tbl := table.Posts
|
||||||
|
stmt := tbl.
|
||||||
|
UPDATE(tbl.UpdatedAt, tbl.Status).
|
||||||
|
SET(
|
||||||
|
tbl.UpdatedAt.SET(TimestampT(time.Now())),
|
||||||
|
tbl.Status.SET(Int16(int16(status))),
|
||||||
|
).
|
||||||
|
WHERE(
|
||||||
|
tbl.ID.EQ(Int64(postID)).AND(
|
||||||
|
tbl.TenantID.EQ(Int64(tenantID)).AND(
|
||||||
|
tbl.UserID.EQ(Int64(userID)),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
||||||
|
|
||||||
|
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return svc.Update(ctx, tenantID, userID, postID, post)
|
||||||
|
}
|
||||||
|
|
||||||
|
// UpdateStage
|
||||||
|
func (svc *Service) UpdateStage(ctx context.Context, tenantID, userID, postID int64, stage fields.PostStage) error {
|
||||||
|
_, span := otel.Start(ctx, "users.service.UpdateStage")
|
||||||
|
defer span.End()
|
||||||
|
span.SetAttributes(
|
||||||
|
attribute.Int64("tenant.id", tenantID),
|
||||||
|
attribute.Int64("user.id", userID),
|
||||||
|
attribute.Int64("post.id", postID),
|
||||||
|
)
|
||||||
|
|
||||||
|
post, err := svc.ForceGetPostByID(ctx, postID)
|
||||||
|
if err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
post.Stage = stage
|
||||||
|
|
||||||
|
tbl := table.Posts
|
||||||
|
stmt := tbl.
|
||||||
|
UPDATE(tbl.UpdatedAt, tbl.Stage).
|
||||||
|
SET(
|
||||||
|
tbl.UpdatedAt.SET(TimestampT(time.Now())),
|
||||||
|
tbl.Stage.SET(Int16(int16(stage))),
|
||||||
|
).
|
||||||
|
WHERE(
|
||||||
|
tbl.ID.EQ(Int64(postID)).AND(
|
||||||
|
tbl.TenantID.EQ(Int64(tenantID)).AND(
|
||||||
|
tbl.UserID.EQ(Int64(userID)),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
|
||||||
|
|
||||||
|
if _, err := stmt.ExecContext(ctx, svc.db); err != nil {
|
||||||
|
span.RecordError(err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -7,14 +7,10 @@ import (
|
|||||||
"backend/app/http/medias"
|
"backend/app/http/medias"
|
||||||
"backend/app/http/posts"
|
"backend/app/http/posts"
|
||||||
"backend/app/http/storages"
|
"backend/app/http/storages"
|
||||||
"backend/database/fields"
|
|
||||||
"backend/database/models/qvyun_v2/public/model"
|
|
||||||
|
|
||||||
_ "git.ipao.vip/rogeecn/atom"
|
_ "git.ipao.vip/rogeecn/atom"
|
||||||
_ "git.ipao.vip/rogeecn/atom/contracts"
|
_ "git.ipao.vip/rogeecn/atom/contracts"
|
||||||
"github.com/pkg/errors"
|
|
||||||
. "github.com/riverqueue/river"
|
. "github.com/riverqueue/river"
|
||||||
"github.com/samber/lo"
|
|
||||||
"github.com/sirupsen/logrus"
|
"github.com/sirupsen/logrus"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -65,43 +61,5 @@ func (w *PostDeleteAssetsJobWorker) NextRetry(job *Job[PostDeleteAssetsJob]) tim
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (w *PostDeleteAssetsJobWorker) Work(ctx context.Context, job *Job[PostDeleteAssetsJob]) error {
|
func (w *PostDeleteAssetsJobWorker) Work(ctx context.Context, job *Job[PostDeleteAssetsJob]) error {
|
||||||
post, err := w.postSvc.ForceGetPostByID(ctx, job.Args.PostID)
|
|
||||||
if err != nil {
|
|
||||||
return errors.Wrapf(err, "failed to get post(%d) by id", job.Args.PostID)
|
|
||||||
}
|
|
||||||
|
|
||||||
mediaIDs := lo.Map(post.Assets.Data, func(asset fields.MediaAsset, _ int) int64 {
|
|
||||||
return asset.Media
|
|
||||||
})
|
|
||||||
|
|
||||||
medias, err := w.mediaSvc.GetMediasByIDs(ctx, post.TenantID, post.UserID, mediaIDs)
|
|
||||||
if err != nil {
|
|
||||||
return errors.Wrapf(err, "failed to get medias by ids(%v)", mediaIDs)
|
|
||||||
}
|
|
||||||
|
|
||||||
storageIds := lo.Map(medias, func(media *model.Medias, _ int) int64 { return media.StorageID })
|
|
||||||
storageMap, err := w.storageSvc.GetFSMapByIDs(ctx, storageIds...)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
// remove assets
|
|
||||||
for _, media := range medias {
|
|
||||||
st, ok := storageMap[media.StorageID]
|
|
||||||
if !ok {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := st.Remove(media.Path); err != nil {
|
|
||||||
return errors.Wrapf(err, "failed to remove media(%d)", media.ID)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// delete all assets
|
|
||||||
ids := lo.Map(medias, func(media *model.Medias, _ int) int64 { return media.ID })
|
|
||||||
if err := w.mediaSvc.DeleteByID(ctx, ids...); err != nil {
|
|
||||||
return errors.Wrapf(err, "failed to delete media(%d)", ids)
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -88,9 +88,8 @@ func (w *PostVideoCutJobWorker) Work(ctx context.Context, job *Job[PostVideoCutJ
|
|||||||
|
|
||||||
videoPath := media.Path
|
videoPath := media.Path
|
||||||
|
|
||||||
// 获取全长度的音频
|
|
||||||
_, ok := lo.Find(post.Assets.Data, func(asset fields.MediaAsset) bool {
|
_, ok := lo.Find(post.Assets.Data, func(asset fields.MediaAsset) bool {
|
||||||
return asset.Type == fields.MediaAssetTypeAudio && asset.Mark != nil && *asset.Mark == "audio-preview"
|
return asset.Type == fields.MediaAssetTypeAudio && asset.Mark != nil && *asset.Mark == "video-preview"
|
||||||
})
|
})
|
||||||
if ok {
|
if ok {
|
||||||
return nil
|
return nil
|
||||||
@@ -142,9 +141,26 @@ func (w *PostVideoCutJobWorker) Work(ctx context.Context, job *Job[PostVideoCutJ
|
|||||||
return errors.Wrapf(err, "attach video(%s) to post(%d) failed", videoPath, post.ID)
|
return errors.Wrapf(err, "attach video(%s) to post(%d) failed", videoPath, post.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
post.Meta.WorkerMark = post.Meta.WorkerMark & 1 << 0
|
post, err = w.postSvc.GetPostByID(ctx, job.Args.PostID)
|
||||||
if err := w.postSvc.UpdateMeta(ctx, job.Args.TenantID, job.Args.UserID, post.ID, post.Meta); err != nil {
|
if err != nil {
|
||||||
return errors.Wrapf(err, "update post(%d) meta failed", post.ID)
|
return errors.Wrapf(err, "get post(%d) failed", job.Args.PostID)
|
||||||
|
}
|
||||||
|
|
||||||
|
marks := lo.Map(post.Assets.Data, func(asset fields.MediaAsset, _ int) string {
|
||||||
|
if asset.Mark != nil {
|
||||||
|
return *asset.Mark
|
||||||
|
}
|
||||||
|
return ""
|
||||||
|
})
|
||||||
|
|
||||||
|
if items := lo.Intersect([]string{"audio-preview", "video-preview", "video", "audio"}, marks); len(items) == 4 {
|
||||||
|
if err := w.postSvc.UpdateStatus(ctx, job.Args.TenantID, job.Args.UserID, post.ID, fields.PostStatusVerified); err != nil {
|
||||||
|
return errors.Wrapf(err, "update post(%d) status failed", post.ID)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := w.postSvc.UpdateStage(ctx, job.Args.TenantID, job.Args.UserID, post.ID, fields.PostStageCompleted); err != nil {
|
||||||
|
return errors.Wrapf(err, "update post(%d) state failed", post.ID)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -145,14 +145,26 @@ func (w *PostVideoExtractAudioJobWorker) Work(ctx context.Context, job *Job[Post
|
|||||||
return errors.Wrapf(err, "attach audio(%s) to post(%d) failed", audioPath, post.ID)
|
return errors.Wrapf(err, "attach audio(%s) to post(%d) failed", audioPath, post.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
if job.Args.Mark == "audio-preview" {
|
// 检查是否可发布
|
||||||
post.Meta.WorkerMark = post.Meta.WorkerMark & 1 << 1
|
post, err = w.postSvc.GetPostByID(ctx, job.Args.PostID)
|
||||||
} else if job.Args.Mark == "audio" {
|
if err != nil {
|
||||||
post.Meta.WorkerMark = post.Meta.WorkerMark & 1 << 2
|
return errors.Wrapf(err, "get post(%d) failed", job.Args.PostID)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := w.postSvc.UpdateMeta(ctx, job.Args.TenantID, job.Args.UserID, post.ID, post.Meta); err != nil {
|
marks := lo.Map(post.Assets.Data, func(asset fields.MediaAsset, _ int) string {
|
||||||
return errors.Wrapf(err, "update post(%d) meta failed", post.ID)
|
if asset.Mark != nil {
|
||||||
|
return *asset.Mark
|
||||||
|
}
|
||||||
|
return ""
|
||||||
|
})
|
||||||
|
if items := lo.Intersect([]string{"audio-preview", "video-preview", "video", "audio"}, marks); len(items) == 4 {
|
||||||
|
if err := w.postSvc.UpdateStatus(ctx, job.Args.TenantID, job.Args.UserID, post.ID, fields.PostStatusVerified); err != nil {
|
||||||
|
return errors.Wrapf(err, "update post(%d) status failed", post.ID)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := w.postSvc.UpdateStage(ctx, job.Args.TenantID, job.Args.UserID, post.ID, fields.PostStageCompleted); err != nil {
|
||||||
|
return errors.Wrapf(err, "update post(%d) state failed", post.ID)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -7,7 +7,8 @@
|
|||||||
"path": "frontend"
|
"path": "frontend"
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"path": "backend"
|
"path": "backend/app",
|
||||||
|
"name":"BACKEND_APP"
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user