140 lines
3.6 KiB
Go
140 lines
3.6 KiB
Go
package spike
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"sync/atomic"
|
||
"time"
|
||
|
||
"cmpp-platform/gateway/internal/queue"
|
||
"cmpp-platform/gateway/internal/tracker"
|
||
)
|
||
|
||
type Simulator struct {
|
||
tracker *tracker.Tracker
|
||
sequence atomic.Uint32
|
||
}
|
||
|
||
type RunResult struct {
|
||
Submitted int
|
||
SubmitResults int
|
||
ReceiptEvents int
|
||
Duration time.Duration
|
||
MessagesPerSec float64
|
||
}
|
||
|
||
func NewSimulator() *Simulator {
|
||
return &Simulator{tracker: tracker.New()}
|
||
}
|
||
|
||
func (s *Simulator) Submit(ctx context.Context, cmd queue.SubmitCommand) (queue.SubmitResult, queue.ReceiptEvent, error) {
|
||
if err := ctx.Err(); err != nil {
|
||
return queue.SubmitResult{}, queue.ReceiptEvent{}, err
|
||
}
|
||
|
||
sequenceID := s.sequence.Add(1)
|
||
s.tracker.TrackSubmit(cmd.MessageID, sequenceID)
|
||
|
||
gatewayMessageID := fmt.Sprintf("gw-%s", cmd.MessageID)
|
||
mapping, err := s.tracker.TrackSubmitResp(sequenceID, gatewayMessageID)
|
||
if err != nil {
|
||
return queue.SubmitResult{}, queue.ReceiptEvent{}, err
|
||
}
|
||
|
||
now := time.Now().UTC()
|
||
result := queue.SubmitResult{
|
||
Envelope: queue.Envelope{
|
||
SchemaVersion: queue.SchemaVersion,
|
||
MessageType: queue.MessageTypeSubmitResult,
|
||
TraceID: cmd.TraceID,
|
||
MessageID: mapping.MessageID,
|
||
ChannelID: cmd.ChannelID,
|
||
CreatedAt: now,
|
||
},
|
||
SequenceID: mapping.SequenceID,
|
||
GatewayMessageID: mapping.GatewayMessageID,
|
||
SubmitStatus: "accepted",
|
||
SubmittedAt: now,
|
||
}
|
||
|
||
receipt := queue.ReceiptEvent{
|
||
Envelope: queue.Envelope{
|
||
SchemaVersion: queue.SchemaVersion,
|
||
MessageType: queue.MessageTypeReceiptEvent,
|
||
TraceID: cmd.TraceID,
|
||
MessageID: mapping.MessageID,
|
||
ChannelID: cmd.ChannelID,
|
||
CreatedAt: now.Add(10 * time.Millisecond),
|
||
},
|
||
SequenceID: mapping.SequenceID,
|
||
GatewayMessageID: mapping.GatewayMessageID,
|
||
ReceiptStatus: "delivered",
|
||
RawStatus: "DELIVRD",
|
||
DeliveredAt: now.Add(10 * time.Millisecond),
|
||
}
|
||
|
||
return result, receipt, nil
|
||
}
|
||
|
||
func (s *Simulator) RunLoad(ctx context.Context, count int) (RunResult, error) {
|
||
start := time.Now()
|
||
result := RunResult{Submitted: count}
|
||
|
||
for i := 0; i < count; i++ {
|
||
cmd := NewSubmitCommand(i)
|
||
_, _, err := s.Submit(ctx, cmd)
|
||
if err != nil {
|
||
return RunResult{}, err
|
||
}
|
||
result.SubmitResults++
|
||
result.ReceiptEvents++
|
||
}
|
||
|
||
result.Duration = time.Since(start)
|
||
if result.Duration > 0 {
|
||
result.MessagesPerSec = float64(count) / result.Duration.Seconds()
|
||
}
|
||
|
||
return result, nil
|
||
}
|
||
|
||
func NewSubmitCommand(index int) queue.SubmitCommand {
|
||
now := time.Now().UTC()
|
||
messageID := fmt.Sprintf("msg-spike-%06d", index)
|
||
return queue.SubmitCommand{
|
||
Envelope: queue.Envelope{
|
||
SchemaVersion: queue.SchemaVersion,
|
||
MessageType: queue.MessageTypeSubmitCommand,
|
||
TraceID: fmt.Sprintf("trace-spike-%06d", index),
|
||
MessageID: messageID,
|
||
ChannelID: "sms-channel-cmpp-spike",
|
||
CreatedAt: now,
|
||
},
|
||
TenantID: "tenant-spike",
|
||
ApplicationID: "app-spike",
|
||
TaskID: "task-spike",
|
||
SubmitID: fmt.Sprintf("submit-spike-%06d", index),
|
||
PhoneNumber: "13800138000",
|
||
Content: "您的验证码为 123456,5 分钟内有效。",
|
||
Signature: "测试平台",
|
||
TemplateID: "tpl-spike",
|
||
BillingUnits: 1,
|
||
Route: queue.Route{
|
||
ChannelCode: "CMCC-CMPP-SPIKE",
|
||
CMPPAccountCode: "cmpp-account-spike",
|
||
Priority: 10,
|
||
RateLimitPerSecond: 500,
|
||
},
|
||
CMPP: queue.CMPP{
|
||
ServiceID: "CMPP",
|
||
SrcID: "106900000000",
|
||
RegisteredDelivery: 1,
|
||
MsgFmt: 15,
|
||
FeeUserType: 2,
|
||
FeeCode: "0",
|
||
FeeType: "01",
|
||
},
|
||
Retry: queue.Retry{Attempt: 0, MaxAttempts: 3},
|
||
}
|
||
}
|