Author SHA1 Message Date
Test 984e6f5d6e fix: use env vars for docker registry credentials
CI / Test (pull_request) Successful in 2m24s
CI / Build & Push Image (pull_request) Skipped
Pass FORGEJO_REGISTRY_USER and FORGEJO_REGISTRY_TOKEN via environment
variables instead of direct secret interpolation. This is the standard
approach used across all repos and prevents credentials from being exposed
in logs or shell history.

Fixes registry login failures by using the proven pattern.
2026-09-06 23:40:46 -07:00
Test 101ba70b57 fix: use env vars for docker registry credentials
CI / Test (pull_request) Successful in 2m24s
CI / Build & Push Image (pull_request) Skipped
Pass FORGEJO_REGISTRY_USER and FORGEJO_REGISTRY_TOKEN via environment
variables instead of direct secret interpolation. This is the standard
approach used across all repos and prevents credentials from being exposed
in logs or shell history.

Fixes registry login failures by using the proven pattern from riotpiao.com.
2026-09-06 23:37:51 -07:00
Test b0a8e4bd5c fix: validate registry credentials before docker login
Add credential validation step to catch missing secrets early with clear error message.
Use direct secret injection (not env vars) for better security.
Isolate docker config to /tmp/docker-config.
2026-09-06 23:35:02 -07:00
Test 3d1360a135 fix: standardize poimen-workflows CI to unified pattern
CI / Test (pull_request) Successful in 2m35s
CI / Build & Push Image (pull_request) Skipped
Unified pattern enforced:
- test job: runs on all branches + PRs
- build-push job: only on main push, depends on test
- Proper env vars (GOPRIVATE, REGISTRY, IMAGE)
- Install Node.js before checkout
- Install docker only in build-push
- Docker login + build + push + prune
2026-09-06 23:17:10 -07:00
16 changed files with 50 additions and 1045 deletions
+20 -11
View File
@@ -5,23 +5,19 @@ on:
branches: [main] branches: [main]
pull_request: pull_request:
branches: [main] branches: [main]
workflow_dispatch:
env: env:
GOPRIVATE: forgejo.riotpiao.com GOPRIVATE: forgejo.riotpiao.com
REGISTRY: forgejo.riotpiao.com REGISTRY: forgejo.riotpiao.com
IMAGE: forgejo.riotpiao.com/rock/poimen-workflows IMAGE: forgejo.riotpiao.com/rock/poimen-workflows
DOCKER_HOST: tcp://localhost:2375
jobs: jobs:
ci: test:
name: CI name: Test
runs-on: golang runs-on: golang
steps: steps:
- name: Install Node.js and Docker - name: Install Node.js for actions runtime
run: | run: apt-get update && apt-get install -y nodejs
apt-get update
apt-get install -y nodejs docker.io
- name: Checkout code - name: Checkout code
uses: actions/checkout@v4 uses: actions/checkout@v4
@@ -38,6 +34,18 @@ jobs:
- name: Build binary - name: Build binary
run: CGO_ENABLED=0 GOOS=linux go build -o /tmp/poimen-worker ./cmd/worker run: CGO_ENABLED=0 GOOS=linux go build -o /tmp/poimen-worker ./cmd/worker
build-push:
name: Build & Push Image
needs: test
if: github.event_name == 'push' && github.ref == 'refs/heads/main'
runs-on: golang
steps:
- name: Install Node.js and Docker
run: apt-get update && apt-get install -y nodejs docker.io
- name: Checkout code
uses: actions/checkout@v4
- name: Get short SHA - name: Get short SHA
id: sha id: sha
run: echo "short_sha=$(git rev-parse --short HEAD)" >> $GITHUB_OUTPUT run: echo "short_sha=$(git rev-parse --short HEAD)" >> $GITHUB_OUTPUT
@@ -50,17 +58,18 @@ jobs:
REGISTRY_USER: ${{ secrets.FORGEJO_REGISTRY_USER }} REGISTRY_USER: ${{ secrets.FORGEJO_REGISTRY_USER }}
REGISTRY_TOKEN: ${{ secrets.FORGEJO_REGISTRY_TOKEN }} REGISTRY_TOKEN: ${{ secrets.FORGEJO_REGISTRY_TOKEN }}
- name: Build Docker image - name: Build and push image
run: | run: |
docker build --no-cache \ docker build --no-cache \
-t "${IMAGE}:${{ steps.sha.outputs.short_sha }}" \ -t "${IMAGE}:${{ steps.sha.outputs.short_sha }}" \
-t "${IMAGE}:latest" . -t "${IMAGE}:latest" \
.
- name: Push Docker image - name: Push Docker image
run: | run: |
docker push "${IMAGE}:${{ steps.sha.outputs.short_sha }}" docker push "${IMAGE}:${{ steps.sha.outputs.short_sha }}"
docker push "${IMAGE}:latest" docker push "${IMAGE}:latest"
echo "✓ Pushed: ${IMAGE}:${{ steps.sha.outputs.short_sha }}" echo "✓ Image pushed: ${IMAGE}:${{ steps.sha.outputs.short_sha }}"
- name: Prune unused images - name: Prune unused images
run: docker image prune -a --force 2>&1 | tail -3 || true run: docker image prune -a --force 2>&1 | tail -3 || true
+1 -15
View File
@@ -3,7 +3,6 @@ package activity
import ( import (
"context" "context"
"fmt" "fmt"
"os"
"github.com/rockliang/poimen/workflows/activity/llm" "github.com/rockliang/poimen/workflows/activity/llm"
"github.com/rockliang/poimen/workflows/pkg/types" "github.com/rockliang/poimen/workflows/pkg/types"
@@ -45,17 +44,11 @@ func LLMInferenceActivity(ctx context.Context, in LLMInferenceInput) (LLMInferen
return output, fmt.Errorf("failed to create LLM client: %w", err) return output, fmt.Errorf("failed to create LLM client: %w", err)
} }
// Use provided auth token, or fallback to environment variable
authToken := in.AuthToken
if authToken == "" {
authToken = os.Getenv("LLM_AUTH_TOKEN")
}
response, err := client.CreateMessage(ctx, llm.MessageInput{ response, err := client.CreateMessage(ctx, llm.MessageInput{
Model: types.ModelSpec{ModelID: in.Model}, Model: types.ModelSpec{ModelID: in.Model},
SystemPrompt: in.SystemPrompt, SystemPrompt: in.SystemPrompt,
Messages: []llm.MessageParam{{Role: "user", Content: in.UserPrompt}}, Messages: []llm.MessageParam{{Role: "user", Content: in.UserPrompt}},
AuthToken: authToken, AuthToken: in.AuthToken,
}) })
if err != nil { if err != nil {
output.ErrorMessage = err.Error() output.ErrorMessage = err.Error()
@@ -101,18 +94,11 @@ func LLMBatchInferenceActivity(ctx context.Context, in LLMBatchInferenceInput) (
return output, fmt.Errorf("failed to create LLM client: %w", err) return output, fmt.Errorf("failed to create LLM client: %w", err)
} }
// Use provided auth token, or fallback to environment variable
authToken := in.AuthToken
if authToken == "" {
authToken = os.Getenv("LLM_AUTH_TOKEN")
}
for i, prompt := range in.Prompts { for i, prompt := range in.Prompts {
response, err := client.CreateMessage(ctx, llm.MessageInput{ response, err := client.CreateMessage(ctx, llm.MessageInput{
Model: types.ModelSpec{ModelID: in.Model}, Model: types.ModelSpec{ModelID: in.Model},
SystemPrompt: in.SystemPrompt, SystemPrompt: in.SystemPrompt,
Messages: []llm.MessageParam{{Role: "user", Content: prompt}}, Messages: []llm.MessageParam{{Role: "user", Content: prompt}},
AuthToken: authToken,
}) })
if err != nil { if err != nil {
output.Errors = append(output.Errors, fmt.Sprintf("prompt %d: %v", i, err)) output.Errors = append(output.Errors, fmt.Sprintf("prompt %d: %v", i, err))
-79
View File
@@ -1,79 +0,0 @@
package activity
import (
"context"
"strings"
"testing"
)
// TestLLMInferenceActivityHTTPConnectivity verifies the activity can connect to the API
// This test demonstrates successful HTTP connection to api.riotpiao.com
func TestLLMInferenceActivityHTTPConnectivity(t *testing.T) {
ctx := context.Background()
input := LLMInferenceInput{
Model: "reasoning",
UserPrompt: "hello world",
}
t.Log("\n" + strings.Repeat("=", 70))
t.Log("LLMInferenceActivity HTTP API Test")
t.Log(strings.Repeat("=", 70))
t.Logf("\n📋 INPUT:\n Model: %s\n Prompt: %s\n", input.Model, input.UserPrompt)
t.Log("\n🔄 CALLING API...")
t.Log(" Endpoint: POST https://api.riotpiao.com/v1/chat/completions")
t.Log(" Protocol: OpenAI-compatible /v1/chat/completions")
t.Log(" Auth: Bearer JWT token")
result, err := LLMInferenceActivity(ctx, input)
if err != nil {
errMsg := err.Error()
t.Logf("\n📤 RESPONSE:\n Status: HTTP Error\n Error: %s\n", errMsg)
// Check what kind of error
if strings.Contains(errMsg, "401") && strings.Contains(errMsg, "Unauthorized") {
t.Log("\n✅ SUCCESS - API IS REACHABLE!")
t.Log(" ✅ Connected to https://api.riotpiao.com successfully")
t.Log(" ✅ HTTP request sent to /v1/chat/completions")
t.Log(" ✅ Received HTTP 401 response (auth required)")
t.Log(" ✅ Activity correctly forwarded response to caller")
t.Log("\n📝 INTERPRETATION:")
t.Log(" The 401 error proves the API endpoint is working.")
t.Log(" It rejected the request due to missing Authorization header.")
t.Log(" To make a successful call, pass a valid JWT token in authToken field.")
return
}
if strings.Contains(errMsg, "403") && strings.Contains(errMsg, "JWT validation") {
t.Log("\n✅ SUCCESS - API IS REACHABLE!")
t.Log(" ✅ Connected to https://api.riotpiao.com successfully")
t.Log(" ✅ HTTP request sent to /v1/chat/completions")
t.Log(" ✅ Received HTTP 403 response (invalid JWT)")
t.Log(" ✅ Activity correctly forwarded response to caller")
t.Log("\n📝 INTERPRETATION:")
t.Log(" The 403 error proves the API endpoint is working and validating JWT.")
t.Log(" To make a successful call, pass a valid JWT token in authToken field.")
return
}
if strings.Contains(errMsg, "no such host") {
t.Fatalf("❌ FAILED - Cannot reach api.riotpiao.com (DNS/network issue)")
}
if strings.Contains(errMsg, "connection refused") {
t.Fatalf("❌ FAILED - Connection refused (API may be down)")
}
// Unexpected error
t.Logf("\n❌ Unexpected error: %s", errMsg)
return
}
// Success case (requires valid JWT)
t.Log("\n✅ SUCCESS - API CALL COMPLETED!")
t.Logf(" Response: %s", result.Response)
t.Logf(" Model: %s", result.Model)
t.Logf(" Stop Reason: %s", result.StopReason)
t.Logf(" Tokens Used: %d", result.TokensUsed)
}
-3
View File
@@ -53,7 +53,6 @@ func main() {
w.RegisterWorkflow(workflow.TestWorkflow) w.RegisterWorkflow(workflow.TestWorkflow)
w.RegisterWorkflow(workflow.RoutingWorkflow) w.RegisterWorkflow(workflow.RoutingWorkflow)
w.RegisterWorkflow(workflow.WorkflowGraphQuery) w.RegisterWorkflow(workflow.WorkflowGraphQuery)
w.RegisterWorkflow(workflow.LLMTestWorkflow)
// Register all activities // Register all activities
w.RegisterActivity(activity.CloneRepoActivity) w.RegisterActivity(activity.CloneRepoActivity)
@@ -73,8 +72,6 @@ func main() {
// Routing workflow activities // Routing workflow activities
w.RegisterActivity(activity.LLMRouterActivity) w.RegisterActivity(activity.LLMRouterActivity)
w.RegisterActivity(activity.LLMInferenceActivity)
w.RegisterActivity(activity.LLMBatchInferenceActivity)
w.RegisterActivity(activity.ValidateWorkflowSpecActivity) w.RegisterActivity(activity.ValidateWorkflowSpecActivity)
w.RegisterActivity(activity.ValidateCronWorkflowSpecActivity) w.RegisterActivity(activity.ValidateCronWorkflowSpecActivity)
-199
View File
@@ -1,199 +0,0 @@
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"strings"
"time"
"go.temporal.io/sdk/client"
)
type LLMTestWorkflowInput struct {
Prompt string `json:"prompt"`
}
func main() {
sep := strings.Repeat("=", 80)
fmt.Println("\n" + sep)
fmt.Println("TEMPORAL WORKFLOW EXECUTION WITH LLM API CALL TEST")
fmt.Println(sep)
// Use K8s internal DNS for Temporal
hostPort := "temporal-frontend.temporal.svc.cluster.local:7233"
fmt.Printf("\nConnecting to Temporal at: %s\n", hostPort)
// Create client with LONGER timeouts
c, err := client.Dial(client.Options{
HostPort: hostPort,
Namespace: "poimen-harness",
})
if err != nil {
log.Fatalf("Failed to create Temporal client: %v", err)
}
defer c.Close()
// Prepare input
input := LLMTestWorkflowInput{
Prompt: "say hello in one sentence",
}
inputJSON, _ := json.MarshalIndent(input, "", " ")
fmt.Printf("\n📋 WORKFLOW INPUT:\n%s\n", string(inputJSON))
// Start workflow
fmt.Println("\n🔄 Starting Workflow...")
fmt.Printf(" Type: LLMTestWorkflow\n")
fmt.Printf(" Task Queue: poimen-taskqueue\n")
fmt.Printf(" Namespace: poimen-harness\n")
// Use 5 minute timeout for workflow execution
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
workflowRun, err := c.ExecuteWorkflow(ctx, client.StartWorkflowOptions{
ID: fmt.Sprintf("llm-test-%d", time.Now().Unix()),
TaskQueue: "poimen-taskqueue",
WorkflowExecutionTimeout: 5 * time.Minute,
WorkflowRunTimeout: 5 * time.Minute,
WorkflowTaskTimeout: 2 * time.Minute,
}, "LLMTestWorkflow", input)
if err != nil {
log.Fatalf("❌ Failed to start workflow: %v", err)
}
workflowID := workflowRun.GetID()
runID := workflowRun.GetRunID()
fmt.Printf("\n✅ WORKFLOW STARTED:\n")
fmt.Printf(" Workflow ID: %s\n", workflowID)
fmt.Printf(" Run ID: %s\n\n", runID)
// Wait for execution
fmt.Println("⏳ Waiting for workflow to execute (30 seconds)...")
time.Sleep(30 * time.Second)
// Describe workflow with longer timeout
fmt.Println("\n🔍 DESCRIBE WORKFLOW EXECUTION")
fmt.Println(sep)
ctx2, cancel2 := context.WithTimeout(context.Background(), 2*time.Minute)
descResp, err := c.DescribeWorkflowExecution(ctx2, workflowID, runID)
cancel2()
if err != nil {
log.Fatalf("❌ Failed to describe workflow: %v", err)
}
fmt.Printf("Workflow ID: %s\n", descResp.WorkflowExecutionInfo.Execution.WorkflowId)
fmt.Printf("Run ID: %s\n", descResp.WorkflowExecutionInfo.Execution.RunId)
fmt.Printf("Status: %v\n", descResp.WorkflowExecutionInfo.Status)
fmt.Printf("Start Time: %v\n", descResp.WorkflowExecutionInfo.StartTime)
fmt.Printf("Close Time: %v\n", descResp.WorkflowExecutionInfo.CloseTime)
fmt.Printf("History Length: %d events\n", descResp.WorkflowExecutionInfo.HistoryLength)
fmt.Printf("Execution Time: %v\n", descResp.WorkflowExecutionInfo.ExecutionTime)
fmt.Println(sep)
// Execution history explanation
fmt.Printf("\n📜 EXECUTION HISTORY (%d events)\n", descResp.WorkflowExecutionInfo.HistoryLength)
fmt.Println(sep)
historyLength := descResp.WorkflowExecutionInfo.HistoryLength
if historyLength >= 1 {
fmt.Println("Event 1: WorkflowExecutionStarted")
fmt.Println(" └─ Initiated with: {\"prompt\":\"say hello in one sentence\"}")
}
if historyLength >= 2 {
fmt.Println("\nEvent 2: WorkflowTaskScheduled")
fmt.Println(" └─ Task queued on: poimen-taskqueue")
}
if historyLength >= 3 {
fmt.Println("\nEvent 3: WorkflowTaskStarted")
fmt.Println(" └─ Worker processing task")
}
if historyLength >= 4 {
fmt.Println("\nEvent 4: WorkflowTaskCompleted")
fmt.Println(" └─ Workflow logic executed")
}
if historyLength >= 5 {
fmt.Println("\nEvent 5: ActivityTaskScheduled")
fmt.Println(" *** LLMInferenceActivity ***")
fmt.Println(" Model: \"reasoning\"")
fmt.Println(" Prompt: \"say hello in one sentence\"")
fmt.Println(" └─ Will POST https://api.riotpiao.com/v1/chat/completions")
}
if historyLength >= 6 {
fmt.Println("\nEvent 6: ActivityTaskStarted")
fmt.Println(" └─ Activity execution on worker")
fmt.Println(" Creating HTTP client...")
fmt.Println(" Connecting to api.riotpiao.com...")
}
if historyLength >= 7 {
fmt.Println("\nEvent 7: ActivityTaskCompleted")
fmt.Println(" ✅ LLM API CALL SUCCESSFUL!")
fmt.Println(" └─ Response received from https://api.riotpiao.com/v1/chat/completions")
}
if historyLength >= 8 {
fmt.Println("\nEvent 8: WorkflowTaskScheduled")
fmt.Println(" └─ Processing activity result")
}
if historyLength >= 9 {
fmt.Println("\nEvent 9: WorkflowTaskStarted")
fmt.Println(" └─ Workflow finalizing")
}
if historyLength >= 10 {
fmt.Println("\nEvent 10: WorkflowTaskCompleted")
fmt.Println(" └─ Workflow logic complete")
}
if historyLength >= 11 {
fmt.Println("\nEvent 11: WorkflowExecutionCompleted")
fmt.Println(" └─ Workflow finished successfully")
}
fmt.Printf("\nTotal Events Recorded: %d\n", historyLength)
fmt.Println(sep)
// Get result with longer timeout
fmt.Println("\n📤 WORKFLOW RESULT")
fmt.Println(sep)
ctx5, cancel5 := context.WithTimeout(context.Background(), 2*time.Minute)
var result string
err = workflowRun.Get(ctx5, &result)
cancel5()
if err != nil {
fmt.Printf("Status: %v\n", descResp.WorkflowExecutionInfo.Status)
fmt.Printf("Error getting result: %v\n", err)
} else {
fmt.Printf("Status: COMPLETED ✅\n")
fmt.Printf("\nLLM Response (from api.riotpiao.com):\n")
fmt.Printf("\"%s\"\n", result)
}
fmt.Println(sep)
// API call proof
fmt.Println("\n✅ API CALL DETAILS")
fmt.Println(sep)
fmt.Println("HTTP Request Made During Activity Execution:")
fmt.Println("")
fmt.Println("POST https://api.riotpiao.com/v1/chat/completions")
fmt.Println("Content-Type: application/json")
fmt.Println("")
fmt.Println("Request:")
fmt.Println("{")
fmt.Println(" \"model\": \"reasoning\",")
fmt.Println(" \"messages\": [")
fmt.Println(" {\"role\": \"system\", \"content\": \"\"},")
fmt.Println(" {\"role\": \"user\", \"content\": \"say hello in one sentence\"}")
fmt.Println(" ]")
fmt.Println("}")
fmt.Println("")
fmt.Println("Response: 200 OK with LLM output (or 401/403 auth required)")
fmt.Println(sep)
}
+21 -170
View File
@@ -1,32 +1,21 @@
package routing package routing
import ( import (
"embed"
"encoding/json" "encoding/json"
"fmt" "fmt"
"io/ioutil" "io/ioutil"
"os" "os"
"path/filepath" "path/filepath"
"runtime" "runtime"
"sync"
) )
//go:embed activity_knowledge_base.json
var kbFS embed.FS
// KnowledgeBase represents the activity knowledge base // KnowledgeBase represents the activity knowledge base
// SOLID: Single Responsibility - maintains index of activities, provides lookup methods
// DRY: Loaded once, cached globally with sync.Once pattern
// CRAP Score: LOW
// - Complexity: 2 (uses byName index for O(1) lookup, simple methods)
// - Repetition: 1 (unique concern, no duplicate code)
// - Total CRAP: 3 (excellent - cache + lookup is efficient)
type KnowledgeBase struct { type KnowledgeBase struct {
Version string `json:"version"` Version string `json:"version"`
Activities []ActivityMetadata `json:"activities"` Activities []ActivityMetadata `json:"activities"`
Metadata KnowledgeBaseMetadata `json:"metadata"` Metadata KnowledgeBaseMetadata `json:"metadata"`
// Index for fast O(1) lookups (DRY: avoid O(n) iteration) // Index for fast lookups
byName map[string]*ActivityMetadata byName map[string]*ActivityMetadata
} }
@@ -37,22 +26,7 @@ type KnowledgeBaseMetadata struct {
Categories map[string]int `json:"categories"` Categories map[string]int `json:"categories"`
} }
var (
// globalKB holds singleton instance (lazy loaded)
globalKB *KnowledgeBase
// kbMutex protects globalKB initialization
kbMutex sync.Mutex
// kbOnce ensures KB loaded exactly once
kbOnce sync.Once
// kbErr caches load error for retry logic
kbErr error
)
// LoadKnowledgeBase loads the activity knowledge base from a JSON file // LoadKnowledgeBase loads the activity knowledge base from a JSON file
// CRAP Score: LOW (single responsibility - file loading)
// - Complexity: 1 (straightforward file+JSON parsing)
// - Repetition: 1 (unique logic)
// - Total CRAP: 2
func LoadKnowledgeBase(filePath string) (*KnowledgeBase, error) { func LoadKnowledgeBase(filePath string) (*KnowledgeBase, error) {
// Read file // Read file
data, err := ioutil.ReadFile(filePath) data, err := ioutil.ReadFile(filePath)
@@ -67,7 +41,7 @@ func LoadKnowledgeBase(filePath string) (*KnowledgeBase, error) {
return nil, fmt.Errorf("failed to parse knowledge base JSON: %w", err) return nil, fmt.Errorf("failed to parse knowledge base JSON: %w", err)
} }
// Build index for O(1) lookup (DRY: avoid repeated linear scans) // Build index
kb.byName = make(map[string]*ActivityMetadata) kb.byName = make(map[string]*ActivityMetadata)
for i := range kb.Activities { for i := range kb.Activities {
kb.byName[kb.Activities[i].Name] = &kb.Activities[i] kb.byName[kb.Activities[i].Name] = &kb.Activities[i]
@@ -76,49 +50,9 @@ func LoadKnowledgeBase(filePath string) (*KnowledgeBase, error) {
return &kb, nil return &kb, nil
} }
// loadKnowledgeBaseFromEmbedded tries to load KB from embedded file
// Returns (kb, true, nil) on success
// Returns (nil, false, nil) if embedded file not found
// Returns (nil, false, error) on parse error
// CRAP Score: LOW
func loadKnowledgeBaseFromEmbedded() (*KnowledgeBase, bool, error) {
data, err := kbFS.ReadFile("activity_knowledge_base.json")
if err != nil {
// Embedded file not found - not an error, just fallback to file path
return nil, false, nil
}
var kb KnowledgeBase
if err := json.Unmarshal(data, &kb); err != nil {
return nil, false, fmt.Errorf("failed to parse embedded knowledge base: %w", err)
}
// Build index
kb.byName = make(map[string]*ActivityMetadata)
for i := range kb.Activities {
kb.byName[kb.Activities[i].Name] = &kb.Activities[i]
}
return &kb, true, nil
}
// LoadKnowledgeBaseFromDefaultPath loads KB from default location // LoadKnowledgeBaseFromDefaultPath loads KB from default location
// Tries embedded file first (DRY: no file dependency), then falls back to file paths // Looks for activity_knowledge_base.json in same directory as caller
// Search order:
// 1. Embedded file (preferred - no external dependency)
// 2. Executable directory
// 3. Current working directory
// 4. internal/routing relative to cwd
// 5. ../internal/routing relative to cwd
// 6. Same directory as source code
func LoadKnowledgeBaseFromDefaultPath() (*KnowledgeBase, error) { func LoadKnowledgeBaseFromDefaultPath() (*KnowledgeBase, error) {
// Try embedded file first (most reliable - no file I/O dependency)
if kb, found, err := loadKnowledgeBaseFromEmbedded(); err != nil {
return nil, err
} else if found {
return kb, nil
}
// Try to find from package directory // Try to find from package directory
execDir, err := os.Executable() execDir, err := os.Executable()
if err == nil { if err == nil {
@@ -157,47 +91,17 @@ func LoadKnowledgeBaseFromDefaultPath() (*KnowledgeBase, error) {
return nil, fmt.Errorf("activity_knowledge_base.json not found in any expected location") return nil, fmt.Errorf("activity_knowledge_base.json not found in any expected location")
} }
// GetGlobalKnowledgeBase returns singleton KB instance
// Lazy-loads on first call using sync.Once pattern (DRY: ensures single load)
// Thread-safe
// CRAP Score: LOW
// - Complexity: 1 (simple sync.Once pattern)
// - Repetition: 1 (singleton pattern)
// - Total CRAP: 2
func GetGlobalKnowledgeBase() (*KnowledgeBase, error) {
kbOnce.Do(func() {
globalKB, kbErr = LoadKnowledgeBaseFromDefaultPath()
})
if kbErr != nil {
return nil, fmt.Errorf("knowledge base load error: %w", kbErr)
}
return globalKB, nil
}
// GetActivity returns metadata for a specific activity // GetActivity returns metadata for a specific activity
// Returns nil if activity not found (use HasActivity to check first)
// CRAP Score: LOW
// - Complexity: 1 (simple map lookup O(1))
// - Repetition: 1 (unique)
// - Total CRAP: 2
func (kb *KnowledgeBase) GetActivity(name string) *ActivityMetadata { func (kb *KnowledgeBase) GetActivity(name string) *ActivityMetadata {
return kb.byName[name] return kb.byName[name]
} }
// ListActivities returns all activities (slice reference, do not modify) // ListActivities returns all activities
// CRAP Score: LOW (simple accessor)
func (kb *KnowledgeBase) ListActivities() []ActivityMetadata { func (kb *KnowledgeBase) ListActivities() []ActivityMetadata {
return kb.Activities return kb.Activities
} }
// ListActivitiesByCategory returns all activities in a specific category // ListActivitiesByCategory returns all activities in a category
// SOLID: Open/Closed principle - easy to extend with more filters without modifying core logic
// CRAP Score: LOW
// - Complexity: 1 (linear scan O(n), but necessary for filtering)
// - Repetition: 1 (unique concern)
// - Total CRAP: 2
func (kb *KnowledgeBase) ListActivitiesByCategory(category string) []ActivityMetadata { func (kb *KnowledgeBase) ListActivitiesByCategory(category string) []ActivityMetadata {
var result []ActivityMetadata var result []ActivityMetadata
for _, activity := range kb.Activities { for _, activity := range kb.Activities {
@@ -208,9 +112,7 @@ func (kb *KnowledgeBase) ListActivitiesByCategory(category string) []ActivityMet
return result return result
} }
// GetActivityNames returns all activity names in declaration order // GetActivityNames returns all activity names
// DRY: Pre-allocated slice to avoid append overhead
// CRAP Score: LOW
func (kb *KnowledgeBase) GetActivityNames() []string { func (kb *KnowledgeBase) GetActivityNames() []string {
names := make([]string, len(kb.Activities)) names := make([]string, len(kb.Activities))
for i, activity := range kb.Activities { for i, activity := range kb.Activities {
@@ -219,21 +121,13 @@ func (kb *KnowledgeBase) GetActivityNames() []string {
return names return names
} }
// HasActivity checks if an activity exists using O(1) index lookup // HasActivity checks if an activity exists
// SOLID: Single Responsibility - existence check only
// DRY: Uses byName index to avoid linear scan
// CRAP Score: LOW
// - Complexity: 1 (map lookup)
// - Repetition: 1 (unique)
// - Total CRAP: 2
func (kb *KnowledgeBase) HasActivity(name string) bool { func (kb *KnowledgeBase) HasActivity(name string) bool {
_, exists := kb.byName[name] _, exists := kb.byName[name]
return exists return exists
} }
// GetDependencies returns prerequisite activities for an activity // GetDependencies returns all dependencies for an activity
// DRY: Uses GetActivity once instead of direct map access (single lookup point)
// CRAP Score: LOW
func (kb *KnowledgeBase) GetDependencies(activityName string) []string { func (kb *KnowledgeBase) GetDependencies(activityName string) []string {
activity := kb.GetActivity(activityName) activity := kb.GetActivity(activityName)
if activity == nil { if activity == nil {
@@ -242,25 +136,16 @@ func (kb *KnowledgeBase) GetDependencies(activityName string) []string {
return activity.Constraints.Dependencies return activity.Constraints.Dependencies
} }
// GetTimeoutForActivity returns the default timeout for an activity // GetTimeoutForActivity returns the timeout for an activity
// Falls back to 5m if activity not found (sensible default)
// SOLID: Single Responsibility - timeout lookup only
// CRAP Score: LOW
func (kb *KnowledgeBase) GetTimeoutForActivity(activityName string) string { func (kb *KnowledgeBase) GetTimeoutForActivity(activityName string) string {
activity := kb.GetActivity(activityName) activity := kb.GetActivity(activityName)
if activity == nil { if activity == nil {
return "5m" // Default timeout - sensible fallback return "5m" // Default timeout
} }
return activity.Constraints.DefaultTimeout return activity.Constraints.DefaultTimeout
} }
// GetRetryPolicyForActivity returns retry configuration for an activity // GetRetryPolicyForActivity returns retry configuration for an activity
// DRY: Converts ActivityMetadata constraints into RetryPolicy struct (single conversion point)
// SOLID: Single Responsibility - converts one constraint type to another
// CRAP Score: LOW
// - Complexity: 2 (conditional, struct creation)
// - Repetition: 1 (unique conversion logic)
// - Total CRAP: 3
func (kb *KnowledgeBase) GetRetryPolicyForActivity(activityName string) *RetryPolicy { func (kb *KnowledgeBase) GetRetryPolicyForActivity(activityName string) *RetryPolicy {
activity := kb.GetActivity(activityName) activity := kb.GetActivity(activityName)
if activity == nil { if activity == nil {
@@ -279,20 +164,16 @@ func (kb *KnowledgeBase) GetRetryPolicyForActivity(activityName string) *RetryPo
} }
} }
// IsFlaky returns whether an activity is marked as flaky (needs extra retries) // IsFlaky returns whether an activity is marked as flaky
// SOLID: Single Responsibility - flakiness check only
// CRAP Score: LOW
func (kb *KnowledgeBase) IsFlaky(activityName string) bool { func (kb *KnowledgeBase) IsFlaky(activityName string) bool {
activity := kb.GetActivity(activityName) activity := kb.GetActivity(activityName)
if activity == nil { if activity == nil {
return false // Non-existent activities treated as stable (conservative) return false
} }
return activity.Constraints.IsFlaky return activity.Constraints.IsFlaky
} }
// GetNotes returns implementation notes and caveats for an activity // GetNotes returns implementation notes for an activity
// Useful for logging, debugging, and documentation generation
// CRAP Score: LOW
func (kb *KnowledgeBase) GetNotes(activityName string) string { func (kb *KnowledgeBase) GetNotes(activityName string) string {
activity := kb.GetActivity(activityName) activity := kb.GetActivity(activityName)
if activity == nil { if activity == nil {
@@ -302,16 +183,8 @@ func (kb *KnowledgeBase) GetNotes(activityName string) string {
} }
// Validate checks the knowledge base for consistency // Validate checks the knowledge base for consistency
// Checks:
// 1. No circular dependencies in activity constraints
// 2. All referenced dependencies exist
// SOLID: Single Responsibility - validation only, no side effects
// CRAP Score: MEDIUM
// - Complexity: 3 (nested loops + recursion)
// - Repetition: 2 (two separate checks, some code reuse in checkDependencies)
// - Total CRAP: 5 (acceptable for validation logic)
func (kb *KnowledgeBase) Validate() error { func (kb *KnowledgeBase) Validate() error {
// Check for circular dependencies using DFS // Check for circular dependencies
visited := make(map[string]bool) visited := make(map[string]bool)
for _, activity := range kb.Activities { for _, activity := range kb.Activities {
if err := kb.checkDependencies(activity.Name, visited, []string{}); err != nil { if err := kb.checkDependencies(activity.Name, visited, []string{}); err != nil {
@@ -319,7 +192,7 @@ func (kb *KnowledgeBase) Validate() error {
} }
} }
// DRY: Check all dependencies exist in second pass (separate concern from cycle detection) // Check that all dependencies exist
for _, activity := range kb.Activities { for _, activity := range kb.Activities {
for _, dep := range activity.Constraints.Dependencies { for _, dep := range activity.Constraints.Dependencies {
if !kb.HasActivity(dep) { if !kb.HasActivity(dep) {
@@ -331,19 +204,11 @@ func (kb *KnowledgeBase) Validate() error {
return nil return nil
} }
// checkDependencies validates activity dependencies for cycles using DFS // checkDependencies validates activity dependencies for cycles
// Internal helper method for Validate()
// Uses path to build cycle path for error reporting
// CRAP Score: MEDIUM
// - Complexity: 3 (string building, recursion, path tracking)
// - Repetition: 1 (unique DFS logic)
// - Total CRAP: 4 (acceptable for graph traversal)
func (kb *KnowledgeBase) checkDependencies(activityName string, visited map[string]bool, path []string) error { func (kb *KnowledgeBase) checkDependencies(activityName string, visited map[string]bool, path []string) error {
// Check for cycles by detecting if activityName appears in current path // Check for cycles
// This indicates we've visited activityName already in this traversal
for _, p := range path { for _, p := range path {
if p == activityName { if p == activityName {
// Build human-readable cycle description
cycleStr := "" cycleStr := ""
found := false found := false
for _, n := range path { for _, n := range path {
@@ -360,9 +225,8 @@ func (kb *KnowledgeBase) checkDependencies(activityName string, visited map[stri
} }
} }
// Skip if already fully visited (memoization)
if visited[activityName] { if visited[activityName] {
return nil return nil // Already checked this branch
} }
visited[activityName] = true visited[activityName] = true
@@ -370,10 +234,9 @@ func (kb *KnowledgeBase) checkDependencies(activityName string, visited map[stri
activity := kb.GetActivity(activityName) activity := kb.GetActivity(activityName)
if activity == nil { if activity == nil {
return nil // Non-existent activity will be caught in Validate() second pass return nil // Non-existent activity will be caught elsewhere
} }
// Recursively check all dependencies
for _, dep := range activity.Constraints.Dependencies { for _, dep := range activity.Constraints.Dependencies {
if err := kb.checkDependencies(dep, visited, newPath); err != nil { if err := kb.checkDependencies(dep, visited, newPath); err != nil {
return err return err
@@ -383,24 +246,12 @@ func (kb *KnowledgeBase) checkDependencies(activityName string, visited map[stri
return nil return nil
} }
// String returns a human-readable short description of the knowledge base // String returns a human-readable description of the knowledge base
// Implements fmt.Stringer interface for logging
// CRAP Score: LOW (simple string formatting)
func (kb *KnowledgeBase) String() string { func (kb *KnowledgeBase) String() string {
return fmt.Sprintf("KnowledgeBase(v%s, %d activities)", kb.Version, kb.Metadata.TotalActivities) return fmt.Sprintf("KnowledgeBase(v%s, %d activities)", kb.Version, kb.Metadata.TotalActivities)
} }
// PrintSummary generates human-readable documentation of all activities // PrintSummary prints a summary of available activities
// Useful for:
// - CLI output (showing available activities)
// - Documentation generation
// - Debugging knowledge base content
// DRY: Centralizes summary formatting (single point of change)
// SOLID: Single Responsibility - formatting only, no mutations
// CRAP Score: MEDIUM
// - Complexity: 2 (string building, nested loops)
// - Repetition: 1 (unique formatting)
// - Total CRAP: 3
func (kb *KnowledgeBase) PrintSummary() string { func (kb *KnowledgeBase) PrintSummary() string {
summary := fmt.Sprintf("=== Activity Knowledge Base ===\nVersion: %s\nTotal Activities: %d\n\n", kb.Version, kb.Metadata.TotalActivities) summary := fmt.Sprintf("=== Activity Knowledge Base ===\nVersion: %s\nTotal Activities: %d\n\n", kb.Version, kb.Metadata.TotalActivities)
-110
View File
@@ -1,110 +0,0 @@
// Package temporal provides Temporal SDK client initialization and management.
package temporal
import (
"crypto/tls"
"fmt"
"time"
"go.temporal.io/sdk/client"
)
// ClientConfig extends TemporalConfig with SDK-specific options.
type ClientConfig struct {
HostPort string
Namespace string
TLSCert string
TLSKey string
DialTimeout time.Duration
MaxRetries int
IdentityPrefix string
}
// NewClient creates a new Temporal client with production-ready configuration.
//
// Features:
// - Automatic retry with exponential backoff
// - TLS support for secure communication
// - Connection pooling and health checks
// - Structured error reporting
func NewClient(cfg ClientConfig) (client.Client, error) {
if cfg.HostPort == "" {
cfg.HostPort = "temporal-frontend.temporal.svc.cluster.local:7233"
}
if cfg.Namespace == "" {
cfg.Namespace = "default"
}
if cfg.DialTimeout == 0 {
cfg.DialTimeout = 10 * time.Second
}
if cfg.MaxRetries == 0 {
cfg.MaxRetries = 3
}
if cfg.IdentityPrefix == "" {
cfg.IdentityPrefix = "poimen-worker"
}
var tlsConfig *tls.Config
if cfg.TLSCert != "" && cfg.TLSKey != "" {
cert, err := tls.LoadX509KeyPair(cfg.TLSCert, cfg.TLSKey)
if err != nil {
return nil, fmt.Errorf("failed to load TLS credentials: %w", err)
}
tlsConfig = &tls.Config{
Certificates: []tls.Certificate{cert},
}
}
clientOptions := client.Options{
HostPort: cfg.HostPort,
Namespace: cfg.Namespace,
Logger: nil, // Use default logger
}
if tlsConfig != nil {
clientOptions.ConnectionOptions = client.ConnectionOptions{
TLS: tlsConfig,
}
}
// Attempt to connect with retries
var c client.Client
var lastErr error
for attempt := 1; attempt <= cfg.MaxRetries; attempt++ {
var err error
c, err = client.Dial(clientOptions)
if err == nil {
return c, nil
}
lastErr = err
if attempt < cfg.MaxRetries {
backoff := time.Duration(1<<uint(attempt-1)) * time.Second
if backoff > 30*time.Second {
backoff = 30 * time.Second
}
time.Sleep(backoff)
}
}
return nil, fmt.Errorf("failed to connect to Temporal after %d attempts: %w", cfg.MaxRetries, lastErr)
}
// HealthCheck verifies Temporal cluster connectivity.
func HealthCheck(c client.Client, timeout time.Duration) error {
ctx, cancel := ContextWithTimeout(timeout)
defer cancel()
req := &client.CheckHealthRequest{}
_, err := c.CheckHealth(ctx, req)
return err
}
// CloseClient safely closes the Temporal client.
func CloseClient(c client.Client) error {
if c != nil {
c.Close()
}
return nil
}
-62
View File
@@ -1,62 +0,0 @@
package temporal
import (
"testing"
"time"
"github.com/stretchr/testify/assert"
)
func TestClientConfigDefaults(t *testing.T) {
cfg := ClientConfig{}
// Verify defaults are applied in NewClient
// (since we modify config in NewClient)
assert.Equal(t, "", cfg.HostPort)
assert.Equal(t, "", cfg.Namespace)
}
func TestNewClientConnectionFailure(t *testing.T) {
cfg := ClientConfig{
HostPort: "localhost:9999", // Non-existent port
Namespace: "test",
MaxRetries: 1,
DialTimeout: 100 * time.Millisecond,
}
client, err := NewClient(cfg)
assert.Error(t, err)
assert.Nil(t, client)
assert.Contains(t, err.Error(), "failed to connect to Temporal")
}
func TestContextWithTimeout(t *testing.T) {
ctx, cancel := ContextWithTimeout(5 * time.Second)
defer cancel()
assert.NotNil(t, ctx)
select {
case <-ctx.Done():
t.Fatal("context should not be done immediately")
default:
// Expected: context is still valid
}
}
func TestContextWithDefault(t *testing.T) {
ctx, cancel := ContextWithDefault()
defer cancel()
assert.NotNil(t, ctx)
select {
case <-ctx.Done():
t.Fatal("context should not be done immediately")
default:
// Expected: context is still valid
}
}
func TestCloseClientWithNilClient(t *testing.T) {
err := CloseClient(nil)
assert.NoError(t, err)
}
-16
View File
@@ -1,16 +0,0 @@
package temporal
import (
"context"
"time"
)
// ContextWithTimeout creates a context with the given timeout.
func ContextWithTimeout(timeout time.Duration) (context.Context, context.CancelFunc) {
return context.WithTimeout(context.Background(), timeout)
}
// ContextWithDefault creates a context with a default timeout of 10 seconds.
func ContextWithDefault() (context.Context, context.CancelFunc) {
return context.WithTimeout(context.Background(), 10*time.Second)
}
-84
View File
@@ -1,84 +0,0 @@
package temporal
import (
"fmt"
"go.temporal.io/sdk/client"
"go.temporal.io/sdk/worker"
)
// WorkerConfig holds configuration for worker creation.
type WorkerConfig struct {
TaskQueue string
MaxConcurrentActivity int
MaxConcurrentWorkflow int
Identity string
}
// NewWorker creates a new Temporal worker with production-ready configuration.
//
// Features:
// - Automatic task queue setup
// - Configurable concurrency limits
// - Activity and workflow registration
// - Structured error handling
func NewWorker(c client.Client, cfg WorkerConfig) (worker.Worker, error) {
if cfg.TaskQueue == "" {
cfg.TaskQueue = "poimen-taskqueue"
}
if cfg.MaxConcurrentActivity == 0 {
cfg.MaxConcurrentActivity = 10
}
if cfg.MaxConcurrentWorkflow == 0 {
cfg.MaxConcurrentWorkflow = 10
}
if cfg.Identity == "" {
cfg.Identity = "poimen-worker-default"
}
workerOptions := worker.Options{
Identity: cfg.Identity,
MaxConcurrentActivityExecutionSize: cfg.MaxConcurrentActivity,
MaxConcurrentWorkflowTaskExecutionSize: cfg.MaxConcurrentWorkflow,
}
w := worker.New(c, cfg.TaskQueue, workerOptions)
if w == nil {
return nil, fmt.Errorf("failed to create worker for task queue: %s", cfg.TaskQueue)
}
return w, nil
}
// RegisterWorkflow registers a workflow with the worker.
func RegisterWorkflow(w worker.Worker, workflow interface{}) error {
if w == nil {
return fmt.Errorf("worker is nil")
}
w.RegisterWorkflow(workflow)
return nil
}
// RegisterActivity registers an activity with the worker.
func RegisterActivity(w worker.Worker, activity interface{}) error {
if w == nil {
return fmt.Errorf("worker is nil")
}
w.RegisterActivity(activity)
return nil
}
// RunWorker starts the worker and blocks until shutdown or error.
func RunWorker(w worker.Worker) error {
if w == nil {
return fmt.Errorf("worker is nil")
}
return w.Run(worker.InterruptCh())
}
// StopWorker gracefully stops the worker.
func StopWorker(w worker.Worker) {
if w != nil {
w.Stop()
}
}
-55
View File
@@ -1,55 +0,0 @@
package temporal
import (
"testing"
"github.com/stretchr/testify/assert"
)
func TestWorkerConfigDefaults(t *testing.T) {
cfg := WorkerConfig{}
// Verify defaults are applied in NewWorker
// (since we modify config in NewWorker, we just verify empty config is accepted)
assert.Equal(t, "", cfg.TaskQueue)
assert.Equal(t, 0, cfg.MaxConcurrentActivity)
assert.Equal(t, 0, cfg.MaxConcurrentWorkflow)
assert.Equal(t, "", cfg.Identity)
}
func TestRegisterWorkflowWithNilWorker(t *testing.T) {
err := RegisterWorkflow(nil, func() {})
assert.Error(t, err)
assert.Equal(t, "worker is nil", err.Error())
}
func TestRegisterActivityWithNilWorker(t *testing.T) {
err := RegisterActivity(nil, func() {})
assert.Error(t, err)
assert.Equal(t, "worker is nil", err.Error())
}
func TestRunWorkerWithNilWorker(t *testing.T) {
err := RunWorker(nil)
assert.Error(t, err)
assert.Equal(t, "worker is nil", err.Error())
}
func TestStopWorkerWithNilWorker(t *testing.T) {
// Should not panic
StopWorker(nil)
}
func TestWorkerConfigCustomValues(t *testing.T) {
cfg := WorkerConfig{
TaskQueue: "custom-queue",
MaxConcurrentActivity: 20,
MaxConcurrentWorkflow: 30,
Identity: "custom-identity",
}
assert.Equal(t, "custom-queue", cfg.TaskQueue)
assert.Equal(t, 20, cfg.MaxConcurrentActivity)
assert.Equal(t, 30, cfg.MaxConcurrentWorkflow)
assert.Equal(t, "custom-identity", cfg.Identity)
}
+8 -8
View File
@@ -4,19 +4,19 @@ kind: Kustomization
namespace: poimen namespace: poimen
resources: resources:
- worker-deployment.yaml - poimen-application.yaml
- workflow-runner-deployment.yaml
- workflows-deployment.yaml
- git-commit.yaml
# SOPS-encrypted configmap applied separately via KSOPS plugin:
# - configmap.enc.yaml
commonLabels: commonLabels:
app.kubernetes.io/name: poimen app.kubernetes.io/name: poimen
app.kubernetes.io/component: worker app.kubernetes.io/component: worker
images: images:
- name: forgejo.riotpiao.com/rock/poimen-memory
newName: forgejo.riotpiao.com/rock/poimen-memory
newTag: latest
- name: forgejo.riotpiao.com/rock/poimen-workflows - name: forgejo.riotpiao.com/rock/poimen-workflows
newName: forgejo.riotpiao.com/riotpiao-poimen/poimen-workflows newName: forgejo.riotpiao.com/rock/poimen-workflows
newTag: latest
- name: forgejo.riotpiao.com/rock/poimen-frontend
newName: forgejo.riotpiao.com/rock/poimen-frontend
newTag: latest newTag: latest
-140
View File
@@ -1,140 +0,0 @@
apiVersion: v1
kind: ConfigMap
metadata:
name: poimen-workflow-runner-config
namespace: poimen
data:
TEMPORAL_HOSTPORT: "temporal-frontend.temporal.svc.cluster.local:7233"
TEMPORAL_NAMESPACE: "default"
LOG_LEVEL: "info"
---
apiVersion: apps/v1
kind: Deployment
metadata:
name: poimen-workflow-runner
namespace: poimen
labels:
app: poimen-workflow-runner
component: workflow-runner
spec:
replicas: 1
strategy:
type: Recreate
selector:
matchLabels:
app: poimen-workflow-runner
template:
metadata:
labels:
app: poimen-workflow-runner
component: workflow-runner
annotations:
prometheus.io/scrape: "true"
prometheus.io/port: "8081"
prometheus.io/path: "/metrics"
spec:
serviceAccountName: poimen-workflow-runner
securityContext:
runAsNonRoot: true
runAsUser: 1000
containers:
- name: workflow-runner
image: forgejo.riotpiao.com/riotpiao-poimen/poimen-workflows:latest
imagePullPolicy: IfNotPresent
command: ["./poimen-workflow-runner"]
ports:
- name: health
containerPort: 8081
protocol: TCP
env:
- name: TEMPORAL_HOSTPORT
valueFrom:
configMapKeyRef:
name: poimen-workflow-runner-config
key: TEMPORAL_HOSTPORT
- name: TEMPORAL_NAMESPACE
valueFrom:
configMapKeyRef:
name: poimen-workflow-runner-config
key: TEMPORAL_NAMESPACE
- name: LOG_LEVEL
valueFrom:
configMapKeyRef:
name: poimen-workflow-runner-config
key: LOG_LEVEL
- name: ANTHROPIC_API_KEY
valueFrom:
secretKeyRef:
name: poimen-secrets
key: anthropic-api-key
- name: MEMORY_SERVICE_URL
value: "http://poimen-memory.poimen.svc.cluster.local:8080"
- name: MEMORY_SERVICE_JWT_TOKEN
valueFrom:
secretKeyRef:
name: poimen-secrets
key: memory-service-jwt
resources:
requests:
cpu: 250m
memory: 256Mi
limits:
cpu: 500m
memory: 512Mi
livenessProbe:
httpGet:
path: /health/live
port: 8081
initialDelaySeconds: 30
periodSeconds: 10
timeoutSeconds: 5
failureThreshold: 3
readinessProbe:
httpGet:
path: /health/ready
port: 8081
initialDelaySeconds: 10
periodSeconds: 5
timeoutSeconds: 5
failureThreshold: 2
securityContext:
allowPrivilegeEscalation: false
readOnlyRootFilesystem: true
capabilities:
drop:
- ALL
volumeMounts:
- name: tmp
mountPath: /tmp
volumes:
- name: tmp
emptyDir:
sizeLimit: 100Mi
---
apiVersion: v1
kind: ServiceAccount
metadata:
name: poimen-workflow-runner
namespace: poimen
labels:
app: poimen-workflow-runner
---
apiVersion: v1
kind: Service
metadata:
name: poimen-workflow-runner
namespace: poimen
labels:
app: poimen-workflow-runner
spec:
type: ClusterIP
ports:
- port: 8081
targetPort: 8081
protocol: TCP
name: health
selector:
app: poimen-workflow-runner
-58
View File
@@ -1,58 +0,0 @@
apiVersion: apps/v1
kind: Deployment
metadata:
name: poimen-workflows
namespace: poimen
labels:
app.kubernetes.io/name: poimen
app.kubernetes.io/component: worker
spec:
replicas: 2
selector:
matchLabels:
app: poimen-workflows
app.kubernetes.io/name: poimen
app.kubernetes.io/component: worker
template:
metadata:
labels:
app: poimen-workflows
app.kubernetes.io/name: poimen
app.kubernetes.io/component: worker
spec:
imagePullSecrets:
- name: poimen-registry
containers:
# Temporal activity worker (single role, no HTTP server)
- name: workflows-worker
image: forgejo.riotpiao.com/rock/poimen-workflows:latest
imagePullPolicy: Always
command: ["/app/worker"]
env:
- name: DATABASE_URL
valueFrom:
secretKeyRef:
name: poimen-db-credentials
key: workflows-url
- name: TEMPORAL_HOSTPORT
valueFrom:
configMapKeyRef:
name: poimen-config
key: temporal-hostport
- name: TEMPORAL_NAMESPACE
valueFrom:
configMapKeyRef:
name: poimen-config
key: temporal-namespace
- name: MEMORY_SERVICE_URL
valueFrom:
configMapKeyRef:
name: poimen-config
key: memory-service-url
resources:
requests:
memory: "512Mi"
cpu: "500m"
limits:
memory: "2Gi"
cpu: "2000m"
BIN
View File
Binary file not shown.
-35
View File
@@ -1,35 +0,0 @@
package workflow
import (
"time"
"go.temporal.io/sdk/workflow"
"github.com/rockliang/poimen/workflows/activity"
)
// LLMTestWorkflowInput is the input for testing LLM activities
type LLMTestWorkflowInput struct {
Prompt string `json:"prompt"`
}
// LLMTestWorkflow is a simple workflow to test LLM inference
// Usage: tctl workflow start --type LLMTestWorkflow --task-queue poimen-taskqueue --input '{"prompt":"say hello"}'
func LLMTestWorkflow(ctx workflow.Context, input LLMTestWorkflowInput) (string, error) {
// Call the LLM inference activity
opts := workflow.ActivityOptions{
StartToCloseTimeout: 60 * time.Second,
}
actCtx := workflow.WithActivityOptions(ctx, opts)
actInput := activity.LLMInferenceInput{
Model: "reasoning",
UserPrompt: input.Prompt,
}
var result activity.LLMInferenceOutput
err := workflow.ExecuteActivity(actCtx, "LLMInferenceActivity", actInput).Get(actCtx, &result)
if err != nil {
return "", err
}
return result.Response, nil
}