2024-06-18 05:36:36 +00:00
|
|
|
package splithttp
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
gotls "crypto/tls"
|
2024-12-12 12:19:18 +00:00
|
|
|
"fmt"
|
2024-11-27 20:19:18 +00:00
|
|
|
"io"
|
2024-06-18 05:36:36 +00:00
|
|
|
"net/http"
|
2024-12-11 14:05:39 +00:00
|
|
|
"net/http/httptrace"
|
2024-06-18 05:36:36 +00:00
|
|
|
"net/url"
|
|
|
|
"strconv"
|
|
|
|
"sync"
|
2024-12-15 05:43:10 +00:00
|
|
|
"sync/atomic"
|
2024-06-18 05:36:36 +00:00
|
|
|
"time"
|
|
|
|
|
2024-07-17 12:10:48 +00:00
|
|
|
"github.com/quic-go/quic-go"
|
|
|
|
"github.com/quic-go/quic-go/http3"
|
2024-06-18 05:36:36 +00:00
|
|
|
"github.com/xtls/xray-core/common"
|
|
|
|
"github.com/xtls/xray-core/common/buf"
|
2024-06-29 18:32:57 +00:00
|
|
|
"github.com/xtls/xray-core/common/errors"
|
2024-06-18 05:36:36 +00:00
|
|
|
"github.com/xtls/xray-core/common/net"
|
2024-12-11 14:05:39 +00:00
|
|
|
"github.com/xtls/xray-core/common/signal/done"
|
2024-06-18 05:36:36 +00:00
|
|
|
"github.com/xtls/xray-core/common/uuid"
|
|
|
|
"github.com/xtls/xray-core/transport/internet"
|
2024-07-11 07:56:20 +00:00
|
|
|
"github.com/xtls/xray-core/transport/internet/browser_dialer"
|
2024-10-30 02:31:05 +00:00
|
|
|
"github.com/xtls/xray-core/transport/internet/reality"
|
2024-06-18 05:36:36 +00:00
|
|
|
"github.com/xtls/xray-core/transport/internet/stat"
|
|
|
|
"github.com/xtls/xray-core/transport/internet/tls"
|
|
|
|
"github.com/xtls/xray-core/transport/pipe"
|
|
|
|
"golang.org/x/net/http2"
|
|
|
|
)
|
|
|
|
|
2024-08-10 05:40:48 +00:00
|
|
|
// defines the maximum time an idle TCP session can survive in the tunnel, so
|
|
|
|
// it should be consistent across HTTP versions and with other transports.
|
|
|
|
const connIdleTimeout = 300 * time.Second
|
|
|
|
|
|
|
|
// consistent with quic-go
|
2024-11-29 00:57:45 +00:00
|
|
|
const quicgoH3KeepAlivePeriod = 10 * time.Second
|
2024-08-10 05:40:48 +00:00
|
|
|
|
|
|
|
// consistent with chrome
|
2024-11-29 00:57:45 +00:00
|
|
|
const chromeH2KeepAlivePeriod = 45 * time.Second
|
2024-08-10 05:40:48 +00:00
|
|
|
|
2024-06-18 05:36:36 +00:00
|
|
|
type dialerConf struct {
|
|
|
|
net.Destination
|
|
|
|
*internet.MemoryStreamConfig
|
|
|
|
}
|
|
|
|
|
|
|
|
var (
|
2024-12-15 05:43:10 +00:00
|
|
|
globalDialerMap map[dialerConf]*XmuxManager
|
2024-06-18 05:36:36 +00:00
|
|
|
globalDialerAccess sync.Mutex
|
|
|
|
)
|
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
func getHTTPClient(ctx context.Context, dest net.Destination, streamSettings *internet.MemoryStreamConfig) (DialerClient, *XmuxClient) {
|
2024-10-30 02:31:05 +00:00
|
|
|
realityConfig := reality.ConfigFromStreamSettings(streamSettings)
|
|
|
|
|
|
|
|
if browser_dialer.HasBrowserDialer() && realityConfig != nil {
|
2024-09-16 12:42:01 +00:00
|
|
|
return &BrowserDialerClient{}, nil
|
2024-07-11 07:56:20 +00:00
|
|
|
}
|
|
|
|
|
2024-06-18 05:36:36 +00:00
|
|
|
globalDialerAccess.Lock()
|
|
|
|
defer globalDialerAccess.Unlock()
|
|
|
|
|
|
|
|
if globalDialerMap == nil {
|
2024-12-15 05:43:10 +00:00
|
|
|
globalDialerMap = make(map[dialerConf]*XmuxManager)
|
2024-09-16 12:42:01 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
key := dialerConf{dest, streamSettings}
|
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
xmuxManager, found := globalDialerMap[key]
|
2024-09-16 12:42:01 +00:00
|
|
|
|
|
|
|
if !found {
|
|
|
|
transportConfig := streamSettings.ProtocolSettings.(*Config)
|
2024-12-15 05:43:10 +00:00
|
|
|
var xmuxConfig XmuxConfig
|
2024-09-16 12:42:01 +00:00
|
|
|
if transportConfig.Xmux != nil {
|
2024-12-15 05:43:10 +00:00
|
|
|
xmuxConfig = *transportConfig.Xmux
|
2024-09-16 12:42:01 +00:00
|
|
|
}
|
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
xmuxManager = NewXmuxManager(xmuxConfig, func() XmuxConn {
|
2024-09-16 12:42:01 +00:00
|
|
|
return createHTTPClient(dest, streamSettings)
|
|
|
|
})
|
2024-12-15 05:43:10 +00:00
|
|
|
globalDialerMap[key] = xmuxManager
|
2024-06-18 05:36:36 +00:00
|
|
|
}
|
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
xmuxClient := xmuxManager.GetXmuxClient(ctx)
|
|
|
|
return xmuxClient.XmuxConn.(DialerClient), xmuxClient
|
2024-09-16 12:42:01 +00:00
|
|
|
}
|
|
|
|
|
2024-12-12 12:19:18 +00:00
|
|
|
func decideHTTPVersion(tlsConfig *tls.Config, realityConfig *reality.Config) string {
|
|
|
|
if realityConfig != nil {
|
|
|
|
return "2"
|
|
|
|
}
|
|
|
|
if tlsConfig == nil {
|
|
|
|
return "1.1"
|
|
|
|
}
|
|
|
|
if len(tlsConfig.NextProtocol) != 1 {
|
|
|
|
return "2"
|
|
|
|
}
|
|
|
|
if tlsConfig.NextProtocol[0] == "http/1.1" {
|
|
|
|
return "1.1"
|
|
|
|
}
|
|
|
|
if tlsConfig.NextProtocol[0] == "h3" {
|
|
|
|
return "3"
|
|
|
|
}
|
|
|
|
return "2"
|
|
|
|
}
|
|
|
|
|
2024-09-16 12:42:01 +00:00
|
|
|
func createHTTPClient(dest net.Destination, streamSettings *internet.MemoryStreamConfig) DialerClient {
|
|
|
|
tlsConfig := tls.ConfigFromStreamSettings(streamSettings)
|
2024-10-30 02:31:05 +00:00
|
|
|
realityConfig := reality.ConfigFromStreamSettings(streamSettings)
|
|
|
|
|
2024-12-12 12:19:18 +00:00
|
|
|
httpVersion := decideHTTPVersion(tlsConfig, realityConfig)
|
|
|
|
if httpVersion == "3" {
|
|
|
|
dest.Network = net.Network_UDP // better to keep this line
|
2024-07-21 08:55:03 +00:00
|
|
|
}
|
2024-06-18 05:36:36 +00:00
|
|
|
|
|
|
|
var gotlsConfig *gotls.Config
|
|
|
|
|
|
|
|
if tlsConfig != nil {
|
|
|
|
gotlsConfig = tlsConfig.GetTLSConfig(tls.WithDestination(dest))
|
|
|
|
}
|
|
|
|
|
2024-09-16 12:42:01 +00:00
|
|
|
transportConfig := streamSettings.ProtocolSettings.(*Config)
|
|
|
|
|
2024-06-18 05:36:36 +00:00
|
|
|
dialContext := func(ctxInner context.Context) (net.Conn, error) {
|
2024-06-20 23:30:51 +00:00
|
|
|
conn, err := internet.DialSystem(ctxInner, dest, streamSettings.SocketSettings)
|
2024-06-18 05:36:36 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2024-10-30 02:31:05 +00:00
|
|
|
if realityConfig != nil {
|
|
|
|
return reality.UClient(conn, realityConfig, ctxInner, dest)
|
|
|
|
}
|
|
|
|
|
2024-06-18 05:36:36 +00:00
|
|
|
if gotlsConfig != nil {
|
|
|
|
if fingerprint := tls.GetFingerprint(tlsConfig.Fingerprint); fingerprint != nil {
|
|
|
|
conn = tls.UClient(conn, gotlsConfig, fingerprint)
|
2024-06-20 23:30:51 +00:00
|
|
|
if err := conn.(*tls.UConn).HandshakeContext(ctxInner); err != nil {
|
2024-06-18 05:36:36 +00:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
conn = tls.Client(conn, gotlsConfig)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return conn, nil
|
|
|
|
}
|
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
var keepAlivePeriod time.Duration
|
|
|
|
if streamSettings.ProtocolSettings.(*Config).Xmux != nil {
|
|
|
|
keepAlivePeriod = time.Duration(streamSettings.ProtocolSettings.(*Config).Xmux.HKeepAlivePeriod) * time.Second
|
|
|
|
}
|
2024-11-29 00:57:45 +00:00
|
|
|
|
2024-09-16 12:42:01 +00:00
|
|
|
var transport http.RoundTripper
|
2024-06-18 05:36:36 +00:00
|
|
|
|
2024-12-12 12:19:18 +00:00
|
|
|
if httpVersion == "3" {
|
2024-11-29 00:57:45 +00:00
|
|
|
if keepAlivePeriod == 0 {
|
|
|
|
keepAlivePeriod = quicgoH3KeepAlivePeriod
|
|
|
|
}
|
|
|
|
if keepAlivePeriod < 0 {
|
|
|
|
keepAlivePeriod = 0
|
|
|
|
}
|
2024-08-10 05:40:48 +00:00
|
|
|
quicConfig := &quic.Config{
|
|
|
|
MaxIdleTimeout: connIdleTimeout,
|
|
|
|
|
|
|
|
// these two are defaults of quic-go/http3. the default of quic-go (no
|
|
|
|
// http3) is different, so it is hardcoded here for clarity.
|
|
|
|
// https://github.com/quic-go/quic-go/blob/b8ea5c798155950fb5bbfdd06cad1939c9355878/http3/client.go#L36-L39
|
|
|
|
MaxIncomingStreams: -1,
|
2024-11-29 00:57:45 +00:00
|
|
|
KeepAlivePeriod: keepAlivePeriod,
|
2024-08-10 05:40:48 +00:00
|
|
|
}
|
2024-09-16 12:42:01 +00:00
|
|
|
transport = &http3.RoundTripper{
|
2024-08-10 05:40:48 +00:00
|
|
|
QUICConfig: quicConfig,
|
2024-07-17 12:10:48 +00:00
|
|
|
TLSClientConfig: gotlsConfig,
|
|
|
|
Dial: func(ctx context.Context, addr string, tlsCfg *gotls.Config, cfg *quic.Config) (quic.EarlyConnection, error) {
|
|
|
|
conn, err := internet.DialSystem(ctx, dest, streamSettings.SocketSettings)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2024-07-19 17:52:34 +00:00
|
|
|
|
2024-08-11 16:58:52 +00:00
|
|
|
var udpConn net.PacketConn
|
2024-07-19 17:52:34 +00:00
|
|
|
var udpAddr *net.UDPAddr
|
|
|
|
|
|
|
|
switch c := conn.(type) {
|
|
|
|
case *internet.PacketConnWrapper:
|
|
|
|
var ok bool
|
|
|
|
udpConn, ok = c.Conn.(*net.UDPConn)
|
|
|
|
if !ok {
|
|
|
|
return nil, errors.New("PacketConnWrapper does not contain a UDP connection")
|
|
|
|
}
|
|
|
|
udpAddr, err = net.ResolveUDPAddr("udp", c.Dest.String())
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
case *net.UDPConn:
|
|
|
|
udpConn = c
|
|
|
|
udpAddr, err = net.ResolveUDPAddr("udp", c.RemoteAddr().String())
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
default:
|
2024-08-11 16:58:52 +00:00
|
|
|
udpConn = &internet.FakePacketConn{c}
|
|
|
|
udpAddr, err = net.ResolveUDPAddr("udp", c.RemoteAddr().String())
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2024-07-17 12:10:48 +00:00
|
|
|
}
|
2024-07-19 17:52:34 +00:00
|
|
|
|
|
|
|
return quic.DialEarly(ctx, udpConn, udpAddr, tlsCfg, cfg)
|
2024-07-17 12:10:48 +00:00
|
|
|
},
|
|
|
|
}
|
2024-12-12 12:19:18 +00:00
|
|
|
} else if httpVersion == "2" {
|
2024-11-29 00:57:45 +00:00
|
|
|
if keepAlivePeriod == 0 {
|
|
|
|
keepAlivePeriod = chromeH2KeepAlivePeriod
|
|
|
|
}
|
|
|
|
if keepAlivePeriod < 0 {
|
|
|
|
keepAlivePeriod = 0
|
|
|
|
}
|
2024-09-16 12:42:01 +00:00
|
|
|
transport = &http2.Transport{
|
2024-06-18 05:36:36 +00:00
|
|
|
DialTLSContext: func(ctxInner context.Context, network string, addr string, cfg *gotls.Config) (net.Conn, error) {
|
|
|
|
return dialContext(ctxInner)
|
|
|
|
},
|
2024-08-10 05:40:48 +00:00
|
|
|
IdleConnTimeout: connIdleTimeout,
|
2024-11-29 00:57:45 +00:00
|
|
|
ReadIdleTimeout: keepAlivePeriod,
|
2024-06-18 05:36:36 +00:00
|
|
|
}
|
|
|
|
} else {
|
|
|
|
httpDialContext := func(ctxInner context.Context, network string, addr string) (net.Conn, error) {
|
|
|
|
return dialContext(ctxInner)
|
|
|
|
}
|
|
|
|
|
2024-09-16 12:42:01 +00:00
|
|
|
transport = &http.Transport{
|
2024-06-18 05:36:36 +00:00
|
|
|
DialTLSContext: httpDialContext,
|
|
|
|
DialContext: httpDialContext,
|
2024-08-10 05:40:48 +00:00
|
|
|
IdleConnTimeout: connIdleTimeout,
|
2024-11-29 00:57:45 +00:00
|
|
|
// chunked transfer download with KeepAlives is buggy with
|
2024-06-18 05:36:36 +00:00
|
|
|
// http.Client and our custom dial context.
|
|
|
|
DisableKeepAlives: true,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-07-11 07:56:20 +00:00
|
|
|
client := &DefaultDialerClient{
|
2024-09-16 12:42:01 +00:00
|
|
|
transportConfig: transportConfig,
|
|
|
|
client: &http.Client{
|
|
|
|
Transport: transport,
|
2024-06-18 05:36:36 +00:00
|
|
|
},
|
2024-12-12 12:19:18 +00:00
|
|
|
httpVersion: httpVersion,
|
2024-06-18 05:36:36 +00:00
|
|
|
uploadRawPool: &sync.Pool{},
|
|
|
|
dialUploadConn: dialContext,
|
|
|
|
}
|
|
|
|
|
|
|
|
return client
|
|
|
|
}
|
|
|
|
|
|
|
|
func init() {
|
|
|
|
common.Must(internet.RegisterTransportDialer(protocolName, Dial))
|
|
|
|
}
|
|
|
|
|
|
|
|
func Dial(ctx context.Context, dest net.Destination, streamSettings *internet.MemoryStreamConfig) (stat.Connection, error) {
|
|
|
|
tlsConfig := tls.ConfigFromStreamSettings(streamSettings)
|
2024-10-30 02:31:05 +00:00
|
|
|
realityConfig := reality.ConfigFromStreamSettings(streamSettings)
|
2024-06-18 05:36:36 +00:00
|
|
|
|
2024-12-12 12:19:18 +00:00
|
|
|
httpVersion := decideHTTPVersion(tlsConfig, realityConfig)
|
|
|
|
if httpVersion == "3" {
|
|
|
|
dest.Network = net.Network_UDP
|
|
|
|
}
|
|
|
|
|
|
|
|
transportConfiguration := streamSettings.ProtocolSettings.(*Config)
|
|
|
|
var requestURL url.URL
|
2024-06-18 05:36:36 +00:00
|
|
|
|
2024-10-30 02:31:05 +00:00
|
|
|
if tlsConfig != nil || realityConfig != nil {
|
2024-06-18 05:36:36 +00:00
|
|
|
requestURL.Scheme = "https"
|
|
|
|
} else {
|
|
|
|
requestURL.Scheme = "http"
|
|
|
|
}
|
|
|
|
requestURL.Host = transportConfiguration.Host
|
2024-12-11 00:58:14 +00:00
|
|
|
if requestURL.Host == "" && tlsConfig != nil {
|
|
|
|
requestURL.Host = tlsConfig.ServerName
|
|
|
|
}
|
|
|
|
if requestURL.Host == "" && realityConfig != nil {
|
|
|
|
requestURL.Host = realityConfig.ServerName
|
|
|
|
}
|
2024-06-18 05:36:36 +00:00
|
|
|
if requestURL.Host == "" {
|
2024-12-11 00:58:14 +00:00
|
|
|
requestURL.Host = dest.Address.String()
|
2024-06-18 05:36:36 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
sessionIdUuid := uuid.New()
|
2024-08-10 21:47:42 +00:00
|
|
|
requestURL.Path = transportConfiguration.GetNormalizedPath() + sessionIdUuid.String()
|
|
|
|
requestURL.RawQuery = transportConfiguration.GetNormalizedQuery()
|
2024-07-29 04:35:17 +00:00
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
httpClient, xmuxClient := getHTTPClient(ctx, dest, streamSettings)
|
2024-06-18 05:36:36 +00:00
|
|
|
|
2024-12-12 12:19:18 +00:00
|
|
|
mode := transportConfiguration.Mode
|
|
|
|
if mode == "" || mode == "auto" {
|
|
|
|
mode = "packet-up"
|
|
|
|
if httpVersion == "2" {
|
|
|
|
mode = "stream-up"
|
|
|
|
}
|
|
|
|
if realityConfig != nil && transportConfiguration.DownloadSettings == nil {
|
|
|
|
mode = "stream-one"
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
errors.LogInfo(ctx, fmt.Sprintf("XHTTP is dialing to %s, mode %s, HTTP version %s, host %s", dest, mode, httpVersion, requestURL.Host))
|
|
|
|
|
2024-11-09 11:05:41 +00:00
|
|
|
requestURL2 := requestURL
|
2024-12-12 12:19:18 +00:00
|
|
|
httpClient2 := httpClient
|
2024-12-15 05:43:10 +00:00
|
|
|
xmuxClient2 := xmuxClient
|
2024-10-31 07:31:19 +00:00
|
|
|
if transportConfiguration.DownloadSettings != nil {
|
2024-11-03 07:25:41 +00:00
|
|
|
globalDialerAccess.Lock()
|
|
|
|
if streamSettings.DownloadSettings == nil {
|
|
|
|
streamSettings.DownloadSettings = common.Must2(internet.ToMemoryStreamConfig(transportConfiguration.DownloadSettings)).(*internet.MemoryStreamConfig)
|
|
|
|
}
|
|
|
|
globalDialerAccess.Unlock()
|
|
|
|
memory2 := streamSettings.DownloadSettings
|
2024-12-11 00:58:14 +00:00
|
|
|
dest2 := *memory2.Destination // just panic
|
|
|
|
tlsConfig2 := tls.ConfigFromStreamSettings(memory2)
|
|
|
|
realityConfig2 := reality.ConfigFromStreamSettings(memory2)
|
2024-12-12 12:19:18 +00:00
|
|
|
httpVersion2 := decideHTTPVersion(tlsConfig2, realityConfig2)
|
|
|
|
if httpVersion2 == "3" {
|
|
|
|
dest2.Network = net.Network_UDP
|
|
|
|
}
|
2024-12-11 00:58:14 +00:00
|
|
|
if tlsConfig2 != nil || realityConfig2 != nil {
|
2024-10-31 07:31:19 +00:00
|
|
|
requestURL2.Scheme = "https"
|
|
|
|
} else {
|
|
|
|
requestURL2.Scheme = "http"
|
|
|
|
}
|
|
|
|
config2 := memory2.ProtocolSettings.(*Config)
|
|
|
|
requestURL2.Host = config2.Host
|
2024-12-11 00:58:14 +00:00
|
|
|
if requestURL2.Host == "" && tlsConfig2 != nil {
|
|
|
|
requestURL2.Host = tlsConfig2.ServerName
|
|
|
|
}
|
|
|
|
if requestURL2.Host == "" && realityConfig2 != nil {
|
|
|
|
requestURL2.Host = realityConfig2.ServerName
|
|
|
|
}
|
2024-10-31 07:31:19 +00:00
|
|
|
if requestURL2.Host == "" {
|
2024-12-11 00:58:14 +00:00
|
|
|
requestURL2.Host = dest2.Address.String()
|
2024-10-31 07:31:19 +00:00
|
|
|
}
|
2024-11-07 03:50:28 +00:00
|
|
|
requestURL2.Path = config2.GetNormalizedPath() + sessionIdUuid.String()
|
2024-10-31 07:31:19 +00:00
|
|
|
requestURL2.RawQuery = config2.GetNormalizedQuery()
|
2024-12-15 05:43:10 +00:00
|
|
|
httpClient2, xmuxClient2 = getHTTPClient(ctx, dest2, memory2)
|
2024-12-12 12:19:18 +00:00
|
|
|
errors.LogInfo(ctx, fmt.Sprintf("XHTTP is downloading from %s, mode %s, HTTP version %s, host %s", dest2, "stream-down", httpVersion2, requestURL2.Host))
|
2024-10-31 07:31:19 +00:00
|
|
|
}
|
|
|
|
|
2024-11-27 20:19:18 +00:00
|
|
|
var writer io.WriteCloser
|
|
|
|
var reader io.ReadCloser
|
|
|
|
var remoteAddr, localAddr net.Addr
|
|
|
|
var err error
|
|
|
|
|
|
|
|
if mode == "stream-one" {
|
|
|
|
requestURL.Path = transportConfiguration.GetNormalizedPath()
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient != nil {
|
|
|
|
xmuxClient.LeftRequests.Add(-1)
|
|
|
|
}
|
2024-11-27 20:19:18 +00:00
|
|
|
writer, reader = httpClient.Open(context.WithoutCancel(ctx), requestURL.String())
|
|
|
|
remoteAddr = &net.TCPAddr{}
|
|
|
|
localAddr = &net.TCPAddr{}
|
|
|
|
} else {
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient2 != nil {
|
|
|
|
xmuxClient2.LeftRequests.Add(-1)
|
|
|
|
}
|
2024-11-27 20:19:18 +00:00
|
|
|
reader, remoteAddr, localAddr, err = httpClient2.OpenDownload(context.WithoutCancel(ctx), requestURL2.String())
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2024-11-09 11:05:41 +00:00
|
|
|
}
|
2024-06-18 05:36:36 +00:00
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient != nil {
|
|
|
|
xmuxClient.OpenUsage.Add(1)
|
2024-11-03 07:25:41 +00:00
|
|
|
}
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient2 != nil && xmuxClient2 != xmuxClient {
|
|
|
|
xmuxClient2.OpenUsage.Add(1)
|
2024-09-16 12:42:01 +00:00
|
|
|
}
|
2024-12-15 05:43:10 +00:00
|
|
|
var once atomic.Int32
|
2024-09-16 12:42:01 +00:00
|
|
|
|
2024-11-09 11:05:41 +00:00
|
|
|
conn := splitConn{
|
2024-11-27 20:19:18 +00:00
|
|
|
writer: writer,
|
2024-11-09 11:05:41 +00:00
|
|
|
reader: reader,
|
|
|
|
remoteAddr: remoteAddr,
|
|
|
|
localAddr: localAddr,
|
|
|
|
onClose: func() {
|
2024-12-15 05:43:10 +00:00
|
|
|
if once.Add(-1) < 0 {
|
2024-11-09 11:05:41 +00:00
|
|
|
return
|
|
|
|
}
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient != nil {
|
|
|
|
xmuxClient.OpenUsage.Add(-1)
|
2024-11-09 11:05:41 +00:00
|
|
|
}
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient2 != nil && xmuxClient2 != xmuxClient {
|
|
|
|
xmuxClient2.OpenUsage.Add(-1)
|
2024-11-09 11:05:41 +00:00
|
|
|
}
|
|
|
|
},
|
|
|
|
}
|
|
|
|
|
2024-11-27 20:19:18 +00:00
|
|
|
if mode == "stream-one" {
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient != nil {
|
|
|
|
xmuxClient.LeftRequests.Add(-1)
|
|
|
|
}
|
2024-11-27 20:19:18 +00:00
|
|
|
return stat.Connection(&conn), nil
|
2024-11-09 11:05:41 +00:00
|
|
|
}
|
|
|
|
if mode == "stream-up" {
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient != nil {
|
|
|
|
xmuxClient.LeftRequests.Add(-1)
|
|
|
|
}
|
2024-11-09 11:05:41 +00:00
|
|
|
conn.writer = httpClient.OpenUpload(ctx, requestURL.String())
|
|
|
|
return stat.Connection(&conn), nil
|
|
|
|
}
|
|
|
|
|
2024-12-12 12:19:18 +00:00
|
|
|
scMaxEachPostBytes := transportConfiguration.GetNormalizedScMaxEachPostBytes()
|
|
|
|
scMinPostsIntervalMs := transportConfiguration.GetNormalizedScMinPostsIntervalMs()
|
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
maxUploadSize := scMaxEachPostBytes.rand()
|
2024-11-09 11:05:41 +00:00
|
|
|
// WithSizeLimit(0) will still allow single bytes to pass, and a lot of
|
|
|
|
// code relies on this behavior. Subtract 1 so that together with
|
|
|
|
// uploadWriter wrapper, exact size limits can be enforced
|
2024-11-17 06:03:25 +00:00
|
|
|
// uploadPipeReader, uploadPipeWriter := pipe.New(pipe.WithSizeLimit(maxUploadSize - 1))
|
|
|
|
uploadPipeReader, uploadPipeWriter := pipe.New(pipe.WithSizeLimit(maxUploadSize - buf.Size))
|
2024-09-16 12:42:01 +00:00
|
|
|
|
2024-11-09 11:05:41 +00:00
|
|
|
conn.writer = uploadWriter{
|
|
|
|
uploadPipeWriter,
|
|
|
|
maxUploadSize,
|
|
|
|
}
|
|
|
|
|
|
|
|
go func() {
|
2024-12-11 14:05:39 +00:00
|
|
|
var seq int64
|
|
|
|
var lastWrite time.Time
|
2024-06-18 05:36:36 +00:00
|
|
|
|
|
|
|
for {
|
2024-12-11 14:05:39 +00:00
|
|
|
wroteRequest := done.New()
|
|
|
|
|
|
|
|
ctx := httptrace.WithClientTrace(ctx, &httptrace.ClientTrace{
|
|
|
|
WroteRequest: func(httptrace.WroteRequestInfo) {
|
|
|
|
wroteRequest.Close()
|
|
|
|
},
|
|
|
|
})
|
|
|
|
|
|
|
|
// this intentionally makes a shallow-copy of the struct so we
|
|
|
|
// can reassign Path (potentially concurrently)
|
|
|
|
url := requestURL
|
|
|
|
url.Path += "/" + strconv.FormatInt(seq, 10)
|
|
|
|
// reassign query to get different padding
|
|
|
|
url.RawQuery = transportConfiguration.GetNormalizedQuery()
|
|
|
|
|
|
|
|
seq += 1
|
|
|
|
|
|
|
|
if scMinPostsIntervalMs.From > 0 {
|
2024-12-15 05:43:10 +00:00
|
|
|
time.Sleep(time.Duration(scMinPostsIntervalMs.rand())*time.Millisecond - time.Since(lastWrite))
|
2024-12-11 14:05:39 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// by offloading the uploads into a buffered pipe, multiple conn.Write
|
|
|
|
// calls get automatically batched together into larger POST requests.
|
|
|
|
// without batching, bandwidth is extremely limited.
|
2024-06-18 05:36:36 +00:00
|
|
|
chunk, err := uploadPipeReader.ReadMultiBuffer()
|
|
|
|
if err != nil {
|
|
|
|
break
|
|
|
|
}
|
|
|
|
|
2024-12-11 14:05:39 +00:00
|
|
|
lastWrite = time.Now()
|
2024-06-18 05:36:36 +00:00
|
|
|
|
2024-12-15 05:43:10 +00:00
|
|
|
if xmuxClient != nil && xmuxClient.LeftRequests.Add(-1) <= 0 {
|
|
|
|
httpClient, xmuxClient = getHTTPClient(ctx, dest, streamSettings)
|
|
|
|
}
|
|
|
|
|
2024-06-18 05:36:36 +00:00
|
|
|
go func() {
|
2024-07-11 07:56:20 +00:00
|
|
|
err := httpClient.SendUploadRequest(
|
|
|
|
context.WithoutCancel(ctx),
|
2024-08-10 21:47:42 +00:00
|
|
|
url.String(),
|
2024-07-11 07:56:20 +00:00
|
|
|
&buf.MultiBufferContainer{MultiBuffer: chunk},
|
|
|
|
int64(chunk.Len()),
|
|
|
|
)
|
2024-12-11 14:05:39 +00:00
|
|
|
wroteRequest.Close()
|
2024-06-18 05:36:36 +00:00
|
|
|
if err != nil {
|
2024-06-29 18:32:57 +00:00
|
|
|
errors.LogInfoInner(ctx, err, "failed to send upload")
|
2024-06-18 05:36:36 +00:00
|
|
|
uploadPipeReader.Interrupt()
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
2024-12-11 14:05:39 +00:00
|
|
|
if _, ok := httpClient.(*DefaultDialerClient); ok {
|
|
|
|
<-wroteRequest.Wait()
|
2024-07-27 12:52:36 +00:00
|
|
|
}
|
2024-06-18 05:36:36 +00:00
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
return stat.Connection(&conn), nil
|
|
|
|
}
|
2024-08-17 11:01:58 +00:00
|
|
|
|
|
|
|
// A wrapper around pipe that ensures the size limit is exactly honored.
|
|
|
|
//
|
|
|
|
// The MultiBuffer pipe accepts any single WriteMultiBuffer call even if that
|
|
|
|
// single MultiBuffer exceeds the size limit, and then starts blocking on the
|
|
|
|
// next WriteMultiBuffer call. This means that ReadMultiBuffer can return more
|
|
|
|
// bytes than the size limit. We work around this by splitting a potentially
|
|
|
|
// too large write up into multiple.
|
|
|
|
type uploadWriter struct {
|
|
|
|
*pipe.Writer
|
|
|
|
maxLen int32
|
|
|
|
}
|
|
|
|
|
|
|
|
func (w uploadWriter) Write(b []byte) (int, error) {
|
2024-11-17 06:03:25 +00:00
|
|
|
/*
|
|
|
|
capacity := int(w.maxLen - w.Len())
|
|
|
|
if capacity > 0 && capacity < len(b) {
|
|
|
|
b = b[:capacity]
|
|
|
|
}
|
|
|
|
*/
|
2024-08-17 11:01:58 +00:00
|
|
|
|
|
|
|
buffer := buf.New()
|
|
|
|
n, err := buffer.Write(b)
|
|
|
|
if err != nil {
|
|
|
|
return 0, err
|
|
|
|
}
|
|
|
|
|
|
|
|
err = w.WriteMultiBuffer([]*buf.Buffer{buffer})
|
|
|
|
if err != nil {
|
|
|
|
return 0, err
|
|
|
|
}
|
|
|
|
return n, nil
|
|
|
|
}
|