diff --git a/internal/client/anthropic.go b/internal/client/anthropic.go index 6a23ef34a..8ff8fb75a 100644 --- a/internal/client/anthropic.go +++ b/internal/client/anthropic.go @@ -2,7 +2,6 @@ package client import ( "context" - "encoding/json" "fmt" "net/http" "strings" @@ -41,10 +40,6 @@ type AnthropicClientInterface interface { APIStyle() protocol.APIStyle SetRecordSink(sink *obs.Sink) Client() *anthropic.Client - - // Prober interface methods - Probe(ctx context.Context, model string) ProbeResult - ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) } // AnthropicClient wraps the Anthropic SDK client @@ -212,234 +207,3 @@ func (c *AnthropicClient) ListModels(ctx context.Context) ([]string, error) { return result, nil } - -// ProbeChatEndpoint tests the messages endpoint with a minimal request -func (c *AnthropicClient) Probe(ctx context.Context, model string) ProbeResult { - startTime := time.Now() - - // Determine system message based on OAuth provider type - systemMessages := []anthropic.TextBlockParam{ - { - Text: "work as `echo`", - }, - } - if c.provider.AuthType == typ.AuthTypeOAuth && c.provider.OAuthDetail != nil && - c.provider.OAuthDetail.GetIssuer() == ai.IssuerClaudeCode { - // Prepend Claude Code system message as the first block - systemMessages = append([]anthropic.TextBlockParam{{ - Text: ClaudeCodeSystemHeader, - }}, systemMessages...) - } - - // Create message request using Anthropic SDK - messageRequest := anthropic.MessageNewParams{ - Model: anthropic.Model(model), - MaxTokens: 100, - System: systemMessages, - Messages: []anthropic.MessageParam{ - anthropic.NewUserMessage(anthropic.NewTextBlock("hi")), - }, - } - - // Make request - resp, err := c.client.Messages.New(ctx, messageRequest) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: err.Error(), - LatencyMs: latencyMs, - } - } - - // Extract response data - responseContent := "" - promptTokens := 0 - completionTokens := 0 - totalTokens := 0 - - if resp != nil { - for _, block := range resp.Content { - if block.Type == "text" { - responseContent += string(block.Text) - } - } - if resp.Usage.InputTokens != 0 { - promptTokens = int(resp.Usage.InputTokens) - completionTokens = int(resp.Usage.OutputTokens) - totalTokens = promptTokens + completionTokens - } - } - - if responseContent == "" { - responseContent = "" - } - - return ProbeResult{ - Success: true, - Message: "Messages endpoint is accessible", - Content: responseContent, - LatencyMs: latencyMs, - PromptTokens: promptTokens, - CompletionTokens: completionTokens, - TotalTokens: totalTokens, - } -} - -// ProbeStream performs a streaming probe with configurable test mode (public interface) -func (c *AnthropicClient) ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - return c.probeStream(ctx, model, message, testMode) -} - -// probeStream performs a streaming probe with configurable test mode -func (c *AnthropicClient) probeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - startTime := time.Now() - - // Determine system message based on OAuth provider type - systemMessages := []anthropic.TextBlockParam{ - { - Text: "work as `echo` if possible", - }, - } - if c.provider.AuthType == typ.AuthTypeOAuth && c.provider.OAuthDetail != nil && - c.provider.OAuthDetail.GetIssuer() == ai.IssuerClaudeCode { - // Prepend Claude Code system message as the first block - systemMessages = append([]anthropic.TextBlockParam{{ - Text: ClaudeCodeSystemHeader, - }}, systemMessages...) - } - - messages := []anthropic.MessageParam{ - anthropic.NewUserMessage(anthropic.NewTextBlock(message)), - } - - params := &anthropic.MessageNewParams{ - Model: anthropic.Model(model), - MaxTokens: 1024, - System: systemMessages, - Messages: messages, - } - - if testMode == ProbeModeTool { - params.Tools = GetProbeToolsAnthropic() - params.ToolChoice = GetProbeToolChoiceAutoAnthropic() - } - - // For simple mode, use non-streaming request - if testMode == ProbeModeSimple { - resp, err := c.client.Messages.New(ctx, *params) - if err != nil { - return nil, err - } - - respJSON, _ := json.Marshal(resp) - return ToProbeResult(string(respJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase+"/v1/messages", false), nil - } - - // For streaming and tool modes, use streaming - stream := c.client.Messages.NewStreaming(ctx, *params) - defer stream.Close() - - var chunks []interface{} - for stream.Next() { - event := stream.Current() - chunks = append(chunks, event) - } - - if err := stream.Err(); err != nil { - return nil, err - } - - chunksJSON, _ := json.Marshal(chunks) - return ToProbeResult(string(chunksJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase+"/v1/messages", true), nil -} - -// ProbeModelsEndpoint tests the models list endpoint -func (c *AnthropicClient) ProbeModelsEndpoint(ctx context.Context) ProbeResult { - startTime := time.Now() - - // Make request to models endpoint - resp, err := c.client.Models.List(ctx, anthropic.ModelListParams{}) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: err.Error(), - LatencyMs: latencyMs, - } - } - - modelsCount := 0 - if resp != nil { - modelsCount = len(resp.Data) - } - - if modelsCount == 0 { - return ProbeResult{ - Success: false, - ErrorMessage: "No models available from provider", - LatencyMs: latencyMs, - } - } - - return ProbeResult{ - Success: true, - Message: "Models endpoint is accessible", - LatencyMs: latencyMs, - ModelsCount: modelsCount, - } -} - -// ProbeOptionsEndpoint tests basic connectivity with an OPTIONS request -func (c *AnthropicClient) ProbeOptionsEndpoint(ctx context.Context) ProbeResult { - startTime := time.Now() - - // Build the options URL - ensure it has /v1 suffix for Anthropic - apiBase := strings.TrimSuffix(c.provider.APIBase, "/") - if !strings.Contains(apiBase, "/v1") { - apiBase = apiBase + "/v1" - } - optionsURL := apiBase - - req, err := http.NewRequestWithContext(ctx, "OPTIONS", optionsURL, nil) - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("Failed to create OPTIONS request: %v", err), - } - } - - // Set authentication headers - req.Header.Set("x-api-key", c.provider.GetAccessToken()) - req.Header.Set("anthropic-version", "2023-06-01") - - client := &http.Client{Timeout: 5 * time.Second} - resp, err := client.Do(req) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("OPTIONS request failed: %v", err), - LatencyMs: latencyMs, - } - } - defer resp.Body.Close() - - // Consider any 2xx status as success for OPTIONS - if resp.StatusCode >= 200 && resp.StatusCode < 300 { - return ProbeResult{ - Success: true, - Message: "OPTIONS request successful", - LatencyMs: latencyMs, - } - } - - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("OPTIONS request failed with status: %d", resp.StatusCode), - LatencyMs: latencyMs, - } -} diff --git a/internal/client/anthropic_probe_test.go b/internal/client/anthropic_probe_test.go index c8abef522..f13dfbd0b 100644 --- a/internal/client/anthropic_probe_test.go +++ b/internal/client/anthropic_probe_test.go @@ -1,131 +1,12 @@ package client import ( - "context" "testing" "github.com/tingly-dev/tingly-box/internal/protocol" "github.com/tingly-dev/tingly-box/internal/typ" ) -// TestAnthropicClient_ProbeChatEndpoint tests the ProbeChatEndpoint method -func TestAnthropicClient_ProbeChatEndpoint(t *testing.T) { - tests := []struct { - name string - provider *typ.Provider - model string - wantErr bool - }{ - { - name: "skip live test - requires valid API key", - provider: &typ.Provider{ - Name: "test-anthropic", - APIBase: "https://api.anthropic.com", - APIStyle: protocol.APIStyleAnthropic, - Token: "sk-test-key", - }, - model: "claude-3-haiku-20240307", - wantErr: true, // Will fail with invalid key - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - client, err := NewAnthropicClient(tt.provider, tt.model, typ.SessionID{}) - if err != nil { - t.Fatalf("NewAnthropicClient() error = %v", err) - } - - result := client.Probe(context.Background(), tt.model) - - if !tt.wantErr && !result.Success { - t.Errorf("ProbeChatEndpoint() failed = %v", result.ErrorMessage) - } - if tt.wantErr && result.Success { - t.Errorf("ProbeChatEndpoint() expected error but succeeded") - } - }) - } -} - -// TestAnthropicClient_ProbeModelsEndpoint tests the ProbeModelsEndpoint method -func TestAnthropicClient_ProbeModelsEndpoint(t *testing.T) { - tests := []struct { - name string - provider *typ.Provider - model string - wantErr bool - }{ - { - name: "skip live test - requires valid API key", - provider: &typ.Provider{ - Name: "test-anthropic", - APIBase: "https://api.anthropic.com", - APIStyle: protocol.APIStyleAnthropic, - Token: "sk-test-key", - }, - model: "claude-3-haiku-20240307", - wantErr: true, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - client, err := NewAnthropicClient(tt.provider, tt.model, typ.SessionID{}) - if err != nil { - t.Fatalf("NewAnthropicClient() error = %v", err) - } - - result := client.ProbeModelsEndpoint(context.Background()) - - if !tt.wantErr && !result.Success { - t.Errorf("ProbeModelsEndpoint() failed = %v", result.ErrorMessage) - } - if tt.wantErr && result.Success { - t.Errorf("ProbeModelsEndpoint() expected error but succeeded") - } - }) - } -} - -// TestAnthropicClient_ProbeOptionsEndpoint tests the ProbeOptionsEndpoint method -func TestAnthropicClient_ProbeOptionsEndpoint(t *testing.T) { - tests := []struct { - name string - provider *typ.Provider - model string - wantErr bool - }{ - { - name: "skip live test - requires valid API key", - provider: &typ.Provider{ - Name: "test-anthropic", - APIBase: "https://api.anthropic.com", - APIStyle: protocol.APIStyleAnthropic, - Token: "sk-test-key", - }, - model: "claude-3-haiku-20240307", - wantErr: true, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - client, err := NewAnthropicClient(tt.provider, tt.model, typ.SessionID{}) - if err != nil { - t.Fatalf("NewAnthropicClient() error = %v", err) - } - - result := client.ProbeOptionsEndpoint(context.Background()) - - if !tt.wantErr && !result.Success { - t.Errorf("ProbeOptionsEndpoint() failed = %v", result.ErrorMessage) - } - // OPTIONS might succeed even with invalid key for some providers - }) - } -} - // TestAnthropicClient_Timeout tests that timeout is properly configured from provider func TestAnthropicClient_Timeout(t *testing.T) { tests := []struct { @@ -171,4 +52,3 @@ func TestAnthropicClient_Timeout(t *testing.T) { }) } } - diff --git a/internal/client/claude_client.go b/internal/client/claude_client.go index 51654d550..31e2c74fb 100644 --- a/internal/client/claude_client.go +++ b/internal/client/claude_client.go @@ -2,10 +2,8 @@ package client import ( "context" - "encoding/json" "fmt" "strings" - "time" "github.com/anthropics/anthropic-sdk-go" anthropicOption "github.com/anthropics/anthropic-sdk-go/option" @@ -129,16 +127,6 @@ func (c *ClaudeClient) ListModels(ctx context.Context) ([]string, error) { } } -// ProbeModelsEndpoint reports the /models endpoint as unavailable without -// issuing a request. Claude Code OAuth tokens cannot access /models, so -// probing would always fail and may spam upstream with rejected calls. -func (c *ClaudeClient) ProbeModelsEndpoint(ctx context.Context) ProbeResult { - return ProbeResult{ - Success: false, - ErrorMessage: "Claude Code OAuth token cannot access /models endpoint", - } -} - func (c *ClaudeClient) Guard(ctx context.Context, req *anthropic.MessageNewParams) (*AnthropicClient, map[string]string) { // Apply thinking transformation for Claude Code OAuth. Thinking can be expressed // either through the thinking union (enabled/adaptive/disabled) or through @@ -289,76 +277,6 @@ func (c *ClaudeClient) Client() *anthropic.Client { return c.AnthropicClient.Client() } -// ProbeChatEndpoint tests the messages endpoint. -func (c *ClaudeClient) Probe(ctx context.Context, model string) ProbeResult { - res, err := c.ProbeStream(ctx, model, "hi", ProbeModeStreaming) - if err != nil { - return ProbeResult{ - Success: false, - Message: err.Error(), - } - } - return *res -} - -// ProbeStream performs a streaming probe with configurable test mode for Claude Code OAuth. -// This uses the ClaudeClient's own MessagesNewStreaming method which properly applies: -// - Beta headers (only for probe operations) -// - Thinking field disabling -// - Session ID injection via Guard pattern -// -// Note: Beta headers are ONLY applied during probe, not during normal message passing. -func (c *ClaudeClient) ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - startTime := time.Now() - - // Build system message - systemMessages := []anthropic.TextBlockParam{ - { - Text: ClaudeCodeSystemHeader, - }, - { - Text: ClaudeCodeSystemBody, - }, - } - - messages := []anthropic.MessageParam{ - anthropic.NewUserMessage(anthropic.NewTextBlock(message)), - } - - params := &anthropic.MessageNewParams{ - Model: anthropic.Model(model), - MaxTokens: 1024, - System: systemMessages, - Messages: messages, - } - - // Disable thinking for Claude Code probe - params.Thinking = anthropic.ThinkingConfigParamUnion{ - OfDisabled: &anthropic.ThinkingConfigDisabledParam{}, - } - - // must use tool && stream - params.Tools = GetProbeToolsAnthropic() - params.ToolChoice = GetProbeToolChoiceAutoAnthropic() - - // For streaming and tool modes, use streaming - stream := c.MessagesNewStreaming(ctx, params) - defer stream.Close() - - var chunks []interface{} - for stream.Next() { - event := stream.Current() - chunks = append(chunks, event) - } - - if err := stream.Err(); err != nil { - return nil, err - } - - chunksJSON, _ := json.Marshal(chunks) - return ToProbeResult(string(chunksJSON), time.Since(startTime).Milliseconds(), c.AnthropicClient.provider.APIBase+"/v1/messages", true), nil -} - // stripBetaClearThinkingEdit removes any clear_thinking_20251015 context-management // edit from the request. The Anthropic API rejects this edit type when thinking is not // enabled or adaptive, so it must be dropped whenever thinking is disabled — otherwise diff --git a/internal/client/claude_client_test.go b/internal/client/claude_client_test.go index 04533d15a..2048035e6 100644 --- a/internal/client/claude_client_test.go +++ b/internal/client/claude_client_test.go @@ -94,7 +94,7 @@ func TestNewClaudeClient_StripV1FromBase(t *testing.T) { } // =================================================================== -// ListModels / ProbeModelsEndpoint +// ListModels // =================================================================== func TestClaudeClient_ListModels_ReturnsError(t *testing.T) { @@ -111,16 +111,6 @@ func TestClaudeClient_ListModels_ReturnsError(t *testing.T) { assert.Equal(t, provider.Name, modelsErr.Provider) } -func TestClaudeClient_ProbeModelsEndpoint_ReturnsUnsupported(t *testing.T) { - provider := newOAuthProvider() - c, err := NewClaudeClient(provider, "", typ.SessionID{Value: "s"}) - require.NoError(t, err) - - result := c.ProbeModelsEndpoint(context.Background()) - assert.False(t, result.Success) - assert.NotEmpty(t, result.ErrorMessage) -} - // =================================================================== // remapBetaToolNames // =================================================================== diff --git a/internal/client/claude_e2e_comparison_test.go b/internal/client/claude_e2e_comparison_test.go index 0bb6d85c2..da7217df1 100644 --- a/internal/client/claude_e2e_comparison_test.go +++ b/internal/client/claude_e2e_comparison_test.go @@ -497,11 +497,5 @@ func TestClaudeRealE2E_ModelsEndpoint_Rejection(t *testing.T) { assert.Error(t, err, "NEW should return error for ListModels") t.Log(" ✅ NEW correctly rejects /models with error") - // Also test ProbeModelsEndpoint - probeResult := newClient.ProbeModelsEndpoint(context.Background()) - t.Logf(" Probe Result: Success=%v, Error=%s", probeResult.Success, probeResult.ErrorMessage) - assert.False(t, probeResult.Success, "ProbeModelsEndpoint should fail") - t.Log(" ✅ NEW ProbeModelsEndpoint correctly returns failure") - t.Log("\n✅ Both implementations properly reject /models endpoint") } diff --git a/internal/client/codex_client.go b/internal/client/codex_client.go index f305bd9c5..a0e2befaa 100644 --- a/internal/client/codex_client.go +++ b/internal/client/codex_client.go @@ -2,11 +2,9 @@ package client import ( "context" - "encoding/json" "fmt" "net/http" "strings" - "time" "github.com/openai/openai-go/v3" "github.com/openai/openai-go/v3/option" @@ -552,87 +550,3 @@ func (c *CodexClient) SetRecordSink(sink *obs.Sink) { func (c *CodexClient) Client() *openai.Client { return c.OpenAIClient.Client() } - -// ProbeChatEndpoint tests the chat endpoint. -// For Codex, this delegates to the embedded OpenAIClient's probeResponsesEndpoint. -func (c *CodexClient) Probe(ctx context.Context, model string) ProbeResult { - startTime := time.Now() - r, err := c.ProbeStream(ctx, model, "hi", ProbeModeTool) - latencyMs := time.Since(startTime).Milliseconds() - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: err.Error(), - LatencyMs: latencyMs, - } - } - - return ProbeResult{ - Success: true, - Message: "Responses endpoint is accessible", - Content: r.Content, - LatencyMs: r.LatencyMs, - PromptTokens: r.PromptTokens, - CompletionTokens: r.CompletionTokens, - TotalTokens: r.TotalTokens, - } -} - -// ProbeStream performs a streaming probe with configurable test mode (public interface) -// Codex uses the Responses API with proper beta headers -func (c *CodexClient) ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - return c.probeChatStream(ctx, model, message, testMode) -} - -// probeChatStream performs a streaming probe using Codex's Responses API -// Codex requires special handling: -// - Uses Responses API (not Chat Completions) -// - Applies beta headers for responses API -// - Proper session ID injection -func (c *CodexClient) probeChatStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - return nil, fmt.Errorf("Codex do not support chat complement") -} - -func (c *CodexClient) ProbeResponsesStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - startTime := time.Now() - - // Build Responses API request - params := responses.ResponseNewParams{ - Model: model, - Instructions: param.NewOpt("work as `echo` if possible"), - Input: responses.ResponseNewParamsInputUnion{ - OfInputItemList: []responses.ResponseInputItemUnionParam{ - responses.ResponseInputItemParamOfMessage( - responses.ResponseInputMessageContentListParam{ - responses.ResponseInputContentParamOfInputText(message), - }, - responses.EasyInputMessageRoleUser, - ), - }, - }, - } - - // Add tools for tool mode - if testMode == ProbeModeTool { - params.Tools = GetProbeToolsResponses() - params.ToolChoice = responses.ResponseNewParamsToolChoiceUnion{ - OfToolChoiceMode: param.NewOpt(responses.ToolChoiceOptionsAuto), - } - } - - // Use ResponsesNewStreaming with proper beta headers - stream := c.ResponsesNewStreaming(ctx, params) - defer stream.Close() - - var chunks []interface{} - for stream.Next() { - chunks = append(chunks, stream.Current()) - } - - if err := stream.Err(); err != nil { - return nil, err - } - - chunksJSON, _ := json.Marshal(chunks) - return ToProbeResult(string(chunksJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase+"/responses", true), nil -} diff --git a/internal/client/google.go b/internal/client/google.go index 05cb5eafb..12a7cd091 100644 --- a/internal/client/google.go +++ b/internal/client/google.go @@ -2,12 +2,8 @@ package client import ( "context" - "encoding/json" - "fmt" "iter" "net/http" - "strings" - "time" "github.com/sirupsen/logrus" "google.golang.org/genai" @@ -147,183 +143,3 @@ func (c *GoogleClient) ListModels(ctx context.Context) ([]string, error) { Reason: "Google genai SDK does not support listing models via API", } } - -// ProbeChatEndpoint tests the chat endpoint with a minimal request -func (c *GoogleClient) Probe(ctx context.Context, model string) ProbeResult { - startTime := time.Now() - - // Create minimal content for probe - contents := []*genai.Content{ - { - Role: "user", - Parts: []*genai.Part{ - {Text: "hi"}, - }, - }, - } - - // Configure generation with minimal tokens - config := &genai.GenerateContentConfig{ - MaxOutputTokens: 1000, - } - - // Make request - resp, err := c.client.Models.GenerateContent(ctx, model, contents, config) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: err.Error(), - LatencyMs: latencyMs, - } - } - - // Extract response data - responseContent := "" - promptTokens := 0 - completionTokens := 0 - totalTokens := 0 - - if resp != nil && len(resp.Candidates) > 0 { - candidate := resp.Candidates[0] - if candidate.Content != nil { - for _, part := range candidate.Content.Parts { - if part.Text != "" { - responseContent += part.Text - } - } - } - if resp.UsageMetadata != nil { - promptTokens = int(resp.UsageMetadata.PromptTokenCount) - completionTokens = int(resp.UsageMetadata.CandidatesTokenCount) - totalTokens = int(resp.UsageMetadata.TotalTokenCount) - } - } - - if responseContent == "" { - responseContent = "" - } - - return ProbeResult{ - Success: true, - Message: "Chat endpoint is accessible", - Content: responseContent, - LatencyMs: latencyMs, - PromptTokens: promptTokens, - CompletionTokens: completionTokens, - TotalTokens: totalTokens, - } -} - -// probeStream performs a streaming probe with configurable test mode -// Note: Google genai SDK has limited tool support in probe context - -// ProbeStream performs a streaming probe with configurable test mode (public interface) -func (c *GoogleClient) ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - return c.probeStream(ctx, model, message, testMode) -} -func (c *GoogleClient) probeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - startTime := time.Now() - - // Create minimal content for probe - contents := []*genai.Content{ - { - Role: "user", - Parts: []*genai.Part{ - {Text: message}, - }, - }, - } - - // Configure generation - config := &genai.GenerateContentConfig{ - MaxOutputTokens: 1024, - } - - // For simple mode, use non-streaming request - if testMode == ProbeModeSimple { - resp, err := c.client.Models.GenerateContent(ctx, model, contents, config) - if err != nil { - return nil, err - } - - respJSON, _ := json.Marshal(resp) - return ToProbeResult(string(respJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase, false), nil - } - - // For streaming and tool modes, use streaming - stream := c.client.Models.GenerateContentStream(ctx, model, contents, config) - - var chunks []interface{} - for resp, err := range stream { - if err != nil { - return nil, err - } - chunks = append(chunks, resp) - } - - chunksJSON, _ := json.Marshal(chunks) - return ToProbeResult(string(chunksJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase, true), nil -} - -// ProbeModelsEndpoint tests the models list endpoint -func (c *GoogleClient) ProbeModelsEndpoint(ctx context.Context) ProbeResult { - // Google genai SDK doesn't provide a models list endpoint - // Return an error result indicating this is not supported - return ProbeResult{ - Success: false, - ErrorMessage: "Google genai SDK does not support listing models via API", - LatencyMs: 0, - } -} - -// ProbeOptionsEndpoint tests basic connectivity with an OPTIONS request -func (c *GoogleClient) ProbeOptionsEndpoint(ctx context.Context) ProbeResult { - startTime := time.Now() - - // Use the API base URL for OPTIONS request - optionsURL := c.provider.APIBase - if !strings.HasSuffix(optionsURL, "/") { - optionsURL += "/" - } - - req, err := http.NewRequestWithContext(ctx, "OPTIONS", optionsURL, nil) - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("Failed to create OPTIONS request: %v", err), - } - } - - // Set authentication header - req.Header.Set("x-goog-api-key", c.provider.GetAccessToken()) - - client := &http.Client{Timeout: 5 * time.Second} - resp, err := client.Do(req) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("OPTIONS request failed: %v", err), - LatencyMs: latencyMs, - } - } - defer resp.Body.Close() - - // Consider any 2xx status as success for OPTIONS - if resp.StatusCode >= 200 && resp.StatusCode < 300 { - return ProbeResult{ - Success: true, - Message: "OPTIONS request successful", - LatencyMs: latencyMs, - } - } - - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("OPTIONS request failed with status: %d", resp.StatusCode), - LatencyMs: latencyMs, - } -} diff --git a/internal/client/kimi_client.go b/internal/client/kimi_client.go index 7d26d6e7b..7c8519b80 100644 --- a/internal/client/kimi_client.go +++ b/internal/client/kimi_client.go @@ -278,38 +278,3 @@ func (c *KimiClient) ResponsesNewStreaming(ctx context.Context, req responses.Re func (c *KimiClient) ListModels(ctx context.Context) ([]string, error) { return c.OpenAIClient.ListModels(ctx) } - -// Probe tests the chat endpoint for Kimi provider. -func (c *KimiClient) Probe(ctx context.Context, model string) ProbeResult { - return c.OpenAIClient.Probe(ctx, model) -} - -// ProbeStream performs a streaming probe with Kimi model normalization. -func (c *KimiClient) ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - // Apply Kimi-specific model name normalization - normalizedModel := c.stripKimiPrefix(model) - return c.OpenAIClient.ProbeStream(ctx, normalizedModel, message, testMode) -} - -// ProbeResponsesStream performs a streaming Responses API probe with Kimi model normalization. -func (c *KimiClient) ProbeResponsesStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - return nil, &ErrKimiNotSupported{ - Operation: "Responses API", - Reason: "Kimi Code API does not support /responses endpoint", - } -} - -// ProbeChatEndpoint tests the chat endpoint with Kimi model normalization. -func (c *KimiClient) ProbeChatEndpoint(ctx context.Context, model string, opts ProbeEndpointOptions) (*ProbeResult, error) { - // Apply Kimi-specific model name normalization - normalizedModel := c.stripKimiPrefix(model) - return c.OpenAIClient.ProbeChatEndpoint(ctx, normalizedModel, opts) -} - -// ProbeResponsesEndpoint tests the Responses endpoint with Kimi model normalization. -func (c *KimiClient) ProbeResponsesEndpoint(ctx context.Context, model string, opts ProbeEndpointOptions) (*ProbeResult, error) { - return nil, &ErrKimiNotSupported{ - Operation: "Responses API", - Reason: "Kimi Code API does not support /responses endpoint", - } -} diff --git a/internal/client/openai.go b/internal/client/openai.go index eeabbc509..ae595f50f 100644 --- a/internal/client/openai.go +++ b/internal/client/openai.go @@ -1,8 +1,6 @@ package client import ( - "bufio" - "bytes" "context" "encoding/json" "fmt" @@ -13,7 +11,6 @@ import ( "github.com/openai/openai-go/v3" "github.com/openai/openai-go/v3/option" - "github.com/openai/openai-go/v3/packages/param" "github.com/openai/openai-go/v3/packages/ssestream" "github.com/openai/openai-go/v3/responses" "github.com/tingly-dev/tingly-box/ai" @@ -42,13 +39,6 @@ type OpenAIClientInterface interface { APIStyle() protocol.APIStyle SetRecordSink(sink *obs.Sink) - // Prober interface methods - Probe(ctx context.Context, model string) ProbeResult - ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) - ProbeResponsesStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) - ProbeChatEndpoint(ctx context.Context, model string, opts ProbeEndpointOptions) (*ProbeResult, error) - ProbeResponsesEndpoint(ctx context.Context, model string, opts ProbeEndpointOptions) (*ProbeResult, error) - // Client returns the underlying OpenAI SDK client (for advanced usage) Client() *openai.Client } @@ -330,448 +320,6 @@ func (c *OpenAIClient) ListModels(ctx context.Context) ([]string, error) { return models, nil } -// ProbeChatEndpoint tests the chat completions endpoint with a minimal request -func (c *OpenAIClient) Probe(ctx context.Context, model string) ProbeResult { - startTime := time.Now() - - // Check if this is a Codex OAuth provider - // Codex OAuth requires the Responses API, not Chat Completions - if c.provider.AuthType == typ.AuthTypeOAuth && - c.provider.OAuthDetail != nil && - c.provider.OAuthDetail.GetIssuer() == ai.IssuerCodex { - return c.probeCodexResponsesEndpoint(ctx, model) - } - - // Create chat completion request using OpenAI SDK - chatRequest := &openai.ChatCompletionNewParams{ - Model: openai.ChatModel(model), - Messages: []openai.ChatCompletionMessageParamUnion{ - openai.SystemMessage("work as `echo`"), - openai.UserMessage("hi"), - }, - } - - // Make request - resp, err := c.client.Chat.Completions.New(ctx, *chatRequest) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: err.Error(), - LatencyMs: latencyMs, - } - } - - // Extract response data - responseContent := "" - promptTokens := 0 - completionTokens := 0 - totalTokens := 0 - - if resp != nil { - if len(resp.Choices) > 0 { - responseContent = resp.Choices[0].Message.Content - } - if resp.Usage.PromptTokens != 0 { - promptTokens = int(resp.Usage.PromptTokens) - completionTokens = int(resp.Usage.CompletionTokens) - totalTokens = int(resp.Usage.TotalTokens) - } - } - - if responseContent == "" { - responseContent = "" - } - - return ProbeResult{ - Success: true, - Message: "Chat endpoint is accessible", - Content: responseContent, - LatencyMs: latencyMs, - PromptTokens: promptTokens, - CompletionTokens: completionTokens, - TotalTokens: totalTokens, - } -} - -// ProbeChatEndpoint explicitly probes the Chat Completions endpoint. -func (c *OpenAIClient) ProbeChatEndpoint(ctx context.Context, model string, opts ProbeEndpointOptions) (*ProbeResult, error) { - message := opts.Message - if message == "" { - message = "Hi" - } - mode := opts.Mode - if mode == "" { - mode = ProbeModeSimple - } - if opts.Stream { - mode = ProbeModeStreaming - } - return c.probeChatEndpoint(ctx, model, message, mode, opts.Stream) -} - -func (c *OpenAIClient) probeChatEndpoint(ctx context.Context, model, message string, testMode ProbeMode, forceStream bool) (*ProbeResult, error) { - startTime := time.Now() - messages := []openai.ChatCompletionMessageParamUnion{ - openai.SystemMessage("work as `echo` if possible"), - openai.UserMessage(message), - } - - params := openai.ChatCompletionNewParams{ - Model: model, - Messages: messages, - } - if testMode == ProbeModeTool { - params.Tools = GetProbeToolsOpenAI() - params.ToolChoice = openai.ChatCompletionToolChoiceOptionUnionParam{ - OfAuto: openai.Opt("auto"), - } - } - - if !forceStream && testMode == ProbeModeSimple { - resp, err := c.client.Chat.Completions.New(ctx, params) - if err != nil { - return nil, err - } - respJSON, _ := json.Marshal(resp) - return ToProbeResult(string(respJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase+"/chat/completions", false), nil - } - - stream := c.client.Chat.Completions.NewStreaming(ctx, params) - defer stream.Close() - - var chunks []interface{} - for stream.Next() { - chunks = append(chunks, stream.Current()) - } - if err := stream.Err(); err != nil { - return nil, err - } - chunksJSON, _ := json.Marshal(chunks) - return ToProbeResult(string(chunksJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase+"/chat/completions", true), nil -} - -// ProbeResponsesEndpoint explicitly probes the Responses API endpoint. -func (c *OpenAIClient) ProbeResponsesEndpoint(ctx context.Context, model string, opts ProbeEndpointOptions) (*ProbeResult, error) { - message := opts.Message - if message == "" { - message = "Hi" - } - mode := opts.Mode - if mode == "" { - mode = ProbeModeSimple - } - return c.probeResponsesEndpoint(ctx, model, message, mode, opts.Stream) -} - -func (c *OpenAIClient) probeResponsesEndpoint(ctx context.Context, model, message string, testMode ProbeMode, stream bool) (*ProbeResult, error) { - startTime := time.Now() - params := responses.ResponseNewParams{ - Model: model, - Instructions: param.NewOpt("work as `echo` if possible"), - Input: responses.ResponseNewParamsInputUnion{ - OfInputItemList: []responses.ResponseInputItemUnionParam{ - responses.ResponseInputItemParamOfMessage( - responses.ResponseInputMessageContentListParam{ - responses.ResponseInputContentParamOfInputText(message), - }, - responses.EasyInputMessageRoleUser, - ), - }, - }, - } - if testMode == ProbeModeTool { - params.Tools = GetProbeToolsResponses() - params.ToolChoice = responses.ResponseNewParamsToolChoiceUnion{ - OfToolChoiceMode: param.NewOpt(responses.ToolChoiceOptionsAuto), - } - } - - if !stream { - resp, err := c.ResponsesNew(ctx, params) - if err != nil { - return nil, err - } - respJSON, _ := json.Marshal(resp) - return ToProbeResult(string(respJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase+"/responses", false), nil - } - - streamResp := c.ResponsesNewStreaming(ctx, params) - defer streamResp.Close() - - var chunks []interface{} - for streamResp.Next() { - chunks = append(chunks, streamResp.Current()) - } - if err := streamResp.Err(); err != nil { - return nil, err - } - chunksJSON, _ := json.Marshal(chunks) - return ToProbeResult(string(chunksJSON), time.Since(startTime).Milliseconds(), c.provider.APIBase+"/responses", true), nil -} - -// probeStream performs a streaming probe with configurable test mode - -// ProbeStream performs a streaming probe with configurable test mode (public interface) -func (c *OpenAIClient) ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - // Keep legacy Codex behavior. Non-Codex legacy probes target Chat explicitly. - if c.provider.AuthType == typ.AuthTypeOAuth && - c.provider.OAuthDetail != nil && - c.provider.OAuthDetail.GetIssuer() == ai.IssuerCodex { - return c.ProbeResponsesStream(ctx, model, message, testMode) - } - return c.probeChatEndpoint(ctx, model, message, testMode, testMode != ProbeModeSimple) -} - -// ProbeResponsesStream performs a streaming probe using Responses API (for Codex) -func (c *OpenAIClient) ProbeResponsesStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) { - return c.probeResponsesEndpoint(ctx, model, message, testMode, true) -} - -// ProbeModelsEndpoint tests the models list endpoint -func (c *OpenAIClient) ProbeModelsEndpoint(ctx context.Context) ProbeResult { - startTime := time.Now() - - // Make request to models endpoint - resp, err := c.client.Models.List(ctx) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: err.Error(), - LatencyMs: latencyMs, - } - } - - modelsCount := 0 - if resp != nil { - modelsCount = len(resp.Data) - } - - if modelsCount == 0 { - return ProbeResult{ - Success: false, - ErrorMessage: "No models available from provider", - LatencyMs: latencyMs, - } - } - - return ProbeResult{ - Success: true, - Message: "Models endpoint is accessible", - LatencyMs: latencyMs, - ModelsCount: modelsCount, - } -} - -// ProbeOptionsEndpoint tests basic connectivity with an OPTIONS request -func (c *OpenAIClient) ProbeOptionsEndpoint(ctx context.Context) ProbeResult { - startTime := time.Now() - - // Use the API base URL for OPTIONS request - optionsURL := c.provider.APIBase - - req, err := http.NewRequestWithContext(ctx, "OPTIONS", optionsURL, nil) - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("Failed to create OPTIONS request: %v", err), - } - } - - // Set authentication header - req.Header.Set("Authorization", "Bearer "+c.provider.GetAccessToken()) - - client := &http.Client{Timeout: 5 * time.Second} - resp, err := client.Do(req) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("OPTIONS request failed: %v", err), - LatencyMs: latencyMs, - } - } - defer resp.Body.Close() - - // Consider any 2xx status as success for OPTIONS - if resp.StatusCode >= 200 && resp.StatusCode < 300 { - return ProbeResult{ - Success: true, - Message: "OPTIONS request successful", - LatencyMs: latencyMs, - } - } - - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("OPTIONS request failed with status: %d", resp.StatusCode), - LatencyMs: latencyMs, - } -} - -// probeResponsesEndpoint tests the Responses API (for Codex OAuth providers) -func (c *OpenAIClient) probeCodexResponsesEndpoint(ctx context.Context, model string) ProbeResult { - startTime := time.Now() - - // Build ChatGPT backend API request format - inputItems := []map[string]interface{}{ - { - "type": "message", - "role": "user", - "content": []map[string]string{ - {"type": "input_text", "text": "Hi"}, - }, - }, - } - - reqBody := map[string]interface{}{ - "model": model, - "instructions": "work as `echo`", - "input": inputItems, - "tools": []interface{}{}, - "tool_choice": "auto", - "stream": true, - "store": false, - "include": []string{}, - } - - bodyBytes, err := json.Marshal(reqBody) - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("Failed to marshal request: %v", err), - } - } - - // Create HTTP request - reqURL := c.provider.APIBase + "/responses" - req, err := http.NewRequestWithContext(ctx, "POST", reqURL, bytes.NewReader(bodyBytes)) - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("Failed to create request: %v", err), - } - } - - // Set required headers - req.Header.Set("Content-Type", "application/json") - req.Header.Set("Authorization", "Bearer "+c.provider.GetAccessToken()) - req.Header.Set("OpenAI-Beta", "responses=experimental") - req.Header.Set("originator", "tingly-box") - - // Add ChatGPT-Account-ID header if available - if accountID := c.provider.OAuthDetail.GetExtraFieldString("account_id"); accountID != "" { - req.Header.Set("ChatGPT-Account-ID", accountID) - } - - // Make the request - resp, err := c.HttpClient.Do(req) - latencyMs := time.Since(startTime).Milliseconds() - - if err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: err.Error(), - LatencyMs: latencyMs, - } - } - defer resp.Body.Close() - - if resp.StatusCode != http.StatusOK { - respBody, _ := io.ReadAll(resp.Body) - return ProbeResult{ - Success: false, - ErrorMessage: string(respBody), - LatencyMs: latencyMs, - } - } - - // Read streaming response and collect all chunks - var responseContent string - tokenUsage := ProbeUsage{} - - scanner := bufio.NewScanner(resp.Body) - for scanner.Scan() { - line := scanner.Text() - - // Skip empty lines and non-data lines - if !strings.HasPrefix(line, "data: ") { - continue - } - - // Extract JSON data from SSE format - jsonData := strings.TrimPrefix(line, "data: ") - - // Check for stream end - if jsonData == "[DONE]" { - break - } - - // Parse SSE chunk - var chunk struct { - Output []struct { - Type string `json:"type"` - Content []struct { - Type string `json:"type"` - Text string `json:"text"` - } `json:"content"` - } `json:"output"` - Usage struct { - InputTokens int `json:"input_tokens"` - OutputTokens int `json:"output_tokens"` - } `json:"usage"` - } - - if err := json.Unmarshal([]byte(jsonData), &chunk); err != nil { - continue - } - - // Extract response content from output - for _, item := range chunk.Output { - if item.Type == "message" { - for _, content := range item.Content { - if content.Type == "output_text" { - responseContent += content.Text - } - } - } - } - - // Extract usage from the last chunk - if chunk.Usage.InputTokens > 0 { - tokenUsage.PromptTokens = chunk.Usage.InputTokens - tokenUsage.CompletionTokens = chunk.Usage.OutputTokens - tokenUsage.TotalTokens = chunk.Usage.InputTokens + chunk.Usage.OutputTokens - } - } - - if err := scanner.Err(); err != nil { - return ProbeResult{ - Success: false, - ErrorMessage: fmt.Sprintf("Failed to read streaming response: %v", err), - LatencyMs: latencyMs, - } - } - - if responseContent == "" { - responseContent = "" - } - - return ProbeResult{ - Success: true, - Message: "Responses endpoint is accessible", - Content: responseContent, - LatencyMs: latencyMs, - PromptTokens: tokenUsage.PromptTokens, - CompletionTokens: tokenUsage.CompletionTokens, - TotalTokens: tokenUsage.TotalTokens, - } -} - // isCodexProvider checks if the current provider is a Codex OAuth provider // Codex OAuth providers require special handling for image generation func (c *OpenAIClient) isCodexProvider() bool { diff --git a/internal/client/openai_probe_test.go b/internal/client/openai_probe_test.go index 33dc9f0be..bec2df995 100644 --- a/internal/client/openai_probe_test.go +++ b/internal/client/openai_probe_test.go @@ -1,133 +1,12 @@ package client import ( - "context" "testing" "github.com/tingly-dev/tingly-box/internal/protocol" "github.com/tingly-dev/tingly-box/internal/typ" ) -// TestOpenAIClient_ProbeChatEndpoint tests the ProbeChatEndpoint method -func TestOpenAIClient_ProbeChatEndpoint(t *testing.T) { - tests := []struct { - name string - provider *typ.Provider - model string - wantErr bool - }{ - { - name: "skip live test - requires valid API key", - // This test would require a real API key to run - // In a real scenario, you might use environment variables or test fixtures - provider: &typ.Provider{ - Name: "test-openai", - APIBase: "https://api.openai.com/v1", - APIStyle: protocol.APIStyleOpenAI, - Token: "sk-test-key", - }, - model: "gpt-3.5-turbo", - wantErr: true, // Will fail with invalid key - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - client, err := NewOpenAIClient(tt.provider, tt.model, typ.SessionID{}) - if err != nil { - t.Fatalf("NewOpenAIClient() error = %v", err) - } - - result := client.Probe(context.Background(), tt.model) - - if !tt.wantErr && !result.Success { - t.Errorf("ProbeChatEndpoint() failed = %v", result.ErrorMessage) - } - if tt.wantErr && result.Success { - t.Errorf("ProbeChatEndpoint() expected error but succeeded") - } - }) - } -} - -// TestOpenAIClient_ProbeModelsEndpoint tests the ProbeModelsEndpoint method -func TestOpenAIClient_ProbeModelsEndpoint(t *testing.T) { - tests := []struct { - name string - provider *typ.Provider - model string - wantErr bool - }{ - { - name: "skip live test - requires valid API key", - provider: &typ.Provider{ - Name: "test-openai", - APIBase: "https://api.openai.com/v1", - APIStyle: protocol.APIStyleOpenAI, - Token: "sk-test-key", - }, - model: "gpt-3.5-turbo", - wantErr: true, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - client, err := NewOpenAIClient(tt.provider, tt.model, typ.SessionID{}) - if err != nil { - t.Fatalf("NewOpenAIClient() error = %v", err) - } - - result := client.ProbeModelsEndpoint(context.Background()) - - if !tt.wantErr && !result.Success { - t.Errorf("ProbeModelsEndpoint() failed = %v", result.ErrorMessage) - } - if tt.wantErr && result.Success { - t.Errorf("ProbeModelsEndpoint() expected error but succeeded") - } - }) - } -} - -// TestOpenAIClient_ProbeOptionsEndpoint tests the ProbeOptionsEndpoint method -func TestOpenAIClient_ProbeOptionsEndpoint(t *testing.T) { - tests := []struct { - name string - provider *typ.Provider - model string - wantErr bool - }{ - { - name: "skip live test - requires valid API key", - provider: &typ.Provider{ - Name: "test-openai", - APIBase: "https://api.openai.com/v1", - APIStyle: protocol.APIStyleOpenAI, - Token: "sk-test-key", - }, - model: "gpt-3.5-turbo", - wantErr: true, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - client, err := NewOpenAIClient(tt.provider, tt.model, typ.SessionID{}) - if err != nil { - t.Fatalf("NewOpenAIClient() error = %v", err) - } - - result := client.ProbeOptionsEndpoint(context.Background()) - - if !tt.wantErr && !result.Success { - t.Errorf("ProbeOptionsEndpoint() failed = %v", result.ErrorMessage) - } - // OPTIONS might succeed even with invalid key for some providers - }) - } -} - // TestOpenAIClient_Timeout tests that timeout is properly configured from provider func TestOpenAIClient_Timeout(t *testing.T) { tests := []struct { @@ -173,4 +52,3 @@ func TestOpenAIClient_Timeout(t *testing.T) { }) } } - diff --git a/internal/client/probe.go b/internal/client/probe.go deleted file mode 100644 index 192fbd361..000000000 --- a/internal/client/probe.go +++ /dev/null @@ -1,90 +0,0 @@ -package client - -import ( - "context" -) - -// ProbeResult represents the result of a probe operation (for both simple and streaming) -type ProbeResult struct { - // Basic fields - Success bool `json:"success"` - Message string `json:"message,omitempty"` - Content string `json:"content,omitempty"` - LatencyMs int64 `json:"latency_ms"` - ModelsCount int `json:"models_count,omitempty"` - ErrorMessage string `json:"error_message,omitempty"` - - // Streaming mode indicator - Stream bool `json:"stream,omitempty"` - - // Token usage - PromptTokens int `json:"prompt_tokens,omitempty"` - CompletionTokens int `json:"completion_tokens,omitempty"` - TotalTokens int `json:"total_tokens,omitempty"` - - // Tool calls (for tool mode) - ToolCalls []ProbeToolCall `json:"tool_calls,omitempty"` - - // Request URL (for debugging) - RequestURL string `json:"request_url,omitempty"` -} - -// ProbeToolCall represents a tool call in probe response -type ProbeToolCall struct { - ID string `json:"id"` - Name string `json:"name"` - Input map[string]interface{} `json:"input"` -} - -// ProbeUsage represents token usage from a probe operation -type ProbeUsage struct { - PromptTokens int - CompletionTokens int - TotalTokens int -} - -// ProbeMode defines the test mode for probeStream -type ProbeMode string - -const ( - ProbeModeSimple ProbeMode = "simple" - ProbeModeStreaming ProbeMode = "streaming" - ProbeModeTool ProbeMode = "tool" -) - -// Prober defines the interface for client probe capabilities -type Prober interface { - // Probe tests the chat/messages endpoint with a minimal request - // Returns a ProbeResult with success status, latency, and any response content - Probe(ctx context.Context, model string) ProbeResult - - // ProbeStream performs a streaming probe with configurable test mode. - // Deprecated for endpoint capability routing: use endpoint-explicit probe methods - // such as OpenAIClient.ProbeChatEndpoint and OpenAIClient.ProbeResponsesEndpoint. - ProbeStream(ctx context.Context, model, message string, testMode ProbeMode) (*ProbeResult, error) -} - -// ToProbeResult creates a ProbeResult with basic fields -func ToProbeResult(content string, latencyMs int64, requestURL string, isStreaming bool) *ProbeResult { - return &ProbeResult{ - Content: content, - LatencyMs: latencyMs, - RequestURL: requestURL, - Stream: isStreaming, - } -} - -// ProbeEndpointType identifies which OpenAI-compatible endpoint a probe must hit. -type ProbeEndpointType string - -const ( - ProbeEndpointChat ProbeEndpointType = "chat" - ProbeEndpointResponses ProbeEndpointType = "responses" -) - -// ProbeEndpointOptions controls endpoint-explicit probing. -type ProbeEndpointOptions struct { - Message string - Stream bool - Mode ProbeMode -} diff --git a/internal/probe/e2e.go b/internal/probe/e2e.go index 43340e7b6..9fcb88c4f 100644 --- a/internal/probe/e2e.go +++ b/internal/probe/e2e.go @@ -271,48 +271,43 @@ func (e *E2EService) resolveSmartRoutingForProbe(rule *typ.Rule) (*loadbalance.S return selectedService, nil } -// getClientForProvider returns a Prober for the given provider via the client pool. -func (e *E2EService) getClientForProvider(provider *typ.Provider, model string) (client.Prober, error) { +// ProbeProviderWithSDK runs an SDK probe by dispatching a minimal request +// through the provider's real-traffic client methods. Public because the +// server's provider onboarding path (testProviderConnectivity) reuses it. +func (e *E2EService) ProbeProviderWithSDK(ctx context.Context, provider *typ.Provider, model, message string, testMode E2EMode) (*E2EData, error) { + mode := testMode + switch provider.APIStyle { - case protocol.APIStyleAnthropic: - c := e.clientPool.GetAnthropicClient(context.Background(), provider, model) - if c == nil { - return nil, fmt.Errorf("failed to get Anthropic client for provider: %s", provider.Name) - } - return c, nil case protocol.APIStyleOpenAI: - c := e.clientPool.GetOpenAIClient(context.Background(), provider, model) - if c == nil { + oc := e.clientPool.GetOpenAIClient(ctx, provider, model) + if oc == nil { return nil, fmt.Errorf("failed to get OpenAI client for provider: %s", provider.Name) } - return c, nil + // Codex OAuth providers only speak the Responses API. + if isCodexOAuth(provider) { + return probeOpenAIResponses(ctx, oc, model, message, mode) + } + return probeOpenAIChat(ctx, oc, model, message, mode) + + case protocol.APIStyleAnthropic: + ac := e.clientPool.GetAnthropicClient(ctx, provider, model) + if ac == nil { + return nil, fmt.Errorf("failed to get Anthropic client for provider: %s", provider.Name) + } + return probeAnthropicMessages(ctx, ac, model, message, mode) + case protocol.APIStyleGoogle: - c := e.clientPool.GetGoogleClient(context.Background(), provider, model) - if c == nil { + gc := e.clientPool.GetGoogleClient(ctx, provider, model) + if gc == nil { return nil, fmt.Errorf("failed to get Google client for provider: %s", provider.Name) } - return c, nil + return probeGoogleGenerate(ctx, gc, model, message, mode) + default: return nil, fmt.Errorf("unsupported API style: %s", provider.APIStyle) } } -// ProbeProviderWithSDK runs a non-streaming SDK probe. Public because the -// server's provider onboarding path (testProviderConnectivity) reuses it. -func (e *E2EService) ProbeProviderWithSDK(ctx context.Context, provider *typ.Provider, model, message string, testMode E2EMode) (*E2EData, error) { - prober, err := e.getClientForProvider(provider, model) - if err != nil { - return nil, err - } - clientMode := client.ProbeMode(testMode) - return prober.ProbeStream(ctx, model, message, clientMode) -} - func (e *E2EService) probeProviderStream(ctx context.Context, provider *typ.Provider, model, message string, testMode E2EMode) (*E2EData, error) { - prober, err := e.getClientForProvider(provider, model) - if err != nil { - return nil, err - } - clientMode := client.ProbeMode(testMode) - return prober.ProbeStream(ctx, model, message, clientMode) + return e.ProbeProviderWithSDK(ctx, provider, model, message, testMode) } diff --git a/internal/probe/lightweight.go b/internal/probe/lightweight.go index b5e24680e..3a40e38f4 100644 --- a/internal/probe/lightweight.go +++ b/internal/probe/lightweight.go @@ -98,39 +98,14 @@ type modelsReport struct { func (l *LightweightService) probeOptionsEndpoint(ctx context.Context, provider *typ.Provider) endpointReport { startTime := time.Now() - var result client.ProbeResult - switch provider.APIStyle { - case protocol.APIStyleOpenAI: - c := l.pool.GetOpenAIClient(context.Background(), provider, "") - if c == nil { - return endpointReport{false, "Failed to create OpenAI client", 0} - } - openaiClient, ok := c.(*client.OpenAIClient) - if !ok { - return endpointReport{false, "OPTIONS probe not implemented for this client type", 0} - } - result = openaiClient.ProbeOptionsEndpoint(ctx) - case protocol.APIStyleAnthropic: - c := l.pool.GetAnthropicClient(context.Background(), provider, "") - if c == nil { - return endpointReport{false, "Failed to create Anthropic client", 0} - } - anthropicClient, ok := c.(*client.AnthropicClient) - if !ok { - return endpointReport{false, "OPTIONS probe not implemented for this client type", 0} - } - result = anthropicClient.ProbeOptionsEndpoint(ctx) - case protocol.APIStyleGoogle: - c := l.pool.GetGoogleClient(context.Background(), provider, "") - if c == nil { - return endpointReport{false, "Failed to create Google client", 0} - } - result = c.ProbeOptionsEndpoint(ctx) + case protocol.APIStyleOpenAI, protocol.APIStyleAnthropic, protocol.APIStyleGoogle: + // supported below default: return endpointReport{false, fmt.Sprintf("Unsupported API style: %s", provider.APIStyle), 0} } + result := probeOptions(ctx, provider) responseTime := time.Since(startTime).Milliseconds() if result.Success { return endpointReport{true, "OPTIONS request successful", responseTime} @@ -210,11 +185,7 @@ func (l *LightweightService) probeChatEndpoint(ctx context.Context, provider *ty probeCtx, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() - result, err := c.ProbeChatEndpoint(probeCtx, "gpt-3.5-turbo", client.ProbeEndpointOptions{ - Message: "Hi", - Stream: false, - Mode: client.ProbeModeSimple, - }) + result, err := probeOpenAIChat(probeCtx, c, "gpt-3.5-turbo", "Hi", E2EModeSimple) responseTime := time.Since(startTime).Milliseconds() if err != nil { @@ -237,11 +208,7 @@ func (l *LightweightService) probeResponsesEndpoint(ctx context.Context, provide probeCtx, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() - result, err := c.ProbeResponsesEndpoint(probeCtx, "gpt-4o", client.ProbeEndpointOptions{ - Message: "Hi", - Stream: false, - Mode: client.ProbeModeSimple, - }) + result, err := probeOpenAIResponses(probeCtx, c, "gpt-4o", "Hi", E2EModeSimple) responseTime := time.Since(startTime).Milliseconds() if err != nil { diff --git a/internal/client/probe_tools.go b/internal/probe/probetools.go similarity index 79% rename from internal/client/probe_tools.go rename to internal/probe/probetools.go index cc6c7b036..cb443860d 100644 --- a/internal/client/probe_tools.go +++ b/internal/probe/probetools.go @@ -1,4 +1,4 @@ -package client +package probe import ( "github.com/anthropics/anthropic-sdk-go" @@ -8,9 +8,9 @@ import ( "github.com/openai/openai-go/v3/shared" ) -// GetProbeToolsAnthropic returns predefined tools in Anthropic format for probe testing -// Uses bash tool to execute simple file system operations -func GetProbeToolsAnthropic() []anthropic.ToolUnionParam { +// getProbeToolsAnthropic returns predefined tools in Anthropic format for probe +// testing. Uses a bash tool to execute simple file system operations. +func getProbeToolsAnthropic() []anthropic.ToolUnionParam { return []anthropic.ToolUnionParam{ { OfTool: &anthropic.ToolParam{ @@ -44,9 +44,9 @@ func GetProbeToolsAnthropic() []anthropic.ToolUnionParam { } } -// GetProbeToolsOpenAI returns predefined tools in OpenAI format for probe testing -// Uses bash tool to execute simple file system operations -func GetProbeToolsOpenAI() []openai.ChatCompletionToolUnionParam { +// getProbeToolsOpenAI returns predefined tools in OpenAI format for probe +// testing. Uses a bash tool to execute simple file system operations. +func getProbeToolsOpenAI() []openai.ChatCompletionToolUnionParam { return []openai.ChatCompletionToolUnionParam{ openai.ChatCompletionFunctionTool(shared.FunctionDefinitionParam{ Name: "bash", @@ -80,9 +80,9 @@ func GetProbeToolsOpenAI() []openai.ChatCompletionToolUnionParam { } } -// GetProbeToolsResponses returns predefined tools in Responses API format for probe testing -// Uses bash tool to execute simple file system operations -func GetProbeToolsResponses() []responses.ToolUnionParam { +// getProbeToolsResponses returns predefined tools in Responses API format for +// probe testing. Uses a bash tool to execute simple file system operations. +func getProbeToolsResponses() []responses.ToolUnionParam { return []responses.ToolUnionParam{ responses.ToolParamOfFunction( "bash", @@ -116,8 +116,8 @@ func GetProbeToolsResponses() []responses.ToolUnionParam { } } -// GetProbeToolChoiceAutoAnthropic returns auto tool choice for testing -func GetProbeToolChoiceAutoAnthropic() anthropic.ToolChoiceUnionParam { +// getProbeToolChoiceAutoAnthropic returns auto tool choice for testing. +func getProbeToolChoiceAutoAnthropic() anthropic.ToolChoiceUnionParam { return anthropic.ToolChoiceUnionParam{ OfAuto: &anthropic.ToolChoiceAutoParam{}, } diff --git a/internal/probe/result.go b/internal/probe/result.go new file mode 100644 index 000000000..de9fcb56c --- /dev/null +++ b/internal/probe/result.go @@ -0,0 +1,46 @@ +package probe + +// ProbeResult is the canonical SDK-level probe result, shared by the E2E and +// lightweight probe strategies. It doubles as the JSON payload returned by the +// probe HTTP endpoints (exposed under the E2EData alias). +type ProbeResult struct { + // Basic fields + Success bool `json:"success"` + Message string `json:"message,omitempty"` + Content string `json:"content,omitempty"` + LatencyMs int64 `json:"latency_ms"` + ModelsCount int `json:"models_count,omitempty"` + ErrorMessage string `json:"error_message,omitempty"` + + // Streaming mode indicator + Stream bool `json:"stream,omitempty"` + + // Token usage + PromptTokens int `json:"prompt_tokens,omitempty"` + CompletionTokens int `json:"completion_tokens,omitempty"` + TotalTokens int `json:"total_tokens,omitempty"` + + // Tool calls (for tool mode) + ToolCalls []ProbeToolCall `json:"tool_calls,omitempty"` + + // Request URL (for debugging) + RequestURL string `json:"request_url,omitempty"` +} + +// ProbeToolCall represents a tool call in a probe response. +type ProbeToolCall struct { + ID string `json:"id"` + Name string `json:"name"` + Input map[string]interface{} `json:"input"` +} + +// toProbeResult builds a ProbeResult carrying the raw (JSON-marshaled) +// upstream response for a successful probe. +func toProbeResult(content string, latencyMs int64, requestURL string, isStreaming bool) *ProbeResult { + return &ProbeResult{ + Content: content, + LatencyMs: latencyMs, + RequestURL: requestURL, + Stream: isStreaming, + } +} diff --git a/internal/probe/sdkprobe.go b/internal/probe/sdkprobe.go new file mode 100644 index 000000000..413907039 --- /dev/null +++ b/internal/probe/sdkprobe.go @@ -0,0 +1,256 @@ +package probe + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "strings" + "time" + + "github.com/anthropics/anthropic-sdk-go" + "github.com/openai/openai-go/v3" + "github.com/openai/openai-go/v3/packages/param" + "github.com/openai/openai-go/v3/responses" + "google.golang.org/genai" + + "github.com/tingly-dev/tingly-box/ai" + "github.com/tingly-dev/tingly-box/internal/client" + "github.com/tingly-dev/tingly-box/internal/protocol" + "github.com/tingly-dev/tingly-box/internal/typ" +) + +// probeEchoInstruction is the system/instruction prompt used by SDK probes to +// keep the upstream response minimal. +const probeEchoInstruction = "work as `echo` if possible" + +// The SDK probe helpers below dispatch a minimal request through each client's +// real-traffic methods (ChatCompletionsNew, ResponsesNew, MessagesNew, +// GenerateContent). Routing probes through the same methods as production +// traffic means provider-specific quirks — Kimi model-name normalization, +// Codex Responses handling — apply identically and cannot drift from the real +// path. The client package therefore no longer owns any probe-specific code. + +// probeOpenAIChat builds and dispatches a minimal Chat Completions probe. +func probeOpenAIChat(ctx context.Context, oc client.OpenAIClientInterface, model, message string, mode E2EMode) (*ProbeResult, error) { + start := time.Now() + params := openai.ChatCompletionNewParams{ + Model: model, + Messages: []openai.ChatCompletionMessageParamUnion{ + openai.SystemMessage(probeEchoInstruction), + openai.UserMessage(message), + }, + } + if mode == E2EModeTool { + params.Tools = getProbeToolsOpenAI() + params.ToolChoice = openai.ChatCompletionToolChoiceOptionUnionParam{OfAuto: openai.Opt("auto")} + } + + url := oc.GetProvider().APIBase + "/chat/completions" + if mode == E2EModeSimple { + resp, err := oc.ChatCompletionsNew(ctx, params) + if err != nil { + return nil, err + } + b, _ := json.Marshal(resp) + return toProbeResult(string(b), time.Since(start).Milliseconds(), url, false), nil + } + + stream := oc.ChatCompletionsNewStreaming(ctx, params) + if stream == nil { + return nil, fmt.Errorf("chat streaming not supported by provider") + } + defer stream.Close() + var chunks []interface{} + for stream.Next() { + chunks = append(chunks, stream.Current()) + } + if err := stream.Err(); err != nil { + return nil, err + } + b, _ := json.Marshal(chunks) + return toProbeResult(string(b), time.Since(start).Milliseconds(), url, true), nil +} + +// probeOpenAIResponses builds and dispatches a minimal Responses API probe. +func probeOpenAIResponses(ctx context.Context, oc client.OpenAIClientInterface, model, message string, mode E2EMode) (*ProbeResult, error) { + start := time.Now() + params := responses.ResponseNewParams{ + Model: model, + Instructions: param.NewOpt(probeEchoInstruction), + Input: responses.ResponseNewParamsInputUnion{ + OfInputItemList: []responses.ResponseInputItemUnionParam{ + responses.ResponseInputItemParamOfMessage( + responses.ResponseInputMessageContentListParam{ + responses.ResponseInputContentParamOfInputText(message), + }, + responses.EasyInputMessageRoleUser, + ), + }, + }, + } + if mode == E2EModeTool { + params.Tools = getProbeToolsResponses() + params.ToolChoice = responses.ResponseNewParamsToolChoiceUnion{ + OfToolChoiceMode: param.NewOpt(responses.ToolChoiceOptionsAuto), + } + } + + url := oc.GetProvider().APIBase + "/responses" + if mode == E2EModeSimple { + resp, err := oc.ResponsesNew(ctx, params) + if err != nil { + return nil, err + } + b, _ := json.Marshal(resp) + return toProbeResult(string(b), time.Since(start).Milliseconds(), url, false), nil + } + + stream := oc.ResponsesNewStreaming(ctx, params) + if stream == nil { + return nil, fmt.Errorf("responses streaming not supported by provider") + } + defer stream.Close() + var chunks []interface{} + for stream.Next() { + chunks = append(chunks, stream.Current()) + } + if err := stream.Err(); err != nil { + return nil, err + } + b, _ := json.Marshal(chunks) + return toProbeResult(string(b), time.Since(start).Milliseconds(), url, true), nil +} + +// probeAnthropicMessages builds and dispatches a minimal Messages probe. +func probeAnthropicMessages(ctx context.Context, ac client.AnthropicClientInterface, model, message string, mode E2EMode) (*ProbeResult, error) { + start := time.Now() + provider := ac.GetProvider() + + system := []anthropic.TextBlockParam{{Text: probeEchoInstruction}} + if provider.AuthType == typ.AuthTypeOAuth && provider.OAuthDetail != nil && + provider.OAuthDetail.GetIssuer() == ai.IssuerClaudeCode { + system = append([]anthropic.TextBlockParam{{Text: client.ClaudeCodeSystemHeader}}, system...) + } + + params := &anthropic.MessageNewParams{ + Model: anthropic.Model(model), + MaxTokens: 1024, + System: system, + Messages: []anthropic.MessageParam{ + anthropic.NewUserMessage(anthropic.NewTextBlock(message)), + }, + } + if mode == E2EModeTool { + params.Tools = getProbeToolsAnthropic() + params.ToolChoice = getProbeToolChoiceAutoAnthropic() + } + + url := provider.APIBase + "/v1/messages" + if mode == E2EModeSimple { + resp, err := ac.MessagesNew(ctx, params) + if err != nil { + return nil, err + } + b, _ := json.Marshal(resp) + return toProbeResult(string(b), time.Since(start).Milliseconds(), url, false), nil + } + + stream := ac.MessagesNewStreaming(ctx, params) + if stream == nil { + return nil, fmt.Errorf("messages streaming not supported by provider") + } + defer stream.Close() + var chunks []interface{} + for stream.Next() { + chunks = append(chunks, stream.Current()) + } + if err := stream.Err(); err != nil { + return nil, err + } + b, _ := json.Marshal(chunks) + return toProbeResult(string(b), time.Since(start).Milliseconds(), url, true), nil +} + +// probeGoogleGenerate builds and dispatches a minimal GenerateContent probe. +func probeGoogleGenerate(ctx context.Context, gc *client.GoogleClient, model, message string, mode E2EMode) (*ProbeResult, error) { + start := time.Now() + contents := []*genai.Content{ + {Role: "user", Parts: []*genai.Part{{Text: message}}}, + } + config := &genai.GenerateContentConfig{MaxOutputTokens: 1024} + url := gc.GetProvider().APIBase + + if mode == E2EModeSimple { + resp, err := gc.GenerateContent(ctx, model, contents, config) + if err != nil { + return nil, err + } + b, _ := json.Marshal(resp) + return toProbeResult(string(b), time.Since(start).Milliseconds(), url, false), nil + } + + var chunks []interface{} + for resp, err := range gc.GenerateContentStream(ctx, model, contents, config) { + if err != nil { + return nil, err + } + chunks = append(chunks, resp) + } + b, _ := json.Marshal(chunks) + return toProbeResult(string(b), time.Since(start).Milliseconds(), url, true), nil +} + +// probeOptions issues a bare OPTIONS request to the provider base URL with the +// auth headers appropriate for its API style. Used by the lightweight probe; +// results are advisory. +func probeOptions(ctx context.Context, provider *typ.Provider) ProbeResult { + start := time.Now() + + url := provider.APIBase + header := http.Header{} + switch provider.APIStyle { + case protocol.APIStyleAnthropic: + apiBase := strings.TrimSuffix(provider.APIBase, "/") + if !strings.Contains(apiBase, "/v1") { + apiBase += "/v1" + } + url = apiBase + header.Set("x-api-key", provider.GetAccessToken()) + header.Set("anthropic-version", "2023-06-01") + case protocol.APIStyleGoogle: + if !strings.HasSuffix(url, "/") { + url += "/" + } + header.Set("x-goog-api-key", provider.GetAccessToken()) + default: + header.Set("Authorization", "Bearer "+provider.GetAccessToken()) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodOptions, url, nil) + if err != nil { + return ProbeResult{Success: false, ErrorMessage: fmt.Sprintf("Failed to create OPTIONS request: %v", err)} + } + req.Header = header + + httpClient := &http.Client{Timeout: 5 * time.Second} + resp, err := httpClient.Do(req) + latencyMs := time.Since(start).Milliseconds() + if err != nil { + return ProbeResult{Success: false, ErrorMessage: fmt.Sprintf("OPTIONS request failed: %v", err), LatencyMs: latencyMs} + } + defer resp.Body.Close() + + if resp.StatusCode >= 200 && resp.StatusCode < 300 { + return ProbeResult{Success: true, Message: "OPTIONS request successful", LatencyMs: latencyMs} + } + return ProbeResult{Success: false, ErrorMessage: fmt.Sprintf("OPTIONS request failed with status: %d", resp.StatusCode), LatencyMs: latencyMs} +} + +// isCodexOAuth reports whether the provider is a Codex OAuth provider, which +// only speaks the Responses API. +func isCodexOAuth(provider *typ.Provider) bool { + return provider.AuthType == typ.AuthTypeOAuth && + provider.OAuthDetail != nil && + provider.OAuthDetail.GetIssuer() == ai.IssuerCodex +} diff --git a/internal/probe/types.go b/internal/probe/types.go index 8d441f3ff..048b54d89 100644 --- a/internal/probe/types.go +++ b/internal/probe/types.go @@ -8,7 +8,6 @@ package probe import ( "fmt" - "github.com/tingly-dev/tingly-box/internal/client" "github.com/tingly-dev/tingly-box/internal/protocol" "github.com/tingly-dev/tingly-box/internal/typ" ) @@ -115,11 +114,10 @@ type E2ERequest struct { Message string `json:"message,omitempty"` } -// E2EData is an alias to client.ProbeResult — the canonical SDK-level -// probe result. Aliased so service-layer Response wrappers and swagger -// registrations can reference a name in this package without re-importing -// internal/client. -type E2EData = client.ProbeResult +// E2EData is an alias to ProbeResult — the canonical SDK-level probe result. +// Aliased so service-layer Response wrappers and swagger registrations can +// keep referring to the historical E2EData name. +type E2EData = ProbeResult // E2EResponseChunk represents a streaming response chunk. type E2EResponseChunk struct { diff --git a/internal/protocol/stream/prime.go b/internal/protocol/stream/prime.go index 553c139ae..1d920bc2e 100644 --- a/internal/protocol/stream/prime.go +++ b/internal/protocol/stream/prime.go @@ -9,10 +9,9 @@ // event it reads is not discarded: it is replayed via // firstEventReplayStream so the handler sees a complete stream. // -// NOTE: this is unrelated to client.ProbeResponsesStream, which issues a -// separate synthetic health-check request. This one pulls the first -// event of the real business stream and replays it — nothing extra is -// sent. +// NOTE: this is unrelated to the probe subsystem, which issues separate +// synthetic health-check requests. This one pulls the first event of the +// real business stream and replays it — nothing extra is sent. package stream import (