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
d06018cd
authored
Jul 19, 2024
by
CalciumIon
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
chore: gopool
parent
261ba5a8
Hide whitespace changes
Inline
Side-by-side
Showing
14 changed files
with
77 additions
and
115 deletions
+77
-115
controller/channel-test.go
+3
-2
go.mod
+1
-0
go.sum
+4
-0
main.go
+3
-2
model/log.go
+2
-1
model/utils.go
+3
-2
relay/channel/ali/adaptor.go
+1
-1
relay/channel/claude/relay-claude.go
+46
-83
relay/channel/ollama/adaptor.go
+1
-1
relay/channel/openai/adaptor.go
+1
-1
relay/channel/openai/relay-openai.go
+9
-8
relay/channel/perplexity/adaptor.go
+1
-1
relay/channel/zhipu/relay-zhipu.go
+1
-12
relay/channel/zhipu_4v/adaptor.go
+1
-1
No files found.
controller/channel-test.go
View file @
d06018cd
...
@@ -5,6 +5,7 @@ import (
...
@@ -5,6 +5,7 @@ import (
"encoding/json"
"encoding/json"
"errors"
"errors"
"fmt"
"fmt"
"github.com/bytedance/gopkg/util/gopool"
"io"
"io"
"math"
"math"
"net/http"
"net/http"
...
@@ -217,7 +218,7 @@ func testAllChannels(notify bool) error {
...
@@ -217,7 +218,7 @@ func testAllChannels(notify bool) error {
if
disableThreshold
==
0
{
if
disableThreshold
==
0
{
disableThreshold
=
10000000
// a impossible value
disableThreshold
=
10000000
// a impossible value
}
}
go
func
()
{
go
pool
.
Go
(
func
()
{
for
_
,
channel
:=
range
channels
{
for
_
,
channel
:=
range
channels
{
isChannelEnabled
:=
channel
.
Status
==
common
.
ChannelStatusEnabled
isChannelEnabled
:=
channel
.
Status
==
common
.
ChannelStatusEnabled
tik
:=
time
.
Now
()
tik
:=
time
.
Now
()
...
@@ -265,7 +266,7 @@ func testAllChannels(notify bool) error {
...
@@ -265,7 +266,7 @@ func testAllChannels(notify bool) error {
common
.
SysError
(
fmt
.
Sprintf
(
"failed to send email: %s"
,
err
.
Error
()))
common
.
SysError
(
fmt
.
Sprintf
(
"failed to send email: %s"
,
err
.
Error
()))
}
}
}
}
}
(
)
})
return
nil
return
nil
}
}
...
...
go.mod
View file @
d06018cd
...
@@ -38,6 +38,7 @@ require (
...
@@ -38,6 +38,7 @@ require (
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.5 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.5 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.5 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.5 // indirect
github.com/aws/smithy-go v1.20.2 // indirect
github.com/aws/smithy-go v1.20.2 // indirect
github.com/bytedance/gopkg v0.0.0-20220118071334-3db87571198b // indirect
github.com/bytedance/sonic v1.9.1 // indirect
github.com/bytedance/sonic v1.9.1 // indirect
github.com/cespare/xxhash/v2 v2.1.2 // indirect
github.com/cespare/xxhash/v2 v2.1.2 // indirect
github.com/chenzhuoyu/base64x v0.0.0-20221115062448-fe3a3abad311 // indirect
github.com/chenzhuoyu/base64x v0.0.0-20221115062448-fe3a3abad311 // indirect
...
...
go.sum
View file @
d06018cd
...
@@ -18,6 +18,8 @@ github.com/aws/aws-sdk-go-v2/service/bedrockruntime v1.7.4 h1:JgHnonzbnA3pbqj76w
...
@@ -18,6 +18,8 @@ github.com/aws/aws-sdk-go-v2/service/bedrockruntime v1.7.4 h1:JgHnonzbnA3pbqj76w
github.com/aws/aws-sdk-go-v2/service/bedrockruntime v1.7.4/go.mod h1:nZspkhg+9p8iApLFoyAqfyuMP0F38acy2Hm3r5r95Cg=
github.com/aws/aws-sdk-go-v2/service/bedrockruntime v1.7.4/go.mod h1:nZspkhg+9p8iApLFoyAqfyuMP0F38acy2Hm3r5r95Cg=
github.com/aws/smithy-go v1.20.2 h1:tbp628ireGtzcHDDmLT/6ADHidqnwgF57XOXZe6tp4Q=
github.com/aws/smithy-go v1.20.2 h1:tbp628ireGtzcHDDmLT/6ADHidqnwgF57XOXZe6tp4Q=
github.com/aws/smithy-go v1.20.2/go.mod h1:krry+ya/rV9RDcV/Q16kpu6ypI4K2czasz0NC3qS14E=
github.com/aws/smithy-go v1.20.2/go.mod h1:krry+ya/rV9RDcV/Q16kpu6ypI4K2czasz0NC3qS14E=
github.com/bytedance/gopkg v0.0.0-20220118071334-3db87571198b h1:LTGVFpNmNHhj0vhOlfgWueFJ32eK9blaIlHR2ciXOT0=
github.com/bytedance/gopkg v0.0.0-20220118071334-3db87571198b/go.mod h1:2ZlV9BaUH4+NXIBF0aMdKKAnHTzqH+iMU4KUjAbL23Q=
github.com/bytedance/sonic v1.5.0/go.mod h1:ED5hyg4y6t3/9Ku1R6dU/4KyJ48DZ4jPhfY1O2AihPM=
github.com/bytedance/sonic v1.5.0/go.mod h1:ED5hyg4y6t3/9Ku1R6dU/4KyJ48DZ4jPhfY1O2AihPM=
github.com/bytedance/sonic v1.9.1 h1:6iJ6NqdoxCDr6mbY8h18oSO+cShGSMRGCEo7F2h0x8s=
github.com/bytedance/sonic v1.9.1 h1:6iJ6NqdoxCDr6mbY8h18oSO+cShGSMRGCEo7F2h0x8s=
github.com/bytedance/sonic v1.9.1/go.mod h1:i736AoUSYt75HyZLoJW9ERYxcy6eaN6h4BZXU064P/U=
github.com/bytedance/sonic v1.9.1/go.mod h1:i736AoUSYt75HyZLoJW9ERYxcy6eaN6h4BZXU064P/U=
...
@@ -198,6 +200,7 @@ golang.org/x/image v0.15.0/go.mod h1:HUYqC05R2ZcZ3ejNQsIHQDQiwWM4JBqmm6MKANTp4LE
...
@@ -198,6 +200,7 @@ golang.org/x/image v0.15.0/go.mod h1:HUYqC05R2ZcZ3ejNQsIHQDQiwWM4JBqmm6MKANTp4LE
golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg=
golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg=
golang.org/x/net v0.21.0 h1:AQyQV4dYCvJ7vGmJyKki9+PBdyvhkSd8EIx/qb0AYv4=
golang.org/x/net v0.21.0 h1:AQyQV4dYCvJ7vGmJyKki9+PBdyvhkSd8EIx/qb0AYv4=
golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44=
golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44=
golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.7.0 h1:YsImfSBoP9QPYL0xyKJPq0gcaJdG3rInoqxTWbfQu9M=
golang.org/x/sync v0.7.0 h1:YsImfSBoP9QPYL0xyKJPq0gcaJdG3rInoqxTWbfQu9M=
golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
...
@@ -206,6 +209,7 @@ golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7w
...
@@ -206,6 +209,7 @@ golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7w
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20210630005230-0f9fa26af87c/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20210630005230-0f9fa26af87c/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20210806184541-e5e7981a1069/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20210806184541-e5e7981a1069/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220110181412-a018aaa089fe/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220704084225-05e143d24a9e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220704084225-05e143d24a9e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
...
...
main.go
View file @
d06018cd
...
@@ -3,6 +3,7 @@ package main
...
@@ -3,6 +3,7 @@ package main
import
(
import
(
"embed"
"embed"
"fmt"
"fmt"
"github.com/bytedance/gopkg/util/gopool"
"github.com/gin-contrib/sessions"
"github.com/gin-contrib/sessions"
"github.com/gin-contrib/sessions/cookie"
"github.com/gin-contrib/sessions/cookie"
"github.com/gin-gonic/gin"
"github.com/gin-gonic/gin"
...
@@ -91,10 +92,10 @@ func main() {
...
@@ -91,10 +92,10 @@ func main() {
go
controller
.
AutomaticallyTestChannels
(
frequency
)
go
controller
.
AutomaticallyTestChannels
(
frequency
)
}
}
if
common
.
IsMasterNode
&&
constant
.
UpdateTask
{
if
common
.
IsMasterNode
&&
constant
.
UpdateTask
{
common
.
SafeGoroutine
(
func
()
{
gopool
.
Go
(
func
()
{
controller
.
UpdateMidjourneyTaskBulk
()
controller
.
UpdateMidjourneyTaskBulk
()
})
})
common
.
SafeGoroutine
(
func
()
{
gopool
.
Go
(
func
()
{
controller
.
UpdateTaskBulk
()
controller
.
UpdateTaskBulk
()
})
})
}
}
...
...
model/log.go
View file @
d06018cd
...
@@ -3,6 +3,7 @@ package model
...
@@ -3,6 +3,7 @@ package model
import
(
import
(
"context"
"context"
"fmt"
"fmt"
"github.com/bytedance/gopkg/util/gopool"
"gorm.io/gorm"
"gorm.io/gorm"
"one-api/common"
"one-api/common"
"strings"
"strings"
...
@@ -87,7 +88,7 @@ func RecordConsumeLog(ctx context.Context, userId int, channelId int, promptToke
...
@@ -87,7 +88,7 @@ func RecordConsumeLog(ctx context.Context, userId int, channelId int, promptToke
common
.
LogError
(
ctx
,
"failed to record log: "
+
err
.
Error
())
common
.
LogError
(
ctx
,
"failed to record log: "
+
err
.
Error
())
}
}
if
common
.
DataExportEnabled
{
if
common
.
DataExportEnabled
{
common
.
SafeGoroutine
(
func
()
{
gopool
.
Go
(
func
()
{
LogQuotaData
(
userId
,
username
,
modelName
,
quota
,
common
.
GetTimestamp
(),
promptTokens
+
completionTokens
)
LogQuotaData
(
userId
,
username
,
modelName
,
quota
,
common
.
GetTimestamp
(),
promptTokens
+
completionTokens
)
})
})
}
}
...
...
model/utils.go
View file @
d06018cd
...
@@ -2,6 +2,7 @@ package model
...
@@ -2,6 +2,7 @@ package model
import
(
import
(
"errors"
"errors"
"github.com/bytedance/gopkg/util/gopool"
"gorm.io/gorm"
"gorm.io/gorm"
"one-api/common"
"one-api/common"
"sync"
"sync"
...
@@ -28,12 +29,12 @@ func init() {
...
@@ -28,12 +29,12 @@ func init() {
}
}
func
InitBatchUpdater
()
{
func
InitBatchUpdater
()
{
go
func
()
{
go
pool
.
Go
(
func
()
{
for
{
for
{
time
.
Sleep
(
time
.
Duration
(
common
.
BatchUpdateInterval
)
*
time
.
Second
)
time
.
Sleep
(
time
.
Duration
(
common
.
BatchUpdateInterval
)
*
time
.
Second
)
batchUpdate
()
batchUpdate
()
}
}
}
(
)
})
}
}
func
addNewRecord
(
type_
int
,
id
int
,
value
int
)
{
func
addNewRecord
(
type_
int
,
id
int
,
value
int
)
{
...
...
relay/channel/ali/adaptor.go
View file @
d06018cd
...
@@ -84,7 +84,7 @@ func (a *Adaptor) DoResponse(c *gin.Context, resp *http.Response, info *relaycom
...
@@ -84,7 +84,7 @@ func (a *Adaptor) DoResponse(c *gin.Context, resp *http.Response, info *relaycom
err
,
usage
=
aliEmbeddingHandler
(
c
,
resp
)
err
,
usage
=
aliEmbeddingHandler
(
c
,
resp
)
default
:
default
:
if
info
.
IsStream
{
if
info
.
IsStream
{
err
,
usage
=
openai
.
O
pen
aiStreamHandler
(
c
,
resp
,
info
)
err
,
usage
=
openai
.
OaiStreamHandler
(
c
,
resp
,
info
)
}
else
{
}
else
{
err
,
usage
=
openai
.
OpenaiHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
)
err
,
usage
=
openai
.
OpenaiHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
)
}
}
...
...
relay/channel/claude/relay-claude.go
View file @
d06018cd
...
@@ -8,12 +8,10 @@ import (
...
@@ -8,12 +8,10 @@ import (
"io"
"io"
"net/http"
"net/http"
"one-api/common"
"one-api/common"
"one-api/constant"
"one-api/dto"
"one-api/dto"
relaycommon
"one-api/relay/common"
relaycommon
"one-api/relay/common"
"one-api/service"
"one-api/service"
"strings"
"strings"
"time"
)
)
func
stopReasonClaude2OpenAI
(
reason
string
)
string
{
func
stopReasonClaude2OpenAI
(
reason
string
)
string
{
...
@@ -332,91 +330,59 @@ func claudeStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
...
@@ -332,91 +330,59 @@ func claudeStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
responseText
:=
""
responseText
:=
""
createdTime
:=
common
.
GetTimestamp
()
createdTime
:=
common
.
GetTimestamp
()
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
(
bufio
.
ScanLines
)
if
atEOF
&&
len
(
data
)
==
0
{
service
.
SetEventStreamHeaders
(
c
)
return
0
,
nil
,
nil
}
for
scanner
.
Scan
()
{
if
i
:=
strings
.
Index
(
string
(
data
),
"
\n
"
);
i
>=
0
{
data
:=
scanner
.
Text
()
return
i
+
1
,
data
[
0
:
i
],
nil
info
.
SetFirstResponseTime
()
if
len
(
data
)
<
6
||
!
strings
.
HasPrefix
(
data
,
"data:"
)
{
continue
}
}
if
atEOF
{
data
=
strings
.
TrimPrefix
(
data
,
"data:"
)
return
len
(
data
),
data
,
nil
data
=
strings
.
TrimSpace
(
data
)
var
claudeResponse
ClaudeResponse
err
:=
json
.
Unmarshal
([]
byte
(
data
),
&
claudeResponse
)
if
err
!=
nil
{
common
.
SysError
(
"error unmarshalling stream response: "
+
err
.
Error
())
continue
}
}
return
0
,
nil
,
nil
})
response
,
claudeUsage
:=
StreamResponseClaude2OpenAI
(
requestMode
,
&
claudeResponse
)
dataChan
:=
make
(
chan
string
,
5
)
if
response
==
nil
{
stopChan
:=
make
(
chan
bool
,
2
)
continue
go
func
()
{
for
scanner
.
Scan
()
{
data
:=
scanner
.
Text
()
if
!
strings
.
HasPrefix
(
data
,
"data: "
)
{
continue
}
data
=
strings
.
TrimPrefix
(
data
,
"data: "
)
if
!
common
.
SafeSendStringTimeout
(
dataChan
,
data
,
constant
.
StreamingTimeout
)
{
// send data timeout, stop the stream
common
.
LogError
(
c
,
"send data timeout, stop the stream"
)
break
}
}
}
stopChan
<-
true
if
requestMode
==
RequestModeCompletion
{
}()
responseText
+=
claudeResponse
.
Completion
isFirst
:=
true
responseId
=
response
.
Id
service
.
SetEventStreamHeaders
(
c
)
}
else
{
c
.
Stream
(
func
(
w
io
.
Writer
)
bool
{
if
claudeResponse
.
Type
==
"message_start"
{
select
{
// message_start, 获取usage
case
data
:=
<-
dataChan
:
responseId
=
claudeResponse
.
Message
.
Id
if
isFirst
{
info
.
UpstreamModelName
=
claudeResponse
.
Message
.
Model
isFirst
=
false
usage
.
PromptTokens
=
claudeUsage
.
InputTokens
info
.
FirstResponseTime
=
time
.
Now
()
}
else
if
claudeResponse
.
Type
==
"content_block_delta"
{
}
responseText
+=
claudeResponse
.
Delta
.
Text
// some implementations may add \r at the end of data
}
else
if
claudeResponse
.
Type
==
"message_delta"
{
data
=
strings
.
TrimSuffix
(
data
,
"
\r
"
)
usage
.
CompletionTokens
=
claudeUsage
.
OutputTokens
var
claudeResponse
ClaudeResponse
usage
.
TotalTokens
=
claudeUsage
.
InputTokens
+
claudeUsage
.
OutputTokens
err
:=
json
.
Unmarshal
([]
byte
(
data
),
&
claudeResponse
)
}
else
if
claudeResponse
.
Type
==
"content_block_start"
{
if
err
!=
nil
{
common
.
SysError
(
"error unmarshalling stream response: "
+
err
.
Error
())
return
true
}
response
,
claudeUsage
:=
StreamResponseClaude2OpenAI
(
requestMode
,
&
claudeResponse
)
if
response
==
nil
{
return
true
}
if
requestMode
==
RequestModeCompletion
{
responseText
+=
claudeResponse
.
Completion
responseId
=
response
.
Id
}
else
{
}
else
{
if
claudeResponse
.
Type
==
"message_start"
{
continue
// message_start, 获取usage
responseId
=
claudeResponse
.
Message
.
Id
info
.
UpstreamModelName
=
claudeResponse
.
Message
.
Model
usage
.
PromptTokens
=
claudeUsage
.
InputTokens
}
else
if
claudeResponse
.
Type
==
"content_block_delta"
{
responseText
+=
claudeResponse
.
Delta
.
Text
}
else
if
claudeResponse
.
Type
==
"message_delta"
{
usage
.
CompletionTokens
=
claudeUsage
.
OutputTokens
usage
.
TotalTokens
=
claudeUsage
.
InputTokens
+
claudeUsage
.
OutputTokens
}
else
if
claudeResponse
.
Type
==
"content_block_start"
{
}
else
{
return
true
}
}
}
//response.Id = responseId
}
response
.
Id
=
responseId
//response.Id = responseId
response
.
Created
=
createdTime
response
.
Id
=
responseId
response
.
Model
=
info
.
UpstreamModelName
response
.
Created
=
createdTime
response
.
Model
=
info
.
UpstreamModelName
err
=
service
.
ObjectData
(
c
,
response
)
err
=
service
.
ObjectData
(
c
,
response
)
if
err
!=
nil
{
if
err
!=
nil
{
common
.
SysError
(
err
.
Error
())
common
.
LogError
(
c
,
"send_stream_response_failed: "
+
err
.
Error
())
}
return
true
case
<-
stopChan
:
return
false
}
}
})
}
if
requestMode
==
RequestModeCompletion
{
if
requestMode
==
RequestModeCompletion
{
usage
,
_
=
service
.
ResponseText2Usage
(
responseText
,
info
.
UpstreamModelName
,
info
.
PromptTokens
)
usage
,
_
=
service
.
ResponseText2Usage
(
responseText
,
info
.
UpstreamModelName
,
info
.
PromptTokens
)
}
else
{
}
else
{
...
@@ -435,10 +401,7 @@ func claudeStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
...
@@ -435,10 +401,7 @@ func claudeStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
}
}
}
}
service
.
Done
(
c
)
service
.
Done
(
c
)
err
:=
resp
.
Body
.
Close
()
resp
.
Body
.
Close
()
if
err
!=
nil
{
return
service
.
OpenAIErrorWrapperLocal
(
err
,
"close_response_body_failed"
,
http
.
StatusInternalServerError
),
nil
}
return
nil
,
usage
return
nil
,
usage
}
}
...
...
relay/channel/ollama/adaptor.go
View file @
d06018cd
...
@@ -64,7 +64,7 @@ func (a *Adaptor) DoRequest(c *gin.Context, info *relaycommon.RelayInfo, request
...
@@ -64,7 +64,7 @@ func (a *Adaptor) DoRequest(c *gin.Context, info *relaycommon.RelayInfo, request
func
(
a
*
Adaptor
)
DoResponse
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
info
*
relaycommon
.
RelayInfo
)
(
usage
*
dto
.
Usage
,
err
*
dto
.
OpenAIErrorWithStatusCode
)
{
func
(
a
*
Adaptor
)
DoResponse
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
info
*
relaycommon
.
RelayInfo
)
(
usage
*
dto
.
Usage
,
err
*
dto
.
OpenAIErrorWithStatusCode
)
{
if
info
.
IsStream
{
if
info
.
IsStream
{
err
,
usage
=
openai
.
O
pen
aiStreamHandler
(
c
,
resp
,
info
)
err
,
usage
=
openai
.
OaiStreamHandler
(
c
,
resp
,
info
)
}
else
{
}
else
{
if
info
.
RelayMode
==
relayconstant
.
RelayModeEmbeddings
{
if
info
.
RelayMode
==
relayconstant
.
RelayModeEmbeddings
{
err
,
usage
=
ollamaEmbeddingHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
,
info
.
RelayMode
)
err
,
usage
=
ollamaEmbeddingHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
,
info
.
RelayMode
)
...
...
relay/channel/openai/adaptor.go
View file @
d06018cd
...
@@ -145,7 +145,7 @@ func (a *Adaptor) DoResponse(c *gin.Context, resp *http.Response, info *relaycom
...
@@ -145,7 +145,7 @@ func (a *Adaptor) DoResponse(c *gin.Context, resp *http.Response, info *relaycom
err
,
usage
=
OpenaiTTSHandler
(
c
,
resp
,
info
)
err
,
usage
=
OpenaiTTSHandler
(
c
,
resp
,
info
)
default
:
default
:
if
info
.
IsStream
{
if
info
.
IsStream
{
err
,
usage
=
O
pen
aiStreamHandler
(
c
,
resp
,
info
)
err
,
usage
=
OaiStreamHandler
(
c
,
resp
,
info
)
}
else
{
}
else
{
err
,
usage
=
OpenaiHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
)
err
,
usage
=
OpenaiHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
)
}
}
...
...
relay/channel/openai/relay-openai.go
View file @
d06018cd
...
@@ -5,6 +5,7 @@ import (
...
@@ -5,6 +5,7 @@ import (
"bytes"
"bytes"
"encoding/json"
"encoding/json"
"fmt"
"fmt"
"github.com/bytedance/gopkg/util/gopool"
"github.com/gin-gonic/gin"
"github.com/gin-gonic/gin"
"io"
"io"
"net/http"
"net/http"
...
@@ -18,8 +19,8 @@ import (
...
@@ -18,8 +19,8 @@ import (
"time"
"time"
)
)
func
O
pen
aiStreamHandler
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
info
*
relaycommon
.
RelayInfo
)
(
*
dto
.
OpenAIErrorWithStatusCode
,
*
dto
.
Usage
)
{
func
OaiStreamHandler
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
info
*
relaycommon
.
RelayInfo
)
(
*
dto
.
OpenAIErrorWithStatusCode
,
*
dto
.
Usage
)
{
has
StreamUsage
:=
false
contain
StreamUsage
:=
false
responseId
:=
""
responseId
:=
""
var
createAt
int64
=
0
var
createAt
int64
=
0
var
systemFingerprint
string
var
systemFingerprint
string
...
@@ -41,7 +42,7 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
...
@@ -41,7 +42,7 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
stopChan
:=
make
(
chan
bool
)
stopChan
:=
make
(
chan
bool
)
defer
close
(
stopChan
)
defer
close
(
stopChan
)
go
func
()
{
go
pool
.
Go
(
func
()
{
for
scanner
.
Scan
()
{
for
scanner
.
Scan
()
{
info
.
SetFirstResponseTime
()
info
.
SetFirstResponseTime
()
ticker
.
Reset
(
time
.
Duration
(
constant
.
StreamingTimeout
)
*
time
.
Second
)
ticker
.
Reset
(
time
.
Duration
(
constant
.
StreamingTimeout
)
*
time
.
Second
)
...
@@ -62,7 +63,7 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
...
@@ -62,7 +63,7 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
}
}
}
}
common
.
SafeSendBool
(
stopChan
,
true
)
common
.
SafeSendBool
(
stopChan
,
true
)
}
(
)
})
select
{
select
{
case
<-
ticker
.
C
:
case
<-
ticker
.
C
:
...
@@ -91,7 +92,7 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
...
@@ -91,7 +92,7 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
model
=
streamResponse
.
Model
model
=
streamResponse
.
Model
if
service
.
ValidUsage
(
streamResponse
.
Usage
)
{
if
service
.
ValidUsage
(
streamResponse
.
Usage
)
{
usage
=
streamResponse
.
Usage
usage
=
streamResponse
.
Usage
has
StreamUsage
=
true
contain
StreamUsage
=
true
}
}
for
_
,
choice
:=
range
streamResponse
.
Choices
{
for
_
,
choice
:=
range
streamResponse
.
Choices
{
responseTextBuilder
.
WriteString
(
choice
.
Delta
.
GetContentString
())
responseTextBuilder
.
WriteString
(
choice
.
Delta
.
GetContentString
())
...
@@ -115,7 +116,7 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
...
@@ -115,7 +116,7 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
model
=
streamResponse
.
Model
model
=
streamResponse
.
Model
if
service
.
ValidUsage
(
streamResponse
.
Usage
)
{
if
service
.
ValidUsage
(
streamResponse
.
Usage
)
{
usage
=
streamResponse
.
Usage
usage
=
streamResponse
.
Usage
has
StreamUsage
=
true
contain
StreamUsage
=
true
}
}
for
_
,
choice
:=
range
streamResponse
.
Choices
{
for
_
,
choice
:=
range
streamResponse
.
Choices
{
responseTextBuilder
.
WriteString
(
choice
.
Delta
.
GetContentString
())
responseTextBuilder
.
WriteString
(
choice
.
Delta
.
GetContentString
())
...
@@ -155,12 +156,12 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
...
@@ -155,12 +156,12 @@ func OpenaiStreamHandler(c *gin.Context, resp *http.Response, info *relaycommon.
}
}
}
}
if
!
has
StreamUsage
{
if
!
contain
StreamUsage
{
usage
,
_
=
service
.
ResponseText2Usage
(
responseTextBuilder
.
String
(),
info
.
UpstreamModelName
,
info
.
PromptTokens
)
usage
,
_
=
service
.
ResponseText2Usage
(
responseTextBuilder
.
String
(),
info
.
UpstreamModelName
,
info
.
PromptTokens
)
usage
.
CompletionTokens
+=
toolCount
*
7
usage
.
CompletionTokens
+=
toolCount
*
7
}
}
if
info
.
ShouldIncludeUsage
&&
!
has
StreamUsage
{
if
info
.
ShouldIncludeUsage
&&
!
contain
StreamUsage
{
response
:=
service
.
GenerateFinalUsageResponse
(
responseId
,
createAt
,
model
,
*
usage
)
response
:=
service
.
GenerateFinalUsageResponse
(
responseId
,
createAt
,
model
,
*
usage
)
response
.
SetSystemFingerprint
(
systemFingerprint
)
response
.
SetSystemFingerprint
(
systemFingerprint
)
service
.
ObjectData
(
c
,
response
)
service
.
ObjectData
(
c
,
response
)
...
...
relay/channel/perplexity/adaptor.go
View file @
d06018cd
...
@@ -58,7 +58,7 @@ func (a *Adaptor) DoRequest(c *gin.Context, info *relaycommon.RelayInfo, request
...
@@ -58,7 +58,7 @@ func (a *Adaptor) DoRequest(c *gin.Context, info *relaycommon.RelayInfo, request
func
(
a
*
Adaptor
)
DoResponse
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
info
*
relaycommon
.
RelayInfo
)
(
usage
*
dto
.
Usage
,
err
*
dto
.
OpenAIErrorWithStatusCode
)
{
func
(
a
*
Adaptor
)
DoResponse
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
info
*
relaycommon
.
RelayInfo
)
(
usage
*
dto
.
Usage
,
err
*
dto
.
OpenAIErrorWithStatusCode
)
{
if
info
.
IsStream
{
if
info
.
IsStream
{
err
,
usage
=
openai
.
O
pen
aiStreamHandler
(
c
,
resp
,
info
)
err
,
usage
=
openai
.
OaiStreamHandler
(
c
,
resp
,
info
)
}
else
{
}
else
{
err
,
usage
=
openai
.
OpenaiHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
)
err
,
usage
=
openai
.
OpenaiHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
)
}
}
...
...
relay/channel/zhipu/relay-zhipu.go
View file @
d06018cd
...
@@ -153,18 +153,7 @@ func streamMetaResponseZhipu2OpenAI(zhipuResponse *ZhipuStreamMetaResponse) (*dt
...
@@ -153,18 +153,7 @@ func streamMetaResponseZhipu2OpenAI(zhipuResponse *ZhipuStreamMetaResponse) (*dt
func
zhipuStreamHandler
(
c
*
gin
.
Context
,
resp
*
http
.
Response
)
(
*
dto
.
OpenAIErrorWithStatusCode
,
*
dto
.
Usage
)
{
func
zhipuStreamHandler
(
c
*
gin
.
Context
,
resp
*
http
.
Response
)
(
*
dto
.
OpenAIErrorWithStatusCode
,
*
dto
.
Usage
)
{
var
usage
*
dto
.
Usage
var
usage
*
dto
.
Usage
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
(
bufio
.
ScanLines
)
if
atEOF
&&
len
(
data
)
==
0
{
return
0
,
nil
,
nil
}
if
i
:=
strings
.
Index
(
string
(
data
),
"
\n\n
"
);
i
>=
0
&&
strings
.
Index
(
string
(
data
),
":"
)
>=
0
{
return
i
+
2
,
data
[
0
:
i
],
nil
}
if
atEOF
{
return
len
(
data
),
data
,
nil
}
return
0
,
nil
,
nil
})
dataChan
:=
make
(
chan
string
)
dataChan
:=
make
(
chan
string
)
metaChan
:=
make
(
chan
string
)
metaChan
:=
make
(
chan
string
)
stopChan
:=
make
(
chan
bool
)
stopChan
:=
make
(
chan
bool
)
...
...
relay/channel/zhipu_4v/adaptor.go
View file @
d06018cd
...
@@ -59,7 +59,7 @@ func (a *Adaptor) DoRequest(c *gin.Context, info *relaycommon.RelayInfo, request
...
@@ -59,7 +59,7 @@ func (a *Adaptor) DoRequest(c *gin.Context, info *relaycommon.RelayInfo, request
func
(
a
*
Adaptor
)
DoResponse
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
info
*
relaycommon
.
RelayInfo
)
(
usage
*
dto
.
Usage
,
err
*
dto
.
OpenAIErrorWithStatusCode
)
{
func
(
a
*
Adaptor
)
DoResponse
(
c
*
gin
.
Context
,
resp
*
http
.
Response
,
info
*
relaycommon
.
RelayInfo
)
(
usage
*
dto
.
Usage
,
err
*
dto
.
OpenAIErrorWithStatusCode
)
{
if
info
.
IsStream
{
if
info
.
IsStream
{
err
,
usage
=
openai
.
O
pen
aiStreamHandler
(
c
,
resp
,
info
)
err
,
usage
=
openai
.
OaiStreamHandler
(
c
,
resp
,
info
)
}
else
{
}
else
{
err
,
usage
=
openai
.
OpenaiHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
)
err
,
usage
=
openai
.
OpenaiHandler
(
c
,
resp
,
info
.
PromptTokens
,
info
.
UpstreamModelName
)
}
}
...
...
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