aurora / internal /chatgpt /handler_response.go
qwen2api's picture
Upload 113 files
6bc074c verified
Raw
History Blame Contribute Delete
21.1 kB
package chatgpt
import (
"bufio"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"os"
"strings"
"sync"
"github.com/gin-gonic/gin"
"aurora/conversion/response/chatgpt"
"aurora/httpclient"
"aurora/internal/accounts"
"aurora/internal/sseparser"
"aurora/typings"
chatgpt_types "aurora/typings/chatgpt"
official_types "aurora/typings/official"
"github.com/bogdanfinn/websocket"
)
type conversationPatchState struct {
response chatgpt_types.ChatGPTResponse
channel string
}
type conversationStreamEvent struct {
response chatgpt_types.ChatGPTResponse
chunk *official_types.ChatCompletionChunk
text string
role string
conversationID string
messageID string
channel string
finishReason string
isStop bool
}
func parseConversationEvent(line string, state *sseparser.PatchState, model string) (conversationStreamEvent, bool) {
var raw map[string]interface{}
if err := json.Unmarshal([]byte(line), &raw); err != nil {
return conversationStreamEvent{}, false
}
if chunk, ok := sseparser.ChunkFromRaw(raw, model); ok {
event := conversationStreamEvent{
chunk: &chunk,
text: sseparser.ChunkContent(chunk),
role: sseparser.ChunkRole(chunk),
conversationID: chunk.ConversationID,
channel: sseparser.ChannelFromValue(raw),
finishReason: sseparser.ChunkFinishReason(chunk),
}
event.isStop = event.finishReason != ""
return event, true
}
var direct chatgpt_types.ChatGPTResponse
if err := json.Unmarshal([]byte(line), &direct); err == nil && sseparser.IsUsableConversationResponse(direct) {
channel := sseparser.ChannelFromValue(raw)
state.Channel = firstNonEmpty(channel, state.Channel)
return conversationStreamEvent{response: direct, messageID: direct.Message.ID, channel: state.Channel}, true
}
if response, ok := sseparser.ResponseFromValue(raw["v"]); ok {
state.Response = response
if channel := sseparser.ChannelFromValue(raw["v"]); channel != "" {
state.Channel = channel
}
return conversationStreamEvent{response: state.Response, messageID: state.Response.Message.ID, channel: state.Channel}, true
}
if text, ok := raw["v"].(string); ok && raw["p"] == nil && raw["o"] == nil {
sseparser.EnsurePatchDefaults(state)
current, _ := state.Response.Message.Content.Parts[0].(string)
state.Response.Message.Content.Parts[0] = current + text
return conversationStreamEvent{response: state.Response, messageID: state.Response.Message.ID, channel: state.Channel}, true
}
if patchPath, ok := raw["p"].(string); ok {
patchOperation, _ := raw["o"].(string)
if patchPath == "" && patchOperation == "patch" {
if batch, ok := raw["v"].([]interface{}); ok {
applied := false
for _, item := range batch {
op, ok := item.(map[string]interface{})
if !ok {
continue
}
subPath, _ := op["p"].(string)
subOp, _ := op["o"].(string)
if sseparser.ApplyPatch(state, subPath, subOp, op["v"]) {
applied = true
}
}
if applied {
return conversationStreamEvent{response: state.Response, messageID: state.Response.Message.ID, channel: state.Channel}, true
}
}
}
if sseparser.ApplyPatch(state, patchPath, patchOperation, raw["v"]) {
return conversationStreamEvent{response: state.Response, messageID: state.Response.Message.ID, channel: state.Channel}, true
}
}
return conversationStreamEvent{}, false
}
// Handler 处理对话响应(简化版)。
func Handler(c *gin.Context, response *http.Response, client httpclient.AuroraHttpClient, account *accounts.Account, uuid string, translated_request chatgpt_types.ChatGPTRequest, stream bool, model string) (string, *ContinueInfo) {
result := HandlerDetailed(c, response, client, account, uuid, translated_request, stream, model)
return result.Text, result.Continue
}
// HandlerDetailed 处理对话响应(详细版)。
func HandlerDetailed(c *gin.Context, response *http.Response, client httpclient.AuroraHttpClient, account *accounts.Account, uuid string, translated_request chatgpt_types.ChatGPTRequest, stream bool, model string) HandlerResult {
return HandlerDetailedWithWebsocket(c, response, client, account, uuid, translated_request, stream, model, nil)
}
// HandlerDetailedWithWebsocket 处理对话响应(带 WebSocket)。
func HandlerDetailedWithWebsocket(c *gin.Context, response *http.Response, client httpclient.AuroraHttpClient, account *accounts.Account, uuid string, translated_request chatgpt_types.ChatGPTRequest, stream bool, model string, wsConn *websocket.Conn) HandlerResult {
return HandlerDetailedWithOptions(c, response, client, account, uuid, translated_request, stream, model, HandlerDetailedOptions{Websocket: wsConn})
}
// HandlerDetailedOptions 是 HandlerDetailedWithOptions 的可选参数。
type HandlerDetailedOptions struct {
Websocket *websocket.Conn
ClientState *ChatClientState
ArtifactDelivery string
ProxyURL string
Tools []official_types.Tool
}
// HandlerDetailedWithOptions 处理对话响应流(最完整版)。
func HandlerDetailedWithOptions(c *gin.Context, response *http.Response, client httpclient.AuroraHttpClient, account *accounts.Account, uuid string, translated_request chatgpt_types.ChatGPTRequest, stream bool, model string, options HandlerDetailedOptions) HandlerResult {
if model == "" {
model = translated_request.Model
}
wsConn := options.Websocket
if options.ClientState != nil {
options.ClientState.ApplyToRequest(&translated_request)
}
max_tokens := false
reader := bufio.NewReader(response.Body)
if wsConn != nil {
// The orchestration layer may establish this connection for a
// non-streaming extended/max request before posting the conversation.
defer wsConn.Close()
} else if stream && client != nil && account != nil {
// Preserve the fallback for streaming callers that did not establish a
// WebSocket before entering the response handler.
if conn, err := DialChatWebsocketWithStateAndProxy(client, account, options.ClientState, options.ProxyURL); err == nil {
wsConn = conn
defer wsConn.Close()
}
}
if stream {
c.Header("Content-Type", "text/event-stream")
c.Header("Cache-Control", "no-cache")
c.Header("Connection", "keep-alive")
c.Header("X-Accel-Buffering", "no")
} else {
c.Header("Content-Type", "application/json")
}
var finish_reason string
var previous_text typings.StringStruct
var original_response chatgpt_types.ChatGPTResponse
var isRole = true
var waitSource = false
var isEnd = false
var imgSource []string
var convId string
var sentinel []map[string]interface{}
var thinkingText string
var activeChannel string
var assistantMessageID string
artifactState := newArtifactAccumulator()
artifactConfig := ArtifactStreamConfig{Delivery: options.ArtifactDelivery}
var patchState sseparser.PatchState
var handoffTopicID string
var currentEvent string
var readingWebsocket bool
var websocketStream io.ReadCloser
emitSentinels := func(items []map[string]interface{}) {
if len(items) == 0 {
return
}
sentinel = append(sentinel, items...)
if !stream {
return
}
for _, item := range items {
chunk := official_types.NewChatCompletionChunk("", model)
if convId != "" {
chunk.ConversationID = convId
}
chunk.Sentinel = item
c.Writer.WriteString("data: " + chunk.String() + "\n\n")
c.Writer.Flush()
}
}
observeArtifacts := func(line string) {
var raw map[string]interface{}
if err := json.Unmarshal([]byte(line), &raw); err != nil {
return
}
if cid := firstConversationID(raw); cid != "" && convId == "" {
convId = cid
}
events := artifactState.ObserveRaw(raw, convId)
emitSentinels(materializeArtifactEvents(client, account, convId, events, artifactConfig))
if artifactState.LastAssistantMsgID != "" {
assistantMessageID = artifactState.LastAssistantMsgID
}
if artifactState.ConversationID != "" && convId == "" {
convId = artifactState.ConversationID
}
}
emitThinking := func(delta string) {
if delta == "" {
return
}
thinkingText += delta
emitSentinels([]map[string]interface{}{{
"event": "thinking",
"kind": "analysis",
"delta": delta,
}})
if stream {
reasoningChunk := official_types.NewReasoningChunk(delta, model)
if convId != "" {
reasoningChunk.ConversationID = convId
}
c.Writer.WriteString("data: " + reasoningChunk.String() + "\n\n")
c.Writer.Flush()
}
}
finalizeArtifacts := func() {
emitSentinels(materializeArtifactEvents(client, account, convId, artifactState.Finalize(), artifactConfig))
}
readLoop:
for {
line, err := reader.ReadString('\n')
if err != nil {
if err == io.EOF && line == "" {
break
}
if err != io.EOF {
return HandlerResult{}
}
}
if eventName, ok := sseparser.EventName(line); ok {
currentEvent = eventName
}
for _, line := range sseparser.DataPayloads(line) {
if strings.HasPrefix(line, "[DONE]") {
if shouldUseWebsocketHandoff(readingWebsocket, handoffTopicID, wsConn, previous_text.Text, imgSource) {
wsReader, err := chatWebsocketStreamReader(wsConn, handoffTopicID)
if err == nil {
websocketStream = wsReader
defer websocketStream.Close()
reader = bufio.NewReader(wsReader)
readingWebsocket = true
currentEvent = ""
continue readLoop
}
}
finalizeArtifacts()
break readLoop
}
observeArtifacts(line)
if topicID, skip := sseparser.HandoffTopicFromPayload(line, currentEvent); skip {
if topicID != "" {
handoffTopicID = topicID
}
currentEvent = ""
continue
}
streamEvent, ok := parseConversationEvent(line, &patchState, model)
if os.Getenv("DEBUG_SSE") != "" {
debugText := streamEvent.text
debugSrc := "chunk"
if streamEvent.response.Message.ID != "" {
debugText = sseparser.FirstStringPart(streamEvent.response.Message.Content.Parts)
debugSrc = "response"
}
raw := strings.TrimSpace(line)
if len(raw) > 200 {
raw = raw[:200] + "..."
}
fmt.Printf("[sse-in] src=%s channel=%q textLen=%d finish=%q parsed=%v raw=%q\n", debugSrc, streamEvent.channel, len(debugText), streamEvent.finishReason, ok, raw)
}
if !ok {
currentEvent = ""
continue
}
if streamEvent.chunk != nil {
if streamEvent.conversationID != "" {
convId = streamEvent.conversationID
}
if streamEvent.chunk.Sentinel != nil {
sentinel = append(sentinel, streamEvent.chunk.Sentinel)
}
deltaText := sseparser.NormalizeContentDelta(previous_text.Text, streamEvent.text)
if streamEvent.channel != "" {
activeChannel = streamEvent.channel
}
if streamEvent.finishReason != "" {
finish_reason = streamEvent.finishReason
if finish_reason == "length" {
max_tokens = true
}
isEnd = true
}
if activeChannel == "analysis" {
emitThinking(streamEvent.text)
if streamEvent.isStop {
if stream {
finalLine := official_types.StopChunkWithConversation(finish_reason, model, convId)
c.Writer.WriteString("data: " + finalLine.String() + "\n\n")
c.Writer.Flush()
}
if max_tokens && convId != "" && assistantMessageID != "" {
finalizeArtifacts()
return HandlerResult{
Text: strings.Join(imgSource, "") + previous_text.Text,
ThinkingText: thinkingText,
ConversationID: convId,
ParentMessageID: assistantMessageID,
Sentinel: sentinel,
ArtifactSignals: artifactState.Signals,
SandboxArtifacts: artifactState.SandboxArtifacts,
PDFArtifacts: artifactState.PDFArtifacts,
GeneratedImageIDs: artifactState.ImageFileIDs,
StopSent: true,
Continue: &ContinueInfo{
ConversationID: convId,
ParentID: assistantMessageID,
},
}
}
finalizeArtifacts()
return HandlerResult{
Text: strings.Join(imgSource, "") + previous_text.Text,
ThinkingText: thinkingText,
ConversationID: convId,
ParentMessageID: assistantMessageID,
Sentinel: sentinel,
ArtifactSignals: artifactState.Signals,
SandboxArtifacts: artifactState.SandboxArtifacts,
PDFArtifacts: artifactState.PDFArtifacts,
GeneratedImageIDs: artifactState.ImageFileIDs,
StopSent: true,
}
}
currentEvent = ""
continue
}
if stream {
outChunk := *streamEvent.chunk
if len(outChunk.Choices) > 0 {
outChunk.Choices[0].Delta.Content = deltaText
if streamEvent.role == "" || !isRole {
outChunk.Choices[0].Delta.Role = ""
}
}
if streamEvent.isStop && outChunk.ConversationID == "" {
outChunk.ConversationID = convId
}
shouldWrite := deltaText != "" ||
(streamEvent.role != "" && isRole) ||
streamEvent.chunk.Sentinel != nil ||
streamEvent.isStop
if shouldWrite {
c.Writer.WriteString("data: " + outChunk.String() + "\n\n")
c.Writer.Flush()
}
if streamEvent.role != "" && isRole {
isRole = false
}
}
if deltaText != "" {
previous_text.Text += deltaText
}
if streamEvent.isStop {
if max_tokens && convId != "" && assistantMessageID != "" {
finalizeArtifacts()
return HandlerResult{
Text: strings.Join(imgSource, "") + previous_text.Text,
ThinkingText: thinkingText,
ConversationID: convId,
ParentMessageID: assistantMessageID,
Sentinel: sentinel,
ArtifactSignals: artifactState.Signals,
SandboxArtifacts: artifactState.SandboxArtifacts,
PDFArtifacts: artifactState.PDFArtifacts,
GeneratedImageIDs: artifactState.ImageFileIDs,
StopSent: true,
Continue: &ContinueInfo{
ConversationID: convId,
ParentID: assistantMessageID,
},
}
}
finalizeArtifacts()
return HandlerResult{
Text: strings.Join(imgSource, "") + previous_text.Text,
ThinkingText: thinkingText,
ConversationID: convId,
ParentMessageID: assistantMessageID,
Sentinel: sentinel,
ArtifactSignals: artifactState.Signals,
SandboxArtifacts: artifactState.SandboxArtifacts,
PDFArtifacts: artifactState.PDFArtifacts,
GeneratedImageIDs: artifactState.ImageFileIDs,
StopSent: true,
}
}
currentEvent = ""
continue
}
original_response = streamEvent.response
if original_response.Error != nil {
c.JSON(500, gin.H{"error": original_response.Error})
return HandlerResult{}
}
sentinel = append(sentinel, sseparser.SentinelsFromResponse(original_response)...)
if original_response.ConversationID != convId {
if convId == "" {
convId = original_response.ConversationID
} else {
continue
}
}
if streamEvent.channel != "" {
activeChannel = streamEvent.channel
}
if original_response.Message.ID != "" && (original_response.Message.Author.Role == "assistant" || original_response.Message.Author.Role == "tool") {
assistantMessageID = original_response.Message.ID
}
if activeChannel == "analysis" {
thinkingDelta := sseparser.NormalizeContentDelta(thinkingText, sseparser.FirstStringPart(original_response.Message.Content.Parts))
emitThinking(thinkingDelta)
currentEvent = ""
continue
}
if !(original_response.Message.Author.Role == "assistant" || (original_response.Message.Author.Role == "tool" && original_response.Message.Content.ContentType != "text")) || original_response.Message.Content.Parts == nil {
continue
}
if original_response.Message.Metadata.MessageType == "" && activeChannel != "final" {
continue
}
if (original_response.Message.Metadata.MessageType != "next" && original_response.Message.Metadata.MessageType != "continue" && activeChannel != "final") || !strings.HasSuffix(original_response.Message.Content.ContentType, "text") {
continue
}
if original_response.Message.EndTurn != nil {
if waitSource {
waitSource = false
}
isEnd = true
}
if len(original_response.Message.Metadata.Citations) != 0 {
r := []rune(original_response.Message.Content.Parts[0].(string))
if waitSource {
if string(r[len(r)-1:]) == "】" {
waitSource = false
} else {
continue
}
}
offset := 0
for _, citation := range original_response.Message.Metadata.Citations {
rl := len(r)
attr := urlAttrMap[citation.Metadata.URL]
if attr == "" {
u, _ := url.Parse(citation.Metadata.URL)
BaseURL := u.Scheme + "://" + u.Host + "/"
attr = getURLAttribution(client, account, BaseURL)
if attr != "" {
urlAttrMap[citation.Metadata.URL] = attr
}
}
original_response.Message.Content.Parts[0] = string(r[:citation.StartIx+offset]) + " ([" + attr + "](" + citation.Metadata.URL + " \"" + citation.Metadata.Title + "\"))" + string(r[citation.EndIx+offset:])
r = []rune(original_response.Message.Content.Parts[0].(string))
offset += len(r) - rl
}
} else if waitSource {
continue
}
response_string := ""
if original_response.Message.Recipient != "all" {
continue
}
if original_response.Message.Content.ContentType == "multimodal_text" {
apiUrl := BaseURL + "/files/"
if FILES_REVERSE_PROXY != "" {
apiUrl = FILES_REVERSE_PROXY
}
imgSource = make([]string, len(original_response.Message.Content.Parts))
var wg sync.WaitGroup
for index, part := range original_response.Message.Content.Parts {
jsonItem, _ := json.Marshal(part)
var dalle_content chatgpt_types.DalleContent
err = json.Unmarshal(jsonItem, &dalle_content)
if err != nil {
continue
}
url := apiUrl + strings.Split(dalle_content.AssetPointer, "//")[1] + "/download"
wg.Add(1)
go GetImageSource(client, &wg, url, dalle_content.Metadata.Dalle.Prompt, account, index, imgSource)
}
wg.Wait()
translated_response := official_types.NewChatCompletionChunk(strings.Join(imgSource, ""), model)
if isRole {
translated_response.Choices[0].Delta.Role = original_response.Message.Author.Role
}
response_string = "data: " + translated_response.String() + "\n\n"
}
if response_string == "" {
response_string = chatgpt.ConvertToString(&original_response, &previous_text, isRole, model)
}
if response_string == "" {
if isEnd {
goto endProcess
} else {
continue
}
}
if response_string == "【" {
waitSource = true
continue
}
endProcess:
isRole = false
if stream {
_, err = c.Writer.WriteString(response_string)
if err != nil {
return HandlerResult{}
}
c.Writer.Flush()
}
if original_response.Message.Metadata.FinishDetails != nil {
if original_response.Message.Metadata.FinishDetails.Type == "max_tokens" {
max_tokens = true
}
finish_reason = original_response.Message.Metadata.FinishDetails.Type
}
if isEnd {
if stream {
final_line := official_types.StopChunkWithConversation(finish_reason, model, convId)
c.Writer.WriteString("data: " + final_line.String() + "\n\n")
c.Writer.Flush()
}
finalizeArtifacts()
return HandlerResult{
Text: strings.Join(imgSource, "") + previous_text.Text,
ThinkingText: thinkingText,
ConversationID: convId,
ParentMessageID: assistantMessageID,
Sentinel: sentinel,
ArtifactSignals: artifactState.Signals,
SandboxArtifacts: artifactState.SandboxArtifacts,
PDFArtifacts: artifactState.PDFArtifacts,
GeneratedImageIDs: artifactState.ImageFileIDs,
StopSent: stream,
}
}
currentEvent = ""
}
if err == io.EOF {
break
}
}
if !max_tokens {
finalizeArtifacts()
return HandlerResult{
Text: strings.Join(imgSource, "") + previous_text.Text,
ThinkingText: thinkingText,
ConversationID: convId,
ParentMessageID: assistantMessageID,
Sentinel: sentinel,
ArtifactSignals: artifactState.Signals,
SandboxArtifacts: artifactState.SandboxArtifacts,
PDFArtifacts: artifactState.PDFArtifacts,
GeneratedImageIDs: artifactState.ImageFileIDs,
}
}
finalizeArtifacts()
return HandlerResult{
Text: strings.Join(imgSource, "") + previous_text.Text,
ThinkingText: thinkingText,
ConversationID: convId,
ParentMessageID: assistantMessageID,
Sentinel: sentinel,
ArtifactSignals: artifactState.Signals,
SandboxArtifacts: artifactState.SandboxArtifacts,
PDFArtifacts: artifactState.PDFArtifacts,
GeneratedImageIDs: artifactState.ImageFileIDs,
Continue: &ContinueInfo{
ConversationID: original_response.ConversationID,
ParentID: original_response.Message.ID,
},
}
}