feat: update medias

This commit is contained in:
Rogee
2025-01-17 14:59:54 +08:00
parent d72f384177
commit b5583bb34a
46 changed files with 1856 additions and 119 deletions

View File

@@ -74,7 +74,7 @@ func (e *PostCreated) Handler(msg *message.Message) ([]*message.Message, error)
// cut video
_, err = job.Insert(context.Background(), jobs.PostVideoCutJob{
PostID: post.ID,
Hash: video.Hash,
MediaID: video.Media,
TenantID: post.TenantID,
UserID: post.UserID,
}, nil)
@@ -85,7 +85,7 @@ func (e *PostCreated) Handler(msg *message.Message) ([]*message.Message, error)
// extract audio
_, err = job.Insert(context.Background(), jobs.PostVideoExtractAudioJob{
PostID: post.ID,
Hash: video.Hash,
MediaID: video.Media,
TenantID: post.TenantID,
UserID: post.UserID,
Mark: "audio-preview",
@@ -96,7 +96,7 @@ func (e *PostCreated) Handler(msg *message.Message) ([]*message.Message, error)
_, err = job.Insert(context.Background(), jobs.PostVideoExtractAudioJob{
PostID: post.ID,
Hash: video.Hash,
MediaID: video.Media,
TenantID: post.TenantID,
UserID: post.UserID,
Mark: "audio",

View File

@@ -52,17 +52,9 @@ func (ctl *Controller) Upload(ctx fiber.Ctx, claim *jwt.Claims, file *multipart.
return uploadedFile, nil
}
uploadedFile, err = storage.Build(defaultStorage).Save(ctx.Context(), uploadedFile)
if err != nil {
return nil, err
}
// save to db
_, err = ctl.svc.Create(ctx.Context(), &model.Medias{
userMediaID, err := ctl.svc.Create(ctx.Context(), *claim.TenantID, claim.UserID, &model.Medias{
CreatedAt: time.Now(),
UpdatedAt: time.Now(),
TenantID: *claim.TenantID,
UserID: claim.UserID,
StorageID: defaultStorage.ID,
Hash: uploadedFile.Hash,
Name: uploadedFile.Name,
@@ -70,6 +62,7 @@ func (ctl *Controller) Upload(ctx fiber.Ctx, claim *jwt.Claims, file *multipart.
Size: uploadedFile.Size,
Path: uploadedFile.Path,
})
uploadedFile.ID = userMediaID
uploadedFile.Preview = ""
return uploadedFile, err

View File

@@ -10,6 +10,7 @@ import (
"backend/providers/otel"
. "github.com/go-jet/jet/v2/postgres"
"github.com/go-jet/jet/v2/qrm"
"github.com/samber/lo"
log "github.com/sirupsen/logrus"
semconv "go.opentelemetry.io/otel/semconv/v1.4.0"
@@ -28,7 +29,7 @@ func (svc *Service) Prepare() error {
}
// Create
func (svc *Service) Create(ctx context.Context, m *model.Medias) (*model.Medias, error) {
func (svc *Service) Create(ctx context.Context, tenantID, userID int64, m *model.Medias) (int64, error) {
_, span := otel.Start(ctx, "medias.service.Create")
defer span.End()
@@ -36,40 +37,50 @@ func (svc *Service) Create(ctx context.Context, m *model.Medias) (*model.Medias,
m.CreatedAt = time.Now()
}
if m.UpdatedAt.IsZero() {
m.UpdatedAt = time.Now()
}
tbl := table.Medias
stmt := tbl.INSERT(tbl.MutableColumns).MODEL(m).RETURNING(tbl.AllColumns)
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
var ret model.Medias
if err := stmt.QueryContext(ctx, svc.db, &ret); err != nil {
return nil, err
var media model.Medias
if err := stmt.QueryContext(ctx, svc.db, &media); err != nil {
return 0, err
}
return &ret, nil
userMediaTbl := table.UserMedias
userMediaStmt := userMediaTbl.INSERT(userMediaTbl.TenantID, userMediaTbl.UserID, userMediaTbl.MediaID).VALUES(tenantID, userID, media.ID).RETURNING(tbl.AllColumns)
span.SetAttributes(semconv.DBStatementKey.String(userMediaStmt.DebugSql()))
var ret model.UserMedias
if err := userMediaStmt.QueryContext(ctx, svc.db, &ret); err != nil {
return 0, err
}
return ret.ID, nil
}
// GetMediasByHash
func (svc *Service) GetMediasByHash(ctx context.Context, tenantID, userID int64, hashes []string) ([]*model.Medias, error) {
_, span := otel.Start(ctx, "medias.service.GetMediasByHash")
func (svc *Service) GetMediasByIDs(ctx context.Context, tenantID, userID int64, ids []int64) ([]*model.Medias, error) {
_, span := otel.Start(ctx, "medias.service.GetMediasByIDs")
defer span.End()
hashExpr := lo.Map(hashes, func(item string, index int) Expression { return String(item) })
idExprs := lo.Map(ids, func(item int64, index int) Expression { return Int64(item) })
tbl := table.Medias
stmt := tbl.
SELECT(tbl.AllColumns).
WHERE(
tbl.TenantID.
EQ(Int64(tenantID)).
AND(
tbl.UserID.EQ(Int64(userID)),
).
AND(
tbl.Hash.IN(hashExpr...),
),
stmt := SELECT(table.Medias.AllColumns).
FROM(
table.Medias.RIGHT_JOIN(
table.UserMedias,
table.UserMedias.TenantID.
EQ(Int64(tenantID)).
AND(
table.UserMedias.UserID.EQ(Int64(userID)),
).
AND(
table.UserMedias.ID.IN(idExprs...),
).
AND(
table.Medias.ID.EQ(table.UserMedias.MediaID),
),
),
)
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
@@ -82,31 +93,20 @@ func (svc *Service) GetMediasByHash(ctx context.Context, tenantID, userID int64,
}), nil
}
func (svc *Service) GetMediaByHash(ctx context.Context, tenantID, userID int64, hash string) (*model.Medias, error) {
func (svc *Service) GetMediaByID(ctx context.Context, tenantID, userID, userMediaID int64) (*model.Medias, error) {
_, span := otel.Start(ctx, "medias.service.GetMediasByHash")
defer span.End()
tbl := table.Medias
stmt := tbl.
SELECT(tbl.AllColumns).
LIMIT(1).
WHERE(
tbl.TenantID.
EQ(Int64(tenantID)).
AND(
tbl.UserID.EQ(Int64(userID)),
).
AND(
tbl.Hash.EQ(String(hash)),
),
)
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))
var ret model.Medias
if err := stmt.QueryContext(ctx, svc.db, &ret); err != nil {
medias, err := svc.GetMediasByIDs(ctx, tenantID, userID, []int64{userMediaID})
if err != nil {
return nil, err
}
return &ret, nil
if len(medias) == 0 {
return nil, qrm.ErrNoRows
}
return medias[0], nil
}
func (svc *Service) DeleteByID(ctx context.Context, id ...int64) error {
@@ -117,7 +117,7 @@ func (svc *Service) DeleteByID(ctx context.Context, id ...int64) error {
_, span := otel.Start(ctx, "medias.service.DeleteByID")
defer span.End()
tbl := table.Medias
tbl := table.UserMedias
stmt := tbl.DELETE().WHERE(tbl.ID.IN(lo.Map(id, func(item int64, _ int) Expression { return Int64(item) })...))
span.SetAttributes(semconv.DBStatementKey.String(stmt.DebugSql()))

View File

@@ -138,13 +138,13 @@ func (ctl *Controller) Create(ctx fiber.Ctx, claim *jwt.Claims, tenantSlug strin
}
// check media assets exists
hashes := lo.Map(body.Assets.Data, func(item fields.MediaAsset, _ int) string { return item.Hash })
medias, err := ctl.mediaSvc.GetMediasByHash(ctx.Context(), tenant.ID, user.ID, hashes)
ids := lo.Map(body.Assets.Data, func(item fields.MediaAsset, _ int) int64 { return item.Media })
medias, err := ctl.mediaSvc.GetMediasByIDs(ctx.Context(), tenant.ID, user.ID, ids)
if err != nil {
return err
}
if len(medias) != len(lo.Uniq(hashes)) {
if len(medias) != len(lo.Uniq(ids)) {
return errorx.BadRequest
}
@@ -218,13 +218,13 @@ func (ctl *Controller) Update(ctx fiber.Ctx, claim *jwt.Claims, hash string, bod
}
// check media assets exists
hashes := lo.Map(body.Assets.Data, func(item fields.MediaAsset, _ int) string { return item.Hash })
medias, err := ctl.mediaSvc.GetMediasByHash(ctx.Context(), *claim.TenantID, post.UserID, hashes)
ids := lo.Map(body.Assets.Data, func(item fields.MediaAsset, _ int) int64 { return item.Media })
medias, err := ctl.mediaSvc.GetMediasByIDs(ctx.Context(), *claim.TenantID, post.UserID, ids)
if err != nil {
return err
}
if len(medias) != len(lo.Uniq(hashes)) {
if len(medias) != len(lo.Uniq(ids)) {
return errorx.BadRequest
}
@@ -247,5 +247,6 @@ func (ctl *Controller) Update(ctx fiber.Ctx, claim *jwt.Claims, hash string, bod
if err := ctl.svc.Update(ctx.Context(), post.TenantID, post.UserID, post.ID, m); err != nil {
return err
}
// todo: trigger event post updated
return nil
}

View File

@@ -70,13 +70,13 @@ func (w *PostDeleteAssetsJobWorker) Work(ctx context.Context, job *Job[PostDelet
return errors.Wrapf(err, "failed to get post(%d) by id", job.Args.PostID)
}
hashes := lo.Map(post.Assets.Data, func(asset fields.MediaAsset, _ int) string {
return asset.Hash
mediaIDs := lo.Map(post.Assets.Data, func(asset fields.MediaAsset, _ int) int64 {
return asset.Media
})
medias, err := w.mediaSvc.GetMediasByHash(ctx, post.TenantID, post.UserID, hashes)
medias, err := w.mediaSvc.GetMediasByIDs(ctx, post.TenantID, post.UserID, mediaIDs)
if err != nil {
return errors.Wrapf(err, "failed to get medias by hashes(%v)", hashes)
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 })

View File

@@ -33,7 +33,7 @@ type PostVideoCutJob struct {
PostID int64
TenantID int64
UserID int64
Hash string
MediaID int64
}
// InsertOpts implements JobArgsWithInsertOpts.
@@ -81,9 +81,9 @@ func (w *PostVideoCutJobWorker) Work(ctx context.Context, job *Job[PostVideoCutJ
return errors.Wrapf(err, "get post(%d) failed", job.Args.PostID)
}
media, err := w.mediaSvc.GetMediaByHash(ctx, job.Args.TenantID, job.Args.UserID, job.Args.Hash)
media, err := w.mediaSvc.GetMediaByID(ctx, job.Args.TenantID, job.Args.UserID, job.Args.MediaID)
if err != nil {
return errors.Wrapf(err, "get media by hash(%s) failed", job.Args.Hash)
return errors.Wrapf(err, "get media by user_media id(%s) failed", job.Args.MediaID)
}
videoPath := media.Path
@@ -108,7 +108,7 @@ func (w *PostVideoCutJobWorker) Work(ctx context.Context, job *Job[PostVideoCutJ
return errors.Wrapf(err, "get preview video(%s) file md5 failed", previewVideoPath)
}
if err := os.Rename(previewVideoPath, strings.Replace(videoPath, job.Args.Hash, fileMd5, 1)); err != nil {
if err := os.Rename(previewVideoPath, strings.Replace(videoPath, media.Hash, fileMd5, 1)); err != nil {
return errors.Wrapf(err, "rename video(%s) file failed", videoPath)
}
@@ -118,10 +118,7 @@ func (w *PostVideoCutJobWorker) Work(ctx context.Context, job *Job[PostVideoCutJ
}
// save to medias
_, err = w.mediaSvc.Create(ctx, &model.Medias{
TenantID: job.Args.TenantID,
UserID: job.Args.UserID,
PostID: post.ID,
mediaID, err := w.mediaSvc.Create(ctx, job.Args.TenantID, job.Args.UserID, &model.Medias{
StorageID: storage.ID,
Hash: fileMd5,
Name: post.Title,
@@ -135,9 +132,9 @@ func (w *PostVideoCutJobWorker) Work(ctx context.Context, job *Job[PostVideoCutJ
assets := []fields.MediaAsset{
{
Type: fields.MediaAssetTypeVideo,
Hash: fileMd5,
Mark: lo.ToPtr("video-preview"),
Type: fields.MediaAssetTypeVideo,
Media: mediaID,
Mark: lo.ToPtr("video-preview"),
},
}

View File

@@ -33,7 +33,7 @@ type PostVideoExtractAudioJob struct {
PostID int64
TenantID int64
UserID int64
Hash string
MediaID int64
Mark string
}
@@ -80,9 +80,9 @@ func (w *PostVideoExtractAudioJobWorker) Work(ctx context.Context, job *Job[Post
return errors.Wrapf(err, "get post(%d) failed", job.Args.PostID)
}
media, err := w.mediaSvc.GetMediaByHash(ctx, job.Args.TenantID, job.Args.UserID, job.Args.Hash)
media, err := w.mediaSvc.GetMediaByID(ctx, job.Args.TenantID, job.Args.UserID, job.Args.MediaID)
if err != nil {
return errors.Wrapf(err, "get media by hash(%s) failed", job.Args.Hash)
return errors.Wrapf(err, "get media by user_media id(%s) failed", job.Args.MediaID)
}
videoPath := media.Path
@@ -111,7 +111,7 @@ func (w *PostVideoExtractAudioJobWorker) Work(ctx context.Context, job *Job[Post
return errors.Wrapf(err, "get audio(%s) file md5 failed", audioPath)
}
if err := os.Rename(audioPath, strings.Replace(audioPath, job.Args.Hash, fileMd5, 1)); err != nil {
if err := os.Rename(audioPath, strings.Replace(audioPath, media.Hash, fileMd5, 1)); err != nil {
return errors.Wrapf(err, "rename audio(%s) file failed", audioPath)
}
@@ -121,10 +121,7 @@ func (w *PostVideoExtractAudioJobWorker) Work(ctx context.Context, job *Job[Post
}
// save to medias
_, err = w.mediaSvc.Create(ctx, &model.Medias{
TenantID: job.Args.TenantID,
UserID: job.Args.UserID,
PostID: post.ID,
mediaID, err := w.mediaSvc.Create(ctx, job.Args.TenantID, job.Args.UserID, &model.Medias{
StorageID: storage.ID,
Hash: fileMd5,
Name: post.Title,
@@ -138,9 +135,9 @@ func (w *PostVideoExtractAudioJobWorker) Work(ctx context.Context, job *Job[Post
assets := []fields.MediaAsset{
{
Type: fields.MediaAssetTypeAudio,
Hash: fileMd5,
Mark: lo.ToPtr(job.Args.Mark),
Type: fields.MediaAssetTypeAudio,
Media: mediaID,
Mark: lo.ToPtr(job.Args.Mark),
},
}

View File

@@ -1,9 +1,9 @@
package fields
type MediaAsset struct {
Type MediaAssetType `json:"type"`
Hash string `json:"hash"`
Mark *string `json:"mark,omitempty"`
Type MediaAssetType `json:"type"`
Media int64 `json:"media"`
Mark *string `json:"mark,omitempty"`
}
// swagger:enum MediaAssetType

View File

@@ -4,11 +4,6 @@
CREATE TABLE medias (
id SERIAL8 PRIMARY KEY,
created_at timestamp NOT NULL default now(),
updated_at timestamp NOT NULL default now(),
tenant_id INT8 NOT NULL,
user_id INT8 NOT NULL,
post_id INT8 NOT NULL,
storage_id INT8 NOT NULL,
hash VARCHAR(32) NOT NULL,
name VARCHAR(255) NOT NULL default '',
@@ -17,13 +12,24 @@ CREATE TABLE medias (
path VARCHAR(255) NOT NULL default ''
);
CREATE INDEX medias_tenant_id_index ON medias (tenant_id);
CREATE INDEX medias_user_id_index ON medias (user_id);
CREATE INDEX medias_post_id_index ON medias (post_id);
CREATE INDEX medias_storage_id_index ON medias (storage_id);
-- index
CREATE UNIQUE INDEX medias_hash_idx ON medias (hash);
-- user medias
CREATE TABLE user_medias (
id SERIAL8 PRIMARY KEY,
created_at timestamp NOT NULL default now(),
updated_at timestamp NOT NULL default now(),
tenant_id INT8 NOT NULL,
user_id INT8 NOT NULL,
media_id INT8 NOT NULL
)
-- +goose StatementEnd
-- +goose Down
-- +goose StatementBegin
DROP TABLE medias;
DROP TABLE user_medias;
-- +goose StatementEnd

View File

@@ -14,10 +14,6 @@ import (
type Medias struct {
ID int64 `sql:"primary_key" json:"id"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
TenantID int64 `json:"tenant_id"`
UserID int64 `json:"user_id"`
PostID int64 `json:"post_id"`
StorageID int64 `json:"storage_id"`
Hash string `json:"hash"`
Name string `json:"name"`

View File

@@ -0,0 +1,21 @@
//
// Code generated by go-jet DO NOT EDIT.
//
// WARNING: Changes to this file may cause incorrect behavior
// and will be lost if the code is regenerated
//
package model
import (
"time"
)
type UserMedias struct {
ID int64 `sql:"primary_key" json:"id"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
TenantID int64 `json:"tenant_id"`
UserID int64 `json:"user_id"`
MediaID int64 `json:"media_id"`
}

View File

@@ -19,10 +19,6 @@ type mediasTable struct {
// Columns
ID postgres.ColumnInteger
CreatedAt postgres.ColumnTimestamp
UpdatedAt postgres.ColumnTimestamp
TenantID postgres.ColumnInteger
UserID postgres.ColumnInteger
PostID postgres.ColumnInteger
StorageID postgres.ColumnInteger
Hash postgres.ColumnString
Name postgres.ColumnString
@@ -71,18 +67,14 @@ func newMediasTableImpl(schemaName, tableName, alias string) mediasTable {
var (
IDColumn = postgres.IntegerColumn("id")
CreatedAtColumn = postgres.TimestampColumn("created_at")
UpdatedAtColumn = postgres.TimestampColumn("updated_at")
TenantIDColumn = postgres.IntegerColumn("tenant_id")
UserIDColumn = postgres.IntegerColumn("user_id")
PostIDColumn = postgres.IntegerColumn("post_id")
StorageIDColumn = postgres.IntegerColumn("storage_id")
HashColumn = postgres.StringColumn("hash")
NameColumn = postgres.StringColumn("name")
MimeTypeColumn = postgres.StringColumn("mime_type")
SizeColumn = postgres.IntegerColumn("size")
PathColumn = postgres.StringColumn("path")
allColumns = postgres.ColumnList{IDColumn, CreatedAtColumn, UpdatedAtColumn, TenantIDColumn, UserIDColumn, PostIDColumn, StorageIDColumn, HashColumn, NameColumn, MimeTypeColumn, SizeColumn, PathColumn}
mutableColumns = postgres.ColumnList{CreatedAtColumn, UpdatedAtColumn, TenantIDColumn, UserIDColumn, PostIDColumn, StorageIDColumn, HashColumn, NameColumn, MimeTypeColumn, SizeColumn, PathColumn}
allColumns = postgres.ColumnList{IDColumn, CreatedAtColumn, StorageIDColumn, HashColumn, NameColumn, MimeTypeColumn, SizeColumn, PathColumn}
mutableColumns = postgres.ColumnList{CreatedAtColumn, StorageIDColumn, HashColumn, NameColumn, MimeTypeColumn, SizeColumn, PathColumn}
)
return mediasTable{
@@ -91,10 +83,6 @@ func newMediasTableImpl(schemaName, tableName, alias string) mediasTable {
//Columns
ID: IDColumn,
CreatedAt: CreatedAtColumn,
UpdatedAt: UpdatedAtColumn,
TenantID: TenantIDColumn,
UserID: UserIDColumn,
PostID: PostIDColumn,
StorageID: StorageIDColumn,
Hash: HashColumn,
Name: NameColumn,

View File

@@ -18,6 +18,7 @@ func UseSchema(schema string) {
TenantUsers = TenantUsers.FromSchema(schema)
Tenants = Tenants.FromSchema(schema)
UserBoughtPosts = UserBoughtPosts.FromSchema(schema)
UserMedias = UserMedias.FromSchema(schema)
UserOauths = UserOauths.FromSchema(schema)
Users = Users.FromSchema(schema)
}

View File

@@ -0,0 +1,90 @@
//
// Code generated by go-jet DO NOT EDIT.
//
// WARNING: Changes to this file may cause incorrect behavior
// and will be lost if the code is regenerated
//
package table
import (
"github.com/go-jet/jet/v2/postgres"
)
var UserMedias = newUserMediasTable("public", "user_medias", "")
type userMediasTable struct {
postgres.Table
// Columns
ID postgres.ColumnInteger
CreatedAt postgres.ColumnTimestamp
UpdatedAt postgres.ColumnTimestamp
TenantID postgres.ColumnInteger
UserID postgres.ColumnInteger
MediaID postgres.ColumnInteger
AllColumns postgres.ColumnList
MutableColumns postgres.ColumnList
}
type UserMediasTable struct {
userMediasTable
EXCLUDED userMediasTable
}
// AS creates new UserMediasTable with assigned alias
func (a UserMediasTable) AS(alias string) *UserMediasTable {
return newUserMediasTable(a.SchemaName(), a.TableName(), alias)
}
// Schema creates new UserMediasTable with assigned schema name
func (a UserMediasTable) FromSchema(schemaName string) *UserMediasTable {
return newUserMediasTable(schemaName, a.TableName(), a.Alias())
}
// WithPrefix creates new UserMediasTable with assigned table prefix
func (a UserMediasTable) WithPrefix(prefix string) *UserMediasTable {
return newUserMediasTable(a.SchemaName(), prefix+a.TableName(), a.TableName())
}
// WithSuffix creates new UserMediasTable with assigned table suffix
func (a UserMediasTable) WithSuffix(suffix string) *UserMediasTable {
return newUserMediasTable(a.SchemaName(), a.TableName()+suffix, a.TableName())
}
func newUserMediasTable(schemaName, tableName, alias string) *UserMediasTable {
return &UserMediasTable{
userMediasTable: newUserMediasTableImpl(schemaName, tableName, alias),
EXCLUDED: newUserMediasTableImpl("", "excluded", ""),
}
}
func newUserMediasTableImpl(schemaName, tableName, alias string) userMediasTable {
var (
IDColumn = postgres.IntegerColumn("id")
CreatedAtColumn = postgres.TimestampColumn("created_at")
UpdatedAtColumn = postgres.TimestampColumn("updated_at")
TenantIDColumn = postgres.IntegerColumn("tenant_id")
UserIDColumn = postgres.IntegerColumn("user_id")
MediaIDColumn = postgres.IntegerColumn("media_id")
allColumns = postgres.ColumnList{IDColumn, CreatedAtColumn, UpdatedAtColumn, TenantIDColumn, UserIDColumn, MediaIDColumn}
mutableColumns = postgres.ColumnList{CreatedAtColumn, UpdatedAtColumn, TenantIDColumn, UserIDColumn, MediaIDColumn}
)
return userMediasTable{
Table: postgres.NewTable(schemaName, tableName, alias, allColumns...),
//Columns
ID: IDColumn,
CreatedAt: CreatedAtColumn,
UpdatedAt: UpdatedAtColumn,
TenantID: TenantIDColumn,
UserID: UserIDColumn,
MediaID: MediaIDColumn,
AllColumns: allColumns,
MutableColumns: mutableColumns,
}
}

View File

@@ -28,6 +28,7 @@ type Uploader struct {
}
type UploadedFile struct {
ID int64 `json:"id"`
Hash string `json:"hash"`
Name string `json:"name"`
Size int64 `json:"size"`