f1a00d128e
- 桌面 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 互通
147 lines
4.2 KiB
Go
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])
|
|
}
|
|
}
|