This commit is contained in:
@@ -107,15 +107,9 @@ func (r *Repository) Start(ctx context.Context, input Attempt) (*model.Integrati
|
||||
if input.IntegrationID == "" {
|
||||
input.IntegrationID = uuid.NewString()
|
||||
}
|
||||
autoAttempt := input.Attempt <= 0 && input.TriggerSeries != nil
|
||||
if input.Attempt <= 0 {
|
||||
input.Attempt = 1
|
||||
if input.TriggerSeries != nil {
|
||||
if err := r.db.WithContext(ctx).Model(&model.IntegrationLog{}).
|
||||
Select("COALESCE(MAX(attempt), 0) + 1").
|
||||
Where("trigger_series = ?", *input.TriggerSeries).Scan(&input.Attempt).Error; err != nil {
|
||||
return nil, pkgerrors.Wrap(pkgerrors.CodeDatabaseError, err, "计算 Integration Log 尝试序号失败")
|
||||
}
|
||||
}
|
||||
}
|
||||
if input.StartedAt == nil {
|
||||
startedAt := r.now().UTC()
|
||||
@@ -128,16 +122,35 @@ func (r *Repository) Start(ctx context.Context, input Attempt) (*model.Integrati
|
||||
resourceType := optionalString(input.ResourceType)
|
||||
log := &model.IntegrationLog{
|
||||
IntegrationID: input.IntegrationID, Provider: input.Provider, Direction: input.Direction,
|
||||
Operation: input.Operation, ExternalID: input.ExternalID, ResourceType: resourceType,
|
||||
ResourceID: input.ResourceID, ResourceKey: input.ResourceKey, TriggerSource: input.TriggerSource,
|
||||
TriggerScene: input.TriggerScene, TriggerSeries: input.TriggerSeries, ScheduledAt: input.ScheduledAt,
|
||||
Operation: input.Operation, ExternalID: sanitizedOptionalText(input.ExternalID), ResourceType: resourceType,
|
||||
ResourceID: input.ResourceID, ResourceKey: sanitizedOptionalText(input.ResourceKey), TriggerSource: input.TriggerSource,
|
||||
TriggerScene: sanitizedOptionalText(input.TriggerScene), TriggerSeries: input.TriggerSeries, ScheduledAt: input.ScheduledAt,
|
||||
StartedAt: input.StartedAt, Attempt: input.Attempt, Result: result,
|
||||
RequestSummary: requestSummary, Metadata: metadata, RequestID: input.RequestID,
|
||||
CorrelationID: input.CorrelationID, AuditEventID: input.AuditEventID,
|
||||
RecoveryStrategy: input.RecoveryStrategy,
|
||||
RecoveryStrategy: sanitizedOptionalText(input.RecoveryStrategy),
|
||||
}
|
||||
if err := r.db.WithContext(ctx).Create(log).Error; err != nil {
|
||||
return nil, pkgerrors.Wrap(pkgerrors.CodeDatabaseError, err, "写入 Integration Log 失败")
|
||||
createAttempt := func(tx *gorm.DB) error {
|
||||
if !autoAttempt {
|
||||
return tx.Create(log).Error
|
||||
}
|
||||
if err := tx.Exec("SELECT pg_advisory_xact_lock(hashtext(?))", *input.TriggerSeries).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if err := tx.Model(&model.IntegrationLog{}).Select("COALESCE(MAX(attempt), 0) + 1").
|
||||
Where("trigger_series = ?", *input.TriggerSeries).Scan(&log.Attempt).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Create(log).Error
|
||||
}
|
||||
var createErr error
|
||||
if autoAttempt {
|
||||
createErr = r.db.WithContext(ctx).Transaction(createAttempt)
|
||||
} else {
|
||||
createErr = createAttempt(r.db.WithContext(ctx))
|
||||
}
|
||||
if createErr != nil {
|
||||
return nil, pkgerrors.Wrap(pkgerrors.CodeDatabaseError, createErr, "写入 Integration Log 失败")
|
||||
}
|
||||
return log, nil
|
||||
}
|
||||
@@ -156,7 +169,7 @@ func (r *Repository) Complete(ctx context.Context, integrationID string, complet
|
||||
if completion.Result == constants.IntegrationResultUnknown && strings.TrimSpace(completion.RecoveryStrategy) == "" {
|
||||
return nil, pkgerrors.New(pkgerrors.CodeInvalidParam, "结果未知必须记录明确恢复策略")
|
||||
}
|
||||
safeProviderMessage := strings.TrimSpace(completion.SafeProviderMessage)
|
||||
safeProviderMessage := sanitizer.SanitizeText(strings.TrimSpace(completion.SafeProviderMessage))
|
||||
if safeProviderMessage != "" && utf8.RuneCountInString(constants.IntegrationSafeMessagePrefix+safeProviderMessage) > constants.IntegrationProviderMessageMaxLength {
|
||||
return nil, pkgerrors.New(pkgerrors.CodeInvalidParam, "Integration Log 安全结果摘要过长")
|
||||
}
|
||||
@@ -187,10 +200,10 @@ func (r *Repository) Complete(ctx context.Context, integrationID string, complet
|
||||
updates["resource_id"] = completion.ResourceID
|
||||
}
|
||||
if completion.ResourceKey != nil {
|
||||
updates["resource_key"] = completion.ResourceKey
|
||||
updates["resource_key"] = sanitizedOptionalText(completion.ResourceKey)
|
||||
}
|
||||
if completion.RecoveryStrategy != "" {
|
||||
updates["recovery_strategy"] = completion.RecoveryStrategy
|
||||
updates["recovery_strategy"] = sanitizer.SanitizeText(completion.RecoveryStrategy)
|
||||
}
|
||||
result := r.db.WithContext(ctx).Model(&model.IntegrationLog{}).
|
||||
Where("integration_id = ? AND result = ?", integrationID, constants.IntegrationResultPending).
|
||||
@@ -237,8 +250,8 @@ func (r *Repository) RecordInbound(ctx context.Context, input InboundAttempt) (*
|
||||
log := &model.IntegrationLog{
|
||||
IntegrationID: input.IntegrationID, IdempotencyKey: &input.IdempotencyKey,
|
||||
Provider: input.Provider, Direction: constants.IntegrationDirectionInbound, Operation: input.Operation,
|
||||
ExternalID: optionalString(input.ExternalID), ResourceType: optionalString(input.ResourceType),
|
||||
ResourceID: input.ResourceID, ResourceKey: input.ResourceKey, StartedAt: &now, Attempt: 1,
|
||||
ExternalID: sanitizedOptionalText(optionalString(input.ExternalID)), ResourceType: optionalString(input.ResourceType),
|
||||
ResourceID: input.ResourceID, ResourceKey: sanitizedOptionalText(input.ResourceKey), StartedAt: &now, Attempt: 1,
|
||||
TriggerSeries: &triggerSeries,
|
||||
Result: constants.IntegrationResultPending, RequestSummary: summary,
|
||||
ContentHash: hex.EncodeToString(hash[:]), RequestID: input.RequestID, CorrelationID: input.CorrelationID,
|
||||
@@ -380,3 +393,11 @@ func optionalString(value string) *string {
|
||||
}
|
||||
return &value
|
||||
}
|
||||
|
||||
func sanitizedOptionalText(value *string) *string {
|
||||
if value == nil {
|
||||
return nil
|
||||
}
|
||||
sanitized := sanitizer.SanitizeText(*value)
|
||||
return &sanitized
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user