Skip to content
Toggle navigation
P
Projects
G
Groups
S
Snippets
Help
phsl
/
new-api
This project
Loading...
Sign in
Toggle navigation
Go to a project
Project
Repository
Issues
0
Merge Requests
0
Pipelines
Wiki
Snippets
Members
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Commit
30668551
authored
Dec 07, 2023
by
Calcium-Ion
Committed by
GitHub
Dec 07, 2023
Browse files
Options
Browse Files
Download
Plain Diff
Merge pull request #20 from Calcium-Ion/optimize/hign--cpu
fix: 修复客户端中断请求,计算补全阻塞问题
parents
4091ffe3
7d2194f9
Show whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
32 additions
and
12 deletions
+32
-12
controller/relay-openai.go
+32
-12
No files found.
controller/relay-openai.go
View file @
30668551
...
@@ -9,10 +9,12 @@ import (
...
@@ -9,10 +9,12 @@ import (
"net/http"
"net/http"
"one-api/common"
"one-api/common"
"strings"
"strings"
"sync"
"time"
)
)
func
openaiStreamHandler
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
relayMode
int
)
(
*
OpenAIErrorWithStatusCode
,
string
)
{
func
openaiStreamHandler
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
relayMode
int
)
(
*
OpenAIErrorWithStatusCode
,
string
)
{
responseText
:=
""
var
responseTextBuilder
strings
.
Builder
scanner
:=
bufio
.
NewScanner
(
resp
.
Body
)
scanner
:=
bufio
.
NewScanner
(
resp
.
Body
)
scanner
.
Split
(
func
(
data
[]
byte
,
atEOF
bool
)
(
advance
int
,
token
[]
byte
,
err
error
)
{
scanner
.
Split
(
func
(
data
[]
byte
,
atEOF
bool
)
(
advance
int
,
token
[]
byte
,
err
error
)
{
if
atEOF
&&
len
(
data
)
==
0
{
if
atEOF
&&
len
(
data
)
==
0
{
...
@@ -26,9 +28,16 @@ func openaiStreamHandler(c *gin.Context, resp *http.Response, relayMode int) (*O
...
@@ -26,9 +28,16 @@ func openaiStreamHandler(c *gin.Context, resp *http.Response, relayMode int) (*O
}
}
return
0
,
nil
,
nil
return
0
,
nil
,
nil
})
})
dataChan
:=
make
(
chan
string
)
dataChan
:=
make
(
chan
string
,
5
)
stopChan
:=
make
(
chan
bool
)
stopChan
:=
make
(
chan
bool
,
2
)
defer
close
(
stopChan
)
defer
close
(
dataChan
)
var
wg
sync
.
WaitGroup
go
func
()
{
go
func
()
{
wg
.
Add
(
1
)
defer
wg
.
Done
()
var
streamItems
[]
string
for
scanner
.
Scan
()
{
for
scanner
.
Scan
()
{
data
:=
scanner
.
Text
()
data
:=
scanner
.
Text
()
if
len
(
data
)
<
6
{
// ignore blank line or wrong format
if
len
(
data
)
<
6
{
// ignore blank line or wrong format
...
@@ -40,29 +49,39 @@ func openaiStreamHandler(c *gin.Context, resp *http.Response, relayMode int) (*O
...
@@ -40,29 +49,39 @@ func openaiStreamHandler(c *gin.Context, resp *http.Response, relayMode int) (*O
dataChan
<-
data
dataChan
<-
data
data
=
data
[
6
:
]
data
=
data
[
6
:
]
if
!
strings
.
HasPrefix
(
data
,
"[DONE]"
)
{
if
!
strings
.
HasPrefix
(
data
,
"[DONE]"
)
{
streamItems
=
append
(
streamItems
,
data
)
}
}
streamResp
:=
"["
+
strings
.
Join
(
streamItems
,
","
)
+
"]"
switch
relayMode
{
switch
relayMode
{
case
RelayModeChatCompletions
:
case
RelayModeChatCompletions
:
var
streamResponse
ChatCompletionsStreamResponseSimple
var
streamResponses
[]
ChatCompletionsStreamResponseSimple
err
:=
json
.
Unmarshal
(
common
.
StringToByteSlice
(
data
),
&
streamResponse
)
err
:=
json
.
Unmarshal
(
common
.
StringToByteSlice
(
streamResp
),
&
streamResponses
)
if
err
!=
nil
{
if
err
!=
nil
{
common
.
SysError
(
"error unmarshalling stream response: "
+
err
.
Error
())
common
.
SysError
(
"error unmarshalling stream response: "
+
err
.
Error
())
continue
// just ignore the error
return
// just ignore the error
}
}
for
_
,
streamResponse
:=
range
streamResponses
{
for
_
,
choice
:=
range
streamResponse
.
Choices
{
for
_
,
choice
:=
range
streamResponse
.
Choices
{
responseText
+=
choice
.
Delta
.
Content
responseTextBuilder
.
WriteString
(
choice
.
Delta
.
Content
)
}
}
}
case
RelayModeCompletions
:
case
RelayModeCompletions
:
var
streamResponse
CompletionsStreamResponse
var
streamResponses
[]
CompletionsStreamResponse
err
:=
json
.
Unmarshal
(
common
.
StringToByteSlice
(
data
),
&
streamResponse
)
err
:=
json
.
Unmarshal
(
common
.
StringToByteSlice
(
streamResp
),
&
streamResponses
)
if
err
!=
nil
{
if
err
!=
nil
{
common
.
SysError
(
"error unmarshalling stream response: "
+
err
.
Error
())
common
.
SysError
(
"error unmarshalling stream response: "
+
err
.
Error
())
continue
return
// just ignore the error
}
}
for
_
,
streamResponse
:=
range
streamResponses
{
for
_
,
choice
:=
range
streamResponse
.
Choices
{
for
_
,
choice
:=
range
streamResponse
.
Choices
{
responseText
+=
choice
.
Text
responseTextBuilder
.
WriteString
(
choice
.
Text
)
}
}
}
}
}
}
if
len
(
dataChan
)
>
0
{
// wait data out
time
.
Sleep
(
2
*
time
.
Second
)
}
}
stopChan
<-
true
stopChan
<-
true
}()
}()
...
@@ -85,7 +104,8 @@ func openaiStreamHandler(c *gin.Context, resp *http.Response, relayMode int) (*O
...
@@ -85,7 +104,8 @@ func openaiStreamHandler(c *gin.Context, resp *http.Response, relayMode int) (*O
if
err
!=
nil
{
if
err
!=
nil
{
return
errorWrapper
(
err
,
"close_response_body_failed"
,
http
.
StatusInternalServerError
),
""
return
errorWrapper
(
err
,
"close_response_body_failed"
,
http
.
StatusInternalServerError
),
""
}
}
return
nil
,
responseText
wg
.
Wait
()
return
nil
,
responseTextBuilder
.
String
()
}
}
func
openaiHandler
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
promptTokens
int
,
model
string
)
(
*
OpenAIErrorWithStatusCode
,
*
Usage
)
{
func
openaiHandler
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
promptTokens
int
,
model
string
)
(
*
OpenAIErrorWithStatusCode
,
*
Usage
)
{
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment