Feat: UDP 框架
This commit is contained in:
202
cmd/internal/peer/udp_client.go
Normal file
202
cmd/internal/peer/udp_client.go
Normal file
@@ -0,0 +1,202 @@
|
||||
package peer
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync/atomic"
|
||||
|
||||
"omnisocketgo/cmd/internal/latencylog"
|
||||
"omnisocketgo/cmd/internal/protocol"
|
||||
"omnisocketgo/cmd/internal/transport"
|
||||
)
|
||||
|
||||
// UDPClient 表示一个通过 UDP 连接到 server 的 peer。
|
||||
type UDPClient struct {
|
||||
id string
|
||||
conn *transport.UDPConn
|
||||
logger latencylog.Logger
|
||||
|
||||
nextID uint64
|
||||
}
|
||||
|
||||
// DialUDP 通过 UDP 连接到 server,并发送 register 消息完成身份注册。
|
||||
func DialUDP(serverAddr, peerID string, opts ...Option) (*UDPClient, error) {
|
||||
options := clientOptions{
|
||||
logger: latencylog.NoopLogger{},
|
||||
}
|
||||
for _, opt := range opts {
|
||||
opt(&options)
|
||||
}
|
||||
if options.logger == nil {
|
||||
options.logger = latencylog.NoopLogger{}
|
||||
}
|
||||
|
||||
udpServerAddr, err := net.ResolveUDPAddr("udp", serverAddr)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("peer: resolve udp server addr %s: %w", serverAddr, err)
|
||||
}
|
||||
|
||||
var localAddr *net.UDPAddr
|
||||
if options.bindIP != "" {
|
||||
ip := net.ParseIP(options.bindIP)
|
||||
if ip == nil {
|
||||
return nil, fmt.Errorf("peer: invalid bind ip %q", options.bindIP)
|
||||
}
|
||||
localAddr = &net.UDPAddr{IP: ip}
|
||||
}
|
||||
|
||||
rawConn, err := net.DialUDP("udp", localAddr, udpServerAddr)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("peer: dial udp server %s: %w", serverAddr, err)
|
||||
}
|
||||
|
||||
conn, err := transport.NewUDPConn(
|
||||
rawConn,
|
||||
nil, // peer 侧已连接模式,不需要指定 peerAddr
|
||||
transport.WithUDPLogger(options.logger, latencylog.NodeRolePeer, peerID),
|
||||
transport.WithUDPTXTimestampDebugLogger(options.txTimestampDebugLogger),
|
||||
)
|
||||
if err != nil {
|
||||
_ = rawConn.Close()
|
||||
return nil, fmt.Errorf("peer: create udp transport conn: %w", err)
|
||||
}
|
||||
|
||||
client := &UDPClient{
|
||||
id: peerID,
|
||||
conn: conn,
|
||||
logger: options.logger,
|
||||
}
|
||||
|
||||
// 发送 register 消息完成身份注册
|
||||
if err := conn.Send(protocol.Message{
|
||||
Type: protocol.MessageTypeRegister,
|
||||
From: peerID,
|
||||
To: protocol.ServerPeerID,
|
||||
}); err != nil {
|
||||
_ = conn.Close()
|
||||
return nil, fmt.Errorf("peer: udp register with server: %w", err)
|
||||
}
|
||||
|
||||
return client, nil
|
||||
}
|
||||
|
||||
// ID 返回当前 client 的 peer 标识。
|
||||
func (c *UDPClient) ID() string {
|
||||
return c.id
|
||||
}
|
||||
|
||||
// SendText 向目标 peer 发送一条文本消息。
|
||||
func (c *UDPClient) SendText(to, body string) error {
|
||||
msg := protocol.Message{
|
||||
Type: protocol.MessageTypeText,
|
||||
ID: c.nextMessageID(),
|
||||
From: c.id,
|
||||
To: to,
|
||||
}
|
||||
latencylog.LogMessageEvent(c.logger, latencylog.NodeRolePeer, c.id, latencylog.EventAAppPrepBegin, msg)
|
||||
|
||||
msg.Body = []byte(body)
|
||||
|
||||
return c.conn.Send(msg)
|
||||
}
|
||||
|
||||
// SendFile 向目标 peer 发送一条文件消息。
|
||||
func (c *UDPClient) SendFile(to, fileName string, body []byte) error {
|
||||
msg := protocol.Message{
|
||||
Type: protocol.MessageTypeFile,
|
||||
ID: c.nextMessageID(),
|
||||
From: c.id,
|
||||
To: to,
|
||||
FileName: fileName,
|
||||
}
|
||||
latencylog.LogMessageEvent(c.logger, latencylog.NodeRolePeer, c.id, latencylog.EventAAppPrepBegin, msg)
|
||||
|
||||
bodyCopy := make([]byte, len(body))
|
||||
copy(bodyCopy, body)
|
||||
|
||||
msg.Body = bodyCopy
|
||||
|
||||
return c.conn.Send(msg)
|
||||
}
|
||||
|
||||
// SendFilePath 从本地文件读取内容并发送给目标 peer。
|
||||
func (c *UDPClient) SendFilePath(to, path string) error {
|
||||
msg := protocol.Message{
|
||||
Type: protocol.MessageTypeFile,
|
||||
ID: c.nextMessageID(),
|
||||
From: c.id,
|
||||
To: to,
|
||||
FileName: filepath.Base(path),
|
||||
}
|
||||
latencylog.LogMessageEvent(c.logger, latencylog.NodeRolePeer, c.id, latencylog.EventAAppPrepBegin, msg)
|
||||
|
||||
body, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return fmt.Errorf("peer: read file %s: %w", path, err)
|
||||
}
|
||||
|
||||
msg.Body = body
|
||||
|
||||
return c.conn.Send(msg)
|
||||
}
|
||||
|
||||
// Receive 读取一条来自 server 的消息。
|
||||
func (c *UDPClient) Receive() (protocol.Message, error) {
|
||||
msg, _, err := c.conn.Receive()
|
||||
if err != nil {
|
||||
return protocol.Message{}, fmt.Errorf("peer: udp receive from server: %w", err)
|
||||
}
|
||||
|
||||
latencylog.LogMessageEvent(c.logger, latencylog.NodeRolePeer, c.id, latencylog.EventBAppRecv, msg)
|
||||
|
||||
return msg, nil
|
||||
}
|
||||
|
||||
// ReceiveLoop 持续接收 server 消息并交给 handler 处理。
|
||||
func (c *UDPClient) ReceiveLoop(handler func(protocol.Message) error) error {
|
||||
return c.conn.ReceiveLoop(func(msg protocol.Message, _ *net.UDPAddr) error {
|
||||
switch msg.Type {
|
||||
case protocol.MessageTypeText, protocol.MessageTypeFile, protocol.MessageTypeError:
|
||||
latencylog.LogMessageEvent(c.logger, latencylog.NodeRolePeer, c.id, latencylog.EventBAppRecv, msg)
|
||||
return handler(msg)
|
||||
default:
|
||||
return fmt.Errorf("peer: unexpected message type from server: %s", msg.Type)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// PersistMessage 将收到的业务消息写入本地磁盘。
|
||||
func (c *UDPClient) PersistMessage(msg protocol.Message, inboxDir string) (string, error) {
|
||||
if !latencylog.IsBusinessMessage(msg) {
|
||||
return "", fmt.Errorf("peer: cannot persist message type %s", msg.Type)
|
||||
}
|
||||
if inboxDir == "" {
|
||||
return "", fmt.Errorf("peer: inbox directory is required")
|
||||
}
|
||||
|
||||
latencylog.LogMessageEvent(c.logger, latencylog.NodeRolePeer, c.id, latencylog.EventBPersistBegin, msg)
|
||||
|
||||
if err := os.MkdirAll(inboxDir, 0o755); err != nil {
|
||||
return "", fmt.Errorf("peer: create inbox dir %s: %w", inboxDir, err)
|
||||
}
|
||||
|
||||
path, err := persistMessageToDisk(msg, inboxDir)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
latencylog.LogMessageEvent(c.logger, latencylog.NodeRolePeer, c.id, latencylog.EventBPersistEnd, msg)
|
||||
|
||||
return path, nil
|
||||
}
|
||||
|
||||
// Close 关闭与 server 的 UDP 连接。
|
||||
func (c *UDPClient) Close() error {
|
||||
return c.conn.Close()
|
||||
}
|
||||
|
||||
func (c *UDPClient) nextMessageID() uint64 {
|
||||
return atomic.AddUint64(&c.nextID, 1)
|
||||
}
|
||||
211
cmd/internal/peer/udp_client_test.go
Normal file
211
cmd/internal/peer/udp_client_test.go
Normal file
@@ -0,0 +1,211 @@
|
||||
package peer
|
||||
|
||||
import (
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"omnisocketgo/cmd/internal/protocol"
|
||||
"omnisocketgo/cmd/internal/server"
|
||||
"omnisocketgo/cmd/internal/transport"
|
||||
)
|
||||
|
||||
// TestUDPDialAndSendText 验证 UDP 客户端可以成功连接、注册并发送文本消息。
|
||||
func TestUDPDialAndSendText(t *testing.T) {
|
||||
hubAddr := startUDPTestHub(t)
|
||||
|
||||
clientA, err := DialUDP(hubAddr.String(), "peer-a")
|
||||
if err != nil {
|
||||
t.Fatalf("DialUDP(peer-a) error = %v", err)
|
||||
}
|
||||
defer clientA.Close()
|
||||
|
||||
clientB, err := DialUDP(hubAddr.String(), "peer-b")
|
||||
if err != nil {
|
||||
t.Fatalf("DialUDP(peer-b) error = %v", err)
|
||||
}
|
||||
defer clientB.Close()
|
||||
|
||||
// 等待注册被处理
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
|
||||
// peer-a 发送文本给 peer-b
|
||||
if err := clientA.SendText("peer-b", "hello from udp"); err != nil {
|
||||
t.Fatalf("SendText() error = %v", err)
|
||||
}
|
||||
|
||||
// peer-b 接收
|
||||
msg := receiveUDPClientMessage(t, clientB)
|
||||
if msg.Type != protocol.MessageTypeText {
|
||||
t.Fatalf("message type = %s, want text", msg.Type)
|
||||
}
|
||||
if string(msg.Body) != "hello from udp" {
|
||||
t.Fatalf("message body = %q, want %q", string(msg.Body), "hello from udp")
|
||||
}
|
||||
}
|
||||
|
||||
// TestUDPClientID 验证 ID() 返回正确的 peer 标识。
|
||||
func TestUDPClientID(t *testing.T) {
|
||||
hubAddr := startUDPTestHub(t)
|
||||
|
||||
client, err := DialUDP(hubAddr.String(), "my-peer-id")
|
||||
if err != nil {
|
||||
t.Fatalf("DialUDP() error = %v", err)
|
||||
}
|
||||
defer client.Close()
|
||||
|
||||
if got := client.ID(); got != "my-peer-id" {
|
||||
t.Fatalf("ID() = %q, want %q", got, "my-peer-id")
|
||||
}
|
||||
}
|
||||
|
||||
// TestUDPClientPersistMessage 验证 UDP 客户端可以将消息持久化到磁盘。
|
||||
func TestUDPClientPersistMessage(t *testing.T) {
|
||||
hubAddr := startUDPTestHub(t)
|
||||
|
||||
client, err := DialUDP(hubAddr.String(), "peer-persist")
|
||||
if err != nil {
|
||||
t.Fatalf("DialUDP() error = %v", err)
|
||||
}
|
||||
defer client.Close()
|
||||
|
||||
inboxDir := t.TempDir()
|
||||
|
||||
// 持久化文本消息
|
||||
textMsg := protocol.Message{
|
||||
Type: protocol.MessageTypeText,
|
||||
ID: 1,
|
||||
From: "sender",
|
||||
To: "peer-persist",
|
||||
Body: []byte("persisted text"),
|
||||
}
|
||||
|
||||
path, err := client.PersistMessage(textMsg, inboxDir)
|
||||
if err != nil {
|
||||
t.Fatalf("PersistMessage(text) error = %v", err)
|
||||
}
|
||||
if !filepath.IsAbs(path) && path == "" {
|
||||
t.Fatalf("PersistMessage(text) returned empty path")
|
||||
}
|
||||
|
||||
// 持久化文件消息
|
||||
fileMsg := protocol.Message{
|
||||
Type: protocol.MessageTypeFile,
|
||||
ID: 2,
|
||||
From: "sender",
|
||||
To: "peer-persist",
|
||||
FileName: "test.bin",
|
||||
Body: []byte{0x01, 0x02, 0x03},
|
||||
}
|
||||
|
||||
filePath, err := client.PersistMessage(fileMsg, inboxDir)
|
||||
if err != nil {
|
||||
t.Fatalf("PersistMessage(file) error = %v", err)
|
||||
}
|
||||
|
||||
content, err := os.ReadFile(filePath)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFile(%s) error = %v", filePath, err)
|
||||
}
|
||||
if len(content) != 3 || content[0] != 0x01 {
|
||||
t.Fatalf("file content mismatch: got %v", content)
|
||||
}
|
||||
}
|
||||
|
||||
// TestUDPClientSendFile 验证 UDP 客户端可以发送文件消息。
|
||||
func TestUDPClientSendFile(t *testing.T) {
|
||||
hubAddr := startUDPTestHub(t)
|
||||
|
||||
clientA, err := DialUDP(hubAddr.String(), "peer-a")
|
||||
if err != nil {
|
||||
t.Fatalf("DialUDP(peer-a) error = %v", err)
|
||||
}
|
||||
defer clientA.Close()
|
||||
|
||||
clientB, err := DialUDP(hubAddr.String(), "peer-b")
|
||||
if err != nil {
|
||||
t.Fatalf("DialUDP(peer-b) error = %v", err)
|
||||
}
|
||||
defer clientB.Close()
|
||||
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
|
||||
fileBody := []byte{0xDE, 0xAD, 0xBE, 0xEF}
|
||||
if err := clientA.SendFile("peer-b", "test.bin", fileBody); err != nil {
|
||||
t.Fatalf("SendFile() error = %v", err)
|
||||
}
|
||||
|
||||
msg := receiveUDPClientMessage(t, clientB)
|
||||
if msg.Type != protocol.MessageTypeFile {
|
||||
t.Fatalf("message type = %s, want file", msg.Type)
|
||||
}
|
||||
if msg.FileName != "test.bin" {
|
||||
t.Fatalf("file name = %q, want %q", msg.FileName, "test.bin")
|
||||
}
|
||||
if len(msg.Body) != 4 {
|
||||
t.Fatalf("body length = %d, want 4", len(msg.Body))
|
||||
}
|
||||
}
|
||||
|
||||
// startUDPTestHub 创建并启动一个测试用 UDPHub。
|
||||
func startUDPTestHub(t *testing.T) *net.UDPAddr {
|
||||
t.Helper()
|
||||
|
||||
addr, err := net.ResolveUDPAddr("udp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatalf("ResolveUDPAddr() error = %v", err)
|
||||
}
|
||||
|
||||
conn, err := net.ListenUDP("udp", addr)
|
||||
if err != nil {
|
||||
t.Fatalf("ListenUDP() error = %v", err)
|
||||
}
|
||||
|
||||
hub, err := server.NewUDPHub(conn)
|
||||
if err != nil {
|
||||
_ = conn.Close()
|
||||
t.Fatalf("NewUDPHub() error = %v", err)
|
||||
}
|
||||
|
||||
go func() {
|
||||
_ = hub.Serve()
|
||||
}()
|
||||
|
||||
t.Cleanup(func() {
|
||||
_ = hub.Close()
|
||||
})
|
||||
|
||||
return conn.LocalAddr().(*net.UDPAddr)
|
||||
}
|
||||
|
||||
// receiveUDPClientMessage 从 UDP 客户端接收一条消息,带超时。
|
||||
func receiveUDPClientMessage(t *testing.T, client *UDPClient) protocol.Message {
|
||||
t.Helper()
|
||||
|
||||
type result struct {
|
||||
msg protocol.Message
|
||||
err error
|
||||
}
|
||||
|
||||
ch := make(chan result, 1)
|
||||
go func() {
|
||||
msg, err := client.Receive()
|
||||
ch <- result{msg: msg, err: err}
|
||||
}()
|
||||
|
||||
select {
|
||||
case r := <-ch:
|
||||
if r.err != nil {
|
||||
t.Fatalf("Receive() error = %v", r.err)
|
||||
}
|
||||
return r.msg
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("Receive() timed out after 2s")
|
||||
return protocol.Message{}
|
||||
}
|
||||
}
|
||||
|
||||
// Ensure transport package is used (needed for WithTXTimestampDebugLogger option).
|
||||
var _ transport.TXTimestampDebugLogger = nil
|
||||
Reference in New Issue
Block a user