Files
poimen-workflows/internal/memory
Test 5ef14ad5ec feat(memory): add Temporal activities integration for memory service
- Implement 12 Temporal activities for memory operations
- Activities: create, update, search, context, diagnose, analyze, document
- Add activity registration and worker setup
- Full retry/timeout configuration with observability
- Include workflow patterns and examples
- All tests passing (23/23)

Documentation:
- MEMORY_INTEGRATION.md: High-level integration guide
- MEMORY_ACTIVITIES.md: Complete activities reference
- REGISTERED_ACTIVITIES.md: Registry and calling conventions
2026-08-29 21:49:24 -07:00
..

Poimen Memory Service Integration

Go client for Poimen Memory Service with Temporal Activities. Provides create, update, retrieve, and context operations for knowledge management with full workflow integration, retry logic, and observability.

Overview

Memory service endpoints:

  • POST /memory/ingest — Create knowledge records (L1/L2/reference)
  • POST /memory/query — Search knowledge (hybrid semantic+lexical)
  • POST /memory/context — Retrieve context (three-tier: signature → vector → reference)
  • GET /memory/vault — Browse vault files
  • GET /health — Health check

Temporal Activities

All operations are Temporal Activities with:

  • Automatic retries (3 attempts by default)
  • Timeout handling (per operation)
  • Heartbeat monitoring
  • Logging + observability
  • Workflow integration

Activity List

Activity Purpose
CreateKnowledgeActivity Create L1/L2/reference records
UpdateKnowledgeActivity Update existing knowledge
SearchKnowledgeActivity Search hybrid (semantic+lexical)
GetContextActivity Retrieve three-tier context
GetVaultActivity Browse vault files
HealthCheckActivity Check service health
LearnFromExecutionActivity Learn from task results
DiagnoseIssueActivity Diagnose tool/task issues
AnalyzeErrorActivity Analyze errors, find solutions
DocumentDecisionActivity Record workflow decisions
SearchAndApplyActivity Search and apply knowledge
RefreshMemoryActivity Periodic memory refresh

Register Activities

In worker setup:

service := memory.NewService(baseURL, token, project)
memory.RegisterMemoryActivities(w, service)

Use in Workflows

// Simple activity call
id, err := memory.ExecuteCreateKnowledge(
	ctx,
	&memory.KnowledgeRecord{
		Level: "L1",
		Content: "...",
	},
	nil, // Use default options
)

// Custom retry policy
options := &memory.ActivityOptions{
	RetryAttempts: 5,
	RetryBackoff: time.Second,
}
recommendations, err := memory.ExecuteDiagnoseIssue(ctx, "kubectl", "pod-crash", options)

Installation

Import package:

import "github.com/poimen/workflows/internal/memory"

Workflow Integration

Example 1: Learning Workflow

// Learn from task execution
func LearningWorkflow(ctx workflow.Context, taskID string) (string, error) {
	// Execute task (placeholder)
	result := fmt.Sprintf("Task %s completed successfully", taskID)

	// Learn from result
	knowledgeID, err := memory.ExecuteLearnFromExecution(
		ctx,
		taskID,
		result,
		[]string{"success", taskID},
		nil, // Default retry policy
	)
	return knowledgeID, err
}

Example 2: Diagnostic Workflow

// Diagnose issue using memory service
func DiagnosticWorkflow(ctx workflow.Context, tool, issue string) ([]string, error) {
	recommendations, err := memory.ExecuteDiagnoseIssue(
		ctx,
		tool,
		issue,
		&memory.ActivityOptions{
			RetryAttempts: 3,
			RetryBackoff: time.Second,
		},
	)
	return recommendations, err
}

Example 3: Error Recovery

// Analyze error and find recovery path
func ErrorRecoveryWorkflow(ctx workflow.Context, errorMsg string) ([]string, error) {
	// Analyze error
	records, err := memory.ExecuteAnalyzeError(ctx, errorMsg, nil)
	if err != nil {
		return nil, err
	}

	// Extract recovery steps
	recovery := make([]string, 0)
	for _, record := range records {
		if record.Level == "L1" { // High confidence
			recovery = append(recovery, record.Content)
		}
	}
	return recovery, nil
}

Example 4: Multi-Step Decision Workflow

// Get context, make decision, document it
func ContextualDecisionWorkflow(ctx workflow.Context, tool, task, decision string) (string, error) {
	// Get context (three-tier retrieval)
	svcCtx, err := memory.ExecuteGetContext(ctx, tool, task, 8192, nil)
	if err != nil {
		return "", err
	}

	// Make decision based on context
	reasoning := fmt.Sprintf("Based on %d lessons (tier %d)", len(svcCtx.Lessons), svcCtx.Tier)

	// Document decision
	docID, err := memory.ExecuteDocumentDecision(ctx, tool, decision, reasoning, nil)
	return docID, err
}

