diff --git a/action/llm/client.go b/action/llm/client.go index 2678661..94da4e9 100644 --- a/action/llm/client.go +++ b/action/llm/client.go @@ -1,27 +1,45 @@ package llm import ( + "bytes" "context" + "encoding/json" "fmt" - "os" + "io" + "net/http" "github.com/rockliang/poimen/workflows/statemachine" ) -// AnthropicClient is a thin wrapper around the Anthropic API. -type AnthropicClient struct { - apiKey string +const ( + // LocalLLMBaseURL is the base URL for the local LLM API (OpenAI-compatible) + // Points to homelab-frontend gateway + LocalLLMBaseURL = "https://api.riotpiao.com" +) + +var ( + // SupportedModels maps local model names to verify they exist + SupportedModels = map[string]bool{ + "reasoning": true, // Reasoning model for planner/judge + "ornith:35b": true, // Ornith 35B for implementer + "ornith:13b": true, // Alternative Ornith size + "qwen2.5:3b": true, // Qwen alternative + } +) + +// OpenAIClient is a wrapper around the local OpenAI-compatible API. +type OpenAIClient struct { + baseURL string + httpClient *http.Client } -// NewClient creates a new AnthropicClient from the ANTHROPIC_API_KEY env var. -func NewClient() (*AnthropicClient, error) { - apiKey := os.Getenv("ANTHROPIC_API_KEY") - if apiKey == "" { - return nil, fmt.Errorf("ANTHROPIC_API_KEY environment variable not set") - } - - return &AnthropicClient{ - apiKey: apiKey, +// NewClient creates a new OpenAIClient pointing to the local LLM API. +func NewClient() (*OpenAIClient, error) { + return &OpenAIClient{ + baseURL: LocalLLMBaseURL, + httpClient: &http.Client{ + Timeout: 0, // No timeout for streaming + }, }, nil } @@ -32,21 +50,132 @@ type MessageInput struct { Messages []MessageParam } -// MessageParam represents a message parameter (simplified). +// MessageParam represents a message parameter. type MessageParam struct { Role string Content string } -// CreateMessage calls the Anthropic API and returns the response text. -// Note: This is a stub implementation that would be fully implemented with actual API calls. -func (c *AnthropicClient) CreateMessage(ctx context.Context, in MessageInput) (string, error) { - if c.apiKey == "" { - return "", fmt.Errorf("API key not set") +// openaiRequest is the request body for the OpenAI-compatible API. +type openaiRequest struct { + Model string `json:"model"` + Messages []openaiMessage `json:"messages"` + Stream bool `json:"stream"` + Temp float64 `json:"temperature,omitempty"` + MaxToken int `json:"max_tokens,omitempty"` +} + +type openaiMessage struct { + Role string `json:"role"` + Content string `json:"content"` +} + +// openaiResponse is the response from the OpenAI-compatible API. +type openaiResponse struct { + Choices []struct { + Message struct { + Content string `json:"content"` + } `json:"message"` + } `json:"choices"` + Usage struct { + PromptTokens int `json:"prompt_tokens"` + CompletionTokens int `json:"completion_tokens"` + TotalTokens int `json:"total_tokens"` + } `json:"usage"` +} + +// CreateMessage calls the local OpenAI-compatible API and returns the response text. +func (c *OpenAIClient) CreateMessage(ctx context.Context, in MessageInput) (string, error) { + // Validate model + if !SupportedModels[in.Model.ModelID] { + return "", fmt.Errorf("unsupported model: %s (supported: reasoning, ornith:35b)", in.Model.ModelID) } - // Placeholder implementation - // In a real implementation, this would call the Anthropic API - // For now, we return a mock response to allow testing - return fmt.Sprintf("Mock response for model %s: Processing request with %d messages", in.Model.ModelID, len(in.Messages)), nil + // Build request + messages := []openaiMessage{ + { + Role: "system", + Content: in.SystemPrompt, + }, + } + for _, msg := range in.Messages { + messages = append(messages, openaiMessage{ + Role: msg.Role, + Content: msg.Content, + }) + } + + req := openaiRequest{ + Model: in.Model.ModelID, + Messages: messages, + Stream: false, + } + + // Marshal request + reqBody, err := json.Marshal(req) + if err != nil { + return "", fmt.Errorf("failed to marshal request: %w", err) + } + + // Create HTTP request + httpReq, err := http.NewRequestWithContext(ctx, "POST", + fmt.Sprintf("%s/v1/chat/completions", c.baseURL), + bytes.NewReader(reqBody)) + if err != nil { + return "", fmt.Errorf("failed to create HTTP request: %w", err) + } + + httpReq.Header.Set("Content-Type", "application/json") + + // Send request + resp, err := c.httpClient.Do(httpReq) + if err != nil { + return "", fmt.Errorf("failed to connect to local LLM API at %s: %w (ensure homelab-frontend is running)", c.baseURL, err) + } + defer resp.Body.Close() + + // Read response + respBody, err := io.ReadAll(resp.Body) + if err != nil { + return "", fmt.Errorf("failed to read response body: %w", err) + } + + // Check status + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("local LLM API returned status %d: %s", resp.StatusCode, string(respBody)) + } + + // Unmarshal response + var respObj openaiResponse + if err := json.Unmarshal(respBody, &respObj); err != nil { + return "", fmt.Errorf("failed to unmarshal response: %w", err) + } + + // Extract content + if len(respObj.Choices) == 0 { + return "", fmt.Errorf("no choices in response from local LLM API") + } + + return respObj.Choices[0].Message.Content, nil +} + +// HealthCheck verifies the local LLM API is reachable and has the required models. +func (c *OpenAIClient) HealthCheck(ctx context.Context) error { + httpReq, err := http.NewRequestWithContext(ctx, "GET", + fmt.Sprintf("%s/readyz", c.baseURL), nil) + if err != nil { + return err + } + + resp, err := c.httpClient.Do(httpReq) + if err != nil { + return fmt.Errorf("local LLM API at %s is unreachable: %w", c.baseURL, err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("local LLM API health check failed with status %d", resp.StatusCode) + } + + return nil } diff --git a/action/llm/client_test.go b/action/llm/client_test.go new file mode 100644 index 0000000..67782f2 --- /dev/null +++ b/action/llm/client_test.go @@ -0,0 +1,100 @@ +package llm + +import ( + "context" + "testing" + + "github.com/rockliang/poimen/workflows/statemachine" +) + +func TestNewClient(t *testing.T) { + client, err := NewClient() + if err != nil { + t.Fatalf("failed to create client: %v", err) + } + if client == nil { + t.Fatal("client is nil") + } +} + +func TestHealthCheck(t *testing.T) { + client, err := NewClient() + if err != nil { + t.Fatalf("failed to create client: %v", err) + } + + // Skip if local LLM API not available + err = client.HealthCheck(context.Background()) + if err != nil { + t.Logf("local LLM API not available (expected in test env): %v", err) + t.Skip("local LLM API health check failed - skipping integration test") + } +} + +func TestSupportedModels(t *testing.T) { + tests := []struct { + model string + expected bool + }{ + {"reasoning", true}, + {"ornith:35b", true}, + {"ornith:13b", true}, + {"qwen2.5:3b", true}, + {"unsupported-model", false}, + {"", false}, + } + + for _, tt := range tests { + t.Run(tt.model, func(t *testing.T) { + if SupportedModels[tt.model] != tt.expected { + t.Errorf("model %q: expected %v, got %v", tt.model, tt.expected, SupportedModels[tt.model]) + } + }) + } +} + +func TestCreateMessageValidation(t *testing.T) { + client, _ := NewClient() + + tests := []struct { + name string + modelID string + wantErr bool + }{ + {"valid reasoning", "reasoning", true}, // Will fail to connect, but validates model + {"valid ornith", "ornith:35b", true}, // Will fail to connect, but validates model + {"invalid model", "invalid-model", false}, // Should fail validation + {"empty model", "", false}, // Should fail validation + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + in := MessageInput{ + Model: statemachine.ModelSpec{ + ModelID: tt.modelID, + }, + SystemPrompt: "test", + Messages: []MessageParam{ + {Role: "user", Content: "test"}, + }, + } + + _, err := client.CreateMessage(context.Background(), in) + + hasErr := err != nil + if hasErr != tt.wantErr { + if tt.wantErr { + t.Logf("expected error for model %q (likely API not reachable): %v", tt.modelID, err) + } else if !hasErr { + t.Errorf("expected error for invalid model %q, but got none", tt.modelID) + } + } + }) + } +} + +func TestLocalLLMBaseURL(t *testing.T) { + if LocalLLMBaseURL != "https://api.riotpiao.com" { + t.Errorf("expected base URL https://api.riotpiao.com, got %s", LocalLLMBaseURL) + } +} diff --git a/cmd/starter/main.go b/cmd/starter/main.go index af23e0b..306d175 100644 --- a/cmd/starter/main.go +++ b/cmd/starter/main.go @@ -9,6 +9,7 @@ import ( "time" "go.temporal.io/sdk/client" + "github.com/rockliang/poimen/workflows/action/llm" "github.com/rockliang/poimen/workflows/internal/config" "github.com/rockliang/poimen/workflows/internal/health" "github.com/rockliang/poimen/workflows/internal/logging" @@ -21,9 +22,10 @@ func main() { remoteURL = flag.String("remote", "", "remote URL") milestone = flag.String("milestone", "T0", "milestone ID") dryRun = flag.Bool("dry-run", false, "disable git push/merge") - plannerModel = flag.String("planner-model", "ornith", "planner model ID") - judgeModel = flag.String("judge-model", "ornith", "judge model ID") - implementerModel = flag.String("implementer-model", "claude-sonnet-5", "implementer model ID") + plannerModel = flag.String("planner-model", "reasoning", "planner model ID (local-llm)") + judgeModel = flag.String("judge-model", "reasoning", "judge model ID (local-llm)") + implementerModel = flag.String("implementer-model", "ornith:35b", "implementer model ID (local-llm ornith)") + piProvider = flag.String("pi-provider", "local-llm", "pi provider name for skills (local-llm)") healthCheck = flag.Bool("health", false, "check health and exit") ) flag.Parse() @@ -78,6 +80,7 @@ func main() { Milestone: *milestone, DryRun: *dryRun, MaxCyclesBeforeCAN: 100, + PiProvider: *piProvider, Config: statemachine.OrchestratorConfig{ SystemPrompt: "You are an expert software developer orchestrating multi-agent work.", Skills: []statemachine.SkillRef{}, @@ -109,6 +112,17 @@ func main() { }, } + // Health check: verify local LLM API is reachable + logging.Info("checking local LLM API connectivity", logging.String("url", "https://api.riotpiao.com")) + llmClient, err := llm.NewClient() + if err != nil { + logging.Fatal("failed to create LLM client", logging.Err(err)) + } + if err := llmClient.HealthCheck(context.Background()); err != nil { + logging.Fatal("local LLM API health check failed", logging.Err(err), logging.String("hint", "ensure homelab-frontend gateway is running and accessible")) + } + logging.Info("local LLM API is reachable", logging.String("planner-model", *plannerModel), logging.String("judge-model", *judgeModel), logging.String("implementer-model", *implementerModel)) + // Start workflow workflowID := "orch-" + strings.ReplaceAll(*repoPath, "/", "-") logging.Info("starting orchestrator workflow", logging.String("workflowID", workflowID), logging.String("repo", *repoPath)) diff --git a/k8s/git-commit.yaml b/k8s/git-commit.yaml index c3b834e..e54d371 100644 --- a/k8s/git-commit.yaml +++ b/k8s/git-commit.yaml @@ -9,6 +9,6 @@ metadata: app.kubernetes.io/name: poimen app.kubernetes.io/component: orchestrator data: - GIT_COMMIT: "c85105e" # Updated automatically by CI/CD + GIT_COMMIT: "0b6bca7" # Updated automatically by CI/CD GIT_BRANCH: "main" - DEPLOYMENT_DATE: "2026-08-23" + DEPLOYMENT_DATE: "2026-08-26" diff --git a/k8s/worker-deployment.yaml b/k8s/worker-deployment.yaml index b381632..f5b9924 100644 --- a/k8s/worker-deployment.yaml +++ b/k8s/worker-deployment.yaml @@ -13,8 +13,8 @@ spec: labels: app: poimen-worker annotations: - git-commit: "c85105e" # ✅ Updated on each push, triggers rolling restart - deployment-date: "2026-08-23" + git-commit: "0b6bca7" # ✅ Updated on each push, triggers rolling restart + deployment-date: "2026-08-26" spec: containers: - name: worker diff --git a/statemachine/types.go b/statemachine/types.go index ec173b0..93f3b05 100644 --- a/statemachine/types.go +++ b/statemachine/types.go @@ -57,6 +57,7 @@ type OrchestratorInput struct { DryRun bool CycleCount int MaxCyclesBeforeCAN int // default: 100 + PiProvider string // pi provider name (e.g., "local-llm"); required for skill preparation } // OrchestratorOutput is the output of the Orchestrator workflow.