diff --git a/internal/biz/integration/integration_config.go b/internal/biz/integration/integration_config.go index 3e6d6ec..537f7bf 100644 --- a/internal/biz/integration/integration_config.go +++ b/internal/biz/integration/integration_config.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" bizpayment "kra/internal/biz/payment" + "net" "sort" "strconv" "strings" @@ -285,6 +286,44 @@ func validateCommunicationIntegrationConfig(kind, provider string, values map[st if integrationInt64(values, "reconnect_interval", 0) <= 0 { return errors.New("rabbitmq reconnect_interval 必须大于 0") } + case IntegrationKindMQ + "/kafka": + brokers := integrationStrings(values, "brokers") + if len(brokers) == 0 { + return errors.New("kafka brokers 不能为空") + } + for _, broker := range brokers { + host, portText, err := net.SplitHostPort(broker) + port, parseErr := strconv.Atoi(portText) + if err != nil || parseErr != nil || strings.TrimSpace(host) == "" || port < 1 || port > 65535 { + return fmt.Errorf("kafka broker %q 必须是有效的 host:port 地址", broker) + } + } + username := integrationText(values, "username") + password := integrationText(values, "password") + if (username == "") != (password == "") { + return errors.New("kafka username 和 password 必须同时配置") + } + startOffset := strings.ToLower(integrationText(values, "start_offset")) + if startOffset != "earliest" && startOffset != "latest" { + return errors.New("kafka start_offset 必须是 earliest 或 latest") + } + minBytes := integrationInt64(values, "min_bytes", 0) + maxBytes := integrationInt64(values, "max_bytes", 0) + if minBytes <= 0 || maxBytes < minBytes { + return errors.New("kafka max_bytes 必须大于等于 min_bytes,且二者必须大于 0") + } + if integrationInt64(values, "max_wait", 0) <= 0 { + return errors.New("kafka max_wait 必须大于 0") + } + if integrationInt64(values, "connect_timeout", 0) <= 0 { + return errors.New("kafka connect_timeout 必须大于 0") + } + if integrationInt64(values, "reconnect_interval", 0) <= 0 { + return errors.New("kafka reconnect_interval 必须大于 0") + } + if integrationBool(values, "tls_skip_verify") && !integrationBool(values, "tls") { + return errors.New("kafka tls_skip_verify 仅能在启用 TLS 时使用") + } case IntegrationKindWebSocket + "/melody": path := integrationText(values, "path") if !strings.HasPrefix(path, "/") { @@ -391,6 +430,34 @@ func integrationText(values map[string]any, key string) string { return strings.TrimSpace(fmt.Sprint(value)) } +func integrationStrings(values map[string]any, key string) []string { + switch items := values[key].(type) { + case []string: + result := make([]string, 0, len(items)) + for _, item := range items { + if item = strings.TrimSpace(item); item != "" { + result = append(result, item) + } + } + return result + case []any: + result := make([]string, 0, len(items)) + for _, item := range items { + if value := strings.TrimSpace(fmt.Sprint(item)); value != "" { + result = append(result, value) + } + } + return result + default: + return nil + } +} + +func integrationBool(values map[string]any, key string) bool { + value, _ := values[key].(bool) + return value +} + func integrationFirst(values map[string]any, keys ...string) string { for _, key := range keys { if value := integrationText(values, key); value != "" { diff --git a/internal/biz/integration/integration_config_definition.go b/internal/biz/integration/integration_config_definition.go index b93d453..59c226a 100644 --- a/internal/biz/integration/integration_config_definition.go +++ b/internal/biz/integration/integration_config_definition.go @@ -110,6 +110,26 @@ var integrationDefinitions = map[string][]IntegrationConfigDefinition{ {Key: "reconnect_interval", Label: "重连退避上限(秒)", Type: "number", Required: true, Description: "网络中断后自动重连的最大退避间隔。"}, }, }, + { + Kind: IntegrationKindMQ, Provider: "kafka", Name: "Kafka", Description: "Apache Kafka 分布式事件队列", + Defaults: map[string]any{"brokers": []string{"127.0.0.1:9092"}, "client_id": "kra", "group_id": "kra", "username": "", "password": "", "tls": false, "tls_skip_verify": false, "start_offset": "earliest", "min_bytes": 1, "max_bytes": 10485760, "max_wait": 1, "connect_timeout": 10, "reconnect_interval": 5, "allow_auto_topic_creation": false}, + Fields: []IntegrationConfigField{ + {Key: "brokers", Label: "Broker 地址", Type: "string-list", Required: true, Placeholder: "127.0.0.1:9092", Description: "每行一个 host:port 地址。"}, + {Key: "client_id", Label: "客户端 ID", Type: "text", Required: true, Placeholder: "kra"}, + {Key: "group_id", Label: "消费组 ID", Type: "text", Required: true, Placeholder: "kra"}, + {Key: "username", Label: "SASL 用户名", Type: "text"}, + {Key: "password", Label: "SASL 密码", Type: "password", Secret: true}, + {Key: "tls", Label: "启用 TLS", Type: "switch"}, + {Key: "tls_skip_verify", Label: "跳过 TLS 证书校验", Type: "switch", Description: "仅用于受控测试环境。"}, + {Key: "start_offset", Label: "初始消费位置", Type: "select", Required: true, Options: []IntegrationConfigOption{{Label: "earliest", Value: "earliest"}, {Label: "latest", Value: "latest"}}}, + {Key: "min_bytes", Label: "最小拉取字节数", Type: "number", Required: true}, + {Key: "max_bytes", Label: "最大拉取字节数", Type: "number", Required: true}, + {Key: "max_wait", Label: "最大拉取等待(秒)", Type: "number", Required: true}, + {Key: "connect_timeout", Label: "连接超时(秒)", Type: "number", Required: true}, + {Key: "reconnect_interval", Label: "重连间隔(秒)", Type: "number", Required: true}, + {Key: "allow_auto_topic_creation", Label: "允许自动创建 Topic", Type: "switch"}, + }, + }, { Kind: IntegrationKindMQ, Provider: "rabbitmq", Name: "RabbitMQ", Description: "RabbitMQ AMQP 消息队列", Defaults: map[string]any{"host": "127.0.0.1", "port": 5672, "username": "guest", "password": "guest", "vhost": "/", "exchange": "kra", "exchange_type": "topic", "queue": "kra", "routing_key": "#", "durable": true, "auto_delete": false, "prefetch_count": 10, "heartbeat": 10, "connect_timeout": 10, "reconnect_interval": 5, "tls": false}, diff --git a/internal/biz/payment/payment.go b/internal/biz/payment/payment.go index b259141..bfe53e7 100644 --- a/internal/biz/payment/payment.go +++ b/internal/biz/payment/payment.go @@ -38,6 +38,7 @@ const ( PaymentModeExternal = "external" PaymentModeInternal = "internal" + DemoBusinessType = "demo_subscription" ) var SupportedPaymentProviders = []string{ @@ -537,6 +538,10 @@ func supportedPaymentProvider(provider string) bool { return false } +func IsDemoPaymentOrder(order *PaymentOrder) bool { + return order != nil && order.BusinessType == DemoBusinessType && strings.HasPrefix(strings.TrimSpace(order.TradeNo), "DEMO-PAY-") +} + func (uc *PaymentUsecase) Query(ctx context.Context, provider, tradeNo string) (*PaymentResult, error) { if !validPaymentText(provider, 64) || !validPaymentText(tradeNo, 128) { return nil, errors.New("查询支付参数不完整") @@ -548,6 +553,9 @@ func (uc *PaymentUsecase) Query(ctx context.Context, provider, tradeNo string) ( if err != nil { return nil, err } + if IsDemoPaymentOrder(order) { + return paymentResultFromOrder(order), nil + } if provider == PaymentInternal || order.PaymentMode == PaymentModeInternal { return paymentResultFromOrder(order), nil } @@ -582,7 +590,14 @@ func (uc *PaymentUsecase) Refund(ctx context.Context, provider, tradeNo string, if uc.orders == nil { return nil, errors.New("支付订单仓储未接入") } - return uc.refundWithOrder(ctx, provider, tradeNo, amount) + order, err := uc.orders.FindPaymentOrder(ctx, provider, tradeNo) + if err != nil { + return nil, err + } + if IsDemoPaymentOrder(order) { + return nil, errors.New("演示订单不执行真实退款") + } + return uc.refundWithOrder(ctx, order, amount) } // TestProvider exercises the configured provider without requiring a business @@ -614,6 +629,9 @@ func (uc *PaymentUsecase) Fulfill(ctx context.Context, provider, tradeNo string) if err != nil { return nil, err } + if IsDemoPaymentOrder(order) { + return nil, errors.New("演示订单不执行真实发货") + } result := paymentResultFromOrder(order) if result == nil { return nil, ErrPaymentOrderNotFound @@ -663,9 +681,6 @@ func (uc *PaymentUsecase) createWithOrder(ctx context.Context, req *PaymentReque ConfirmationID: paymentConfirmationID(req.Provider, req.TradeNo), RequestFingerprint: paymentOrderFingerprint(req, extra), Extra: extra, - ClientIP: req.ClientIP, - UserAgent: req.UserAgent, - DeviceID: req.DeviceID, } order, created, err := uc.orders.CreatePaymentOrder(ctx, order) if err != nil { @@ -866,17 +881,17 @@ func (uc *PaymentUsecase) fulfillOrder(ctx context.Context, order *PaymentOrder, return attachOrderResult(result, order), nil } -func (uc *PaymentUsecase) refundWithOrder(ctx context.Context, provider, tradeNo string, amount int64) (*PaymentResult, error) { - current, err := uc.orders.FindPaymentOrder(ctx, provider, tradeNo) - if err != nil { - return nil, err +func (uc *PaymentUsecase) refundWithOrder(ctx context.Context, current *PaymentOrder, amount int64) (*PaymentResult, error) { + if current == nil { + return nil, ErrPaymentOrderNotFound } + provider, tradeNo := current.Provider, current.TradeNo handler := uc.fulfillments.Handler(current.BusinessType) authorizer, ok := handler.(PaymentRefundAuthorizer) if !ok { return nil, fmt.Errorf("支付业务未实现退款授权: %s", current.BusinessType) } - if err = authorizer.AuthorizeRefund(ctx, current, amount); err != nil { + if err := authorizer.AuthorizeRefund(ctx, current, amount); err != nil { return nil, err } order, token, err := uc.orders.BeginPaymentRefund(ctx, provider, tradeNo, amount, 10*time.Minute) diff --git a/internal/biz/payment/payment_order.go b/internal/biz/payment/payment_order.go index f05c377..6c05be1 100644 --- a/internal/biz/payment/payment_order.go +++ b/internal/biz/payment/payment_order.go @@ -88,9 +88,6 @@ type PaymentOrder struct { RefundedAt *time.Time FulfillmentLeaseUntil *time.Time RefundLeaseUntil *time.Time - ClientIP string - UserAgent string - DeviceID string } type PaymentProviderUpdate struct { diff --git a/internal/biz/payment/payment_test.go b/internal/biz/payment/payment_test.go index 74d5e67..6a05784 100644 --- a/internal/biz/payment/payment_test.go +++ b/internal/biz/payment/payment_test.go @@ -157,6 +157,33 @@ func testBizPaymentOrder() *PaymentOrder { return &PaymentOrder{Provider: PaymentAlipay, TradeNo: "order-1", BusinessType: "game_item", BusinessID: "item-1", Subject: "item", Amount: 100, Currency: "CNY", PaymentStatus: PaymentStatusPending, FulfillmentStatus: FulfillmentStatusPending, RefundStatus: RefundStatusNone, ConfirmationID: "11111111-1111-1111-1111-111111111111"} } +func TestDemoPaymentOperationsNeverReachExternalProviderOrFulfillment(t *testing.T) { + repo := &paymentRepoStub{query: successPaymentResult(), refund: successPaymentResult()} + order := testBizPaymentOrder() + order.TradeNo = "DEMO-PAY-01" + order.BusinessType = DemoBusinessType + order.PaymentStatus = PaymentStatusPaid + orders := &paymentOrderRepoStub{order: order} + registry := NewPaymentFulfillmentRegistry() + handler := &paymentFulfillmentStub{} + if err := registry.Register(handler); err != nil { + t.Fatal(err) + } + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + uc := NewPaymentUsecase(repo, orders, nil, nil, registry, logger) + + result, err := uc.Query(context.Background(), order.Provider, order.TradeNo) + if err != nil || result == nil || repo.queryCalls != 0 { + t.Fatalf("demo query result=%#v calls=%d err=%v", result, repo.queryCalls, err) + } + if _, err = uc.Refund(context.Background(), order.Provider, order.TradeNo, 10); err == nil || repo.refundedReq != nil { + t.Fatalf("demo refund reached provider: req=%#v err=%v", repo.refundedReq, err) + } + if _, err = uc.Fulfill(context.Background(), order.Provider, order.TradeNo); err == nil || handler.calls != 0 { + t.Fatalf("demo fulfillment reached handler: calls=%d err=%v", handler.calls, err) + } +} + func successPaymentResult() *PaymentResult { return &PaymentResult{ Provider: PaymentAlipay, diff --git a/internal/data/integration/migrations.go b/internal/data/integration/migrations.go index ff0a2de..622c677 100644 --- a/internal/data/integration/migrations.go +++ b/internal/data/integration/migrations.go @@ -18,12 +18,14 @@ func Migrations() []migration.Step { }}, {ID: "202608210001_communication_integration_defaults", Migrate: ensureCommunicationIntegrationConfigs}, {ID: "202608200004_payment_defaults", Migrate: ensurePaymentIntegrationConfigs}, + {ID: "202608290004_kafka_integration_default", Migrate: ensureCommunicationIntegrationConfigs}, } } func ensureCommunicationIntegrationConfigs(db *gorm.DB) error { defaults := []struct{ kind, provider string }{ {integrationbiz.IntegrationKindMQ, "emqx"}, + {integrationbiz.IntegrationKindMQ, "kafka"}, {integrationbiz.IntegrationKindMQ, "rabbitmq"}, {integrationbiz.IntegrationKindWebSocket, "melody"}, } diff --git a/internal/data/payment/demo_seed.go b/internal/data/payment/demo_seed.go index 89f3287..2b5470b 100644 --- a/internal/data/payment/demo_seed.go +++ b/internal/data/payment/demo_seed.go @@ -25,9 +25,10 @@ func SeedDemoOrders(db *gorm.DB, now time.Time) (*DemoSeedResult, error) { if !db.Migrator().HasTable(&paymentOrderPO{}) || !db.Migrator().HasTable(&paymentEventPO{}) { return nil, fmt.Errorf("payment tables are not migrated") } - now = now.UTC().Truncate(time.Second) if now.IsZero() { now = time.Now().UTC().Truncate(time.Second) + } else { + now = now.UTC().Truncate(time.Second) } type demoOrder struct { @@ -37,7 +38,7 @@ func SeedDemoOrders(db *gorm.DB, now time.Time) (*DemoSeedResult, error) { makeOrder := func(index int, provider, subject, status, fulfillment, refund string, amount int64, currency string, created time.Time) paymentOrderPO { tradeNo := fmt.Sprintf("%s%02d", DemoTradePrefix, index) return paymentOrderPO{ - TradeNo: tradeNo, Provider: provider, BusinessType: "demo_subscription", BusinessID: fmt.Sprintf("DEMO-LICENSE-%02d", index), + TradeNo: tradeNo, Provider: provider, BusinessType: bizpayment.DemoBusinessType, BusinessID: fmt.Sprintf("DEMO-LICENSE-%02d", index), Subject: subject, PaymentMode: bizpayment.PaymentModeExternal, OriginalAmount: amount, Amount: amount, Currency: currency, PaymentStatus: status, FulfillmentStatus: fulfillment, RefundStatus: refund, ConfirmationID: uuid.NewString(), RequestFingerprint: fmt.Sprintf("demo-fingerprint-%02d", index), diff --git a/internal/data/payment/payment_event.go b/internal/data/payment/payment_event.go index aa6a637..788d1a4 100644 --- a/internal/data/payment/payment_event.go +++ b/internal/data/payment/payment_event.go @@ -2,6 +2,7 @@ package payment import ( "context" + "errors" bizpayment "kra/internal/biz/payment" "time" @@ -48,6 +49,10 @@ func appendPaymentEvent(ctx context.Context, tx *gorm.DB, order *paymentOrderPO, if tx == nil || order == nil || event == nil { return nil } + eventType := trimTo(event.Type, 64) + if eventType == "" { + return errors.New("支付事件类型为空") + } operation := bizpayment.PaymentOperationFromContext(ctx) source := trimTo(event.Source, 64) if source == "" { @@ -81,7 +86,7 @@ func appendPaymentEvent(ctx context.Context, tx *gorm.DB, order *paymentOrderPO, deviceID = trimTo(operation.DeviceID, 128) } return tx.Create(&paymentEventPO{ - Provider: order.Provider, TradeNo: order.TradeNo, Type: trimTo(event.Type, 64), Source: source, + Provider: order.Provider, TradeNo: order.TradeNo, Type: eventType, Source: source, Status: trimTo(event.Status, 64), ProviderStatus: trimTo(event.ProviderStatus, 64), Message: message, EventID: trimTo(event.EventID, 128), PayloadHash: trimTo(event.PayloadHash, 64), Amount: event.Amount, Currency: trimTo(event.Currency, 16), OperatorID: operatorID, OperatorName: operatorName, diff --git a/internal/data/payment/payment_order.go b/internal/data/payment/payment_order.go index 23c9da2..ca79c4e 100644 --- a/internal/data/payment/payment_order.go +++ b/internal/data/payment/payment_order.go @@ -152,8 +152,7 @@ func (r *paymentOrderRepo) CreatePaymentOrder(ctx context.Context, order *bizpay } return appendPaymentEvent(ctx, tx, po, &bizpayment.PaymentEvent{ Type: bizpayment.PaymentEventOrderCreated, Source: "client", Status: po.PaymentStatus, - Amount: po.Amount, Currency: po.Currency, ClientIP: order.ClientIP, - UserAgent: order.UserAgent, DeviceID: order.DeviceID, + Amount: po.Amount, Currency: po.Currency, }) }) if err == nil { @@ -637,8 +636,9 @@ func defaultVersion(value uint64) uint64 { func trimTo(value string, max int) string { value = strings.TrimSpace(value) - if len(value) > max { - return value[:max] + runes := []rune(value) + if len(runes) > max { + return string(runes[:max]) } return value } diff --git a/internal/data/payment/payment_order_test.go b/internal/data/payment/payment_order_test.go index 4ade4db..2767331 100644 --- a/internal/data/payment/payment_order_test.go +++ b/internal/data/payment/payment_order_test.go @@ -5,11 +5,13 @@ import ( bizpayment "kra/internal/biz/payment" "testing" "time" + + "github.com/google/uuid" ) func newPaymentOrderRepoForTest(t *testing.T) *paymentOrderRepo { t.Helper() - db, err := openWithDriver("sqlite", "file:"+t.Name()+"?mode=memory&cache=shared") + db, err := openWithDriver("sqlite", "file:"+t.Name()+"-"+uuid.NewString()+"?mode=memory&cache=shared") if err != nil { t.Fatal(err) } @@ -119,11 +121,9 @@ func TestPaymentOrderRepositorySummarizesAndListsEvents(t *testing.T) { repo := newPaymentOrderRepoForTest(t) ctx := bizpayment.WithPaymentOperation(context.Background(), bizpayment.PaymentOperation{ Source: "admin", Reason: "客户重复购买", OperatorID: 7, OperatorName: "operator", - ClientIP: "127.0.0.1", DeviceID: "device-1", + ClientIP: "127.0.0.1", UserAgent: "payment-test/1.0", DeviceID: "device-1", }) order := testPaymentOrder() - order.ClientIP = "10.0.0.8" - order.UserAgent = "payment-client/1.0" if _, _, err := repo.CreatePaymentOrder(ctx, order); err != nil { t.Fatal(err) } @@ -168,3 +168,9 @@ func TestPaymentOrderRepositorySummarizesAndListsEvents(t *testing.T) { t.Fatalf("refund audit event missing operation context: %#v", events) } } + +func TestTrimToPreservesUTF8(t *testing.T) { + if got := trimTo(" 退款处理失败 ", 4); got != "退款处理" { + t.Fatalf("trimTo() = %q, want %q", got, "退款处理") + } +} diff --git a/internal/data/system/surface_upgrade_test.go b/internal/data/system/surface_upgrade_test.go index ffb8e66..6141304 100644 --- a/internal/data/system/surface_upgrade_test.go +++ b/internal/data/system/surface_upgrade_test.go @@ -3,11 +3,12 @@ package system import ( "testing" + "github.com/google/uuid" platformmodule "kra/pkg/module" ) func TestEnsureAdminSurfaceAndPolicyInheritance(t *testing.T) { - db, err := openWithDriver("sqlite", "file:"+t.Name()+"?mode=memory&cache=shared") + db, err := openWithDriver("sqlite", "file:"+t.Name()+"-"+uuid.NewString()+"?mode=memory&cache=shared") if err != nil { t.Fatal(err) } diff --git a/internal/integration/mq/emqx.go b/internal/integration/mq/emqx.go index 38abf32..908ded2 100644 --- a/internal/integration/mq/emqx.go +++ b/internal/integration/mq/emqx.go @@ -16,10 +16,13 @@ import ( const ( ProviderEMQX = platformmq.ProviderEMQX + ProviderKafka = platformmq.ProviderKafka ProviderRabbitMQ = platformmq.ProviderRabbitMQ retryTick = time.Second ) +var providers = []string{ProviderEMQX, ProviderKafka, ProviderRabbitMQ} + // Reloadable owns process-wide message clients and the logical subscription // declarations used to restore them after reconnects or configuration reloads. type Reloadable struct { @@ -65,12 +68,11 @@ func New(store *runtimeconfig.Store, logger *slog.Logger) (*Reloadable, func(), logger: logger, } if store != nil { - r.apply(ProviderEMQX, storeConfig(store, ProviderEMQX)) - r.apply(ProviderRabbitMQ, storeConfig(store, ProviderRabbitMQ)) - r.stop = append(r.stop, - store.Subscribe("mq", ProviderEMQX, func(config runtimeconfig.Config) { r.apply(ProviderEMQX, config) }), - store.Subscribe("mq", ProviderRabbitMQ, func(config runtimeconfig.Config) { r.apply(ProviderRabbitMQ, config) }), - ) + for _, provider := range providers { + r.apply(provider, storeConfig(store, provider)) + provider := provider + r.stop = append(r.stop, store.Subscribe("mq", provider, func(config runtimeconfig.Config) { r.apply(provider, config) })) + } } go r.retryLoop() cleanup := func() { @@ -88,7 +90,7 @@ func storeConfig(store *runtimeconfig.Store, provider string) runtimeconfig.Conf } // TestConfig creates a short-lived provider client and closes it immediately. -// For RabbitMQ this also checks the configured exchange and queue topology. +// RabbitMQ also checks topology; Kafka reads cluster metadata. func TestConfig(ctx context.Context, provider string, raw json.RawMessage) error { if ctx != nil { select { @@ -218,6 +220,24 @@ func newProviderClient(provider string, raw json.RawMessage) (platformmq.Client, ReconnectInterval: configSeconds(values, "reconnect_interval"), TLS: configBool(values, "tls"), }) + case ProviderKafka: + return platformmq.NewKafka(platformmq.KafkaConfig{ + Enabled: true, + Brokers: configStrings(values, "brokers"), + ClientID: configText(values, "client_id"), + GroupID: configText(values, "group_id"), + Username: configText(values, "username"), + Password: configText(values, "password"), + TLS: configBool(values, "tls"), + TLSSkipVerify: configBool(values, "tls_skip_verify"), + StartOffset: configText(values, "start_offset"), + MinBytes: configInt(values, "min_bytes"), + MaxBytes: configInt(values, "max_bytes"), + MaxWait: configSeconds(values, "max_wait"), + ConnectTimeout: configSeconds(values, "connect_timeout"), + ReconnectInterval: configSeconds(values, "reconnect_interval"), + AllowAutoTopicCreation: configBool(values, "allow_auto_topic_creation"), + }) default: return nil, fmt.Errorf("unsupported message provider %q", provider) } @@ -246,6 +266,24 @@ func configInt(values map[string]any, key string) int { } } +func configStrings(values map[string]any, key string) []string { + items, ok := values[key].([]any) + if ok { + result := make([]string, 0, len(items)) + for _, item := range items { + if value := strings.TrimSpace(fmt.Sprint(item)); value != "" { + result = append(result, value) + } + } + return result + } + stringsValue, ok := values[key].([]string) + if ok { + return append([]string(nil), stringsValue...) + } + return nil +} + func configSeconds(values map[string]any, key string) time.Duration { seconds := configInt(values, key) if seconds <= 0 { @@ -281,7 +319,7 @@ func (r *Reloadable) retryOnce() { } r.ensureStateLocked() now := time.Now() - for _, provider := range []string{ProviderEMQX, ProviderRabbitMQ} { + for _, provider := range providers { config, exists := r.configs[provider] if !exists || !config.Enabled { continue @@ -631,7 +669,7 @@ func (r *Reloadable) UnsubscribeFrom(ctx context.Context, provider string, topic func (r *Reloadable) Client(provider string) platformmq.Client { provider = strings.ToLower(strings.TrimSpace(provider)) - if provider != ProviderEMQX && provider != ProviderRabbitMQ { + if provider != ProviderEMQX && provider != ProviderKafka && provider != ProviderRabbitMQ { return nil } return &namedClient{owner: r, provider: provider} diff --git a/internal/service/payment/payment.go b/internal/service/payment/payment.go index 7bde468..0f290b6 100644 --- a/internal/service/payment/payment.go +++ b/internal/service/payment/payment.go @@ -189,6 +189,7 @@ func (s *PaymentService) Create(ctx context.Context, req *dto.PaymentRequest) (* if req == nil { return nil, errors.New("支付请求为空") } + ctx = operationContext(ctx, "client", "客户端创建支付订单", 0, "", req.ClientIP, req.UserAgent, req.DeviceID) result, err := s.uc.Create(ctx, &paymentbiz.PaymentRequest{ Provider: req.Provider, TradeNo: req.TradeNo, Subject: req.Subject, Amount: req.Amount, OriginalAmount: req.OriginalAmount, Currency: req.Currency, diff --git a/pkg/mq/mq.go b/pkg/mq/mq.go index eba48c6..14ad160 100644 --- a/pkg/mq/mq.go +++ b/pkg/mq/mq.go @@ -12,6 +12,7 @@ var ErrUnavailable = errors.New("message broker unavailable") const ( ProviderEMQX = "emqx" + ProviderKafka = "kafka" ProviderRabbitMQ = "rabbitmq" ) diff --git a/pkg/mq/subscription.go b/pkg/mq/subscription.go index 6dec36f..350af86 100644 --- a/pkg/mq/subscription.go +++ b/pkg/mq/subscription.go @@ -12,7 +12,7 @@ func NormalizeSubscriptionSet(set SubscriptionSet) (SubscriptionSet, error) { if set.Owner == "" { return SubscriptionSet{}, fmt.Errorf("mq subscription owner is empty") } - if set.Provider != ProviderEMQX && set.Provider != ProviderRabbitMQ { + if set.Provider != ProviderEMQX && set.Provider != ProviderKafka && set.Provider != ProviderRabbitMQ { return SubscriptionSet{}, fmt.Errorf("unsupported mq provider %q", set.Provider) } if len(set.Topics) == 0 { diff --git a/web/src/view/systemTools/integration/config.vue b/web/src/view/systemTools/integration/config.vue index 7d9ae0d..c1aa3e9 100644 --- a/web/src/view/systemTools/integration/config.vue +++ b/web/src/view/systemTools/integration/config.vue @@ -265,6 +265,12 @@ const TARGETS = { description: 'RabbitMQ 消息队列', icon: Promotion }, + 'mq/kafka': { + name: 'Kafka', + protocol: 'Kafka', + description: 'Apache Kafka 分布式事件队列', + icon: Promotion + }, 'websocket/melody': { name: 'WebSocket', protocol: 'WS', @@ -276,7 +282,7 @@ const TARGETS = { const TARGET_ORDER = Object.keys(TARGETS) const RECONNECT_INTERVAL_KEY = 'reconnect_interval' const RECONNECT_DEFAULT_SECONDS = 5 -const MQ_RECONNECT_TARGETS = new Set(['mq/emqx', 'mq/rabbitmq']) +const MQ_RECONNECT_TARGETS = new Set(['mq/emqx', 'mq/kafka', 'mq/rabbitmq']) const RECONNECT_INTERVAL_FIELD = { key: RECONNECT_INTERVAL_KEY, label: '重连间隔(秒)', @@ -291,6 +297,9 @@ const NUMBER_CONSTRAINTS = { reconnect_interval: { min: 1, integer: true }, prefetch_count: { min: 0 }, heartbeat: { min: 0 }, + min_bytes: { min: 1, integer: true }, + max_bytes: { min: 1, integer: true }, + max_wait: { min: 1 }, max_message_size: { min: 0 }, message_buffer_size: { min: 0 } } diff --git a/web/src/view/systemTools/payment/orders.vue b/web/src/view/systemTools/payment/orders.vue index e878daa..d933f73 100644 --- a/web/src/view/systemTools/payment/orders.vue +++ b/web/src/view/systemTools/payment/orders.vue @@ -78,7 +78,7 @@