Files
poimen-memory/docs/LLM_INFERENCE_ACTIVITY_INTEGRATION.md
rock 4c275525e9 Implement LLMInferenceActivity integration for Temporal workflows
Workflow Input Structure:
  ├─ question: User content for reasoning
  ├─ project: Project ID for scoping
  ├─ operations: Flags for link_entities, infer_facts, reason_query, summarize
  └─ llm_activity: Configuration for LLMInferenceActivity
       ├─ model: Selected based on complexity (reasoning|ornith:35b|qwen2.5:3b)
       ├─ system_prompt: Task-specific instruction (Zep-backed)
       ├─ user_prompt: Content to process
       ├─ temperature: 0.7 (reasoning) or 0.5 (validation)
       └─ max_tokens: 2048 (reasoning) or 512 (validation)

Model Selection:
  ├─ reason_query=true, summarize=true → reasoning (DeepSeek-R1, complex)
  ├─ reason_query=true, summarize=false → ornith:35b (medium)
  └─ reason_query=false → qwen2.5:3b (fast, <100ms)

System Prompts (handlers/llm_prompts.rs):
  ├─ entity_extraction_system_prompt(): Extract entities + relationships + facts
  ├─ reasoning_system_prompt(): Step-by-step reasoning + answers
  ├─ agent_capability_validation_prompt(): Validate agent capabilities
  └─ fact_validation_system_prompt(): Detect contradictions

Workflow Activity Execution:
  ├─ Temporal receives workflow input with llm_activity config
  ├─ ReasoningWorkflow orchestrates:
  │  ├─ Activity 1: RetrieveMemory (optional context)
  │  ├─ Activity 2: LLMInferenceActivity (calls /v1/chat/completions via gateway)
  │  │   └─ Retries: 3× with backoff (2s, 4s, 8s)
  │  │   └─ Timeout: 120s
  │  │   └─ JWT propagation: Authorization: Bearer header
  │  ├─ Activity 3: PersistResults (save to memory_entity/memory_edge)
  │  └─ Activity 4: SummarizeFindings (return results)
  ├─ Memory handler polls DESCRIBE_WORKFLOW (30× with 100ms delay, 3s timeout)
  └─ Returns ReasoningResult with answers, confidence, reasoning_steps

Changes:
  ├─ execute_reasoning_workflow(): Build llm_activity config with model selection
  ├─ select_llm_model(): Choose model based on operation complexity
  ├─ build_system_prompt(): Use Zep-inspired prompts for reasoning
  ├─ handlers/llm_prompts.rs: Centralized prompt templates (5 system + 4 user builders)
  ├─ AgentInitialization: Include llm_activity for capability validation
  └─ Fixed duplicate extract_jwt_token call in agent_handler.rs

Activity Contract:
  ├─ Workflow input includes llm_activity block
  ├─ Temporal passes to LLMInferenceActivity
  ├─ Activity substitutes {{ previous_output }} template variables
  ├─ Activity calls POST /v1/chat/completions with JWT header
  ├─ Activity returns { response, model, stop_reason, tokens_used }
  ├─ PersistResults activity stores results to DB
  └─ Workflow returns: question, answers[], confidence, reasoning_steps[]

Tests Added:
  + 14 new tests in llm_prompts.rs (prompt validation, user prompt builders)

Compilation: 
2026-09-05 00:52:30 -07:00

13 KiB
Raw Permalink Blame History

LLM Inference Activity Integration

Overview

Memory service integrates with Temporal's LLMInferenceActivity for LLM-powered reasoning. The activity handles:

  • Template variable substitution ({{ previous_output.field }})
  • LLM API calls via gateway (POST /v1/chat/completions)
  • Retry logic (exponential backoff: 2s, 4s, 8s)
  • JWT token propagation (Bearer header)
  • Timeout management (120s per call)

Architecture

Memory Handler (this service)
  ↓
SynthesisClient.execute_workflow()
  ↓
POST https://api.riotpiao.com/workflow (with JWT)
  ↓
Temporal Workflow Executor
  ↓
