Files
admin f1a00d128e 原生数据面协议层:桌面 engine 进度节流 / pion 日志 + iOS 数据面路由
- 桌面 engine(pion 数据面,承 ab57afd 之后的精修):
  - 新增 logging.go——把 pion 内部日志路由到宿主 OnLog(真机无 stderr,连接失败无从查);仅 ice/mdns 作用域放 Debug、余 Info+、Trace 丢弃,按 session 聚合
  - wire.go / session.go / engine.go:进度回调 ~10Hz 节流(progressEmitThrottleMs,高吞吐下每片一回调打满宿主主线程;终态由 emitProgressNow 强发最终值),与 web store push 同量级;engine_test.go 跟进
- web 数据面路由:p2p.ts 翻 IOS_NATIVE=true(iOS 走原生 libwebrtc 引擎)+ 新增 p2pIos.ts(iOS p2p 后端)+ p2pNative.ts / net/ios.ts 跟进
- 数据面设计文档 NATIVE-TRANSFER.md:U1(gomobile+pion)真机证伪 → 翻案 U2(libwebrtc)的依据与实测证据(§7/§8)
- 线协议三端一致(cdrop-file ordered / meta+chunk(64KB)+done+ack / 16MB-4MB 水位 / ack 追平完成),桌面 pion ↔ iOS libwebrtc ↔ web JS 互通
2026-06-28 19:32:37 +08:00

147 lines
4.2 KiB
Go

package engine
import (
"bytes"
"os"
"path/filepath"
"testing"
"time"
)
// sigQueue 是测试里每个引擎的入站信令 FIFO:单消费者保序处理,确保 offer/answer 先于
// 其后 trickle 的 ICE 候选被处理(pion 的 AddICECandidate 在 remote description 未设时报错)。
type sigQueue struct {
ch chan [2]string // [fromPeer, payloadJSON]
eng *Engine
}
func newSigQueue() *sigQueue { return &sigQueue{ch: make(chan [2]string, 256)} }
func (q *sigQueue) run(t *testing.T, done <-chan struct{}) {
go func() {
for {
select {
case <-done:
return
case m := <-q.ch:
if err := q.eng.HandleSignal(m[0], m[1]); err != nil {
t.Logf("HandleSignal(%s) err: %v", m[0], err)
}
}
}
}()
}
// peerEvents 把一端的出站信令投到对端队列,并暴露 saved/failed 通道供断言。
type peerEvents struct {
t *testing.T
selfName string
out *sigQueue // 对端队列
saved chan string
failed chan string
}
func (p *peerEvents) OnProgress(string, int64) {}
func (p *peerEvents) OnState(_ string, st string) {
if st == "failed" {
select {
case p.failed <- st:
default:
}
}
}
func (p *peerEvents) OnSignal(_ string, _ string, payloadJSON string) {
// 对端看到的 fromPeer 即本端名字。
p.out.ch <- [2]string{p.selfName, payloadJSON}
}
func (p *peerEvents) OnSaved(_ string, path string) {
select {
case p.saved <- path:
default:
}
}
func (p *peerEvents) OnIcePair(_ string, _ string, _ string) {}
func (p *peerEvents) OnLog(line string) { p.t.Log(line) }
// TestEngineTransfer 端到端:两个进程内 pion peer(A 发 / B 收)经 host 候选直连,传一个
// 3MB 已知内容的文件,断言落盘字节与源一致。覆盖线协议(meta/chunk/done/ack)、水位回压、
// 接收端落盘改名。
func TestEngineTransfer(t *testing.T) {
dir := t.TempDir()
srcPath := filepath.Join(dir, "payload.bin")
dlDir := filepath.Join(dir, "downloads")
if err := os.MkdirAll(dlDir, 0o755); err != nil {
t.Fatal(err)
}
// 3MB 确定性内容:触发多片 + 至少一次回压窗口。
const size = 3 * 1024 * 1024
want := make([]byte, size)
for i := range want {
want[i] = byte(i*7 + 3)
}
if err := os.WriteFile(srcPath, want, 0o644); err != nil {
t.Fatal(err)
}
qA, qB := newSigQueue(), newSigQueue()
evA := &peerEvents{t: t, selfName: "A", out: qB, saved: make(chan string, 1), failed: make(chan string, 1)}
evB := &peerEvents{t: t, selfName: "B", out: qA, saved: make(chan string, 1), failed: make(chan string, 1)}
engA := New(Config{DownloadDir: dlDir}, evA) // 发送端(不落盘)
engB := New(Config{DownloadDir: dlDir}, evB) // 接收端
qA.eng, qB.eng = engA, engB
done := make(chan struct{})
defer close(done)
qA.run(t, done)
qB.run(t, done)
const sid = "sess-1"
if err := engB.StartIncoming(sid, "A", ""); err != nil { // 先武装接收端
t.Fatalf("StartIncoming: %v", err)
}
if err := engA.StartOutgoing(sid, "B", srcPath, ""); err != nil {
t.Fatalf("StartOutgoing: %v", err)
}
var savedPath string
select {
case savedPath = <-evB.saved:
case <-evB.failed:
t.Fatal("receiver reported failed")
case <-evA.failed:
t.Fatal("sender reported failed")
case <-time.After(20 * time.Second):
t.Fatal("timeout waiting for transfer to complete")
}
got, err := os.ReadFile(savedPath)
if err != nil {
t.Fatalf("read saved file: %v", err)
}
if len(got) != len(want) {
t.Fatalf("size mismatch: got %d want %d", len(got), len(want))
}
if !bytes.Equal(got, want) {
t.Fatal("content mismatch between source and received file")
}
// 给发送端 waitAck/收尾一点时间收束,避免泄漏 goroutine 噪声(非断言项)。
time.Sleep(200 * time.Millisecond)
}
// TestParseICEServers 验证 urls 单串 / 数组两种形态都能解析(与 web getICEServers 输出兼容)。
func TestParseICEServers(t *testing.T) {
out := parseICEServers(`[{"urls":"stun:a:1"},{"urls":["turns:b:2","turn:b:3"],"username":"u","credential":"c"}]`)
if len(out) != 2 {
t.Fatalf("want 2 servers, got %d", len(out))
}
if len(out[0].URLs) != 1 || out[0].URLs[0] != "stun:a:1" {
t.Errorf("server0 urls wrong: %+v", out[0].URLs)
}
if len(out[1].URLs) != 2 || out[1].Username != "u" {
t.Errorf("server1 wrong: %+v", out[1])
}
}