Usage

Client (Low-Level)

package main

import (
	"context"
	"fmt"
	"log"
	
	"github.com/poimen/workflows/internal/memory"
)

func main() {
	// Create client
	client := memory.NewClient(
		"http://localhost:8080",
		"your-jwt-token",
	)

	ctx := context.Background()

	// Ingest knowledge
	resp, err := client.Ingest(ctx, &memory.IngestRequest{
		Project: "poimen",
		Source:  "workflow://task-123",
		Kind:    "L1",
		Text:    "Pod CrashLoopBackOff: check logs with kubectl logs",
		Metadata: map[string]interface{}{
			"topic": "kubernetes",
			"task_id": "debug-pod",
		},
	})
	if err != nil {
		log.Fatal(err)
	}
	fmt.Printf("Created: %s (SHA256: %s)\n", resp.ID, resp.SHA256)

	// Search knowledge
	query, err := client.Query(ctx, &memory.QueryRequest{
		Project: "poimen",
		Query:   "fix pod crash loop",
		Limit:   5,
		Floor:   0.6, // minimum relevance
	})
	if err != nil {
		log.Fatal(err)
	}
	for _, r := range query.Results {
		fmt.Printf("%s (score: %.2f): %s\n", r.Level, r.Score, r.Text)
	}

	// Get context (three-tier retrieval)
	ctxResp, err := client.Context(ctx, &memory.ContextRequest{
		Project:         "poimen",
		Tool:            "kubectl",
		Task:            "debug-pod",
		SignatureSource: "error_log",
		Budget:          8192,
	})
	if err != nil {
		log.Fatal(err)
	}
	fmt.Printf("Context tier: %d\n", ctxResp.Tier)
	for _, lesson := range ctxResp.Lessons {
		fmt.Printf("- [Tier %d] %s: %.2f\n", lesson.Tier, lesson.Level, lesson.Score)
	}

	// Browse vault
	vault, err := client.Vault(ctx, "poimen")
	if err != nil {
		log.Fatal(err)
	}
	fmt.Printf("Total records: %d\n", vault.TotalRecords)
	for _, f := range vault.Files {
		fmt.Printf("- %s (%s, %d records)\n", f.Path, f.Level, f.RecordCount)
	}
}

Service (High-Level)

package main

import (
	"context"
	"log"
	
	"github.com/poimen/workflows/internal/memory"
)

func main() {
	// Create service
	svc := memory.NewService(
		"http://localhost:8080",
		"your-jwt-token",
		"poimen", // project
	)

	ctx := context.Background()

	// Create knowledge
	id, err := svc.CreateKnowledge(ctx, &memory.KnowledgeRecord{
		Level:   "L1",
		Title:   "Pod Debugging",
		Content: "To debug CrashLoopBackOff: kubectl logs <pod>",
		Source:  "workflow://debug-task",
	})
	if err != nil {
		log.Fatal(err)
	}
	log.Printf("Created knowledge: %s\n", id)

	// Update knowledge (re-ingest with same ID)
	id, err = svc.UpdateKnowledge(ctx, &memory.KnowledgeRecord{
		ID:      id,
		Level:   "L2",
		Content: "Advanced debugging: check events, describe pod, check node status",
	})
	if err != nil {
		log.Fatal(err)
	}
	log.Printf("Updated knowledge: %s\n", id)

	// Retrieve knowledge
	records, err := svc.RetrieveKnowledge(ctx, "kubernetes pod debugging", &memory.RetrievalOptions{
		LevelFilter: []string{"L1", "L2"},
		Limit:       10,
		Floor:       0.7,
	})
	if err != nil {
		log.Fatal(err)
	}
	for _, rec := range records {
		log.Printf("- %s: %s\n", rec.ID, rec.Content)
	}

	// Retrieve context
	svcCtx, err := svc.RetrieveContext(ctx, "kubectl", "debug-pod", 8192)
	if err != nil {
		log.Fatal(err)
	}
	log.Printf("Context tier: %d (%d lessons, %d skills)\n", 
		svcCtx.Tier, len(svcCtx.Lessons), len(svcCtx.Skills))
	for _, skill := range svcCtx.Skills {
		log.Printf("  - %s: %s\n", skill.Name, skill.Why)
	}

	// Get vault
	files, err := svc.GetVault(ctx)
	if err != nil {
		log.Fatal(err)
	}
	log.Printf("Vault has %d files\n", len(files))

	// Check health
	if svc.IsHealthy(ctx) {
		log.Println("Memory service is healthy")
	}
}

API Reference

Client Methods

