mirror of
https://github.com/MHSanaei/3x-ui.git
synced 2026-10-04 13:12:07 +03:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 815c9c5772 | |||
| 05eb06f333 |
@@ -96,9 +96,7 @@ func TestAudit3ManagerEnsureActualSendersWithPersistentTraffic(t *testing.T) {
|
|||||||
if _, err := stream.Write(frame.Bytes()); err != nil {
|
if _, err := stream.Write(frame.Bytes()); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if err := stream.Close(); err != nil {
|
closeUniStream(t, stream)
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
response, err := p.client.AcceptUniStream(ctx)
|
response, err := p.client.AcceptUniStream(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
@@ -164,9 +162,7 @@ func TestAudit3ManagerEnsureActualSendersWithPersistentTraffic(t *testing.T) {
|
|||||||
if _, err := auth.Write(authBytes); err != nil {
|
if _, err := auth.Write(authBytes); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if err := auth.Close(); err != nil {
|
closeUniStream(t, auth)
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
waitForClientCongestionSender(t, server, client, served)
|
waitForClientCongestionSender(t, server, client, served)
|
||||||
var serverConn *quic.Conn
|
var serverConn *quic.Conn
|
||||||
|
|||||||
@@ -57,9 +57,7 @@ func audit3LogsStart(t *testing.T, level, marker, relayAddr string) (*Server, *c
|
|||||||
if _, err := auth.Write(frame); err != nil {
|
if _, err := auth.Write(frame); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if err := auth.Close(); err != nil {
|
closeUniStream(t, auth)
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
_, _ = authenticatedServerConnection(t, s, id)
|
_, _ = authenticatedServerConnection(t, s, id)
|
||||||
return s, c, id, password, token
|
return s, c, id, password, token
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -10,12 +10,14 @@ import (
|
|||||||
"crypto/x509"
|
"crypto/x509"
|
||||||
"crypto/x509/pkix"
|
"crypto/x509/pkix"
|
||||||
"encoding/pem"
|
"encoding/pem"
|
||||||
|
"errors"
|
||||||
"io"
|
"io"
|
||||||
"math/big"
|
"math/big"
|
||||||
"net"
|
"net"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
serverquic "github.com/apernet/quic-go"
|
||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
"github.com/quic-go/quic-go"
|
"github.com/quic-go/quic-go"
|
||||||
)
|
)
|
||||||
@@ -196,12 +198,51 @@ func testServerTCPConnectE2E(t *testing.T, controller string) {
|
|||||||
t.Fatalf("expected active email alice@example.com, got %v", activeEmails)
|
t.Fatalf("expected active email alice@example.com, got %v", activeEmails)
|
||||||
}
|
}
|
||||||
|
|
||||||
deltas := server.CollectClientTraffic()
|
waitForClientTraffic(t, server, "alice@example.com", int64(len(testMsg)))
|
||||||
if len(deltas) == 0 {
|
}
|
||||||
t.Fatalf("expected traffic deltas, got none")
|
|
||||||
|
// closeUniStream tolerates only the server's STOP_SENDING: it cancels the read side of a
|
||||||
|
// uni stream once the command is parsed, which can land before the client's FIN.
|
||||||
|
func closeUniStream(t *testing.T, stream interface {
|
||||||
|
Close() error
|
||||||
|
Context() context.Context
|
||||||
|
},
|
||||||
|
) {
|
||||||
|
t.Helper()
|
||||||
|
err := stream.Close()
|
||||||
|
if err == nil {
|
||||||
|
return
|
||||||
}
|
}
|
||||||
if deltas[0].Email != "alice@example.com" || deltas[0].Up < int64(len(testMsg)) || deltas[0].Down < int64(len(testMsg)) {
|
cause := context.Cause(stream.Context())
|
||||||
t.Fatalf("unexpected traffic deltas: %+v", deltas[0])
|
var clientErr *quic.StreamError
|
||||||
|
var serverErr *serverquic.StreamError
|
||||||
|
if (errors.As(cause, &clientErr) && clientErr.Remote) || (errors.As(cause, &serverErr) && serverErr.Remote) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
t.Fatalf("close uni stream: %v (cause %v)", err, cause)
|
||||||
|
}
|
||||||
|
|
||||||
|
// waitForClientTraffic accumulates drained deltas because the up and down counters are
|
||||||
|
// bumped on different relay goroutines, so the echo can arrive before the upload is counted.
|
||||||
|
func waitForClientTraffic(t *testing.T, server *Server, email string, minBytes int64) {
|
||||||
|
t.Helper()
|
||||||
|
var up, down int64
|
||||||
|
deadline := time.Now().Add(4 * time.Second)
|
||||||
|
for {
|
||||||
|
for _, delta := range server.CollectClientTraffic() {
|
||||||
|
if delta.Email != email {
|
||||||
|
t.Fatalf("unexpected traffic delta for %q: %+v", delta.Email, delta)
|
||||||
|
}
|
||||||
|
up += delta.Up
|
||||||
|
down += delta.Down
|
||||||
|
}
|
||||||
|
if up >= minBytes && down >= minBytes {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if time.Now().After(deadline) {
|
||||||
|
t.Fatalf("traffic for %s = up %d, down %d; want both >= %d", email, up, down, minBytes)
|
||||||
|
}
|
||||||
|
time.Sleep(5 * time.Millisecond)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -356,13 +397,7 @@ func testServerUDPDatagramE2E(t *testing.T, controller string) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// 4. Verify traffic
|
// 4. Verify traffic
|
||||||
deltas := server.CollectClientTraffic()
|
waitForClientTraffic(t, server, "bob@example.com", int64(len(udpMsg)))
|
||||||
if len(deltas) == 0 {
|
|
||||||
t.Fatalf("expected traffic deltas, got none")
|
|
||||||
}
|
|
||||||
if deltas[0].Email != "bob@example.com" || deltas[0].Up < int64(len(udpMsg)) || deltas[0].Down < int64(len(udpMsg)) {
|
|
||||||
t.Fatalf("unexpected traffic deltas: %+v", deltas[0])
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestServerUDPStreamE2E(t *testing.T) {
|
func TestServerUDPStreamE2E(t *testing.T) {
|
||||||
@@ -433,9 +468,7 @@ func testServerUDPStreamE2E(t *testing.T, controller string) {
|
|||||||
if _, err := authStream.Write(authPayload); err != nil {
|
if _, err := authStream.Write(authPayload); err != nil {
|
||||||
t.Fatalf("write authentication payload failed: %v", err)
|
t.Fatalf("write authentication payload failed: %v", err)
|
||||||
}
|
}
|
||||||
if err := authStream.Close(); err != nil {
|
closeUniStream(t, authStream)
|
||||||
t.Fatalf("close authentication stream failed: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
target := &Address{Type: AddrTypeIPv4, IP: net.ParseIP("8.8.8.8"), Port: 53}
|
target := &Address{Type: AddrTypeIPv4, IP: net.ParseIP("8.8.8.8"), Port: 53}
|
||||||
udpMsg := bytes.Repeat([]byte("s"), 8500)
|
udpMsg := bytes.Repeat([]byte("s"), 8500)
|
||||||
@@ -458,9 +491,7 @@ func testServerUDPStreamE2E(t *testing.T, controller string) {
|
|||||||
if _, err := packetStream.Write(frame.Bytes()); err != nil {
|
if _, err := packetStream.Write(frame.Bytes()); err != nil {
|
||||||
t.Fatalf("write packet frame failed: %v", err)
|
t.Fatalf("write packet frame failed: %v", err)
|
||||||
}
|
}
|
||||||
if err := packetStream.Close(); err != nil {
|
closeUniStream(t, packetStream)
|
||||||
t.Fatalf("close packet stream failed: %v", err)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
replyReassembler := newPacketReassembler(maxUdpRelayPacketSize)
|
replyReassembler := newPacketReassembler(maxUdpRelayPacketSize)
|
||||||
|
|||||||
@@ -57,9 +57,7 @@ func startLifecycleTestServer(t *testing.T, relayAddr, email string) (*Server, *
|
|||||||
if _, err := stream.Write(auth); err != nil {
|
if _, err := stream.Write(auth); err != nil {
|
||||||
t.Fatalf("write authentication: %v", err)
|
t.Fatalf("write authentication: %v", err)
|
||||||
}
|
}
|
||||||
if err := stream.Close(); err != nil {
|
closeUniStream(t, stream)
|
||||||
t.Fatalf("close authentication stream: %v", err)
|
|
||||||
}
|
|
||||||
return server, client, clientID, password
|
return server, client, clientID, password
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user