Files
NixMsg/internal/app/message/app.go
T

172 lines
4.5 KiB
Go

package message
import (
"sync"
"time"
"git.asio.asia/nixevol/NixMsg/internal/app/port"
"git.asio.asia/nixevol/NixMsg/internal/auth"
"git.asio.asia/nixevol/NixMsg/internal/config"
"git.asio.asia/nixevol/NixMsg/internal/metrics"
"git.asio.asia/nixevol/NixMsg/internal/protocol"
"git.asio.asia/nixevol/NixMsg/internal/store"
)
// 消息状态(DEVELOPMENT 7.1)。
const (
StateScheduled = "scheduled"
StateDispatched = "dispatched"
StateCompleted = "completed"
)
// talk_grants.kind。
const (
GrantKindPassword = "password"
GrantKindReply = "reply"
)
// 请求频率桶默认突发容量(DEVELOPMENT 6.10;配置无单独字段)。
const defaultRequestBurst = 100
// Limits 是消息子系统所需配置上限。
type Limits struct {
MaxBodyBytes int
MaxMetaBytes int
MaxFrameBytes int
MaxTTLSeconds int64
MaxScheduleSeconds int64
RequestsPerSecond float64
RequestBurst int
MaxPendingPerSender int
MaxPendingPerReceiver int
GraceSeconds int64
AckTimeoutSeconds int64
DeliveryWindow int
ReceiptWindow int
RecordRetentionDays int
ReceiptRetentionDays int
IdempotencyHours int
}
// LimitsFromConfig 从平台配置构造 Limits。
func LimitsFromConfig(c config.LimitsConfig) Limits {
return Limits{
MaxBodyBytes: c.MaxBodyBytes,
MaxMetaBytes: c.MaxMetaBytes,
MaxFrameBytes: c.MaxFrameBytes,
MaxTTLSeconds: int64(c.MaxTTLSeconds),
MaxScheduleSeconds: int64(c.MaxScheduleSeconds),
RequestsPerSecond: float64(c.RequestsPerSecond),
RequestBurst: defaultRequestBurst,
MaxPendingPerSender: c.MaxPendingPerSender,
MaxPendingPerReceiver: c.MaxPendingPerReceiver,
GraceSeconds: int64(c.GraceSeconds),
AckTimeoutSeconds: int64(c.AckTimeoutSeconds),
DeliveryWindow: c.DeliveryWindow,
ReceiptWindow: c.ReceiptWindow,
}
}
// LimitsFromFullConfig 附带保留天数等顶层配置。
func LimitsFromFullConfig(cfg config.Config) Limits {
lim := LimitsFromConfig(cfg.Limits)
lim.RecordRetentionDays = cfg.RecordRetentionDays
lim.ReceiptRetentionDays = cfg.ReceiptRetentionDays
lim.IdempotencyHours = cfg.IdempotencyHours
return lim
}
// App 实现 Service:提交、分发、推送、确认、撤回、回执、清理与启动恢复。
type App struct {
db *store.DB
lim Limits
hash auth.HashPool
locks auth.LoginLocks
nowFn func() time.Time
rates *rateLimiter
down port.Downlink
conns ConnRegistry
met *metrics.Registry
mu sync.Mutex
largeSem chan struct{}
largeHeld map[string]bool
pendingRevoke []revokeJob
repushTimers map[string]*time.Timer
}
// Option 配置 App。
type Option func(*App)
// WithNow 注入时钟(测试用)。
func WithNow(now func() time.Time) Option {
return func(a *App) { a.nowFn = now }
}
// WithLocks 注入对话密码锁定计数器;nil 表示不锁定。
func WithLocks(locks auth.LoginLocks) Option {
return func(a *App) { a.locks = locks }
}
// WithDownlink 注入下行发布器(未接线时测试用 RecordingDownlink)。
func WithDownlink(d port.Downlink) Option {
return func(a *App) { a.down = d }
}
// WithConnRegistry 注入连接查询。
func WithConnRegistry(c ConnRegistry) Option {
return func(a *App) { a.conns = c }
}
// WithMetrics 注入 Prometheus 注册表(投递耗时直方图)。
func WithMetrics(m *metrics.Registry) Option {
return func(a *App) { a.met = m }
}
// New 创建消息服务实现。
func New(db *store.DB, lim Limits, hash auth.HashPool, opts ...Option) *App {
if lim.RequestBurst <= 0 {
lim.RequestBurst = defaultRequestBurst
}
if lim.DeliveryWindow <= 0 {
lim.DeliveryWindow = defaultDeliveryWindow
}
if lim.ReceiptWindow <= 0 {
lim.ReceiptWindow = defaultReceiptWindow
}
if lim.AckTimeoutSeconds <= 0 {
lim.AckTimeoutSeconds = 300
}
if lim.GraceSeconds < 0 {
lim.GraceSeconds = 60
}
a := &App{
db: db,
lim: lim,
hash: hash,
nowFn: time.Now,
rates: newRateLimiter(lim.RequestsPerSecond, lim.RequestBurst),
largeSem: make(chan struct{}, maxLargeInflight),
largeHeld: make(map[string]bool),
}
for _, opt := range opts {
opt(a)
}
return a
}
func (a *App) now() time.Time {
return a.nowFn()
}
func (a *App) protocolLimits() protocol.Limits {
return protocol.Limits{
MaxBodyBytes: a.lim.MaxBodyBytes,
MaxMetaBytes: a.lim.MaxMetaBytes,
MaxFrameBytes: a.lim.MaxFrameBytes,
}
}
var _ Service = (*App)(nil)