ReasoningWorkflow (defined in homelab-frontend)
  ├─ Activity 1: RetrieveMemory (fetch context)
  ├─ Activity 2: LLMInferenceActivity
  │    ├─ model: "reasoning" | "ornith:35b" | "qwen2.5:3b"
  │    ├─ system_prompt: "You are a knowledge extraction expert..."
  │    ├─ user_prompt: "Extract entities from: {{ previous_output.memory }}"
  │    ├─ auth_token: "{{ header.authorization }}" (from Memory call)
  │    ├─ temperature: 0.7
  │    ├─ max_tokens: 2048
  │    ↓
  │    Calls: POST api.riotpiao.com/v1/chat/completions
  │      + Header: Authorization: Bearer {jwt}
  │      + Retries: 3× with backoff (2s, 4s, 8s)
  │      + Timeout: 120s
  │    ↓
  │    Returns: { response, model, stop_reason, tokens_used }
  │
  ├─ Activity 3: PersistResults
  │    ├─ Extract entities from LLM response
  │    ├─ Insert into memory_entity table
  │    ├─ Insert into memory_edge table
  │    └─ Link in temporal_workflow_links table
  │
  └─ Activity 4: SummarizeFindings
       └─ Return reasoning result

Workflow completes
  ↓
Poll DESCRIBE_WORKFLOW
  ↓
Return result to Memory handler
  ↓
Return to user

Request/Response Flow

1. Memory Handler Initiates Reasoning

// From handlers/unified_synthesis.rs

let workflow_input = json!({
    "question": "Extract entities from this text...",
    "project": "poimen",
    "operations": {
        "link_entities": true,
        "infer_facts": true,
        "reason_query": true,
        "summarize": false
    }
});

let workflow_req = WorkflowBuilder::new("ReasoningWorkflow")
    .with_input(workflow_input)
    .build();

// Calls: POST /workflow with JWT
let response = client.execute_workflow(workflow_req).await?;
// Response: { "data": { "workflow_id": "...", "run_id": "..." } }

2. Workflow START_WORKFLOW Request

{
  "action": "START_WORKFLOW",
  "namespace": "poimen",
  "payload": {
    "workflow_id": "reasoning-abc123",
    "workflow_type": "ReasoningWorkflow",
    "task_queue": "synthesis",
    "input": {
      "question": "Extract entities from this text...",
      "project": "poimen",
      "operations": { ... }
    }
  }
}

3. LLMInferenceActivity Input (Internal)

