mirror of
https://github.com/shtorm-7/sing-box-extended.git
synced 2026-09-25 09:20:27 +00:00
xhttp: Update xhttp core
This commit is contained in:
@@ -12,6 +12,7 @@ import (
|
||||
"reflect"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"unsafe"
|
||||
|
||||
"github.com/sagernet/quic-go/http3"
|
||||
@@ -114,7 +115,7 @@ func (c *DefaultDialerClient) OpenStream(ctx context.Context, url string, sessio
|
||||
reqCtx, cancel := context.WithCancel(context.WithoutCancel(ctx))
|
||||
req, _ := http.NewRequestWithContext(reqCtx, method, url, body)
|
||||
FillStreamRequest(req, sessionId, "", c.options)
|
||||
wrc = &WaitReadCloser{Wait: make(chan struct{}), Cancel: cancel}
|
||||
wrc = &WaitReadCloser{wait: done.New(), Cancel: cancel}
|
||||
go func() {
|
||||
var resp *http.Response
|
||||
resp, err = c.client.Do(req)
|
||||
@@ -215,42 +216,40 @@ func (c *DefaultDialerClient) PostPacket(ctx context.Context, url string, sessio
|
||||
}
|
||||
|
||||
type WaitReadCloser struct {
|
||||
Wait chan struct{}
|
||||
wait *done.Instance
|
||||
Cancel context.CancelFunc
|
||||
io.ReadCloser
|
||||
reader atomic.Pointer[io.ReadCloser]
|
||||
}
|
||||
|
||||
func (w *WaitReadCloser) Set(rc io.ReadCloser) {
|
||||
w.ReadCloser = rc
|
||||
defer func() {
|
||||
if recover() != nil {
|
||||
rc.Close()
|
||||
w.reader.Store(&rc)
|
||||
if w.wait.Done() {
|
||||
if p := w.reader.Swap(nil); p != nil {
|
||||
(*p).Close()
|
||||
}
|
||||
}()
|
||||
close(w.Wait)
|
||||
}
|
||||
w.wait.Close()
|
||||
}
|
||||
|
||||
func (w *WaitReadCloser) Read(b []byte) (int, error) {
|
||||
<-w.Wait
|
||||
if w.ReadCloser == nil {
|
||||
return 0, io.ErrClosedPipe
|
||||
rc := w.reader.Load()
|
||||
if rc == nil {
|
||||
<-w.wait.Wait()
|
||||
if rc = w.reader.Load(); rc == nil {
|
||||
return 0, io.ErrClosedPipe
|
||||
}
|
||||
}
|
||||
return w.ReadCloser.Read(b)
|
||||
return (*rc).Read(b)
|
||||
}
|
||||
|
||||
func (w *WaitReadCloser) Close() error {
|
||||
if w.Cancel != nil {
|
||||
w.Cancel()
|
||||
}
|
||||
if w.ReadCloser != nil {
|
||||
return w.ReadCloser.Close()
|
||||
w.wait.Close()
|
||||
if p := w.reader.Swap(nil); p != nil {
|
||||
return (*p).Close()
|
||||
}
|
||||
defer func() {
|
||||
if recover() != nil && w.ReadCloser != nil {
|
||||
w.ReadCloser.Close()
|
||||
}
|
||||
}()
|
||||
close(w.Wait)
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package xhttp
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/base64"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -73,14 +74,17 @@ func FillStreamRequest(request *http.Request, sessionId string, seqStr string, o
|
||||
|
||||
func FillPacketRequest(request *http.Request, sessionId string, seqStr string, payload buf.MultiBuffer, options *option.V2RayXHTTPBaseOptions) error {
|
||||
dataPlacement := options.GetNormalizedUplinkDataPlacement()
|
||||
data := make([]byte, payload.Len())
|
||||
payload.Copy(data)
|
||||
buf.ReleaseMulti(payload)
|
||||
if dataPlacement == option.PlacementBody || dataPlacement == option.PlacementAuto {
|
||||
request.Header = options.GetRequestHeader()
|
||||
request.Body = io.NopCloser(&buf.MultiBufferContainer{MultiBuffer: payload})
|
||||
request.ContentLength = int64(payload.Len())
|
||||
request.Body = io.NopCloser(bytes.NewReader(data))
|
||||
request.ContentLength = int64(len(data))
|
||||
request.GetBody = func() (io.ReadCloser, error) {
|
||||
return io.NopCloser(bytes.NewReader(data)), nil
|
||||
}
|
||||
} else {
|
||||
data := make([]byte, payload.Len())
|
||||
payload.Copy(data)
|
||||
buf.ReleaseMulti(payload)
|
||||
switch dataPlacement {
|
||||
case option.PlacementHeader:
|
||||
request.Header = GetRequestHeaderWithPayload(data, options)
|
||||
|
||||
Reference in New Issue
Block a user