diff --git a/backend/internal/middleware/rabbitmq/rabbitMQ.go b/backend/internal/middleware/rabbitmq/rabbitMQ.go index 4dd0d1a..5647835 100644 --- a/backend/internal/middleware/rabbitmq/rabbitMQ.go +++ b/backend/internal/middleware/rabbitmq/rabbitMQ.go @@ -14,8 +14,8 @@ import ( ) type RabbitMQ struct { - conn *amqp.Connection - ch *amqp.Channel + Conn *amqp.Connection + Ch *amqp.Channel } func NewRabbitMQ(cfg *config.RabbitMQConfig) (*RabbitMQ, error) { @@ -23,39 +23,39 @@ func NewRabbitMQ(cfg *config.RabbitMQConfig) (*RabbitMQ, error) { return nil, errors.New("rabbitmq config is nil") } url := "amqp://" + cfg.Username + ":" + cfg.Password + "@" + cfg.Host + ":" + strconv.Itoa(cfg.Port) + "/" - conn, err := amqp.Dial(url) + Conn, err := amqp.Dial(url) if err != nil { return nil, err } - ch, err := conn.Channel() + Ch, err := Conn.Channel() if err != nil { return nil, err } - return &RabbitMQ{conn: conn, ch: ch}, nil + return &RabbitMQ{Conn: Conn, Ch: Ch}, nil } func (r *RabbitMQ) Close() error { - if r == nil || r.ch == nil || r.conn == nil { + if r == nil || r.Ch == nil || r.Conn == nil { return nil } - if err := r.ch.Close(); err != nil { + if err := r.Ch.Close(); err != nil { return err } - if err := r.conn.Close(); err != nil { + if err := r.Conn.Close(); err != nil { return err } return nil } func (r *RabbitMQ) DeclareTopic(exchange string, queue string, bindingKey string) error { - if r == nil || r.ch == nil { + if r == nil || r.Ch == nil { return errors.New("rabbitmq is not initialized") } if exchange == "" || queue == "" || bindingKey == "" { return errors.New("exchange/queue/bindingKey is required") } - if err := r.ch.ExchangeDeclare( + if err := r.Ch.ExchangeDeclare( exchange, "topic", true, @@ -67,7 +67,7 @@ func (r *RabbitMQ) DeclareTopic(exchange string, queue string, bindingKey string return err } - q, err := r.ch.QueueDeclare( + q, err := r.Ch.QueueDeclare( queue, true, false, @@ -79,7 +79,7 @@ func (r *RabbitMQ) DeclareTopic(exchange string, queue string, bindingKey string return err } - return r.ch.QueueBind( + return r.Ch.QueueBind( q.Name, bindingKey, exchange, @@ -89,7 +89,7 @@ func (r *RabbitMQ) DeclareTopic(exchange string, queue string, bindingKey string } func (r *RabbitMQ) PublishJSON(ctx context.Context, exchange string, routingKey string, payload any) error { - if r == nil || r.ch == nil { + if r == nil || r.Ch == nil { return errors.New("rabbitmq is not initialized") } if exchange == "" || routingKey == "" { @@ -99,7 +99,7 @@ func (r *RabbitMQ) PublishJSON(ctx context.Context, exchange string, routingKey if err != nil { return err } - return r.ch.PublishWithContext(ctx, exchange, routingKey, false, false, amqp.Publishing{ + return r.Ch.PublishWithContext(ctx, exchange, routingKey, false, false, amqp.Publishing{ ContentType: "application/json", DeliveryMode: amqp.Persistent, Timestamp: time.Now(), diff --git a/backend/internal/middleware/rabbitmq/timelineMQ.go b/backend/internal/middleware/rabbitmq/timelineMQ.go new file mode 100644 index 0000000..f3d1fdc --- /dev/null +++ b/backend/internal/middleware/rabbitmq/timelineMQ.go @@ -0,0 +1,55 @@ +package rabbitmq + +import ( + "context" + "errors" + "time" +) + +type TimelineMQ struct { + *RabbitMQ +} + +const ( + timelineExchange = "video.timeline.events" + timelineQueue = "video.timeline.update.queue" + timelineBindingKey = "video.timeline.*" + timelinePublishRK = "video.timeline.publish" +) + +type TimelineEvent struct { + EventID string `json:"event_id"` + ViedoID uint `json:"video_id"` + CreateTime int64 `json:"create_time"` + OccurredAt time.Time `json:"occurred_at"` +} + +func NewTimelineMQ(base *RabbitMQ) (*TimelineMQ, error) { + if base == nil { + return nil, errors.New("rabbitmq base is nil") + } + if err := base.DeclareTopic(timelineExchange, timelineQueue, timelineBindingKey); err != nil { + return nil, err + } + return &TimelineMQ{RabbitMQ: base}, nil +} + +func (t *TimelineMQ) PublishVideo(ctx context.Context, videoID uint, createTime time.Time) error { + if t == nil || t.RabbitMQ == nil { + return errors.New("like mq is not initialized") + } + if videoID == 0 { + return errors.New("videoID are required") + } + id, err := newEventID(16) + if err != nil { + return err + } + timeline := TimelineEvent{ + EventID: id, + ViedoID: videoID, + CreateTime: createTime.UnixMilli(), + OccurredAt: time.Now(), + } + return t.PublishJSON(ctx, timelineExchange, timelinePublishRK, timeline) +} diff --git a/backend/internal/video/video_entity.go b/backend/internal/video/video_entity.go index 9259540..9487427 100644 --- a/backend/internal/video/video_entity.go +++ b/backend/internal/video/video_entity.go @@ -38,3 +38,11 @@ type UpdateLikesCountRequest struct { ID uint `json:"id"` LikesCount int64 `json:"likes_count"` } + +type OutboxMsg struct { + ID uint `gorm:"primaryKey"` + VideoID uint `gorm:"index"` + EventType string `gorm:"type:varchar(50)"` + CreateTime time.Time `gorm:"autoCreateTime"` + Status string `gorm:"type:varchar(50);index"` +}