package message import ( "context" "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/protocol" "git.asio.asia/nixevol/NixMsg/internal/store" ) // 消息状态(DEVELOPMENT 7.1)。 const ( StateScheduled = "scheduled" StateDispatched = "dispatched" StateCompleted = "completed" ) // 投递状态。 const ( DeliveryPending = "pending" ) // talk_grants.kind。 const ( GrantKindPassword = "password" GrantKindReply = "reply" ) // 请求频率桶默认突发容量(DEVELOPMENT 6.10;配置无单独字段)。 const defaultRequestBurst = 100 // Limits 是提交所需的配置上限(来自 config.LimitsConfig)。 type Limits struct { MaxBodyBytes int MaxMetaBytes int MaxFrameBytes int MaxTTLSeconds int64 MaxScheduleSeconds int64 RequestsPerSecond float64 RequestBurst int MaxPendingPerSender int MaxPendingPerReceiver int GraceSeconds int64 } // LimitsFromConfig 从平台配置构造 Limits。 func LimitsFromConfig(c config.LimitsConfig) Limits { burst := defaultRequestBurst return Limits{ MaxBodyBytes: c.MaxBodyBytes, MaxMetaBytes: c.MaxMetaBytes, MaxFrameBytes: c.MaxFrameBytes, MaxTTLSeconds: int64(c.MaxTTLSeconds), MaxScheduleSeconds: int64(c.MaxScheduleSeconds), RequestsPerSecond: float64(c.RequestsPerSecond), RequestBurst: burst, MaxPendingPerSender: c.MaxPendingPerSender, MaxPendingPerReceiver: c.MaxPendingPerReceiver, GraceSeconds: int64(c.GraceSeconds), } } // App 实现 Service 的提交路径(M1);其余方法暂返回未实现或空操作。 type App struct { db *store.DB lim Limits hash auth.HashPool locks auth.LoginLocks nowFn func() time.Time rates *rateLimiter } // 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 } } // New 创建消息服务实现。hash 用于校验对话密码;locks 可为 nil。 func New(db *store.DB, lim Limits, hash auth.HashPool, opts ...Option) *App { if lim.RequestBurst <= 0 { lim.RequestBurst = defaultRequestBurst } a := &App{ db: db, lim: lim, hash: hash, nowFn: time.Now, rates: newRateLimiter(lim.RequestsPerSecond, lim.RequestBurst), } 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, } } func (a *App) Ack(context.Context, string, *protocol.Ack) (AckResult, error) { return AckResult{}, ErrNotImplemented } func (a *App) Recall(context.Context, string, *protocol.Recall) (protocol.RecallData, error) { return protocol.RecallData{}, ErrNotImplemented } func (a *App) Status(context.Context, string, *protocol.Status) (any, error) { return nil, ErrNotImplemented } func (a *App) ReceiptAck(context.Context, string, *protocol.ReceiptAck) error { return ErrNotImplemented } func (a *App) DispatchDue(context.Context, int64, int) (int, error) { return 0, nil } func (a *App) PushPending(context.Context, string, port.ConnID) error { return nil } func (a *App) OnPublishDropped(context.Context, string, port.ConnID, []byte) error { return nil } func (a *App) CleanupOnce(context.Context, int64) error { return nil } func (a *App) RecoverOnStart(context.Context) error { return nil } func (a *App) WakePush(string) {} var _ Service = (*App)(nil)