feat: 补齐模块接口契约与管理 API 文档
This commit is contained in:
@@ -0,0 +1,73 @@
|
||||
// Package group 定义 DEVELOPMENT 第 6.7 节群操作的接口。
|
||||
package group
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
// ErrNotImplemented 表示假实现未提供业务能力。
|
||||
var ErrNotImplemented = errors.New("group: not implemented")
|
||||
|
||||
// MemberFail 是建群/加人时部分失败的项。
|
||||
type MemberFail struct {
|
||||
ID string `json:"id"`
|
||||
Code string `json:"code"`
|
||||
}
|
||||
|
||||
// CreateResult 是建群结果。
|
||||
type CreateResult struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
OwnerID string `json:"owner_id"`
|
||||
Failed []MemberFail `json:"failed,omitempty"`
|
||||
}
|
||||
|
||||
// AddResult 是加人结果。
|
||||
type AddResult struct {
|
||||
Failed []MemberFail `json:"failed,omitempty"`
|
||||
}
|
||||
|
||||
// ListItem 是 group.list 一项。
|
||||
type ListItem struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
OwnerID string `json:"owner_id"`
|
||||
MemberCount int `json:"member_count"`
|
||||
}
|
||||
|
||||
// MemberItem 是 group.get 成员项。
|
||||
type MemberItem struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
Online bool `json:"online"`
|
||||
}
|
||||
|
||||
// GetResult 是群详情。
|
||||
type GetResult struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
OwnerID string `json:"owner_id"`
|
||||
Members []MemberItem `json:"members"`
|
||||
NextCursor string `json:"next_cursor,omitempty"`
|
||||
}
|
||||
|
||||
// Service 是群子系统契约(第 6.7 节)。
|
||||
type Service interface {
|
||||
Create(ctx context.Context, actorID string, req *protocol.GroupCreate) (CreateResult, error)
|
||||
Add(ctx context.Context, actorID string, req *protocol.GroupAdd) (AddResult, error)
|
||||
Remove(ctx context.Context, actorID string, req *protocol.GroupRemove) error
|
||||
Leave(ctx context.Context, actorID string, req *protocol.GroupLeave) error
|
||||
Transfer(ctx context.Context, actorID string, req *protocol.GroupTransfer) error
|
||||
Rename(ctx context.Context, actorID string, req *protocol.GroupRename) error
|
||||
Dissolve(ctx context.Context, actorID string, req *protocol.GroupDissolve) error
|
||||
List(ctx context.Context, actorID string, req *protocol.GroupList) (items []ListItem, nextCursor string, err error)
|
||||
Get(ctx context.Context, actorID string, req *protocol.GroupGet) (GetResult, error)
|
||||
|
||||
// AdminCreate 后台建群(不要求对话密码)。
|
||||
AdminCreate(ctx context.Context, name string, ownerID string, memberIDs []string) (CreateResult, error)
|
||||
// AdminAddMembers 后台加人。
|
||||
AdminAddMembers(ctx context.Context, groupID string, memberIDs []string) (AddResult, error)
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package group
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
func TestStubCreateNotImplemented(t *testing.T) {
|
||||
s := NewStub()
|
||||
_, err := s.Create(context.Background(), "a", &protocol.GroupCreate{Name: "g"})
|
||||
if !errors.Is(err, ErrNotImplemented) {
|
||||
t.Fatalf("got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStubListEmpty(t *testing.T) {
|
||||
s := NewStub()
|
||||
items, cursor, err := s.List(context.Background(), "a", &protocol.GroupList{})
|
||||
if err != nil || len(items) != 0 || cursor != "" {
|
||||
t.Fatalf("items=%v cursor=%q err=%v", items, cursor, err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
package group
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
// Stub 是测试用假实现。
|
||||
type Stub struct{}
|
||||
|
||||
func NewStub() *Stub { return &Stub{} }
|
||||
|
||||
func (s *Stub) Create(context.Context, string, *protocol.GroupCreate) (CreateResult, error) {
|
||||
return CreateResult{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Add(context.Context, string, *protocol.GroupAdd) (AddResult, error) {
|
||||
return AddResult{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Remove(context.Context, string, *protocol.GroupRemove) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Leave(context.Context, string, *protocol.GroupLeave) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Transfer(context.Context, string, *protocol.GroupTransfer) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Rename(context.Context, string, *protocol.GroupRename) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Dissolve(context.Context, string, *protocol.GroupDissolve) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) List(context.Context, string, *protocol.GroupList) ([]ListItem, string, error) {
|
||||
return nil, "", nil
|
||||
}
|
||||
|
||||
func (s *Stub) Get(context.Context, string, *protocol.GroupGet) (GetResult, error) {
|
||||
return GetResult{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) AdminCreate(context.Context, string, string, []string) (CreateResult, error) {
|
||||
return CreateResult{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) AdminAddMembers(context.Context, string, []string) (AddResult, error) {
|
||||
return AddResult{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
var _ Service = (*Stub)(nil)
|
||||
@@ -0,0 +1,63 @@
|
||||
// Package identity 定义注册、self.*、对话密码授权、停用与删除级联的接口。
|
||||
package identity
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
// ErrNotImplemented 表示假实现未提供业务能力。
|
||||
var ErrNotImplemented = errors.New("identity: not implemented")
|
||||
|
||||
// RegisterRequest 对应 POST /api/client/register 与后台开通的公共字段。
|
||||
type RegisterRequest struct {
|
||||
RegistrationCode string
|
||||
ID string
|
||||
LoginPassword string
|
||||
Name string
|
||||
TalkPassword string
|
||||
Remark string
|
||||
DefaultDelayMs int64
|
||||
Source string // admin | self
|
||||
RemoteIP string
|
||||
}
|
||||
|
||||
// RegisterResult 是注册/开通结果。
|
||||
type RegisterResult struct {
|
||||
ID string
|
||||
LoginPassword string // 仅当请求留空由服务器生成时非空
|
||||
}
|
||||
|
||||
// SelfInfo 是 self.get 返回。
|
||||
type SelfInfo struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
DefaultDelayMs int64 `json:"default_delay_ms"`
|
||||
TalkPasswordSet bool `json:"talk_password_set"`
|
||||
}
|
||||
|
||||
// Service 是身份子系统契约。
|
||||
type Service interface {
|
||||
// Register 处理自助注册或供管理开通复用(第 6.9 节)。
|
||||
Register(ctx context.Context, req RegisterRequest) (RegisterResult, error)
|
||||
|
||||
SelfGet(ctx context.Context, endpointID string) (SelfInfo, error)
|
||||
SelfUpdate(ctx context.Context, endpointID string, req *protocol.SelfUpdate) error
|
||||
SelfSetTalkPassword(ctx context.Context, endpointID string, talkPassword string) error
|
||||
SelfChangeLoginPassword(ctx context.Context, endpointID string, oldPassword, newPassword string) (sessionToken string, err error)
|
||||
SelfLogout(ctx context.Context, endpointID string) error
|
||||
|
||||
// UnlockTalk 校验并写入对话密码授权(第 6.6 节 unlock)。
|
||||
UnlockTalk(ctx context.Context, senderID, targetID, talkPassword string) error
|
||||
// HasTalkGrant 查询发送方对目标是否有有效授权。
|
||||
HasTalkGrant(ctx context.Context, senderID, targetID string) (bool, error)
|
||||
|
||||
// Disable 停用端并作废相关消息/令牌(第 7.6 节)。
|
||||
Disable(ctx context.Context, endpointID string) error
|
||||
// Enable 重新启用。
|
||||
Enable(ctx context.Context, endpointID string) error
|
||||
// Delete 删除端并做级联清理。
|
||||
Delete(ctx context.Context, endpointID string) error
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
package identity
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestStubRegisterNotImplemented(t *testing.T) {
|
||||
s := NewStub()
|
||||
_, err := s.Register(context.Background(), RegisterRequest{Source: "self"})
|
||||
if !errors.Is(err, ErrNotImplemented) {
|
||||
t.Fatalf("got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStubHasTalkGrantFalse(t *testing.T) {
|
||||
s := NewStub()
|
||||
ok, err := s.HasTalkGrant(context.Background(), "a", "b")
|
||||
if err != nil || ok {
|
||||
t.Fatalf("ok=%v err=%v", ok, err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
package identity
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
// Stub 是测试用假实现。
|
||||
type Stub struct{}
|
||||
|
||||
func NewStub() *Stub { return &Stub{} }
|
||||
|
||||
func (s *Stub) Register(context.Context, RegisterRequest) (RegisterResult, error) {
|
||||
return RegisterResult{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) SelfGet(context.Context, string) (SelfInfo, error) {
|
||||
return SelfInfo{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) SelfUpdate(context.Context, string, *protocol.SelfUpdate) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) SelfSetTalkPassword(context.Context, string, string) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) SelfChangeLoginPassword(context.Context, string, string, string) (string, error) {
|
||||
return "", ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) SelfLogout(context.Context, string) error { return ErrNotImplemented }
|
||||
|
||||
func (s *Stub) UnlockTalk(context.Context, string, string, string) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) HasTalkGrant(context.Context, string, string) (bool, error) {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
func (s *Stub) Disable(context.Context, string) error { return ErrNotImplemented }
|
||||
func (s *Stub) Enable(context.Context, string) error { return ErrNotImplemented }
|
||||
func (s *Stub) Delete(context.Context, string) error { return ErrNotImplemented }
|
||||
|
||||
var _ Service = (*Stub)(nil)
|
||||
@@ -0,0 +1,54 @@
|
||||
// Package message 定义消息提交、分发、推送、确认、撤回、回执、清理与启动恢复的接口。
|
||||
package message
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/app/port"
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
// ErrNotImplemented 表示假实现未提供业务能力。
|
||||
var ErrNotImplemented = errors.New("message: not implemented")
|
||||
|
||||
// SubmitResult 是发送提交的结果(对应 send 的 resp data)。
|
||||
type SubmitResult struct {
|
||||
ID string
|
||||
SendAtMs int64
|
||||
State string // scheduled | dispatched | ...
|
||||
}
|
||||
|
||||
// AckResult 是确认结果。
|
||||
type AckResult struct {
|
||||
Result string // accepted | recalled | expired | dropped | rejected 等最终态
|
||||
}
|
||||
|
||||
// Service 是消息子系统对外契约。方法签名供 M 线实现;T0.4 不写状态机。
|
||||
type Service interface {
|
||||
// Submit 处理端发送请求(第 7.3 节)。
|
||||
Submit(ctx context.Context, senderID string, conn port.ConnInfo, req *protocol.Send) (SubmitResult, error)
|
||||
// Ack 处理确认(第 7.6 节)。
|
||||
Ack(ctx context.Context, endpointID string, req *protocol.Ack) (AckResult, error)
|
||||
// Recall 处理撤回。
|
||||
Recall(ctx context.Context, senderID string, req *protocol.Recall) (protocol.RecallData, error)
|
||||
// Status 查询自己发出的消息状态。
|
||||
Status(ctx context.Context, senderID string, req *protocol.Status) (any, error)
|
||||
// ReceiptAck 确认回执已收下。
|
||||
ReceiptAck(ctx context.Context, endpointID string, req *protocol.ReceiptAck) error
|
||||
|
||||
// DispatchDue 分发已到点的 scheduled 消息(第 7.4 节);由调度循环调用。
|
||||
DispatchDue(ctx context.Context, nowMs int64, limit int) (dispatched int, err error)
|
||||
// PushPending 向已握手连接推送 pending 投递与回执(第 7.5 节)。
|
||||
PushPending(ctx context.Context, endpointID string, connID port.ConnID) error
|
||||
// OnPublishDropped 下行未写入发送队列时,清推送标记并安排重推。
|
||||
OnPublishDropped(ctx context.Context, endpointID string, connID port.ConnID, payload []byte) error
|
||||
|
||||
// CleanupOnce 执行一轮过期投递与记录清理(第 7.5 / 7.6 节)。
|
||||
CleanupOnce(ctx context.Context, nowMs int64) error
|
||||
// RecoverOnStart 启动恢复:清残留 pushed_conn、宽限、补发定时等(第 7.8 节)。
|
||||
RecoverOnStart(ctx context.Context) error
|
||||
|
||||
// WakePush 唤醒某端推送循环(提交/分发后由内部或其它模块调用)。
|
||||
WakePush(endpointID string)
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package message
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/app/port"
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
func TestStubSubmitNotImplemented(t *testing.T) {
|
||||
s := NewStub()
|
||||
_, err := s.Submit(context.Background(), "a", port.ConnInfo{}, &protocol.Send{})
|
||||
if !errors.Is(err, ErrNotImplemented) {
|
||||
t.Fatalf("got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStubRecoverNoop(t *testing.T) {
|
||||
s := NewStub()
|
||||
if err := s.RecoverOnStart(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,53 @@
|
||||
package message
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/app/port"
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
// Stub 是测试用假实现:方法返回 ErrNotImplemented 或固定空结果。
|
||||
type Stub struct{}
|
||||
|
||||
func NewStub() *Stub { return &Stub{} }
|
||||
|
||||
func (s *Stub) Submit(context.Context, string, port.ConnInfo, *protocol.Send) (SubmitResult, error) {
|
||||
return SubmitResult{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Ack(context.Context, string, *protocol.Ack) (AckResult, error) {
|
||||
return AckResult{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Recall(context.Context, string, *protocol.Recall) (protocol.RecallData, error) {
|
||||
return protocol.RecallData{}, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) Status(context.Context, string, *protocol.Status) (any, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) ReceiptAck(context.Context, string, *protocol.ReceiptAck) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
|
||||
func (s *Stub) DispatchDue(context.Context, int64, int) (int, error) {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func (s *Stub) PushPending(context.Context, string, port.ConnID) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Stub) OnPublishDropped(context.Context, string, port.ConnID, []byte) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Stub) CleanupOnce(context.Context, int64) error { return nil }
|
||||
|
||||
func (s *Stub) RecoverOnStart(context.Context) error { return nil }
|
||||
|
||||
func (s *Stub) WakePush(string) {}
|
||||
|
||||
var _ Service = (*Stub)(nil)
|
||||
@@ -0,0 +1,81 @@
|
||||
// Package port 定义 broker 与 app 之间的契约。
|
||||
//
|
||||
// broker(连接 N)实现 Downlink / ConnControl;app 各模块通过这些接口下行发布或踢线。
|
||||
// broker 在连接生命周期与上行帧到达时调用 UplinkHandler。
|
||||
// 本包不依赖 mochi,也不包含业务状态机。
|
||||
package port
|
||||
|
||||
import (
|
||||
"context"
|
||||
)
|
||||
|
||||
// ConnID 是连接代号(同一端编号新旧连接 ClientID 相同,用独立代号区分)。
|
||||
type ConnID string
|
||||
|
||||
// Transport 区分接入方式。
|
||||
type Transport string
|
||||
|
||||
const (
|
||||
TransportTCP Transport = "tcp"
|
||||
TransportWS Transport = "ws"
|
||||
)
|
||||
|
||||
// ConnInfo 描述一条已建立的 MQTT 连接(登录校验通过之后)。
|
||||
type ConnInfo struct {
|
||||
ConnID ConnID
|
||||
EndpointID string
|
||||
Transport Transport
|
||||
RemoteIP string
|
||||
// SessionToken 非空表示本次用密码登录后新签发的令牌(握手响应里交给端)。
|
||||
SessionToken string
|
||||
// MaxPacketSize 来自 CONNECT;0 表示未声明。
|
||||
MaxPacketSize uint32
|
||||
}
|
||||
|
||||
// HandshakeInfo 是握手完成时的补充信息。
|
||||
type HandshakeInfo struct {
|
||||
ConnInfo
|
||||
MaxReceiveBytes int // 0 表示不限(仍受 MaxPacketSize 约束)
|
||||
Client string
|
||||
}
|
||||
|
||||
// DisconnectReason 说明断开原因,便于 app 区分当前连接与被顶号的旧连接。
|
||||
type DisconnectReason string
|
||||
|
||||
const (
|
||||
DisconnectNormal DisconnectReason = "normal"
|
||||
DisconnectTakenOver DisconnectReason = "taken_over"
|
||||
DisconnectKicked DisconnectReason = "kicked"
|
||||
DisconnectFatal DisconnectReason = "fatal"
|
||||
DisconnectIdle DisconnectReason = "idle"
|
||||
)
|
||||
|
||||
// UplinkHandler 由 app 实现,broker 在钩子里调用。
|
||||
type UplinkHandler interface {
|
||||
// OnSessionEstablished 在 MQTT 会话建立、登录已通过后调用(握手前)。
|
||||
OnSessionEstablished(ctx context.Context, conn ConnInfo) error
|
||||
// OnHandshakeComplete 在 hello 成功处理后调用;此后该连接算在线并可推送。
|
||||
OnHandshakeComplete(ctx context.Context, hs HandshakeInfo) error
|
||||
// OnDisconnect 在连接断开时调用;是否为当前连接由 app 按 ConnID 判断。
|
||||
OnDisconnect(ctx context.Context, conn ConnInfo, reason DisconnectReason)
|
||||
// HandleUplink 处理上行应用帧原始 JSON(已从 MQTT 发布拷贝)。
|
||||
HandleUplink(ctx context.Context, conn ConnInfo, payload []byte) error
|
||||
}
|
||||
|
||||
// PublishOpts 控制下行发布。
|
||||
type PublishOpts struct {
|
||||
QoS byte // 0 或 1;msg/receipt/revoked/resp/fatal 用 1,presence/group_event 用 0
|
||||
}
|
||||
|
||||
// Downlink 由 broker 实现,供 app 向下行主题发布。
|
||||
type Downlink interface {
|
||||
// PublishDown 向 nix/c/{endpointID}/down 发布一帧。
|
||||
// connID 非空时仅在该连接仍是当前连接时发布;空表示发给该端当前连接。
|
||||
PublishDown(ctx context.Context, endpointID string, connID ConnID, payload []byte, opts PublishOpts) error
|
||||
}
|
||||
|
||||
// ConnControl 由 broker 实现,供 app/admin 踢线或发 fatal 后断开。
|
||||
type ConnControl interface {
|
||||
// Disconnect 断开指定连接;connID 为空则断开该端当前连接。
|
||||
Disconnect(ctx context.Context, endpointID string, connID ConnID, reason DisconnectReason) error
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
package port
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestStubDownlinkPublish(t *testing.T) {
|
||||
d := &StubDownlink{}
|
||||
if err := d.PublishDown(context.Background(), "a", "c1", []byte(`{}`), PublishOpts{QoS: 1}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if d.Published != 1 {
|
||||
t.Fatalf("published=%d", d.Published)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStubUplinkHandlerCompile(t *testing.T) {
|
||||
var h UplinkHandler = StubUplinkHandler{}
|
||||
if err := h.OnSessionEstablished(context.Background(), ConnInfo{EndpointID: "a"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,40 @@
|
||||
package port
|
||||
|
||||
import "context"
|
||||
|
||||
// StubDownlink 吞掉所有下行发布。
|
||||
type StubDownlink struct {
|
||||
Published int
|
||||
}
|
||||
|
||||
func (s *StubDownlink) PublishDown(_ context.Context, _ string, _ ConnID, _ []byte, _ PublishOpts) error {
|
||||
s.Published++
|
||||
return nil
|
||||
}
|
||||
|
||||
// StubConnControl 记录断开请求。
|
||||
type StubConnControl struct {
|
||||
Calls []string
|
||||
}
|
||||
|
||||
func (s *StubConnControl) Disconnect(_ context.Context, endpointID string, _ ConnID, _ DisconnectReason) error {
|
||||
s.Calls = append(s.Calls, endpointID)
|
||||
return nil
|
||||
}
|
||||
|
||||
// StubUplinkHandler 空实现,供 broker 接线测试。
|
||||
type StubUplinkHandler struct{}
|
||||
|
||||
func (StubUplinkHandler) OnSessionEstablished(context.Context, ConnInfo) error { return nil }
|
||||
|
||||
func (StubUplinkHandler) OnHandshakeComplete(context.Context, HandshakeInfo) error { return nil }
|
||||
|
||||
func (StubUplinkHandler) OnDisconnect(context.Context, ConnInfo, DisconnectReason) {}
|
||||
|
||||
func (StubUplinkHandler) HandleUplink(context.Context, ConnInfo, []byte) error { return nil }
|
||||
|
||||
var (
|
||||
_ Downlink = (*StubDownlink)(nil)
|
||||
_ ConnControl = (*StubConnControl)(nil)
|
||||
_ UplinkHandler = StubUplinkHandler{}
|
||||
)
|
||||
@@ -0,0 +1,53 @@
|
||||
// Package presence 定义在线查询、目录与上下线订阅的接口。
|
||||
package presence
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/app/port"
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
// ErrNotImplemented 表示假实现未提供业务能力。
|
||||
var ErrNotImplemented = errors.New("presence: not implemented")
|
||||
|
||||
// StatusItem 是 presence.get 单项。
|
||||
type StatusItem struct {
|
||||
ID string `json:"id"`
|
||||
Online bool `json:"online"`
|
||||
SinceMs int64 `json:"since_ms"`
|
||||
// NotFound 为 true 时该项对应未知编号(协议用 not_found 表达,实现可汇总)。
|
||||
NotFound bool `json:"-"`
|
||||
}
|
||||
|
||||
// DirectoryItem 是目录一项。
|
||||
type DirectoryItem struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
Online bool `json:"online"`
|
||||
OnlineSinceMs *int64 `json:"online_since_ms"`
|
||||
OfflineSinceMs *int64 `json:"offline_since_ms"`
|
||||
TalkPasswordSet bool `json:"talk_password_set"`
|
||||
}
|
||||
|
||||
// Service 是在线与目录子系统契约(第 6.5 节)。
|
||||
type Service interface {
|
||||
// Get 查询若干编号在线状态。
|
||||
Get(ctx context.Context, ids []string) ([]StatusItem, error)
|
||||
// Directory 分页目录。
|
||||
Directory(ctx context.Context, req *protocol.DirectoryList) (items []DirectoryItem, nextCursor string, err error)
|
||||
// Watch 覆盖本连接的上下线订阅;断线由 ClearWatch 清空。
|
||||
Watch(ctx context.Context, connID port.ConnID, endpointID string, req *protocol.PresenceWatch) error
|
||||
// ClearWatch 在连接断开时清空订阅。
|
||||
ClearWatch(connID port.ConnID)
|
||||
|
||||
// SetOnline 握手完成时标记在线并通知订阅者。
|
||||
SetOnline(ctx context.Context, endpointID string, connID port.ConnID, atMs int64) error
|
||||
// SetOffline 当前连接断开时标记离线并通知订阅者。
|
||||
SetOffline(ctx context.Context, endpointID string, connID port.ConnID, atMs int64) error
|
||||
// IsOnline 查询端是否有已握手连接。
|
||||
IsOnline(endpointID string) bool
|
||||
// CurrentConn 返回端的当前连接代号;无则空。
|
||||
CurrentConn(endpointID string) (port.ConnID, bool)
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
package presence
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestStubGetOffline(t *testing.T) {
|
||||
s := NewStub()
|
||||
items, err := s.Get(context.Background(), []string{"a", "b"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(items) != 2 || items[0].Online || items[1].Online {
|
||||
t.Fatalf("%+v", items)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStubIsOnlineFalse(t *testing.T) {
|
||||
if NewStub().IsOnline("x") {
|
||||
t.Fatal("expected offline")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package presence
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/app/port"
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
// Stub 是测试用假实现:全部视为离线。
|
||||
type Stub struct{}
|
||||
|
||||
func NewStub() *Stub { return &Stub{} }
|
||||
|
||||
func (s *Stub) Get(_ context.Context, ids []string) ([]StatusItem, error) {
|
||||
out := make([]StatusItem, 0, len(ids))
|
||||
for _, id := range ids {
|
||||
out = append(out, StatusItem{ID: id, Online: false, SinceMs: 0})
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (s *Stub) Directory(context.Context, *protocol.DirectoryList) ([]DirectoryItem, string, error) {
|
||||
return nil, "", nil
|
||||
}
|
||||
|
||||
func (s *Stub) Watch(context.Context, port.ConnID, string, *protocol.PresenceWatch) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Stub) ClearWatch(port.ConnID) {}
|
||||
|
||||
func (s *Stub) SetOnline(context.Context, string, port.ConnID, int64) error { return nil }
|
||||
|
||||
func (s *Stub) SetOffline(context.Context, string, port.ConnID, int64) error { return nil }
|
||||
|
||||
func (s *Stub) IsOnline(string) bool { return false }
|
||||
|
||||
func (s *Stub) CurrentConn(string) (port.ConnID, bool) { return "", false }
|
||||
|
||||
var _ Service = (*Stub)(nil)
|
||||
@@ -0,0 +1,95 @@
|
||||
// Package auth 定义密码哈希池、会话令牌、API 令牌与登录锁定计数的接口。
|
||||
// 本包在 T0.4 只提供契约与测试假实现;真正的 argon2 池与令牌逻辑由平台 P 实现。
|
||||
package auth
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"time"
|
||||
)
|
||||
|
||||
// ErrNotImplemented 表示假实现尚未提供业务能力。
|
||||
var ErrNotImplemented = errors.New("auth: not implemented")
|
||||
|
||||
// PasswordKind 区分哈希用途(便于日志与限流,不影响算法)。
|
||||
type PasswordKind string
|
||||
|
||||
const (
|
||||
PasswordLogin PasswordKind = "login"
|
||||
PasswordTalk PasswordKind = "talk"
|
||||
PasswordAdmin PasswordKind = "admin"
|
||||
)
|
||||
|
||||
// HashPool 是 argon2id 并发池(DEVELOPMENT 第 12 节)。
|
||||
type HashPool interface {
|
||||
// Hash 计算 PHC 格式哈希。
|
||||
Hash(ctx context.Context, kind PasswordKind, password string) (phc string, err error)
|
||||
// Verify 常量时间比较;ok 为 false 时 err 仍可为 nil(密码不匹配)。
|
||||
Verify(ctx context.Context, kind PasswordKind, password, phc string) (ok bool, err error)
|
||||
// QueueLen 返回等待哈希的任务数(指标用)。
|
||||
QueueLen() int
|
||||
}
|
||||
|
||||
// SessionTokens 管理端会话令牌(nst_ 前缀,库中存 SHA-256)。
|
||||
type SessionTokens interface {
|
||||
// Issue 生成新令牌明文,返回明文与 SHA-256 十六进制(或原始哈希字节由实现约定)。
|
||||
Issue(ctx context.Context) (token string, hash []byte, err error)
|
||||
// HashToken 对已有令牌做 SHA-256(校验用,不走 argon2)。
|
||||
HashToken(token string) []byte
|
||||
// LooksLikeSessionToken 判断密码字段是否以 nst_ 开头。
|
||||
LooksLikeSessionToken(credential string) bool
|
||||
}
|
||||
|
||||
// APITokenInfo 是列表项(不含令牌明文)。
|
||||
type APITokenInfo struct {
|
||||
ID int64
|
||||
Name string
|
||||
Enabled bool
|
||||
CreatedAt time.Time
|
||||
LastUsedAt *time.Time
|
||||
}
|
||||
|
||||
// APITokens 管理后台 API 令牌(nxm_ 前缀)。
|
||||
type APITokens interface {
|
||||
Issue(ctx context.Context) (token string, hash []byte, err error)
|
||||
HashToken(token string) []byte
|
||||
LooksLikeAPIToken(credential string) bool
|
||||
}
|
||||
|
||||
// LockKind 区分锁定计数器类型。
|
||||
type LockKind string
|
||||
|
||||
const (
|
||||
// LockLoginEndpointIP:编号 + IP,5 分钟内 10 次错 → 锁该组合 5 分钟。
|
||||
LockLoginEndpointIP LockKind = "login_endpoint_ip"
|
||||
// LockLoginEndpoint:编号总数,1 小时内 50 次错 → 暂停该编号密码登录 1 小时。
|
||||
LockLoginEndpoint LockKind = "login_endpoint"
|
||||
// LockTalkPair:发送方 + 对方对话密码。
|
||||
LockTalkPair LockKind = "talk_pair"
|
||||
// LockTalkTarget:对方对话密码总数。
|
||||
LockTalkTarget LockKind = "talk_target"
|
||||
// LockAdminIP:管理员登录 / 错误 API 令牌。
|
||||
LockAdminIP LockKind = "admin_ip"
|
||||
// LockRegisterIP:注册安全码错误。
|
||||
LockRegisterIP LockKind = "register_ip"
|
||||
)
|
||||
|
||||
// LockKey 是一次锁定查询/计数的键。
|
||||
type LockKey struct {
|
||||
Kind LockKind
|
||||
EndpointID string // 端编号;管理员/注册场景可空
|
||||
PeerID string // 对话密码对方;可空
|
||||
IP string // 来源 IP;按编号总数等场景可空
|
||||
}
|
||||
|
||||
// LoginLocks 是内存锁定计数器(重启清零)。
|
||||
type LoginLocks interface {
|
||||
// Check 若当前已锁定返回 locked=true 与剩余时间。
|
||||
Check(key LockKey) (locked bool, retryAfter time.Duration)
|
||||
// Fail 记录一次失败;若因此触发锁定,返回 locked=true。
|
||||
Fail(key LockKey) (locked bool, retryAfter time.Duration)
|
||||
// ClearEndpoint 清除某端编号相关的登录锁定(两种都清),对应管理 unlock。
|
||||
ClearEndpoint(endpointID string)
|
||||
// Clear 清除精确键。
|
||||
Clear(key LockKey)
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
package auth
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestStubHashPoolRoundTrip(t *testing.T) {
|
||||
p := NewStubHashPool()
|
||||
h, err := p.Hash(context.Background(), PasswordLogin, "secret")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ok, err := p.Verify(context.Background(), PasswordLogin, "secret", h)
|
||||
if err != nil || !ok {
|
||||
t.Fatalf("verify: ok=%v err=%v", ok, err)
|
||||
}
|
||||
ok, _ = p.Verify(context.Background(), PasswordLogin, "wrong", h)
|
||||
if ok {
|
||||
t.Fatal("expected mismatch")
|
||||
}
|
||||
}
|
||||
|
||||
func TestStubSessionTokenPrefix(t *testing.T) {
|
||||
s := NewStubSessionTokens()
|
||||
tok, hash, err := s.Issue(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !s.LooksLikeSessionToken(tok) {
|
||||
t.Fatalf("token %q should look like session", tok)
|
||||
}
|
||||
if len(hash) != 32 {
|
||||
t.Fatalf("hash len %d", len(hash))
|
||||
}
|
||||
}
|
||||
|
||||
func TestStubAPITokenPrefix(t *testing.T) {
|
||||
s := NewStubAPITokens()
|
||||
tok, _, err := s.Issue(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !s.LooksLikeAPIToken(tok) {
|
||||
t.Fatalf("token %q should look like api", tok)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStubLoginLocksFailRecord(t *testing.T) {
|
||||
l := NewStubLoginLocks()
|
||||
locked, _ := l.Fail(LockKey{Kind: LockLoginEndpointIP, EndpointID: "a", IP: "1.2.3.4"})
|
||||
if locked {
|
||||
t.Fatal("stub should not lock")
|
||||
}
|
||||
l.ClearEndpoint("a")
|
||||
if len(l.Cleared) != 1 || l.Cleared[0] != "a" {
|
||||
t.Fatalf("cleared=%v", l.Cleared)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,105 @@
|
||||
package auth
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// StubHashPool 是测试用假哈希池:明文加前缀,不做 argon2。
|
||||
type StubHashPool struct {
|
||||
mu sync.Mutex
|
||||
queue int
|
||||
}
|
||||
|
||||
func NewStubHashPool() *StubHashPool { return &StubHashPool{} }
|
||||
|
||||
func (p *StubHashPool) Hash(_ context.Context, _ PasswordKind, password string) (string, error) {
|
||||
return "stub$" + password, nil
|
||||
}
|
||||
|
||||
func (p *StubHashPool) Verify(_ context.Context, _ PasswordKind, password, phc string) (bool, error) {
|
||||
return phc == "stub$"+password, nil
|
||||
}
|
||||
|
||||
func (p *StubHashPool) QueueLen() int {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
return p.queue
|
||||
}
|
||||
|
||||
// StubSessionTokens 假会话令牌。
|
||||
type StubSessionTokens struct{}
|
||||
|
||||
func NewStubSessionTokens() *StubSessionTokens { return &StubSessionTokens{} }
|
||||
|
||||
func (s *StubSessionTokens) Issue(_ context.Context) (string, []byte, error) {
|
||||
tok := "nst_stub_session_token_000000000000"
|
||||
sum := sha256.Sum256([]byte(tok))
|
||||
return tok, sum[:], nil
|
||||
}
|
||||
|
||||
func (s *StubSessionTokens) HashToken(token string) []byte {
|
||||
sum := sha256.Sum256([]byte(token))
|
||||
return sum[:]
|
||||
}
|
||||
|
||||
func (s *StubSessionTokens) LooksLikeSessionToken(credential string) bool {
|
||||
return strings.HasPrefix(credential, "nst_")
|
||||
}
|
||||
|
||||
// StubAPITokens 假 API 令牌。
|
||||
type StubAPITokens struct{}
|
||||
|
||||
func NewStubAPITokens() *StubAPITokens { return &StubAPITokens{} }
|
||||
|
||||
func (s *StubAPITokens) Issue(_ context.Context) (string, []byte, error) {
|
||||
tok := "nxm_stub_api_token_0000000000000000"
|
||||
sum := sha256.Sum256([]byte(tok))
|
||||
return tok, sum[:], nil
|
||||
}
|
||||
|
||||
func (s *StubAPITokens) HashToken(token string) []byte {
|
||||
sum := sha256.Sum256([]byte(token))
|
||||
return sum[:]
|
||||
}
|
||||
|
||||
func (s *StubAPITokens) LooksLikeAPIToken(credential string) bool {
|
||||
return strings.HasPrefix(credential, "nxm_")
|
||||
}
|
||||
|
||||
// StubLoginLocks 假锁定:永不锁定,记录调用便于测试。
|
||||
type StubLoginLocks struct {
|
||||
mu sync.Mutex
|
||||
Fails []LockKey
|
||||
Cleared []string
|
||||
}
|
||||
|
||||
func NewStubLoginLocks() *StubLoginLocks { return &StubLoginLocks{} }
|
||||
|
||||
func (l *StubLoginLocks) Check(LockKey) (bool, time.Duration) { return false, 0 }
|
||||
|
||||
func (l *StubLoginLocks) Fail(key LockKey) (bool, time.Duration) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
l.Fails = append(l.Fails, key)
|
||||
return false, 0
|
||||
}
|
||||
|
||||
func (l *StubLoginLocks) ClearEndpoint(endpointID string) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
l.Cleared = append(l.Cleared, endpointID)
|
||||
}
|
||||
|
||||
func (l *StubLoginLocks) Clear(LockKey) {}
|
||||
|
||||
// 编译期检查:假实现满足接口。
|
||||
var (
|
||||
_ HashPool = (*StubHashPool)(nil)
|
||||
_ SessionTokens = (*StubSessionTokens)(nil)
|
||||
_ APITokens = (*StubAPITokens)(nil)
|
||||
_ LoginLocks = (*StubLoginLocks)(nil)
|
||||
)
|
||||
Reference in New Issue
Block a user