Files
break 5e552d99bc 收口审计治理与套餐任务进展
Constraint: 在线热修前必须保存当前迭代分支全部有效代码进展
Confidence: medium
Scope-risk: broad
Directive: 后续修改需保持审计事件与业务事务边界一致
Tested: git diff --cached --check
Not-tested: 未运行全量测试,提交用于切换分支前保存既有工作
2026-08-05 14:30:54 +08:00

170 lines
7.1 KiB
Go

package wecom
import (
"context"
"fmt"
"time"
"github.com/bytedance/sonic"
"gorm.io/gorm"
systemconfigapp "github.com/break/junhong_cmp_fiber/internal/application/systemconfig"
"github.com/break/junhong_cmp_fiber/internal/model"
"github.com/break/junhong_cmp_fiber/internal/model/dto"
"github.com/break/junhong_cmp_fiber/pkg/constants"
"github.com/break/junhong_cmp_fiber/pkg/errors"
"github.com/break/junhong_cmp_fiber/pkg/middleware"
)
// DirectoryMember 是企微 Adapter 交给 Application 的最小成员快照。
type DirectoryMember struct {
UserID string
Name string
DepartmentIDs []int64
}
// DirectoryProvider 定义拉取应用可见成员的外部端口。
type DirectoryProvider interface {
ListVisibleMembers(ctx context.Context, applicationID uint) ([]DirectoryMember, error)
}
// MemberRepository 定义可见成员快照同步和分页查询边界。
type MemberRepository interface {
ReplaceVisible(ctx context.Context, tx *gorm.DB, applicationID uint, members []model.WeComMember, syncedAt time.Time) error
ListVisible(ctx context.Context, applicationID uint, page, pageSize int, keyword string) ([]model.WeComMember, int64, error)
}
// DirectoryService 同步并分页查询企业微信应用可见成员。
type DirectoryService struct {
db *gorm.DB
applications interface {
GetEnabled(ctx context.Context, applicationID uint) (*model.WeComApplication, error)
}
provider DirectoryProvider
members MemberRepository
audit systemconfigapp.AuditWriter
now func() time.Time
}
// NewDirectoryService 创建企业微信通讯录同步用例。
func NewDirectoryService(db *gorm.DB, applications interface {
GetEnabled(ctx context.Context, applicationID uint) (*model.WeComApplication, error)
}, provider DirectoryProvider, members MemberRepository, audit systemconfigapp.AuditWriter) *DirectoryService {
return &DirectoryService{db: db, applications: applications, provider: provider, members: members, audit: audit, now: time.Now}
}
// Sync 拉取并替换指定应用当前可见成员快照。
func (s *DirectoryService) Sync(ctx context.Context, applicationID uint) (*dto.WeComMemberSyncResponse, error) {
if s == nil || s.db == nil || s.applications == nil || s.provider == nil || s.members == nil || s.audit == nil || applicationID == 0 {
return nil, errors.New(errors.CodeServiceUnavailable, "企业微信通讯录服务未配置")
}
if !canManageWeComDirectory(ctx) {
return nil, errors.New(errors.CodeForbidden)
}
application, err := s.applications.GetEnabled(ctx, applicationID)
if err != nil {
return nil, err
}
remoteMembers, err := s.provider.ListVisibleMembers(ctx, applicationID)
if err != nil {
s.recordFailure(ctx, application, "同步企业微信应用可见成员失败")
return nil, err
}
syncedAt := s.now().UTC()
members := make([]model.WeComMember, 0, len(remoteMembers))
for _, member := range remoteMembers {
departments, err := sonic.Marshal(member.DepartmentIDs)
if err != nil {
return nil, errors.Wrap(errors.CodeInternalError, err, "序列化企业微信成员部门失败")
}
members = append(members, model.WeComMember{
ApplicationID: applicationID, CorpID: application.CorpID, UserID: member.UserID,
Name: member.Name, DepartmentIDs: departments, Visible: true, SyncedAt: syncedAt,
CreatedAt: syncedAt, UpdatedAt: syncedAt,
})
}
err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := s.members.ReplaceVisible(ctx, tx, applicationID, members, syncedAt); err != nil {
return err
}
requestID := ""
if value := middleware.GetRequestIDFromContext(ctx); value != nil {
requestID = *value
}
resourceID := fmt.Sprintf("%d", applicationID)
after := applicationAuditSnapshot(application)
after["synced_count"] = len(members)
after["synced_at"] = syncedAt
return s.audit.WriteConfigChange(ctx, tx, systemconfigapp.ChangeAudit{
OperatorID: middleware.GetUserIDFromContext(ctx), OperationType: constants.AuditOperationWeComMembersSync,
Description: "同步企业微信应用可见成员", ConfigKey: fmt.Sprintf("wecom.application.%d.members", applicationID),
ResourceID: &resourceID, DisplayName: application.Name, Identity: applicationAuditIdentity(application),
AfterData: after, RequestID: requestID, CorrelationID: requestID,
})
})
if err != nil {
s.recordFailure(ctx, application, "保存企业微信应用可见成员快照失败")
return nil, err
}
return &dto.WeComMemberSyncResponse{ApplicationID: applicationID, SyncedCount: len(members), SyncedAt: syncedAt}, nil
}
func (s *DirectoryService) recordFailure(ctx context.Context, application *model.WeComApplication, description string) {
if application == nil {
return
}
resourceID := fmt.Sprintf("%d", application.ID)
recordConfigFailure(ctx, s.db, s.audit, systemconfigapp.ChangeAudit{
OperatorID: middleware.GetUserIDFromContext(ctx), OperationType: constants.AuditOperationWeComMembersSync,
Description: description, ConfigKey: fmt.Sprintf("wecom.application.%d.members", application.ID),
ResourceID: &resourceID, DisplayName: application.Name, Identity: applicationAuditIdentity(application),
BeforeData: applicationAuditSnapshot(application), Result: constants.AuditResultFailed,
ErrorCode: fmt.Sprintf("%d", errors.CodeInternalError), ErrorSummary: description,
})
}
// List 分页返回本地最近一次同步的应用可见成员。
func (s *DirectoryService) List(ctx context.Context, applicationID uint, request dto.WeComMemberListRequest) (*dto.WeComMemberListResponse, error) {
if s == nil || s.applications == nil || s.members == nil || applicationID == 0 {
return nil, errors.New(errors.CodeServiceUnavailable, "企业微信通讯录服务未配置")
}
if !canManageWeComDirectory(ctx) {
return nil, errors.New(errors.CodeForbidden)
}
if request.Page <= 0 {
request.Page = constants.DefaultPage
}
if request.PageSize <= 0 {
request.PageSize = constants.DefaultPageSize
}
if request.PageSize > constants.MaxPageSize {
return nil, errors.New(errors.CodeInvalidParam)
}
if _, err := s.applications.GetEnabled(ctx, applicationID); err != nil {
return nil, err
}
members, total, err := s.members.ListVisible(ctx, applicationID, request.Page, request.PageSize, request.Keyword)
if err != nil {
return nil, err
}
items := make([]dto.WeComMemberResponse, 0, len(members))
for _, member := range members {
var departmentIDs []int64
if len(member.DepartmentIDs) > 0 {
if err := sonic.Unmarshal(member.DepartmentIDs, &departmentIDs); err != nil {
return nil, errors.Wrap(errors.CodeInternalError, err, "解析企业微信成员部门失败")
}
}
items = append(items, dto.WeComMemberResponse{
ApplicationID: member.ApplicationID, CorpID: member.CorpID, UserID: member.UserID,
Name: member.Name, DepartmentIDs: departmentIDs, SyncedAt: member.SyncedAt,
})
}
return &dto.WeComMemberListResponse{Items: items, Total: total, Page: request.Page, PageSize: request.PageSize}, nil
}
func canManageWeComDirectory(ctx context.Context) bool {
userType := middleware.GetUserTypeFromContext(ctx)
return userType == constants.UserTypeSuperAdmin || userType == constants.UserTypePlatform
}