Skip to content

Commit 3e7dc35

Browse files
committed
fix: use sync.Map in request to packet conn server
1 parent ae4b773 commit 3e7dc35

1 file changed

Lines changed: 18 additions & 7 deletions

File tree

transport/internet/request/assembler/packetconn/req2packet.go

Lines changed: 18 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"context"
66
"crypto/rand"
77
"io"
8+
"sync"
89
"time"
910

1011
"github.com/golang-collections/go-datastructures/queue"
@@ -102,7 +103,7 @@ copyFromChan:
102103
waitTimer.Stop()
103104
go func() {
104105
reader, writer := io.Pipe()
105-
defer writer.Close()
106+
defer writer.Close()
106107
streamingRespOpt := &pipedStreamingRespOption{writer}
107108
go func() {
108109
for {
@@ -176,7 +177,7 @@ func (r *requestToPacketConnClientSession) Close() error {
176177

177178
func newRequestToPacketConnServer(ctx context.Context, config *ServerConfig) *requestToPacketConnServer {
178179
return &requestToPacketConnServer{
179-
sessionMap: make(map[string]*requestToPacketConnServerSession),
180+
sessionMap: sync.Map{},
180181
ctx: ctx,
181182
config: config,
182183
}
@@ -185,7 +186,7 @@ func newRequestToPacketConnServer(ctx context.Context, config *ServerConfig) *re
185186
type requestToPacketConnServer struct {
186187
packetSessionReceiver request.SessionReceiver
187188

188-
sessionMap map[string]*requestToPacketConnServerSession
189+
sessionMap sync.Map
189190

190191
ctx context.Context
191192
config *ServerConfig
@@ -203,7 +204,15 @@ func (r *requestToPacketConnServer) OnRoundTrip(ctx context.Context, req request
203204
return request.Response{}, newError("nil session id")
204205
}
205206
sessionID := string(SessionID)
206-
session, found := r.sessionMap[sessionID]
207+
var session *requestToPacketConnServerSession
208+
sessionAny, found := r.sessionMap.Load(sessionID)
209+
if found {
210+
var ok bool
211+
session, ok = sessionAny.(*requestToPacketConnServerSession)
212+
if !ok {
213+
return request.Response{}, newError("failed to cast session")
214+
}
215+
}
207216
if !found {
208217
ctxWithFinish, finish := context.WithCancel(ctx)
209218
session = &requestToPacketConnServerSession{
@@ -218,8 +227,10 @@ func (r *requestToPacketConnServer) OnRoundTrip(ctx context.Context, req request
218227
maxWriteDuration: int(r.config.MaxWriteDurationMs),
219228
maxSimultaneousWriteConnection: int(r.config.MaxSimultaneousWriteConnection),
220229
}
221-
r.sessionMap[sessionID] = session
222-
err = r.packetSessionReceiver.OnNewSession(ctx, session)
230+
_, loaded := r.sessionMap.LoadOrStore(sessionID, session)
231+
if !loaded {
232+
err = r.packetSessionReceiver.OnNewSession(ctx, session)
233+
}
223234
}
224235
if err != nil {
225236
return request.Response{}, err
@@ -228,7 +239,7 @@ func (r *requestToPacketConnServer) OnRoundTrip(ctx context.Context, req request
228239
}
229240

230241
func (r *requestToPacketConnServer) removeSessionID(sessionID []byte) {
231-
delete(r.sessionMap, string(sessionID))
242+
r.sessionMap.Delete(string(sessionID))
232243
}
233244

234245
type requestToPacketConnServerSession struct {

0 commit comments

Comments
 (0)