Temporal constructs this (not Memory's responsibility):

{
  "type": "llm-inference",
  "model": "reasoning",
  "system_prompt": "You are a knowledge extraction expert. Extract all entities, relationships, and facts.",
  "user_prompt": "Extract from: {{ workflow.input.question }}",
  "temperature": 0.7,
  "max_tokens": 2048,
  "auth_token": "{{ workflow.auth_context.jwt }}"
}

4. LLMInferenceActivity Execution

Activity backend (homelab-frontend):

  1. Template substitution:

    user_prompt: "Extract from: Extract entities from this text..."
    auth_token: "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9..."
    
  2. LLM API call:

    POST https://api.riotpiao.com/v1/chat/completions
    Authorization: Bearer eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9...
    Content-Type: application/json
    
    {
      "model": "reasoning",
      "messages": [
        { "role": "system", "content": "You are a knowledge extraction expert..." },
        { "role": "user", "content": "Extract from: Extract entities from this text..." }
      ],
      "temperature": 0.7,
      "max_tokens": 2048
    }
    
  3. LLM Response:

    {
      "choices": [
        {
          "message": {
            "content": "Entities found:\n1. Entity: 'Kubernetes' (Technology)\n2. Entity: 'Docker' (Technology)..."
          }
        }
      ],
      "usage": { "prompt_tokens": 50, "completion_tokens": 200, "total_tokens": 250 }
    }
    
  4. Activity Output:

    {
      "response": "Entities found:\n1. Entity: 'Kubernetes' (Technology)\n2. Entity: 'Docker' (Technology)...",
      "model": "reasoning",
      "stop_reason": "stop_sequence",
      "tokens_used": 250
    }
    
  5. Retry Logic (if LLM call fails):

    Attempt 1: Failed (network timeout)
      → Wait 2 seconds
    Attempt 2: Failed (rate limited, 429)
      → Wait 4 seconds
    Attempt 3: Failed (model overloaded)
      → Workflow error recorded
      → Fallback: proceed with best-effort result or fail workflow
    

5. PersistResults Activity (Custom)

Temporal's custom activity in homelab-frontend:

{
  "type": "persist-results",
  "input": {
    "workflow_id": "reasoning-abc123",
    "llm_response": "Entities found:\n1. Kubernetes (Technology)...",
    "project": "poimen"
  }
}

Activity implementation:

1. Parse LLM response
2. Extract entities/facts
3. INSERT INTO memory_entity (name, summary, entity_type, contributed_by)
4. INSERT INTO memory_edge (source, relation, target)
5. INSERT INTO temporal_workflow_links (workflow_id, run_id, entity_id)
6. Return: { "entities_count": 2, "edges_count": 3 }

6. Workflow DESCRIBE_WORKFLOW Poll

Memory handler polls periodically:

{
  "action": "DESCRIBE_WORKFLOW",
  "namespace": "poimen",
  "payload": {
    "workflow_id": "reasoning-abc123"
  }
}

Response (while running):

{
  "data": {
    "workflow_id": "reasoning-abc123",
    "status": "RUNNING",
    "last_update": "2025-01-30T10:05:00Z"
  }
}

Response (when complete):

{
  "data": {
    "workflow_id": "reasoning-abc123",
    "status": "COMPLETED",
    "result": {
      "question": "Extract entities from this text...",
      "answers": [
        "Entities: Kubernetes, Docker",
        "Relationships: Kubernetes uses Docker"
      ],
      "confidence": 0.92,
      "reasoning_steps": [
        "Extracted all entities using NER",
        "Identified entity types",
        "Built relationship graph"
      ],
      "entities_persisted": 2,
      "edges_persisted": 3
    }
  }
}

JWT Token Flow

Header Propagation

Memory Handler Request:

POST /memory/synthesis
Authorization: Bearer eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJzdWIiOiJ1c2VyMTIzIn0...
Content-Type: application/json

{
  "project": "poimen",
  "content": "Extract entities...",
  "reason_query": true
}

Extract in Handler:

let jwt = extract_jwt_token(&req)?;  // "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9..."

let client = SynthesisClient::new(
    "https://api.riotpiao.com".to_string(),
    jwt,  // ← Stored in client
);

POST /workflow with JWT:

POST https://api.riotpiao.com/workflow
Authorization: Bearer eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9...

{
  "action": "START_WORKFLOW",
  "namespace": "poimen",
  "payload": { ... }
}

Temporal Workflow with JWT:

ReasoningWorkflow receives:
  - workflow input (question, project, operations)
  - auth context (JWT from request header)

LLMInferenceActivity:
  auth_token = "{{ workflow.auth_context.jwt }}"
  
Activity calls LLM with:
  Authorization: Bearer eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9...

JWT Validation at LLM:

homelab-frontend proxy checks:
  1. Signature valid (signed by Authentik)
  2. Not expired
  3. Has "llm:inference" capability
  
If valid:
  → Forward to LLM backend (reasoning/ollama/etc)
  
If invalid:
  → 401 Unauthorized
  → Activity retry or fail

Model Selection

Available Models

Model Use Case Speed Cost Max Tokens
reasoning Complex analysis, entity extraction Slow (500-1000ms) Free 4096
ornith:35b General reasoning Medium (300-500ms) Free 2048
ornith:13b Fast reasoning Fast (100-200ms) Free 2048
qwen2.5:3b Quick tasks Fastest (<100ms) Free 1024

Selection Strategy

// From handlers/unified_synthesis.rs

let model = match body.operations.reason_query {
    true => match body.operations.summarize {
        true => "reasoning",      // Complex: extract + reason + summarize
        false => "ornith:35b",    // Medium: extract + reason
    },
    false => "qwen2.5:3b",        // Quick: only linking/inference (no reasoning)
};

Error Handling

Retry Behavior

LLMInferenceActivity automatically retries:

Attempt 1: Failed
  Error: ConnectionError (network issue)
  Backoff: 2 seconds
  
Attempt 2: Failed
  Error: HTTPError 429 (rate limited)
  Backoff: 4 seconds
  
Attempt 3: Failed
  Error: HTTPError 500 (backend overload)
  → Workflow error recorded
  → No further retries
  
Result:
  {
    "success": false,
    "error": "Max retries exceeded after 3 attempts",
    "last_error": "HTTPError 500 from LLM backend"
  }

Terminal Errors (No Retry)

"Model not found: xyz"
  → Immediate failure (no retry)
  → Activity returns error
  → Workflow fails
  
"Context length exceeded"
  → Immediate failure (no retry)
  → Activity returns error
  → Workflow fails
  
"Invalid auth token"
  → Immediate failure (retry won't help)
  → Activity returns 401
  → Workflow fails

Workflow Error Handling

// From handlers/workflow_poller.rs

match poll_workflow_until_complete(...).await {
    Ok(("COMPLETED", Some(result))) => {
        // Parse result into ReasoningResult
        Ok(ReasoningResult { ... })
    }
    Ok(("COMPLETED", None)) => {
        // Workflow succeeded but no result
        Err("Workflow completed without result")
    }
    Ok((status, _)) => {
        // Unexpected status
        Err(format!("Unexpected workflow status: {}", status))
    }
    Err(e) => {
        // Workflow failed or polling timeout
        Err(e)
    }
}

Template Variables

LLMInferenceActivity supports Handlebars-style templates:

{
  "user_prompt": "Extract entities from: {{ previous_output.memory }}"
}

Available Variables

{{ workflow.input.field }}      // Input from ReasoningWorkflow
{{ previous_output.field }}     // Output from prior activity
{{ workflow.auth_context.jwt }} // JWT from request header
{{ workflow.execution_id }}     // Workflow execution ID

Example Substitution

Before: "Extract from: {{ previous_output.memory }}"
Memory variable: "Kubernetes is a container orchestrator"

After: "Extract from: Kubernetes is a container orchestrator"

Performance Considerations

Latency Budget

ReasoningWorkflow Latency Breakdown:

RetrieveMemory activity:     ~50-100ms
  ↓
LLMInferenceActivity:        ~500-1000ms (reasoning model)
  ├─ Template substitution:  ~10ms
  ├─ LLM API call:           ~400-900ms
  └─ Response parsing:       ~5ms
  ↓
PersistResults activity:     ~100-200ms
  ├─ Parse LLM response:     ~10ms
  ├─ Extract entities:       ~20ms
  └─ DB inserts:             ~70-170ms
  ↓
SummarizeFindings activity:  ~50ms
  ↓
Total: ~750-1350ms (1.3 seconds typical)

Memory handler poll overhead:
  ├─ 30 polls × 100ms delay: 3000ms
  └─ So total time: ~4-5 seconds (with polling)

Optimization

  1. Use faster model for quick tasks:

    if body.content.len() < 500 {
        model = "qwen2.5:3b";  // Fast
    } else {
        model = "reasoning";   // Accurate
    }
    
  2. Cache frequent queries:

    If same question asked twice:
      1st time: Call workflow → 4-5 seconds
      2nd time: Cache hit → <1ms
    
  3. Batch processing (if needed):

    Use LLMBatchInferenceActivity:
      - 3 prompts: ~1500ms (serial)
      - vs 3 separate calls: ~4500ms (sequential)
    

Production Checklist

JWT token extraction working SynthesisClient.execute_workflow() wired Workflow polling implemented (30 retries, 100ms interval, 3s timeout) Error handling for workflow failures Model selection strategy chosen Database persistence (memory_entity, memory_edge, temporal_workflow_links) Retry logic understood (Activity retries handled by Temporal) Token validation (Authentik checks signature + expiration)

TODO (Phase 6.5):

  • Verify LLMInferenceActivity input format with homelab-frontend
  • Test end-to-end workflow execution
  • Monitor actual latency (should be ~4-5s with polling)
  • Add metrics/tracing for workflow lifecycle
  • Document production SLA (e.g., 99% success within 10 seconds)
  • Set up alerting for workflow failures

References

  • Temporal Workflow API: /Users/rockliang/workplace/homelab-frontend/API.md (lines 789-1000)
  • LLMInferenceActivity: Supports template variables, retry, JWT propagation
  • Memory Service Integration: This document
  • Production Code: handlers/unified_synthesis.rs, handlers/workflow_poller.rs

Status: Architecture complete, ready for integration testing and production deployment.