Ingest(ctx, req) → IngestResponse, error

Create knowledge record.

Request:

&IngestRequest{
	Project: "poimen",
	Source: "workflow://task-id",
	Kind: "L1",  // L1|L2|reference
	Text: "knowledge content",
	Metadata: map[string]interface{}{...},
}

Response:

{
	ID: "chunk-abc123",
	SHA256: "de12cd34ef56...",
	QueueStatus: "pending",  // Async processing
	IdempotencyID: "sess-123:0",
}

Query(ctx, req) → QueryResponse, error

Search knowledge (hybrid semantic + lexical).

Request:

&QueryRequest{
	Project: "poimen",
	Query: "fix kubernetes pod crash",
	LevelFilter: []string{"L1", "L2"},  // Optional
	Floor: 0.6,  // Minimum relevance
	Limit: 10,
	Scope: "all",  // learned|reference|all
}

Response:

{
	Query: "...",
	Results: []QueryResult{
		{
			ID: "chunk-abc123",
			Level: "L1",
			Score: 0.992,
			SemanticScore: 1.0,
			LexicalScore: 0.98,
			Text: "...",
			Breadcrumb: "kubernetes.md > Troubleshooting",
			Source: "transcript://session-123",
		},
		...
	},
	TotalHits: 127,
	SearchTimeMS: 145,
}

Context(ctx, req) → ContextResponse, error

Retrieve context for tool/task (three-tier retrieval: signature → vector → reference).

Request:

&ContextRequest{
	Project: "poimen",
	Tool: "kubectl",
	Task: "debug-pod",
	SignatureSource: "failure_log",  // Where to find signature
	Scope: "tool_context",
	Budget: 8192,  // Max response bytes
}

Response:

{
	Tier: 1,  // Highest tier with results
	Lessons: []ContextLesson{
		{
			Tier: 1,
			Level: "L1",
			Score: 1.0,
			Text: "Pod in CrashLoopBackOff: check logs",
			MatchedKind: "signature",
			SeenCount: 23,
			LastSeen: "2025-01-28T15:30:00Z",
		},
		...
	},
	Skills: []ContextSkill{
		{
			Name: "diagnose-pod-failure",
			Why: "Tier-1 signature matched",
		},
	},
	Budget: {
		Requested: 8192,
		Used: 4156,
		Dropped: 0,
		Degradation: nil,
	},
}

Vault(ctx, project) → VaultResponse, error

Browse vault files.

Response:

{
	Project: "poimen",
	Files: []VaultFile{
		{
			Path: "kubernetes/debugging.md",
			Title: "Debugging",
			Level: "L1",
			UpdatedAt: "2025-01-28T10:00:00Z",
			RecordCount: 23,
		},
		...
	},
	TotalRecords: 542,
}

Health(ctx) → bool, error

Check service health.

Service Methods

Service provides higher-level operations:

  • CreateKnowledge(ctx, record) → id, error
  • UpdateKnowledge(ctx, record) → id, error
  • RetrieveKnowledge(ctx, query, opts) → []KnowledgeRecord, error
  • RetrieveContext(ctx, tool, task, budget) → *ServiceContext, error
  • GetVault(ctx) → []VaultInfo, error
  • IsHealthy(ctx) → bool

Error Handling

// All operations return (result, error)
resp, err := client.Ingest(ctx, req)
if err != nil {
	// Possible errors:
	// - Request marshal/network errors
	// - 401 Unauthorized: Missing/invalid JWT
	// - 403 Forbidden: Token lacks capability
	// - 429 Too Many Requests: Rate limit exceeded
	// - 409 Conflict: Duplicate (same idempotency key within 24h)
	// - 503 Service Unavailable: Database unreachable
	log.Fatalf("ingest failed: %v", err)
}

Authentication

Pass JWT bearer token to NewClient/NewService:

// Get token from Authentik
token := "eyJ0eXAiOiJKV1QiLCJhbGc..."

client := memory.NewClient(baseURL, token)

Token must have capability:

  • memory:read — for Query, Context, Vault
  • memory:write — for Ingest

Rate Limits

Per JWT identity:

  • Ingest: 100/hour
  • Query: 1000/hour
  • Context: 100/hour

Exceed limit → 429 Too Many Requests.

Deployment

Memory service endpoints (k8s):

  • Service: memory-service.poimen.svc.cluster.local:8080
  • Ingress: https://memory.riotpiao.com (external)

Environment:

baseURL := "http://memory-service.poimen.svc.cluster.local:8080"
token := os.Getenv("MEMORY_SERVICE_TOKEN")
svc := memory.NewService(baseURL, token, "poimen")

Testing

Run tests:

go test ./internal/memory -v

Mock server example in client_test.go and service_test.go.