From 7d092437d6dace79356750c01182b04f5f56d328 Mon Sep 17 00:00:00 2001 From: zhx Date: Tue, 21 Jul 2026 14:09:14 +0800 Subject: [PATCH] feat(agency): notify hosts when removed --- .../configs/config.docker.yaml | 1 + .../configs/config.tencent.example.yaml | 1 + services/activity-service/configs/config.yaml | 1 + .../app/agency_host_removed_notice.go | 132 ++++++++++++++++++ services/activity-service/internal/app/mq.go | 12 ++ .../internal/config/config.go | 33 +++-- .../storage/mysql/host/applications.go | 54 +++++++ 7 files changed, 220 insertions(+), 14 deletions(-) create mode 100644 services/activity-service/internal/app/agency_host_removed_notice.go diff --git a/services/activity-service/configs/config.docker.yaml b/services/activity-service/configs/config.docker.yaml index 856c3eca..e2f46850 100644 --- a/services/activity-service/configs/config.docker.yaml +++ b/services/activity-service/configs/config.docker.yaml @@ -117,6 +117,7 @@ rocketmq: task_consumer_group: "hyapp-activity-task-user-outbox" invite_activity_consumer_group: "hyapp-activity-invite-user-outbox" message_action_consumer_group: "hyapp-activity-message-action-user-outbox" + agency_host_removed_notice_consumer_group: "hyapp-activity-agency-host-removed-notice" consumer_max_reconsume_times: 16 activity_template_outbox: enabled: true diff --git a/services/activity-service/configs/config.tencent.example.yaml b/services/activity-service/configs/config.tencent.example.yaml index 24f2438b..1346b1a3 100644 --- a/services/activity-service/configs/config.tencent.example.yaml +++ b/services/activity-service/configs/config.tencent.example.yaml @@ -118,6 +118,7 @@ rocketmq: task_consumer_group: "hyapp-activity-task-user-outbox" invite_activity_consumer_group: "hyapp-activity-invite-user-outbox" message_action_consumer_group: "hyapp-activity-message-action-user-outbox" + agency_host_removed_notice_consumer_group: "hyapp-activity-agency-host-removed-notice" consumer_max_reconsume_times: 16 activity_template_outbox: enabled: true diff --git a/services/activity-service/configs/config.yaml b/services/activity-service/configs/config.yaml index dba33362..cb28a059 100644 --- a/services/activity-service/configs/config.yaml +++ b/services/activity-service/configs/config.yaml @@ -118,6 +118,7 @@ rocketmq: task_consumer_group: "hyapp-activity-task-user-outbox" invite_activity_consumer_group: "hyapp-activity-invite-user-outbox" message_action_consumer_group: "hyapp-activity-message-action-user-outbox" + agency_host_removed_notice_consumer_group: "hyapp-activity-agency-host-removed-notice" consumer_max_reconsume_times: 16 activity_template_outbox: enabled: true diff --git a/services/activity-service/internal/app/agency_host_removed_notice.go b/services/activity-service/internal/app/agency_host_removed_notice.go new file mode 100644 index 00000000..a3a12649 --- /dev/null +++ b/services/activity-service/internal/app/agency_host_removed_notice.go @@ -0,0 +1,132 @@ +package app + +import ( + "context" + "encoding/json" + "fmt" + "strconv" + "strings" + + "hyapp/pkg/appcode" + "hyapp/pkg/rocketmqx" + "hyapp/pkg/usermq" + "hyapp/pkg/xerr" + "hyapp/services/activity-service/internal/config" + messageservice "hyapp/services/activity-service/internal/service/message" +) + +const userOutboxEventAgencyHostRemoved = "AgencyHostRemoved" + +type agencyHostRemovedUserSnapshot struct { + UserID int64 `json:"user_id"` + DisplayUserID string `json:"display_user_id"` + Username string `json:"username"` +} + +type agencyHostRemovedEvent struct { + AppCode string + EventID string + MembershipID int64 + AgencyID int64 + HostUserID int64 + Operator agencyHostRemovedUserSnapshot + RemovedAtMS int64 +} + +func newAgencyHostRemovedNoticeUserOutboxConsumer(cfg config.Config, services *serviceBundle) (*rocketmqx.Consumer, error) { + consumer, err := rocketmqx.NewConsumer(userOutboxAgencyHostRemovedNoticeConsumerConfig(cfg.RocketMQ)) + if err != nil { + return nil, err + } + if err := consumer.Subscribe(cfg.RocketMQ.UserOutbox.Topic, usermq.TagUserOutboxEvent, func(ctx context.Context, message rocketmqx.ConsumedMessage) error { + event, ok, err := agencyHostRemovedEventFromUserMessage(message.Body) + if err != nil || !ok { + return err + } + return createAgencyHostRemovedSystemNotice(appcode.WithContext(ctx, event.AppCode), services.message, event) + }); err != nil { + _ = consumer.Shutdown() + return nil, err + } + return consumer, nil +} + +func agencyHostRemovedEventFromUserMessage(body []byte) (agencyHostRemovedEvent, bool, error) { + message, err := usermq.DecodeUserOutboxMessage(body) + if err != nil { + return agencyHostRemovedEvent{}, false, err + } + if message.EventType != userOutboxEventAgencyHostRemoved { + return agencyHostRemovedEvent{}, false, nil + } + var payload struct { + MembershipID int64 `json:"membership_id"` + AgencyID int64 `json:"agency_id"` + HostUserID int64 `json:"host_user_id"` + Operator agencyHostRemovedUserSnapshot `json:"operator"` + RemovedAtMS int64 `json:"removed_at_ms"` + } + if err := json.Unmarshal([]byte(message.PayloadJSON), &payload); err != nil { + // 目标事件损坏必须交回 MQ 重试并告警;确认消息会永久丢失主播可见的移除事实。 + return agencyHostRemovedEvent{}, true, err + } + payload.Operator.DisplayUserID = strings.TrimSpace(payload.Operator.DisplayUserID) + payload.Operator.Username = strings.TrimSpace(payload.Operator.Username) + if payload.MembershipID <= 0 || payload.AgencyID <= 0 || payload.HostUserID <= 0 || + payload.Operator.UserID <= 0 || payload.Operator.DisplayUserID == "" || payload.RemovedAtMS <= 0 || + message.AggregateType != "agency_membership" || message.AggregateID != payload.MembershipID { + return agencyHostRemovedEvent{}, true, xerr.New(xerr.InvalidArgument, "agency host removed event is invalid") + } + return agencyHostRemovedEvent{ + AppCode: appcode.Normalize(message.AppCode), + EventID: message.EventID, + MembershipID: payload.MembershipID, + AgencyID: payload.AgencyID, + HostUserID: payload.HostUserID, + Operator: payload.Operator, + RemovedAtMS: payload.RemovedAtMS, + }, true, nil +} + +func createAgencyHostRemovedSystemNotice(ctx context.Context, service *messageservice.Service, event agencyHostRemovedEvent) error { + if service == nil { + return xerr.New(xerr.Unavailable, "message service is not configured") + } + // 昵称为空时只展示短号,避免形成双空格;正常资料同时展示“当前靓号/短 ID + 昵称”。 + operatorLabel := strings.Join(nonEmptyStrings(event.Operator.DisplayUserID, event.Operator.Username), " ") + content := fmt.Sprintf("你已被 %s agency移除", operatorLabel) + _, err := service.CreateSystemNotice(ctx, messageservice.NoticeCommand{ + TargetUserID: event.HostUserID, + Producer: "user-service", + ProducerEventID: event.EventID, + ProducerEventType: userOutboxEventAgencyHostRemoved, + AggregateType: "agency_membership", + AggregateID: strconv.FormatInt(event.MembershipID, 10), + TemplateID: "agency_host_removed", + TemplateVersion: "v1", + Title: "Agency移除通知", + Summary: content, + Body: content, + Priority: 50, + SentAtMS: event.RemovedAtMS, + Metadata: map[string]any{ + "agency_id": event.AgencyID, + "membership_id": event.MembershipID, + "operator_user_id": event.Operator.UserID, + "operator_display_user_id": event.Operator.DisplayUserID, + "operator_username": event.Operator.Username, + "removed_at_ms": event.RemovedAtMS, + }, + }) + return err +} + +func nonEmptyStrings(values ...string) []string { + items := make([]string, 0, len(values)) + for _, value := range values { + if value = strings.TrimSpace(value); value != "" { + items = append(items, value) + } + } + return items +} diff --git a/services/activity-service/internal/app/mq.go b/services/activity-service/internal/app/mq.go index 6a83b720..a22cb532 100644 --- a/services/activity-service/internal/app/mq.go +++ b/services/activity-service/internal/app/mq.go @@ -113,6 +113,14 @@ func buildMQConsumers(cfg config.Config, services *serviceBundle) ([]*rocketmqx. } mqConsumers = append(mqConsumers, consumer) } + if cfg.RocketMQ.UserOutbox.Enabled { + consumer, err := newAgencyHostRemovedNoticeUserOutboxConsumer(cfg, services) + if err != nil { + shutdownConsumers(mqConsumers) + return nil, err + } + mqConsumers = append(mqConsumers, consumer) + } if cfg.RedPacketBroadcastWorker.Enabled && cfg.RocketMQ.WalletOutbox.Enabled { consumer, err := newRedPacketWalletConsumer(cfg, services) if err != nil { @@ -548,6 +556,10 @@ func userOutboxMessageActionConsumerConfig(cfg config.RocketMQConfig) rocketmqx. return rocketMQConsumerConfig(cfg, cfg.UserOutbox.MessageActionConsumerGroup, cfg.UserOutbox.ConsumerMaxReconsumeTimes) } +func userOutboxAgencyHostRemovedNoticeConsumerConfig(cfg config.RocketMQConfig) rocketmqx.ConsumerConfig { + return rocketMQConsumerConfig(cfg, cfg.UserOutbox.AgencyHostRemovedNoticeConsumerGroup, cfg.UserOutbox.ConsumerMaxReconsumeTimes) +} + func walletOutboxConsumerConfig(cfg config.RocketMQConfig, group string) rocketmqx.ConsumerConfig { return rocketMQConsumerConfig(cfg, group, cfg.WalletOutbox.ConsumerMaxReconsumeTimes) } diff --git a/services/activity-service/internal/config/config.go b/services/activity-service/internal/config/config.go index 62bf169f..337bcacb 100644 --- a/services/activity-service/internal/config/config.go +++ b/services/activity-service/internal/config/config.go @@ -261,13 +261,14 @@ type WalletOutboxMQConfig struct { // UserOutboxMQConfig 控制 activity 对 user_outbox topic 的独立消费组。 type UserOutboxMQConfig struct { - Enabled bool `yaml:"enabled"` - Topic string `yaml:"topic"` - ConsumerGroup string `yaml:"consumer_group"` - TaskConsumerGroup string `yaml:"task_consumer_group"` - InviteActivityConsumerGroup string `yaml:"invite_activity_consumer_group"` - MessageActionConsumerGroup string `yaml:"message_action_consumer_group"` - ConsumerMaxReconsumeTimes int32 `yaml:"consumer_max_reconsume_times"` + Enabled bool `yaml:"enabled"` + Topic string `yaml:"topic"` + ConsumerGroup string `yaml:"consumer_group"` + TaskConsumerGroup string `yaml:"task_consumer_group"` + InviteActivityConsumerGroup string `yaml:"invite_activity_consumer_group"` + MessageActionConsumerGroup string `yaml:"message_action_consumer_group"` + AgencyHostRemovedNoticeConsumerGroup string `yaml:"agency_host_removed_notice_consumer_group"` + ConsumerMaxReconsumeTimes int32 `yaml:"consumer_max_reconsume_times"` } // MessageActionOutboxMQConfig 控制确认消息 outbox 投递到 notice-service 的 MQ producer。 @@ -411,13 +412,14 @@ func defaultRocketMQConfig() RocketMQConfig { ConsumerMaxReconsumeTimes: 16, }, UserOutbox: UserOutboxMQConfig{ - Enabled: false, - Topic: "hyapp_user_outbox", - ConsumerGroup: "hyapp-activity-user-region-broadcast", - TaskConsumerGroup: "hyapp-activity-task-user-outbox", - InviteActivityConsumerGroup: "hyapp-activity-invite-user-outbox", - MessageActionConsumerGroup: "hyapp-activity-message-action-user-outbox", - ConsumerMaxReconsumeTimes: 16, + Enabled: false, + Topic: "hyapp_user_outbox", + ConsumerGroup: "hyapp-activity-user-region-broadcast", + TaskConsumerGroup: "hyapp-activity-task-user-outbox", + InviteActivityConsumerGroup: "hyapp-activity-invite-user-outbox", + MessageActionConsumerGroup: "hyapp-activity-message-action-user-outbox", + AgencyHostRemovedNoticeConsumerGroup: "hyapp-activity-agency-host-removed-notice", + ConsumerMaxReconsumeTimes: 16, }, MessageActionOutbox: MessageActionOutboxMQConfig{ Enabled: false, @@ -689,6 +691,9 @@ func normalizeRocketMQConfig(cfg RocketMQConfig) (RocketMQConfig, error) { if cfg.UserOutbox.MessageActionConsumerGroup = strings.TrimSpace(cfg.UserOutbox.MessageActionConsumerGroup); cfg.UserOutbox.MessageActionConsumerGroup == "" { cfg.UserOutbox.MessageActionConsumerGroup = defaults.UserOutbox.MessageActionConsumerGroup } + if cfg.UserOutbox.AgencyHostRemovedNoticeConsumerGroup = strings.TrimSpace(cfg.UserOutbox.AgencyHostRemovedNoticeConsumerGroup); cfg.UserOutbox.AgencyHostRemovedNoticeConsumerGroup == "" { + cfg.UserOutbox.AgencyHostRemovedNoticeConsumerGroup = defaults.UserOutbox.AgencyHostRemovedNoticeConsumerGroup + } if cfg.UserOutbox.ConsumerMaxReconsumeTimes <= 0 { cfg.UserOutbox.ConsumerMaxReconsumeTimes = defaults.UserOutbox.ConsumerMaxReconsumeTimes } diff --git a/services/user-service/internal/storage/mysql/host/applications.go b/services/user-service/internal/storage/mysql/host/applications.go index f21118aa..14898747 100644 --- a/services/user-service/internal/storage/mysql/host/applications.go +++ b/services/user-service/internal/storage/mysql/host/applications.go @@ -2,12 +2,18 @@ package host import ( "context" + "database/sql" + "encoding/json" + "hyapp/pkg/appcode" + "hyapp/pkg/idgen" "hyapp/pkg/xerr" hostdomain "hyapp/services/user-service/internal/domain/host" hostservice "hyapp/services/user-service/internal/service/host" ) +const userOutboxEventAgencyHostRemoved = "AgencyHostRemoved" + // ApplyToAgency 创建待处理申请;用户行和 Agency 行在同一事务内锁定。 func (r *Repository) ApplyToAgency(ctx context.Context, command hostservice.ApplyToAgencyCommand) (hostdomain.AgencyApplication, error) { tx, err := r.db.BeginTx(ctx, nil) @@ -320,9 +326,57 @@ func (r *Repository) KickAgencyHost(ctx context.Context, command hostservice.Kic return hostdomain.KickAgencyHostResult{}, err } } + // 系统消息消费的是和成员关系结束同事务提交的 user_outbox 快照。HTTP 返回后再同步调用 + // activity-service 会让“关系已移除、通知调用失败”变成不可重试的半成功状态,也无法覆盖 MQ 重投。 + if err := insertAgencyHostRemovedOutbox(ctx, tx, membership, command.OperatorUserID, command.NowMs); err != nil { + return hostdomain.KickAgencyHostResult{}, err + } if err := tx.Commit(); err != nil { return hostdomain.KickAgencyHostResult{}, err } return hostdomain.KickAgencyHostResult{Membership: membership, HostProfile: hostProfile}, nil } + +type agencyHostRemovedPayload struct { + MembershipID int64 `json:"membership_id"` + AgencyID int64 `json:"agency_id"` + HostUserID int64 `json:"host_user_id"` + Operator roleInvitationUserSnapshot `json:"operator"` + RemovedAtMS int64 `json:"removed_at_ms"` +} + +func insertAgencyHostRemovedOutbox(ctx context.Context, tx *sql.Tx, membership hostdomain.AgencyMembership, operatorUserID int64, nowMs int64) error { + // current_display_user_id 是用户当前对外展示号:存在有效靓号时为靓号,否则为永久短 ID。 + // 在关系事务内固化操作者资料,避免异步消费时昵称或靓号变化导致通知内容不再对应移除当时的事实。 + operator, err := roleInvitationUserSnapshotByID(ctx, tx, operatorUserID) + if err != nil { + return err + } + payloadBytes, err := json.Marshal(agencyHostRemovedPayload{ + MembershipID: membership.MembershipID, + AgencyID: membership.AgencyID, + HostUserID: membership.HostUserID, + Operator: operator, + RemovedAtMS: nowMs, + }) + if err != nil { + return err + } + _, err = tx.ExecContext(ctx, ` + INSERT INTO user_outbox ( + app_code, event_id, event_type, aggregate_type, aggregate_id, + status, worker_id, lock_until_ms, retry_count, next_retry_at_ms, + last_error, payload_json, created_at_ms, updated_at_ms + ) VALUES (?, ?, ?, ?, ?, 'pending', '', 0, 0, 0, '', CAST(? AS JSON), ?, ?)`, + appcode.FromContext(ctx), + idgen.New("uout"), + userOutboxEventAgencyHostRemoved, + hostdomain.ResultTypeMembership, + membership.MembershipID, + string(payloadBytes), + nowMs, + nowMs, + ) + return err +}