Gin + GORM + pure-Go SQLite. Users/auth (JWT), API key management with quotas, proxy gateway with weighted channel failover and health checks, usage/billing ledger, cross-protocol conversion (Anthropic Messages / OpenAI Chat Completions / OpenAI Responses), and channel/model admin API. Channels declare native API formats and auto-convert the rest. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
417 lines
12 KiB
Go
417 lines
12 KiB
Go
package convert
|
|
|
|
import (
|
|
"encoding/json"
|
|
|
|
"openteam/server/internal/proxy/claude"
|
|
"openteam/server/internal/proxy/openai"
|
|
"openteam/server/internal/proxy/stream"
|
|
)
|
|
|
|
// Usage is the streaming token usage snapshot.
|
|
type Usage struct {
|
|
Input int64
|
|
Output int64
|
|
CacheRead int64
|
|
CacheCreation int64
|
|
}
|
|
|
|
// Translator converts SSE events from an upstream stream into client frames.
|
|
type Translator interface {
|
|
// Feed handles one upstream SSE event, returning client frames to write.
|
|
Feed(ev stream.SSEEvent) ([]stream.SSEEvent, error)
|
|
// Finish is called at end-of-stream, returning final frames.
|
|
Finish() ([]stream.SSEEvent, error)
|
|
// Usage returns the latest known usage.
|
|
Usage() *Usage
|
|
}
|
|
|
|
// claudeToChatTranslator converts a Claude stream to OpenAI chat chunks.
|
|
type claudeToChatTranslator struct {
|
|
model string
|
|
usage *Usage
|
|
started bool
|
|
finishSent bool
|
|
toolCallIndex int
|
|
toolCallID string
|
|
toolCallName string
|
|
}
|
|
|
|
func (t *claudeToChatTranslator) Feed(ev stream.SSEEvent) ([]stream.SSEEvent, error) {
|
|
if ev.Done {
|
|
return nil, nil
|
|
}
|
|
var e claude.StreamEvent
|
|
if err := json.Unmarshal([]byte(ev.Data), &e); err != nil {
|
|
return nil, nil
|
|
}
|
|
var out []stream.SSEEvent
|
|
|
|
switch e.Type {
|
|
case "message_start":
|
|
if e.Message != nil {
|
|
var msg struct {
|
|
Model string `json:"model"`
|
|
}
|
|
_ = json.Unmarshal(e.Message, &msg)
|
|
t.model = msg.Model
|
|
}
|
|
chunk, _ := json.Marshal(openai.ChatChunk{
|
|
ID: "chatcmpl-stream", Object: "chat.completion.chunk", Model: t.model,
|
|
Choices: []openai.ChatChunkChoice{{Index: 0, Delta: openai.ChatDelta{Role: "assistant"}}},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(chunk)})
|
|
t.started = true
|
|
|
|
case "content_block_start":
|
|
var cb struct {
|
|
Index int `json:"index"`
|
|
Block json.RawMessage `json:"content_block"`
|
|
}
|
|
_ = json.Unmarshal([]byte(ev.Data), &cb)
|
|
var block struct {
|
|
Type string `json:"type"`
|
|
ID string `json:"id"`
|
|
Name string `json:"name"`
|
|
}
|
|
_ = json.Unmarshal(cb.Block, &block)
|
|
if block.Type == "tool_use" {
|
|
t.toolCallIndex = cb.Index
|
|
t.toolCallID = block.ID
|
|
t.toolCallName = block.Name
|
|
tc, _ := json.Marshal([]map[string]any{{
|
|
"index": cb.Index, "id": block.ID, "type": "function",
|
|
"function": map[string]any{"name": block.Name, "arguments": ""},
|
|
}})
|
|
chunk, _ := json.Marshal(openai.ChatChunk{
|
|
ID: "chatcmpl-stream", Object: "chat.completion.chunk", Model: t.model,
|
|
Choices: []openai.ChatChunkChoice{{Index: 0, Delta: openai.ChatDelta{ToolCalls: tc}}},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(chunk)})
|
|
}
|
|
|
|
case "content_block_delta":
|
|
var d struct {
|
|
Index int `json:"index"`
|
|
Delta json.RawMessage `json:"delta"`
|
|
}
|
|
_ = json.Unmarshal([]byte(ev.Data), &d)
|
|
var delta struct {
|
|
Type string `json:"type"`
|
|
Text string `json:"text"`
|
|
PartialJSON string `json:"partial_json"`
|
|
}
|
|
_ = json.Unmarshal(d.Delta, &delta)
|
|
if delta.Type == "text_delta" && delta.Text != "" {
|
|
chunk, _ := json.Marshal(openai.ChatChunk{
|
|
ID: "chatcmpl-stream", Object: "chat.completion.chunk", Model: t.model,
|
|
Choices: []openai.ChatChunkChoice{{Index: 0, Delta: openai.ChatDelta{Content: delta.Text}}},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(chunk)})
|
|
} else if delta.Type == "input_json_delta" && delta.PartialJSON != "" {
|
|
tc, _ := json.Marshal([]map[string]any{{
|
|
"index": d.Index, "function": map[string]any{"arguments": delta.PartialJSON},
|
|
}})
|
|
chunk, _ := json.Marshal(openai.ChatChunk{
|
|
ID: "chatcmpl-stream", Object: "chat.completion.chunk", Model: t.model,
|
|
Choices: []openai.ChatChunkChoice{{Index: 0, Delta: openai.ChatDelta{ToolCalls: tc}}},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(chunk)})
|
|
}
|
|
|
|
case "message_delta":
|
|
var d struct {
|
|
Delta json.RawMessage `json:"delta"`
|
|
Usage json.RawMessage `json:"usage"`
|
|
}
|
|
_ = json.Unmarshal([]byte(ev.Data), &d)
|
|
if len(d.Usage) > 0 {
|
|
var u claude.Usage
|
|
if json.Unmarshal(d.Usage, &u) == nil {
|
|
t.usage = &Usage{
|
|
Input: u.InputTokens, Output: u.OutputTokens,
|
|
CacheRead: u.CacheReadInputTokens, CacheCreation: u.CacheCreationInputTokens,
|
|
}
|
|
}
|
|
}
|
|
if len(d.Delta) > 0 && !t.finishSent {
|
|
var delta struct {
|
|
StopReason string `json:"stop_reason"`
|
|
}
|
|
_ = json.Unmarshal(d.Delta, &delta)
|
|
if delta.StopReason != "" {
|
|
reason := mapClaudeStopReason(delta.StopReason)
|
|
chunk, _ := json.Marshal(openai.ChatChunk{
|
|
ID: "chatcmpl-stream", Object: "chat.completion.chunk", Model: t.model,
|
|
Choices: []openai.ChatChunkChoice{{Index: 0, Delta: openai.ChatDelta{}, FinishReason: &reason}},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(chunk)})
|
|
t.finishSent = true
|
|
}
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (t *claudeToChatTranslator) Finish() ([]stream.SSEEvent, error) {
|
|
if !t.finishSent {
|
|
reason := "stop"
|
|
chunk, _ := json.Marshal(openai.ChatChunk{
|
|
ID: "chatcmpl-stream", Object: "chat.completion.chunk", Model: t.model,
|
|
Choices: []openai.ChatChunkChoice{{Index: 0, Delta: openai.ChatDelta{}, FinishReason: &reason}},
|
|
})
|
|
t.finishSent = true
|
|
return []stream.SSEEvent{{Data: string(chunk)}, {Data: "[DONE]"}}, nil
|
|
}
|
|
return []stream.SSEEvent{{Data: "[DONE]"}}, nil
|
|
}
|
|
|
|
func (t *claudeToChatTranslator) Usage() *Usage { return t.usage }
|
|
|
|
// chatToClaudeTranslator converts an OpenAI chat stream to Claude events.
|
|
type chatToClaudeTranslator struct {
|
|
usage *Usage
|
|
started bool
|
|
openBlock bool
|
|
blockType string
|
|
toolIndex int
|
|
finishSent bool
|
|
}
|
|
|
|
func (t *chatToClaudeTranslator) Feed(ev stream.SSEEvent) ([]stream.SSEEvent, error) {
|
|
if ev.Done {
|
|
return nil, nil
|
|
}
|
|
var chunk openai.ChatChunk
|
|
if err := json.Unmarshal([]byte(ev.Data), &chunk); err != nil {
|
|
return nil, nil
|
|
}
|
|
if chunk.Usage != nil {
|
|
t.usage = &Usage{Input: chunk.Usage.PromptTokens, Output: chunk.Usage.CompletionTokens}
|
|
}
|
|
var out []stream.SSEEvent
|
|
if len(chunk.Choices) == 0 {
|
|
return out, nil
|
|
}
|
|
ch := chunk.Choices[0]
|
|
|
|
if !t.started {
|
|
msg, _ := json.Marshal(map[string]any{
|
|
"id": "msg_stream", "type": "message", "role": "assistant",
|
|
"model": chunk.Model, "content": []any{},
|
|
})
|
|
start, _ := json.Marshal(map[string]any{"type": "message_start", "message": json.RawMessage(msg)})
|
|
out = append(out, stream.SSEEvent{Data: string(start)})
|
|
t.started = true
|
|
}
|
|
|
|
if ch.Delta.Content != "" {
|
|
if !t.openBlock || t.blockType != "text" {
|
|
start, _ := json.Marshal(map[string]any{
|
|
"type": "content_block_start", "index": 0,
|
|
"content_block": map[string]any{"type": "text", "text": ""},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(start)})
|
|
t.openBlock = true
|
|
t.blockType = "text"
|
|
t.toolIndex = 0
|
|
}
|
|
delta, _ := json.Marshal(map[string]any{
|
|
"type": "content_block_delta", "index": 0,
|
|
"delta": map[string]any{"type": "text_delta", "text": ch.Delta.Content},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(delta)})
|
|
}
|
|
|
|
if len(ch.Delta.ToolCalls) > 0 {
|
|
var calls []struct {
|
|
Index *int `json:"index"`
|
|
ID string `json:"id"`
|
|
Function struct {
|
|
Name string `json:"name"`
|
|
Arguments string `json:"arguments"`
|
|
} `json:"function"`
|
|
}
|
|
_ = json.Unmarshal(ch.Delta.ToolCalls, &calls)
|
|
for _, call := range calls {
|
|
idx := 0
|
|
if call.Index != nil {
|
|
idx = *call.Index
|
|
}
|
|
if !t.openBlock || t.blockType != "tool_use" || idx != t.toolIndex {
|
|
start, _ := json.Marshal(map[string]any{
|
|
"type": "content_block_start", "index": idx,
|
|
"content_block": map[string]any{
|
|
"type": "tool_use", "id": call.ID, "name": call.Function.Name, "input": map[string]any{},
|
|
},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(start)})
|
|
t.openBlock = true
|
|
t.blockType = "tool_use"
|
|
t.toolIndex = idx
|
|
}
|
|
if call.Function.Arguments != "" {
|
|
delta, _ := json.Marshal(map[string]any{
|
|
"type": "content_block_delta", "index": idx,
|
|
"delta": map[string]any{"type": "input_json_delta", "partial_json": call.Function.Arguments},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(delta)})
|
|
}
|
|
}
|
|
}
|
|
|
|
if ch.FinishReason != nil && !t.finishSent {
|
|
reason := mapChatStopReasonToClaude(*ch.FinishReason)
|
|
md, _ := json.Marshal(map[string]any{
|
|
"type": "message_delta",
|
|
"delta": map[string]any{"stop_reason": reason, "stop_sequence": nil},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(md)})
|
|
t.finishSent = true
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (t *chatToClaudeTranslator) Finish() ([]stream.SSEEvent, error) {
|
|
if !t.finishSent {
|
|
md, _ := json.Marshal(map[string]any{
|
|
"type": "message_delta",
|
|
"delta": map[string]any{"stop_reason": "end_turn", "stop_sequence": nil},
|
|
})
|
|
t.finishSent = true
|
|
return []stream.SSEEvent{{Data: string(md)}, {Data: "{\"type\":\"message_stop\"}"}}, nil
|
|
}
|
|
return []stream.SSEEvent{{Data: "{\"type\":\"message_stop\"}"}}, nil
|
|
}
|
|
|
|
func (t *chatToClaudeTranslator) Usage() *Usage { return t.usage }
|
|
|
|
// claudeToResponsesTranslator converts a Claude stream to Responses events.
|
|
type claudeToResponsesTranslator struct {
|
|
usage *Usage
|
|
model string
|
|
completed bool
|
|
}
|
|
|
|
func (t *claudeToResponsesTranslator) Feed(ev stream.SSEEvent) ([]stream.SSEEvent, error) {
|
|
if ev.Done {
|
|
return nil, nil
|
|
}
|
|
var e claude.StreamEvent
|
|
if err := json.Unmarshal([]byte(ev.Data), &e); err != nil {
|
|
return nil, nil
|
|
}
|
|
var out []stream.SSEEvent
|
|
switch e.Type {
|
|
case "message_start":
|
|
if e.Message != nil {
|
|
var msg struct {
|
|
Model string `json:"model"`
|
|
}
|
|
_ = json.Unmarshal(e.Message, &msg)
|
|
t.model = msg.Model
|
|
}
|
|
created, _ := json.Marshal(map[string]any{
|
|
"type": "response.created",
|
|
"response": map[string]any{"id": "resp_stream", "object": "response", "status": "in_progress", "model": t.model, "output": []any{}},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(created)})
|
|
case "content_block_delta":
|
|
var d struct {
|
|
Index int `json:"index"`
|
|
Delta json.RawMessage `json:"delta"`
|
|
}
|
|
_ = json.Unmarshal([]byte(ev.Data), &d)
|
|
var delta struct {
|
|
Type string `json:"type"`
|
|
Text string `json:"text"`
|
|
}
|
|
_ = json.Unmarshal(d.Delta, &delta)
|
|
if delta.Type == "text_delta" && delta.Text != "" {
|
|
item, _ := json.Marshal(map[string]any{
|
|
"type": "response.output_text.delta", "item_id": "msg_stream", "output_index": 0,
|
|
"delta": delta.Text,
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(item)})
|
|
}
|
|
case "message_delta":
|
|
var d struct {
|
|
Usage json.RawMessage `json:"usage"`
|
|
}
|
|
_ = json.Unmarshal([]byte(ev.Data), &d)
|
|
if len(d.Usage) > 0 {
|
|
var u claude.Usage
|
|
if json.Unmarshal(d.Usage, &u) == nil {
|
|
t.usage = &Usage{
|
|
Input: u.InputTokens, Output: u.OutputTokens,
|
|
CacheRead: u.CacheReadInputTokens, CacheCreation: u.CacheCreationInputTokens,
|
|
}
|
|
}
|
|
}
|
|
if !t.completed {
|
|
item, _ := json.Marshal(map[string]any{
|
|
"type": "response.completed",
|
|
"response": map[string]any{
|
|
"id": "resp_stream", "object": "response", "status": "completed", "model": t.model,
|
|
"output": []map[string]any{{"type": "message", "role": "assistant", "content": []any{}}},
|
|
},
|
|
})
|
|
out = append(out, stream.SSEEvent{Data: string(item)})
|
|
t.completed = true
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (t *claudeToResponsesTranslator) Finish() ([]stream.SSEEvent, error) {
|
|
if !t.completed {
|
|
item, _ := json.Marshal(map[string]any{
|
|
"type": "response.completed",
|
|
"response": map[string]any{
|
|
"id": "resp_stream", "object": "response", "status": "completed", "model": t.model,
|
|
"output": []map[string]any{{"type": "message", "role": "assistant", "content": []any{}}},
|
|
},
|
|
})
|
|
t.completed = true
|
|
return []stream.SSEEvent{{Data: string(item)}, {Data: "[DONE]"}}, nil
|
|
}
|
|
return []stream.SSEEvent{{Data: "[DONE]"}}, nil
|
|
}
|
|
|
|
func (t *claudeToResponsesTranslator) Usage() *Usage { return t.usage }
|
|
|
|
func mapChatStopReasonToClaude(reason string) string {
|
|
switch reason {
|
|
case "stop":
|
|
return "end_turn"
|
|
case "length":
|
|
return "max_tokens"
|
|
case "tool_calls":
|
|
return "tool_use"
|
|
case "content_filter":
|
|
return "refusal"
|
|
default:
|
|
return "end_turn"
|
|
}
|
|
}
|
|
|
|
// NewTranslator returns the stream translator for a client protocol +
|
|
// upstream provider pair, or nil for passthrough (no translation needed).
|
|
func NewTranslator(client ClientProtocol, provider string) Translator {
|
|
switch client {
|
|
case ClientOpenAIChat:
|
|
if provider == "anthropic" {
|
|
return &claudeToChatTranslator{}
|
|
}
|
|
case ClientOpenAIResponses:
|
|
if provider == "anthropic" {
|
|
return &claudeToResponsesTranslator{}
|
|
}
|
|
case ClientAnthropic:
|
|
if provider == "openai" || provider == "compatible" {
|
|
return &chatToClaudeTranslator{}
|
|
}
|
|
}
|
|
return nil
|
|
}
|