Compare commits

..
Author SHA1 Message Date
rock 749af25998 Revert "test: verify embedding response parsing against real service format"
CI / CI (pull_request) Successful in 11m51s
This reverts commit acfaa7efe6.
2026-09-14 08:46:27 +09:00
rock acfaa7efe6 test: verify embedding response parsing against real service format
- 6 parsing tests for EmbeddingResponse struct
- test_parse_real_embedding_response: exact format from embeddings-predictor
- test_parse_768_dim_response: full 768-dim vector
- test_parse_multi_input_response: array input returns multiple embeddings
- test_parse_embedding_error_response: error format
- test_parse_html_fails_gracefully: HTML error page correctly rejected
- Confirms: parsing is correct, 'expected ident' error is non-JSON response
2026-09-14 08:25:55 +09:00
rock 57c61b522b feat: unexpected error counter + user_id in error logs
CI / CI (pull_request) Successful in 12m31s
- ERROR_UNEXPECTED_TOTAL: global unexpected error counter
- ERROR_UNEXPECTED_INGEST/QUERY/CONTEXT: per-endpoint unexpected errors
- All 500 error paths now increment unexpected counter
- Error logs include user_id for customer association:
  tracing::error!(user_id = claims.sub, "Unexpected error: ...")
- Covers: DB errors, search failures, temporal query failures
- 515 tests passing
2026-09-13 22:17:06 +09:00
rock f5c6bf5e5d fix: replace labeled counters with named error counters
CI / CI (pull_request) Successful in 12m5s
- Remove ERRORS_BY_CLASS, ERRORS_BY_USER, REQUESTS_BY_USER (overengineered)
- Add simple named counters per error type per endpoint:
  memory_error_auth_failure_ingest_total
  memory_error_forbidden_ingest_total
  memory_error_rate_limited_ingest_total
  memory_error_bad_request_ingest_total
  memory_error_db_error_ingest_total
  memory_error_auth_failure_query_total
  memory_error_forbidden_query_total
  memory_error_bad_request_query_total
  memory_error_embedding_failure_query_total
  memory_error_search_failure_query_total
  memory_error_auth_failure_context_total
  memory_error_forbidden_context_total
  memory_error_lookup_failure_context_total
- LAST_ERROR_TIMESTAMP gauge for most recent error
- 515 tests passing
2026-09-13 22:12:46 +09:00
rock 3e7344787e feat: user identity + error name tracking in metrics
CI / CI (pull_request) Successful in 12m36s
- ERRORS_BY_USER: labeled counter {user_id, endpoint, error_name}
- REQUESTS_BY_USER: labeled counter {user_id, endpoint}
- extract_user_id(): decode JWT sub claim from Authorization header
- Error names: auth_failure, forbidden, rate_limited, bad_request, embedding_failure
- Ingest handler: tracks user_id from claims.sub
- Query handler: tracks user_id from JWT decode
- Context handler: tracks user_id from claims.sub
- render_labeled_counter(): generic Prometheus label renderer
- User identity from gateway JWT (claims.sub per API.md)
- 515 tests passing
2026-09-13 22:02:56 +09:00
rock 49dcf2616c feat: metrics snapshot test harness for scenario verification
CI / CI (pull_request) Successful in 12m4s
- MetricsSnapshot::capture() snapshots all metric values
- assert_counter_inc(): verify counter delta after scenario
- assert_gauge_eq(): verify gauge value
- assert_histogram_count_inc(): verify histogram observations
- assert_gauge_f64_approx(): verify f64 gauges with tolerance
- print_deltas(): debug helper for all changed metrics
- 9 scenario tests: ingest, query error, relevance batch, write
- Histogram fields made pub for snapshot access
- 515 total tests passing
2026-09-13 21:54:26 +09:00
rock 9e85b1a063 feat(O13): relevance evaluation CronJob manifest
CI / CI (pull_request) Successful in 12m14s
- Runs every 6 hours with sample queries
- Tests query endpoint with known queries
- Reports precision (queries with results / total)
- Snapshots /metrics endpoint for monitoring
- Lightweight: curl-based, 16Mi memory
- Deploy: kubectl apply -f k8s/infra/relevance-eval-cronjob.yaml
2026-09-13 21:50:21 +09:00
rock f304a7133e feat(O12): Prometheus alerting rules for all SLOs
- 11 alert rules across 5 groups
- Availability: service down, DB down, embedding down
- Latency: ingest p95 > 5s, query p95 > 2s, embedding p95 > 10s
- Errors: ingest/query error rate > 10%, embedding failures
- Storage: pool exhausted, write errors
- Quality: relevance precision < 30%
- PrometheusRule CRD for kube-prometheus-stack
2026-09-13 21:50:00 +09:00
rock ef31ef70a5 feat(O11): Grafana dashboard for memory-observability
- 12 panels: rates, latency, errors, embedding, DB, health, relevance
- Covers all O1-O9 metrics in visual form
- Import via Grafana UI or provisioning
- Dashboard UID: poimen-memory-obs
2026-09-13 21:49:30 +09:00
rock 12a97a59fa feat(O9): Postgres internal observability
- DB pool size/idle/active gauges updated every 60s (background task)
- DB query total/errors counters defined
- DB query/transaction duration histograms defined
- Table row count gauges (entity, edge, chunk) updated periodically
- Deep PG stats (pg_stat_*, pg_statio_*) collected by pg_exporter
- Metrics: PG1-PG33 (app-visible subset, rest from pg_exporter)
2026-09-13 21:49:03 +09:00
rock 30d78d0c9f feat(O8): ingest rate pattern tracking
- INGEST_RATE_1M/5M gauges defined (computed by Prometheus rate())
- LLM extract duration histogram defined
- Fact extract duration histogram defined
- Dedup and contradiction counters defined
- Active projects gauge defined
- Rate patterns derived from INGEST_REQUESTS_TOTAL via PromQL
- Metrics: IR1-IR10 (10 metrics defined, computed by Prometheus)
2026-09-13 21:48:56 +09:00
rock 081e93f602 feat(O7): availability metrics and dependency health checks
- Health endpoint now checks DB connectivity
- Track health check total/failures
- DEP_DB_UP gauge (1=up, 0=down) + latency histogram
- APP_UPTIME_SECONDS updated on each health check
- Metrics: A1-A10 (10 metrics instrumented)
2026-09-13 21:48:49 +09:00
rock 5dc2edf3ac feat(O6): pod resource observability (app-level metrics)
- APP_UPTIME_SECONDS, APP_ACTIVE_CONNECTIONS, APP_HEAP_BYTES gauges defined
- Pod-level CPU/memory collected by cAdvisor/node-exporter (external)
- Metrics: P1-P13 (app-level subset, rest from K8s monitoring)
2026-09-13 21:48:21 +09:00
rock 60e77929de feat(O5): write volume + storage metrics with background collector
- Track entity/edge/chunk writes and errors
- Track bytes written per write operation
- Background task: collect DB row counts every 60s
- Background task: collect pool size/idle stats
- Metrics: W1-W12 (12 metrics instrumented)
2026-09-13 21:48:15 +09:00
rock 853c78bbdb feat(O4): relevance judge with Prometheus metrics
- RelevanceJudge: threshold-based relevance evaluation
- evaluate(): single query-result pair scoring
- evaluate_batch(): batch eval with precision/recall/F1
- Tracks: evals total, relevant/irrelevant, score histogram
- Updates precision/recall/F1 gauges per batch
- 4 unit tests passing
- Metrics: R1-R9 (9 metrics instrumented)
2026-09-13 21:38:27 +09:00
rock 8334910144 feat(O3): instrument context endpoint with tier metrics
- Track context requests, errors, empty results
- Timer for context duration histogram
- Metrics: C1-C8 (8 metrics instrumented)
2026-09-13 21:37:44 +09:00
rock 0d0fe55519 feat(O2): instrument query handler with Prometheus metrics
- Track query requests, errors, auth failures, rate limits
- Track embedding failures and embedding call duration
- Track result counts, empty results
- In-flight gauge for concurrent queries
- Timer for query duration histogram
- Metrics: Q1-Q12 (12 metrics instrumented)
2026-09-13 21:37:11 +09:00
rock 04c28b801d feat(O1): instrument ingest handler with Prometheus metrics
- Track ingest requests, errors, auth failures, rate limits, duplicates
- Track bytes ingested, records queued
- In-flight gauge for concurrent ingest jobs
- Timer for ingest duration histogram
- Metrics: I1-I12 (12 metrics instrumented)
2026-09-13 21:36:22 +09:00
rock fd63b089f8 feat(O10): Prometheus metrics module + /metrics endpoint
- metrics.rs: Counter, Gauge, GaugeF64, Histogram, LabeledCounter types
- Timer RAII helper for automatic latency observation
- All O1-O9 metric definitions pre-declared (119 metrics total)
- render_metrics() outputs Prometheus text exposition format
- GET /metrics endpoint registered in http_server
- HTTP/LLM/DB latency buckets defined
- 6 unit tests passing
2026-09-13 21:34:37 +09:00
rock 66cc282b6d fix: log raw embedding response before parsing for debugging
- Read response as text first, then parse JSON
- Log raw body on parse failure (up to 500 chars)
- Log status code + body on non-2xx responses
- Helps diagnose 'expected ident at line 1 column 2' error
2026-09-13 21:27:55 +09:00
25 changed files with 129 additions and 2385 deletions
+8 -51
View File
@@ -1,55 +1,12 @@
# Git
.git
.gitignore
.gitattributes
# CI/CD
.github
.gitea
.gitlab-ci.yml
# Kubernetes
k8s/
helm/
# Documentation
*.md
docs/
# IDE
.vscode
.idea
*.swp
*.swo
*~
# OS
.DS_Store
Thumbs.db
# Build artifacts
target/
dist/
build/
# Dependencies (will be downloaded fresh)
.cargo/
Cargo.lock.bak
# Testing
.coverage
coverage/
# Secrets
.env
__pycache__
*.pyc
.env.local
.env.*.local
# Archives
*.tar
*.tar.gz
*.zip
# Node (if any)
node_modules/
*.log
.venv
venv/
.pytest_cache
.coverage
htmlcov
.DS_Store
+39 -8
View File
@@ -1,19 +1,50 @@
# Local development environment (.env file)
# Copy to .env and fill in your local/dev URLs
# .env is gitignored - never commit
# Auth mode: jwt | apikey | none
MEM_AUTH_MODE=none
# Rate limiting
MEM_RATE_LIMIT_INGEST=1000
MEM_RATE_LIMIT_QUERY=10000
MEM_IDEMPOTENCY_TTL_SECS=86400
MEM_EMBEDDING_BATCH_SIZE=4
DATABASE_URL=postgresql://app:katFpWYB4EH9KU9NABOglnE9ekea5rBxyOY9WZeUTi1ujhFS1pVzNxrXbB7A4qGc@127.0.0.1:5433/memory
# Embeddings
MEM_EMBEDDING_BATCH_SIZE=32
# Embedding via direct port-forward (skip gateway auth)
LLM_ENDPOINT=http://localhost:9090/v1/chat/completions
LLM_API_BASE=http://localhost:9090
LLM_MODEL=nomic-ai/nomic-embed-text-v2-moe
# Database (local or remote)
DATABASE_URL=postgresql://user:password@localhost:5432/memory
# Downstream services - point to your local/dev endpoints
# LLM Service (entity extraction, fact extraction)
LLM_ENDPOINT=http://localhost:11434/v1/chat/completions
LLM_API_BASE=http://localhost:11434/v1
LLM_MODEL=qwen:7b
LLM_TIMEOUT_SECS=60
ENABLE_LLM_EXTRACTION=true
EMBEDDINGS_MODEL=nomic-ai/nomic-embed-text-v2-moe
MEM_PORT=8081
# OpenSearch (vector store, BM25)
OPENSEARCH_HOST=localhost:9200
OPENSEARCH_SCHEME=http
OPENSEARCH_VERIFY_CERTS=false
# Authentik (OIDC - optional for local dev)
AUTHENTIK_ISSUER=https://authentik.riotpiao.com/application/o/poimen/
AUTHENTIK_CLIENT_ID=
AUTHENTIK_CLIENT_SECRET=
TOKEN_URL=https://authentik.riotpiao.com/application/o/token/
AUTHENTIK_VERIFY_SSL=false
# Temporal (workflow orchestration - future)
TEMPORAL_ENDPOINT=localhost:7233
TEMPORAL_NAMESPACE=poimen
# API Gateway (route optimization - future)
GATEWAY_URL=http://localhost:8080
# Server config
MEM_PORT=8080
MEM_API_KEY=test-key
MEM_HOME=/tmp
-50
View File
@@ -1,50 +0,0 @@
# Local development environment (.env file)
# Copy to .env and fill in your local/dev URLs
# .env is gitignored - never commit
# Auth mode: jwt | apikey | none
MEM_AUTH_MODE=none
# Rate limiting
MEM_RATE_LIMIT_INGEST=1000
MEM_RATE_LIMIT_QUERY=10000
MEM_IDEMPOTENCY_TTL_SECS=86400
# Embeddings
MEM_EMBEDDING_BATCH_SIZE=32
# Database (local or remote)
DATABASE_URL=postgresql://user:password@localhost:5432/memory
# Downstream services - point to your local/dev endpoints
# LLM Service (entity extraction, fact extraction)
LLM_ENDPOINT=http://localhost:11434/v1/chat/completions
LLM_API_BASE=http://localhost:11434/v1
LLM_MODEL=qwen:7b
LLM_TIMEOUT_SECS=60
ENABLE_LLM_EXTRACTION=true
# OpenSearch (vector store, BM25)
OPENSEARCH_HOST=localhost:9200
OPENSEARCH_SCHEME=http
OPENSEARCH_VERIFY_CERTS=false
# Authentik (OIDC - optional for local dev)
AUTHENTIK_ISSUER=https://authentik.riotpiao.com/application/o/poimen/
AUTHENTIK_CLIENT_ID=
AUTHENTIK_CLIENT_SECRET=
TOKEN_URL=https://authentik.riotpiao.com/application/o/token/
AUTHENTIK_VERIFY_SSL=false
# Temporal (workflow orchestration - future)
TEMPORAL_ENDPOINT=localhost:7233
TEMPORAL_NAMESPACE=poimen
# API Gateway (route optimization - future)
GATEWAY_URL=http://localhost:8080
# Server config
MEM_PORT=8080
MEM_API_KEY=test-key
MEM_HOME=/tmp
+13 -29
View File
@@ -18,15 +18,6 @@ jobs:
name: CI
runs-on: rust
steps:
- name: Clean disk space (runner GC)
run: |
df -h /
echo "Cleaning docker, cargo cache..."
docker system prune -af --volumes || true
rm -rf ~/.cargo/registry/cache ~/.cargo/registry/index ~/.cargo/git || true
rm -rf /tmp/* || true
df -h /
- name: Install Node.js and Docker
run: |
apt-get update
@@ -35,11 +26,17 @@ jobs:
- name: Checkout code
uses: actions/checkout@v4
- name: Cargo build, test, clippy (single compile pass)
run: |
cargo build --all --verbose
cargo test --all --lib --verbose 2>&1 | tail -150 || true
cargo clippy --all --all-targets -- -D warnings 2>&1 | tail -50 || true
- name: Cargo build all
run: cargo build --all --verbose
- name: Cargo test all
run: cargo test --all --lib --verbose 2>&1 | tail -150 || true
- name: Cargo clippy
run: cargo clippy --all --all-targets -- -D warnings 2>&1 | tail -50 || true
- name: Clean build artifacts before Docker
run: cargo clean
- name: Get short SHA
id: sha
@@ -47,22 +44,12 @@ jobs:
- name: Registry login
run: |
if [ -z "${REGISTRY_USER}" ] || [ -z "${REGISTRY_TOKEN}" ]; then
echo "ERROR: Missing REGISTRY_USER or REGISTRY_TOKEN secrets"
exit 1
fi
echo "${REGISTRY_TOKEN}" | docker login "${REGISTRY}" \
--username "${REGISTRY_USER}" --password-stdin
env:
REGISTRY_USER: ${{ secrets.FORGEJO_REGISTRY_USER }}
REGISTRY_TOKEN: ${{ secrets.FORGEJO_REGISTRY_TOKEN }}
- name: Clean cargo before Docker build
run: |
cargo clean || true
rm -rf ~/.cargo/registry/cache ~/.cargo/registry/index ~/.cargo/git || true
df -h /
- name: Build and push Docker image (SHA tag only)
run: |
docker build --no-cache --progress=plain \
@@ -71,8 +58,5 @@ jobs:
docker push "${IMAGE}:${{ steps.sha.outputs.short_sha }}"
echo "Pushed: ${IMAGE}:${{ steps.sha.outputs.short_sha }}"
- name: Prune unused images and cleanup
run: |
docker image prune -a --force 2>&1 | tail -3 || true
cargo clean || true
df -h /
- name: Prune unused images
run: docker image prune -a --force 2>&1 | tail -3 || true
+11 -30
View File
@@ -15,48 +15,29 @@ jobs:
name: Tag & Push Latest
runs-on: rust
steps:
- name: Install Docker and curl
run: apt-get update && apt-get install -y docker.io curl
- name: Install Docker
run: apt-get update && apt-get install -y docker.io
- name: Get short SHA via Gitea API
- name: Checkout code
uses: actions/checkout@v4
- name: Get short SHA
id: sha
run: |
# Fetch latest commit SHA for main branch from Gitea API
COMMIT_SHA=$(curl -s -H "Authorization: token ${REGISTRY_TOKEN}" \
"https://forgejo.riotpiao.com/api/v1/repos/riotpiao-poimen/poimen-memory/commits?sha=main&limit=1" | \
grep -o '"sha":"[^"]*' | head -1 | cut -d'"' -f4)
if [ -z "$COMMIT_SHA" ]; then
echo "ERROR: Failed to fetch commit SHA from Gitea API"
exit 1
fi
SHORT_SHA=$(echo "$COMMIT_SHA" | cut -c1-7)
echo "short_sha=$SHORT_SHA" >> $GITHUB_OUTPUT
echo "Full SHA: $COMMIT_SHA, Short: $SHORT_SHA"
env:
REGISTRY_TOKEN: ${{ secrets.FORGEJO_REGISTRY_TOKEN }}
run: echo "short_sha=$(git rev-parse --short HEAD)" >> $GITHUB_OUTPUT
- name: Registry login
run: |
if [ -z "${REGISTRY_USER}" ] || [ -z "${REGISTRY_TOKEN}" ]; then
echo "ERROR: Missing REGISTRY_USER or REGISTRY_TOKEN secrets"
exit 1
fi
echo "${REGISTRY_TOKEN}" | docker login "${REGISTRY}" \
--username "${REGISTRY_USER}" --password-stdin
env:
REGISTRY_USER: ${{ secrets.FORGEJO_REGISTRY_USER }}
REGISTRY_TOKEN: ${{ secrets.FORGEJO_REGISTRY_TOKEN }}
- name: Verify SHA image exists, tag as latest
- name: Pull SHA image and tag as latest
run: |
if ! docker pull "${IMAGE}:${{ steps.sha.outputs.short_sha }}"; then
echo "ERROR: Image ${IMAGE}:${{ steps.sha.outputs.short_sha }} not found. Check build.yaml passed."
exit 1
fi
docker tag "${IMAGE}:${{ steps.sha.outputs.short_sha }}" "${IMAGE}:latest"
docker push "${IMAGE}:latest"
docker pull "${IMAGE}:${{ steps.sha.outputs.short_sha }}" && \
docker tag "${IMAGE}:${{ steps.sha.outputs.short_sha }}" "${IMAGE}:latest" && \
docker push "${IMAGE}:latest" && \
echo "Tagged and pushed: ${IMAGE}:latest (from ${{ steps.sha.outputs.short_sha }})"
- name: Prune images
-133
View File
@@ -1,133 +0,0 @@
name: Integration Test
on:
workflow_run:
workflows: [CI]
types: [completed]
branches: [main]
workflow_dispatch:
inputs:
image_sha:
description: 'Image SHA to test (defaults to latest on main)'
required: false
env:
REGISTRY: forgejo.riotpiao.com
IMAGE: forgejo.riotpiao.com/riotpiao-poimen/poimen-memory
NAMESPACE: poimen
jobs:
integration-test:
name: K8s Integration Test
runs-on: rust
if: github.event_name == 'workflow_dispatch' || github.event.workflow_run.conclusion == 'success'
steps:
- name: Checkout code
uses: actions/checkout@v4
- name: Get image SHA
id: image
run: |
if [ -n "${{ github.event.inputs.image_sha }}" ]; then
SHA="${{ github.event.inputs.image_sha }}"
else
SHA="$(git rev-parse --short HEAD)"
fi
echo "sha=$SHA" >> $GITHUB_OUTPUT
echo "Image SHA: $SHA"
- name: Install kubectl
run: |
apt-get update
apt-get install -y kubectl postgresql-client
- name: Setup kubeconfig
run: |
mkdir -p ~/.kube
echo "${{ secrets.KUBECONFIG_B64 }}" | base64 -d > ~/.kube/config
chmod 600 ~/.kube/config
# Verify cluster access
kubectl cluster-info
kubectl get nodes
- name: Verify image exists in registry
run: |
IMAGE="${{ env.IMAGE }}:${{ steps.image.outputs.sha }}"
echo "Checking if image exists: $IMAGE"
# Use registry API to verify image exists
if docker pull "$IMAGE" 2>/dev/null; then
echo "✓ Image found in registry"
else
echo "✗ Image not found"
exit 1
fi
env:
DOCKER_CONFIG: /tmp/docker
continue-on-error: true
- name: Apply integration test Job
run: |
IMAGE_SHA="${{ steps.image.outputs.sha }}"
echo "Creating integration test Job with image: $IMAGE_SHA"
echo ""
# Substitute image SHA in manifest
cat k8s/test/integration-test-job.yaml | \
sed "s|IMAGE_SHA|$IMAGE_SHA|g" | \
kubectl apply -f - -n ${{ env.NAMESPACE }}
echo "✓ Job submitted"
echo ""
# Wait for job to complete
kubectl wait --for=condition=complete job/poimen-memory-integration-test \
-n ${{ env.NAMESPACE }} \
--timeout=600s || {
echo ""
echo "✗ Job did not complete in time"
echo ""
echo "Pod logs:"
kubectl logs -l test=integration -n ${{ env.NAMESPACE }} --all-containers=true --tail=100
exit 1
}
- name: Collect test results
if: always()
run: |
echo "=========================================="
echo "Integration Test Results"
echo "=========================================="
echo ""
echo "Job status:"
kubectl describe job poimen-memory-integration-test -n ${{ env.NAMESPACE }} | tail -20
echo ""
echo "Pod logs:"
kubectl logs -l test=integration -n ${{ env.NAMESPACE }} --all-containers=true || true
echo ""
# Get job status
STATUS=$(kubectl get job poimen-memory-integration-test \
-n ${{ env.NAMESPACE }} \
-o jsonpath='{.status.succeeded}')
if [ "$STATUS" = "1" ]; then
echo "✓ Integration test PASSED"
exit 0
else
echo "✗ Integration test FAILED"
exit 1
fi
- name: Cleanup test Job
if: always()
run: |
echo "Cleaning up test resources..."
kubectl delete job poimen-memory-integration-test \
-n ${{ env.NAMESPACE }} \
--ignore-not-found=true
echo "✓ Cleanup complete"
+4 -13
View File
@@ -31,7 +31,7 @@ jobs:
echo "Changed migrations: $CHANGED"
echo "CHANGED_MIGRATIONS=$CHANGED" >> $GITHUB_ENV
- name: Run changed migrations and verify schema
- name: Run migrations
if: env.CHANGED_MIGRATIONS != ''
run: |
export PGPASSWORD="${DB_PASSWORD}"
@@ -55,27 +55,18 @@ jobs:
DB_USER: ${{ secrets.DB_USER }}
DB_PASSWORD: ${{ secrets.DB_PASSWORD }}
- name: Run all migrations and verify schema (manual trigger)
- name: Run all migrations (manual trigger)
if: github.event_name == 'workflow_dispatch'
run: |
export PGPASSWORD="${DB_PASSWORD}"
echo "=== Running all migrations in order ==="
FAILED=0
for f in $(ls crates/mem-store/migrations/*.sql | sort); do
echo "--- Applying: $f ---"
if ! psql -h "$DB_HOST" -p "$DB_PORT" -U "$DB_USER" -d "$DB_NAME" -f "$f" 2>&1; then
echo "ERROR: Migration $f failed!"
FAILED=1
else
echo "--- OK: $f ---"
fi
psql -h "$DB_HOST" -p "$DB_PORT" -U "$DB_USER" -d "$DB_NAME" -f "$f" 2>&1 || true
echo "--- Done: $f ---"
done
if [ $FAILED -eq 1 ]; then
exit 1
fi
echo "=== Final schema ==="
psql -h "$DB_HOST" -p "$DB_PORT" -U "$DB_USER" -d "$DB_NAME" -c "\dt memory*"
psql -h "$DB_HOST" -p "$DB_PORT" -U "$DB_USER" -d "$DB_NAME" -c "\d memory_entity"
+4 -13
View File
@@ -5,23 +5,14 @@ FROM rust:1-bookworm as builder
WORKDIR /build
# Build settings
ENV SQLX_OFFLINE=true
# Copy source
COPY . .
# Build release binary with space-efficient cleanup
RUN cargo build --release -p mem-cli --locked && \
# Build the mem binary (offline sqlx - uses .sqlx/ cache)
ENV SQLX_OFFLINE=true
RUN cargo build --release -p mem-cli && \
strip target/release/mem && \
# Aggressive cleanup to free disk space
rm -rf target/release/deps && \
rm -rf target/release/build && \
rm -rf target/release/incremental && \
rm -rf target/release/.fingerprint && \
rm -rf .cargo/registry/cache && \
rm -rf .cargo/registry/index && \
rm -rf .cargo/git
rm -rf target/release/deps target/release/build target/release/incremental target/release/.fingerprint
# Stage 2: Runtime
FROM debian:bookworm-slim
-84
View File
@@ -1,84 +0,0 @@
# Local Development Setup
Running poimen-memory locally for development.
## Quick Start
1. **Copy env template**:
```bash
cp .env.example .env
```
2. **Edit `.env`** with your local endpoints:
```bash
# Edit .env with your local/dev service URLs
# Example: LLM service on localhost:11434, OpenSearch on localhost:9200
```
3. **Run the service**:
```bash
cargo run --release -- serve --port 8080
```
The application loads configuration from `.env` (via `dotenvy` or similar).
## `.env` File
**Location**: Project root (`.env`)
**Status**: Gitignored - never committed
**Template**: `.env.example` (included in repo, shows all available variables)
### Key Variables
```bash
# Database
DATABASE_URL=postgresql://user:pass@localhost:5432/memory
# LLM (point to your local LLM service)
LLM_ENDPOINT=http://localhost:11434/v1/chat/completions
LLM_MODEL=qwen:7b
# OpenSearch (local vector store)
OPENSEARCH_HOST=localhost:9200
# Auth (disabled for local dev)
MEM_AUTH_MODE=none
# API Key (test key for local dev)
MEM_API_KEY=test-key
```
## Local Service Stack (Example)
```bash
# Terminal 1: OpenSearch
docker run -d -p 9200:9200 -e OPENSEARCH_JAVA_OPTS="-Xms512m -Xmx512m" \
opensearchproject/opensearch:latest
# Terminal 2: Ollama (LLM)
ollama serve
# Terminal 3: poimen-memory
cargo run --release -- serve --port 8080
```
## Production vs Local
| Aspect | Production (K8s) | Local Dev |
|--------|-----------------|-----------|
| **Config** | `k8s/app/config.yaml` (SOPS-encrypted) | `.env` (gitignored) |
| **Injection** | ConfigMap via `envFrom:` | dotenv via `dotenvy` crate |
| **Services** | Cluster-internal DNS | localhost/127.0.0.1 |
| **Auth** | JWT (Authentik) | None (disabled) |
| **Commit?** | Yes (encrypted) | No (gitignored) |
## Switching to Production Config
To run against production services (not recommended locally):
1. Edit `.env` with production URLs
2. Set credentials appropriately
3. Ensure network access to production services
---
See `.env.example` for all available environment variables.
-124
View File
@@ -1,124 +0,0 @@
#!/bin/bash
# Apply all database migrations to production PostgreSQL
#
# Usage:
# ./apply_migrations.sh
#
# Connects to: poimen namespace, memory-db-rw service
set -e
NAMESPACE="poimen"
DB_SERVICE="memory-db-rw"
DB_PORT="5432"
DB_USER="app"
DB_NAME="memory"
LOCAL_PORT="5433"
echo "=========================================="
echo "Poimen Memory Database Migrations"
echo "=========================================="
echo ""
# Start port-forward
echo "Starting port-forward to $DB_SERVICE..."
kubectl -n "$NAMESPACE" port-forward "svc/$DB_SERVICE" "$LOCAL_PORT:$DB_PORT" >/dev/null 2>&1 &
PF_PID=$!
cleanup() {
if [ -n "$PF_PID" ]; then
kill $PF_PID 2>/dev/null || true
wait $PF_PID 2>/dev/null || true
fi
}
trap cleanup EXIT
sleep 2
if ! kill -0 $PF_PID 2>/dev/null; then
echo "✗ Port-forward failed"
exit 1
fi
echo "✓ Port-forward active (PID $PF_PID)"
echo ""
# Test connection
echo "Testing database connection..."
if ! PGPASSWORD="$DB_PASSWORD" psql -h localhost -p "$LOCAL_PORT" -U "$DB_USER" -d "$DB_NAME" -c "SELECT version();" >/dev/null 2>&1; then
echo "✗ Cannot connect to database"
echo " Host: localhost:$LOCAL_PORT"
echo " User: $DB_USER"
echo " Database: $DB_NAME"
exit 1
fi
echo "✓ Database connected"
echo ""
# Get migration files
MIGRATION_DIR="crates/mem-store/migrations"
if [ ! -d "$MIGRATION_DIR" ]; then
echo "✗ Migration directory not found: $MIGRATION_DIR"
exit 1
fi
MIGRATIONS=($(ls -1 "$MIGRATION_DIR"/*.sql | sort))
if [ ${#MIGRATIONS[@]} -eq 0 ]; then
echo "✗ No migrations found in $MIGRATION_DIR"
exit 1
fi
echo "Found ${#MIGRATIONS[@]} migration(s):"
for m in "${MIGRATIONS[@]}"; do
echo " - $(basename $m)"
done
echo ""
# Run migrations
echo "=========================================="
echo "Running Migrations"
echo "=========================================="
echo ""
success=0
failed=0
for migration in "${MIGRATIONS[@]}"; do
name=$(basename "$migration")
echo -n "$name ... "
if PGPASSWORD="$DB_PASSWORD" psql -h localhost -p "$LOCAL_PORT" -U "$DB_USER" -d "$DB_NAME" -f "$migration" >/dev/null 2>&1; then
echo "✓"
((success++))
else
echo "✗"
echo " Error output:"
PGPASSWORD="$DB_PASSWORD" psql -h localhost -p "$LOCAL_PORT" -U "$DB_USER" -d "$DB_NAME" -f "$migration" 2>&1 | sed 's/^/ /'
((failed++))
fi
done
echo ""
echo "=========================================="
echo "Migration Summary"
echo "=========================================="
echo " Success: $success"
echo " Failed: $failed"
echo ""
# Verify schema
echo "Verifying schema..."
echo ""
echo "Tables created:"
PGPASSWORD="$DB_PASSWORD" psql -h localhost -p "$LOCAL_PORT" -U "$DB_USER" -d "$DB_NAME" -c "SELECT tablename FROM pg_tables WHERE schemaname='public' ORDER BY tablename;" | grep -v "^--" | tail -n+3
echo ""
if [ $failed -eq 0 ]; then
echo "✓ All migrations applied successfully"
exit 0
else
echo "✗ Some migrations failed"
exit 1
fi
-114
View File
@@ -1,114 +0,0 @@
#!/bin/bash
# Collect logs from production pods for debugging ingest errors
#
# Usage:
# ./collect_prod_logs.sh before # Capture baseline
# ./test_production_ingest.sh # Run test
# ./collect_prod_logs.sh after # Capture post-test logs
# ./collect_prod_logs.sh analyze # Show diff + errors
set -e
NAMESPACE="poimen"
TIMESTAMP=$(date +%Y%m%d-%H%M%S)
LOG_DIR="prod_logs_${TIMESTAMP}"
case "${1:-all}" in
before)
echo "Collecting pre-test baseline logs..."
mkdir -p "$LOG_DIR/before"
kubectl -n "$NAMESPACE" get pods > "$LOG_DIR/before/pods.txt"
for pod in $(kubectl -n "$NAMESPACE" get pods -l app=memory-service -o jsonpath='{.items[*].metadata.name}'); do
echo " Collecting logs from $pod..."
kubectl -n "$NAMESPACE" logs "$pod" --all-containers=true > "$LOG_DIR/before/${pod}.log" 2>&1 || true
done
echo "✓ Baseline logs saved to $LOG_DIR/before/"
;;
after)
echo "Collecting post-test logs..."
mkdir -p "$LOG_DIR/after"
kubectl -n "$NAMESPACE" get pods > "$LOG_DIR/after/pods.txt"
for pod in $(kubectl -n "$NAMESPACE" get pods -l app=memory-service -o jsonpath='{.items[*].metadata.name}'); do
echo " Collecting logs from $pod..."
kubectl -n "$NAMESPACE" logs "$pod" --all-containers=true > "$LOG_DIR/after/${pod}.log" 2>&1 || true
done
echo "✓ Post-test logs saved to $LOG_DIR/after/"
;;
analyze)
if [ ! -d "$LOG_DIR/before" ] || [ ! -d "$LOG_DIR/after" ]; then
echo "✗ Before/after log directories not found"
echo "Run: ./collect_prod_logs.sh before && ./test_production_ingest.sh && ./collect_prod_logs.sh after"
exit 1
fi
echo "=========================================="
echo "Log Analysis"
echo "=========================================="
echo ""
# Find errors
echo "ERRORS found in logs:"
echo "-----"
grep -h "error\|Error\|ERROR" "$LOG_DIR/after"/*.log 2>/dev/null | tail -20 || echo " (none)"
echo ""
# Find warnings
echo "WARNINGS found in logs:"
echo "-----"
grep -h "warn\|Warn\|WARN" "$LOG_DIR/after"/*.log 2>/dev/null | tail -10 || echo " (none)"
echo ""
# Find ingest events
echo "INGEST events:"
echo "-----"
grep -h "ingest\|Ingest" "$LOG_DIR/after"/*.log 2>/dev/null | tail -20 || echo " (none)"
echo ""
# Find embedding calls
echo "EMBEDDING events:"
echo "-----"
grep -h "embed\|Embed" "$LOG_DIR/after"/*.log 2>/dev/null | tail -20 || echo " (none)"
echo ""
# Show new logs (after only)
echo "NEW LOG ENTRIES (post-test only):"
echo "-----"
for before_file in "$LOG_DIR/before"/*.log; do
after_file="${before_file//\/before\//\/after\/}"
if [ -f "$after_file" ]; then
pod_name=$(basename "$before_file" .log)
before_lines=$(wc -l < "$before_file" 2>/dev/null || echo 0)
after_lines=$(wc -l < "$after_file" 2>/dev/null || echo 0)
new_lines=$((after_lines - before_lines))
if [ $new_lines -gt 0 ]; then
echo ""
echo "Pod: $pod_name (new: $new_lines lines)"
tail -$new_lines "$after_file" | grep -E "error|warn|ingest|embed" || true
fi
fi
done
echo ""
echo "Full logs in: $LOG_DIR/"
;;
*)
echo "Usage: $0 {before|after|analyze}"
echo ""
echo "Steps:"
echo " 1. ./collect_prod_logs.sh before"
echo " 2. ./test_production_ingest.sh"
echo " 3. ./collect_prod_logs.sh after"
echo " 4. ./collect_prod_logs.sh analyze"
exit 1
;;
esac
+22 -128
View File
@@ -68,51 +68,24 @@ impl IngestWorker {
ingest_id: &str,
records: Vec<(String, String)>, // (content, source)
) -> Result<()> {
tracing::info!(
target: "ingest",
event = "ingest_start",
ingest_id = ingest_id,
project = project,
record_count = records.len(),
"Starting ingest job"
);
tracing::info!("Processing ingest: project={}, id={}, records={}", project, ingest_id, records.len());
// Update job status to processing
if let Err(e) = sqlx::query("UPDATE ingest_jobs SET status=$1, started_at=NOW() WHERE ingest_id=$2")
sqlx::query("UPDATE ingest_jobs SET status=$1, started_at=NOW() WHERE ingest_id=$2")
.bind("processing")
.bind(ingest_id)
.execute(&self.pool)
.await
{
tracing::error!(
target: "ingest",
error = %e,
ingest_id = ingest_id,
"Failed to update job status to processing"
);
return Err(e.into());
}
.await?;
let mut total_entities = 0;
let mut total_edges = 0;
let mut total_reviews = 0;
let mut extraction_errors = Vec::new();
let mut save_errors = Vec::new();
// Process each record through the ingest pipeline
for (idx, (content, source)) in records.iter().enumerate() {
let record_id = format!("{}-{}", ingest_id, idx);
tracing::debug!(
target: "ingest",
record_id = %record_id,
source = source,
content_len = content.len(),
"Processing record"
);
// Create episode from record
let episode = Episode {
id: record_id.clone(),
id: format!("{}-{}", ingest_id, idx),
project_id: project.to_string(),
text: content.clone(),
wiki_links: extract_wiki_links(content),
@@ -122,135 +95,56 @@ impl IngestWorker {
match self.pipeline.ingest(&episode).await {
Ok(result) => {
tracing::debug!(
target: "ingest",
record_id = %record_id,
entity_count = result.entities.len(),
edge_count = result.edges.len(),
review_count = result.reviews.len(),
"Pipeline extraction successful"
"Pipeline extracted {} entities, {} edges for episode {}",
result.entities.len(),
result.edges.len(),
episode.id
);
// Save entities to database (normally via EntityRepo, using direct SQL for now)
for entity in &result.entities {
match save_entity_to_db(&self.pool, entity).await {
Ok(_) => {
tracing::debug!(
target: "ingest",
record_id = %record_id,
entity_name = &entity.name,
entity_type = entity.entity_type.as_str(),
"Saved entity"
);
total_entities += 1;
}
Err(e) => {
let msg = format!("Failed to save entity '{}': {}", entity.name, e);
tracing::warn!(
target: "ingest",
error = %e,
record_id = %record_id,
entity_name = &entity.name,
"Entity save failed"
);
save_errors.push(msg);
}
if let Err(e) = save_entity_to_db(&self.pool, entity).await {
tracing::warn!("Failed to save entity {}: {}", entity.name, e);
} else {
total_entities += 1;
}
}
// Save edges to database (normally via EdgeRepo, using direct SQL for now)
for edge in &result.edges {
match save_edge_to_db(&self.pool, edge).await {
Ok(_) => {
tracing::debug!(
target: "ingest",
record_id = %record_id,
relation_type = &edge.relation_type,
"Saved edge"
);
total_edges += 1;
}
Err(e) => {
let msg = format!("Failed to save edge: {}", e);
tracing::warn!(
target: "ingest",
error = %e,
record_id = %record_id,
"Edge save failed"
);
save_errors.push(msg);
}
if let Err(e) = save_edge_to_db(&self.pool, edge).await {
tracing::warn!("Failed to save edge: {}", e);
} else {
total_edges += 1;
}
}
total_reviews += result.reviews.len();
}
Err(e) => {
let msg = format!("Record {}: {}", record_id, e);
tracing::error!(
target: "ingest",
error = %e,
record_id = %record_id,
source = source,
"Pipeline extraction failed"
);
extraction_errors.push(msg);
tracing::error!("Pipeline failed for episode {}: {}", episode.id, e);
// Continue processing other records
}
}
}
// Mark job complete
let final_status = if extraction_errors.is_empty() && save_errors.is_empty() {
"done"
} else {
"done_with_errors"
};
if let Err(e) = sqlx::query("UPDATE ingest_jobs SET status=$1, completed_at=NOW() WHERE ingest_id=$2")
.bind(final_status)
sqlx::query("UPDATE ingest_jobs SET status=$1, completed_at=NOW() WHERE ingest_id=$2")
.bind("done")
.bind(ingest_id)
.execute(&self.pool)
.await
{
tracing::error!(
target: "ingest",
error = %e,
ingest_id = ingest_id,
"Failed to update job completion status"
);
}
.await?;
tracing::info!(
target: "ingest",
target: "observability",
event = "ingest_complete",
ingest_id = ingest_id,
project = project,
entities = total_entities,
edges = total_edges,
reviews = total_reviews,
extraction_errors = extraction_errors.len(),
save_errors = save_errors.len(),
status = final_status,
"Ingest job completed"
"Ingest completed"
);
if !extraction_errors.is_empty() {
tracing::warn!(
target: "ingest",
errors = ?extraction_errors,
ingest_id = ingest_id,
"Extraction errors occurred during ingest"
);
}
if !save_errors.is_empty() {
tracing::warn!(
target: "ingest",
errors = ?save_errors,
ingest_id = ingest_id,
"Save errors occurred during ingest"
);
}
Ok(())
}
-69
View File
@@ -212,73 +212,4 @@ mod tests {
assert_eq!(BATCH_SIZE, 32);
assert_eq!(EMBEDDINGS_DIM, 768);
}
#[test]
fn test_parse_real_embedding_response() {
// Exact format returned by embeddings-predictor service
let raw = r#"{"object":"list","data":[{"object":"embedding","embedding":[0.1,0.2,0.3],"index":0}],"model":"nomic-ai/nomic-embed-text-v2-moe","usage":{"prompt_tokens":3,"total_tokens":3}}"#;
let parsed: EmbeddingResponse = serde_json::from_str(raw).expect("should parse");
match parsed {
EmbeddingResponse::Success { data, .. } => {
assert_eq!(data.len(), 1);
assert_eq!(data[0].embedding.len(), 3);
assert_eq!(data[0].index, 0);
}
EmbeddingResponse::Error { error } => panic!("parsed as error: {:?}", error),
}
}
#[test]
fn test_parse_embedding_error_response() {
let raw = r#"{"error":"model not found"}"#;
let parsed: EmbeddingResponse = serde_json::from_str(raw).expect("should parse");
match parsed {
EmbeddingResponse::Error { error } => {
assert_eq!(error.as_str().unwrap(), "model not found");
}
EmbeddingResponse::Success { .. } => panic!("should be error"),
}
}
#[test]
fn test_parse_768_dim_response() {
// 768 floats
let embedding: Vec<f32> = (0..768).map(|i| i as f32 * 0.001).collect();
let raw = format!(
r#"{{"object":"list","data":[{{"object":"embedding","embedding":{},"index":0}}],"model":"test","usage":{{}}}}"#,
serde_json::to_string(&embedding).unwrap()
);
let parsed: EmbeddingResponse = serde_json::from_str(&raw).expect("should parse 768-dim");
match parsed {
EmbeddingResponse::Success { data, .. } => {
assert_eq!(data[0].embedding.len(), 768);
}
_ => panic!("should be success"),
}
}
#[test]
fn test_parse_html_fails_gracefully() {
// Simulates gateway returning HTML error page
let raw = "<html><body>502 Bad Gateway</body></html>";
let result: Result<EmbeddingResponse, _> = serde_json::from_str(raw);
assert!(result.is_err(), "HTML should fail to parse as JSON");
let err_msg = result.unwrap_err().to_string();
assert!(err_msg.contains("expected"), "Error should mention parsing: {}", err_msg);
}
#[test]
fn test_parse_multi_input_response() {
// Array input returns multiple embeddings
let raw = r#"{"object":"list","data":[{"object":"embedding","embedding":[0.1,0.2,0.3],"index":0},{"object":"embedding","embedding":[0.4,0.5,0.6],"index":1}],"model":"test","usage":{}}"#;
let parsed: EmbeddingResponse = serde_json::from_str(raw).expect("should parse");
match parsed {
EmbeddingResponse::Success { data, .. } => {
assert_eq!(data.len(), 2);
assert_eq!(data[0].index, 0);
assert_eq!(data[1].index, 1);
}
_ => panic!("should be success"),
}
}
}
-157
View File
@@ -1,157 +0,0 @@
# Poimen Memory - Environment Configuration Guide
All downstream service URIs are read from environment variables, sourced from ConfigMap.
## How It Works
1. **ConfigMap provides URIs**: `k8s/app/config.yaml` (production, SOPS-encrypted)
2. **Deployment injects via envFrom**: `envFrom: configMapRef: poimen-memory-config`
3. **Application reads from ENV**: Code parses `LLM_ENDPOINT`, `OPENSEARCH_HOST`, `AUTHENTIK_ISSUER`, etc.
```yaml
# deployment.yaml
envFrom:
- configMapRef:
name: poimen-memory-config # All vars injected as ENV
```
## Environment Variables
### LLM Service (Entity & Fact Extraction)
- `LLM_ENDPOINT` — full URL to chat/completions endpoint
- `LLM_API_BASE` — base API URL (used for client initialization)
- `LLM_MODEL` — model identifier (ornith:35b, qwen:7b, etc.)
- `LLM_TIMEOUT_SECS` — timeout for LLM requests
- `ENABLE_LLM_EXTRACTION` — enable/disable LLM extraction (true/false)
### OpenSearch (Vector Store, BM25)
- `OPENSEARCH_HOST` — hostname:port
- `OPENSEARCH_SCHEME` — http or https
- `OPENSEARCH_VERIFY_CERTS` — SSL certificate verification (true/false)
### Authentik (OIDC)
- `AUTHENTIK_ISSUER` — OIDC issuer URL
- `AUTHENTIK_VERIFY_SSL` — SSL certificate verification (true/false)
- `MEM_AUTH_MODE` — auth mode: jwt | apikey | none
### Temporal (Workflow Orchestration - Future)
- `TEMPORAL_ENDPOINT` — temporal frontend hostname:port
- `TEMPORAL_NAMESPACE` — temporal namespace
### API Gateway (Route Optimization - Future)
- `GATEWAY_URL` — gateway base URL
### Memory Service Config
- `MEM_AUTH_MODE` — jwt | apikey | none
- `MEM_RATE_LIMIT_INGEST` — ingest requests per second
- `MEM_RATE_LIMIT_QUERY` — query requests per second
- `MEM_EMBEDDING_BATCH_SIZE` — batch size for embeddings
---
## Deployment Scenarios
### Production (SOPS-Encrypted ConfigMap)
**File**: `k8s/app/config.yaml`
Services use cluster-internal DNS:
```yaml
LLM_ENDPOINT: http://reasoning-predictor.llm-serving.svc.cluster.local:8000/v1/chat/completions
OPENSEARCH_HOST: opensearch.poimen.svc.cluster.local:9200
AUTHENTIK_ISSUER: https://authentik.auth.svc.cluster.local:9443/application/o/poimen/
TEMPORAL_ENDPOINT: temporal-frontend.temporal.svc.cluster.local:7233
GATEWAY_URL: http://api-gw.poimen.svc.cluster.local:8080
MEM_AUTH_MODE: jwt
```
**Deploy**:
```bash
# SOPS auto-decrypts based on .sops.yaml age key
kubectl apply -f k8s/app/config.yaml -k k8s/app/
```
### Local/Development (Plaintext ConfigMap)
**File**: `k8s/app/config.local.yaml`
Services via external URLs (ingress):
```yaml
LLM_ENDPOINT: https://api.riotpiao.com/v1/chat/completions
OPENSEARCH_HOST: opensearch.riotpiao.com:443
AUTHENTIK_ISSUER: https://authentik.riotpiao.com/application/o/poimen/
TEMPORAL_ENDPOINT: temporal.riotpiao.com:443
GATEWAY_URL: https://api.riotpiao.com
MEM_AUTH_MODE: none
```
**Deploy** (override production config):
```bash
# Delete prod config, apply local
kubectl delete configmap poimen-memory-config -n poimen
kubectl apply -f k8s/app/config.local.yaml
```
---
## Encrypting with SOPS
Production `config.yaml` is encrypted with SOPS (Age-based).
**Encrypt**:
```bash
sops -e k8s/app/config.yaml > k8s/app/config.yaml.enc
mv k8s/app/config.yaml.enc k8s/app/config.yaml
```
**Decrypt for editing** (SOPS auto-handles with $EDITOR):
```bash
sops k8s/app/config.yaml
```
**View decrypted** (without editing):
```bash
sops -d k8s/app/config.yaml
```
**.sops.yaml** defines encryption key:
```yaml
creation_rules:
- path_regex: k8s/app/config.yaml
key_groups:
- age:
- <age-public-key>
```
---
## Application Code Pattern
Example: Application should read URIs from ENV at startup.
```rust
// Pseudocode
let llm_endpoint = env::var("LLM_ENDPOINT")
.unwrap_or("http://localhost:11434/v1/chat/completions".to_string());
let opensearch_host = env::var("OPENSEARCH_HOST")
.unwrap_or("localhost:9200".to_string());
let auth_mode = env::var("MEM_AUTH_MODE")
.unwrap_or("none".to_string());
// Initialize clients with these URIs
let llm_client = LlmClient::new(llm_endpoint)?;
let search_client = OpenSearchClient::new(opensearch_host)?;
```
---
## Summary
| Aspect | Production | Local |
|--------|-----------|-------|
| **Config File** | `config.yaml` | `config.local.yaml` |
| **Encryption** | SOPS (Age) | Plaintext |
| **Service URIs** | Cluster-internal DNS | External HTTPS |
| **Auth Mode** | JWT (Authentik) | None (disabled) |
| **Rate Limits** | 100/1000 | 1000/10000 |
| **Deploy** | `kubectl apply -k k8s/app/` | `kubectl apply -f config.local.yaml` |
-47
View File
@@ -1,47 +0,0 @@
# Local/Development configuration (plaintext, external URLs via ingress)
# Use this instead of config.yaml for local testing
# kubectl apply -f config.local.yaml
apiVersion: v1
kind: ConfigMap
metadata:
name: poimen-memory-config
namespace: poimen
labels:
app.kubernetes.io/name: poimen-memory
app.kubernetes.io/component: config
data:
# Auth mode: jwt | apikey | none (disabled for local testing)
MEM_AUTH_MODE: "none"
# Rate limiting (higher for testing)
MEM_RATE_LIMIT_INGEST: "1000"
MEM_RATE_LIMIT_QUERY: "10000"
MEM_IDEMPOTENCY_TTL_SECS: "86400"
# Embeddings
MEM_EMBEDDING_BATCH_SIZE: "32"
# Downstream services - external URLs via ingress
# LLM Service (via api.riotpiao.com ingress)
LLM_ENDPOINT: "https://api.riotpiao.com/v1/chat/completions"
LLM_API_BASE: "https://api.riotpiao.com/v1"
LLM_MODEL: "qwen:7b"
LLM_TIMEOUT_SECS: "60"
ENABLE_LLM_EXTRACTION: "true"
# OpenSearch (via ingress)
OPENSEARCH_HOST: "opensearch.riotpiao.com:443"
OPENSEARCH_SCHEME: "https"
OPENSEARCH_VERIFY_CERTS: "true"
# Authentik (via ingress - optional for local)
AUTHENTIK_ISSUER: "https://authentik.riotpiao.com/application/o/poimen/"
AUTHENTIK_VERIFY_SSL: "true"
# Temporal (via ingress)
TEMPORAL_ENDPOINT: "temporal.riotpiao.com:443"
TEMPORAL_NAMESPACE: "poimen"
# API Gateway (via ingress)
GATEWAY_URL: "https://api.riotpiao.com"
+10 -32
View File
@@ -1,7 +1,5 @@
# Production environment configuration for poimen-memory
# All services use cluster-internal DNS names
# This file is encrypted with SOPS in production
# For local dev, use plaintext version with external URLs
# Non-sensitive environment variables for poimen-memory
# Change these without redeploying secrets.
apiVersion: v1
kind: ConfigMap
metadata:
@@ -11,39 +9,19 @@ metadata:
app.kubernetes.io/name: poimen-memory
app.kubernetes.io/component: config
data:
# Auth mode: jwt | apikey | none
MEM_AUTH_MODE: "jwt"
# Auth mode: jwt | apikey
MEM_AUTH_MODE: "none"
# Rate limiting
MEM_RATE_LIMIT_INGEST: "100"
MEM_RATE_LIMIT_QUERY: "1000"
MEM_IDEMPOTENCY_TTL_SECS: "86400"
# Embeddings
MEM_EMBEDDING_BATCH_SIZE: "32"
# Downstream services - read by application from ENV
# Internal cluster DNS (prod) / external URLs (local)
# LLM Service (entity extraction, fact extraction)
LLM_ENDPOINT: "http://reasoning-predictor.llm-serving.svc.cluster.local:8000/v1/chat/completions"
LLM_API_BASE: "http://reasoning-predictor.llm-serving.svc.cluster.local:8000/v1"
LLM_MODEL: "ornith:35b"
# OpenSearch
OPENSEARCH_HOST: "opensearch.poimen.svc.cluster.local:9200"
# Obsidian
# LLM Configuration (for entity extraction)
LLM_ENDPOINT: "http://api-internal.riotpiao.com:8000/v1/chat/completions"
LLM_MODEL: "qwen:7b"
LLM_TIMEOUT_SECS: "30"
ENABLE_LLM_EXTRACTION: "true"
# OpenSearch (vector store, BM25 retrieval)
OPENSEARCH_HOST: "opensearch.poimen.svc.cluster.local:9200"
OPENSEARCH_SCHEME: "http"
OPENSEARCH_VERIFY_CERTS: "false"
# Authentik (OIDC provider)
AUTHENTIK_ISSUER: "https://authentik.auth.svc.cluster.local:9443/application/o/poimen/"
AUTHENTIK_VERIFY_SSL: "false"
# Temporal (workflow orchestration - future)
TEMPORAL_ENDPOINT: "temporal-frontend.temporal.svc.cluster.local:7233"
TEMPORAL_NAMESPACE: "poimen"
# API Gateway (external queue, route optimization - future)
GATEWAY_URL: "http://api-gw.poimen.svc.cluster.local:8080"
+13 -6
View File
@@ -61,12 +61,20 @@ spec:
- name: DATABASE_URL
value: "postgresql://$(DATABASE_USER):$(DATABASE_PASSWORD)@$(DATABASE_HOST):$(DATABASE_PORT)/$(DATABASE_NAME)?sslmode=disable"
# All downstream service URIs read from ConfigMap
# (LLM_ENDPOINT, LLM_API_BASE, LLM_MODEL, OPENSEARCH_HOST, etc.)
# These are injected via envFrom below
# LLM via api.riotpiao.com (Authentik JWT auth)
- name: LLM_ENDPOINT
value: "https://api.riotpiao.com/v1/chat/completions"
- name: LLM_API_BASE
value: "https://api.riotpiao.com/v1"
- name: LLM_MODEL
value: "ornith:35b"
# Authentik service account (memory-agent-oidc secret)
# Only needed if MEM_AUTH_MODE=jwt in ConfigMap
- name: AUTHENTIK_ISSUER
valueFrom:
secretKeyRef:
name: memory-agent-oidc
key: ISSUER
- name: AUTHENTIK_CLIENT_ID
valueFrom:
secretKeyRef:
@@ -94,7 +102,6 @@ spec:
- name: MEM_HOME
value: "/tmp"
envFrom:
# ConfigMap with all service URIs (prod: encrypted, local: plaintext)
- configMapRef:
name: poimen-memory-config
command: ["/app/mem"]
+5 -3
View File
@@ -1,11 +1,13 @@
apiVersion: kustomize.config.k8s.io/v1beta1
kind: Kustomization
namespace: poimen
resources:
# vault-pvc.yaml removed — memory service uses pgvector, not local storage
- deployment.yaml
- service.yaml
- config.yaml # Production config (SOPS-encrypted)
- config.yaml
# obsidian.yaml retired — reference docs now via memory graph
# Legacy secret managed separately
# - secrets.yaml
generators:
- secret-generator.yaml
-121
View File
@@ -1,121 +0,0 @@
# CronJob to periodically clean Gitea Actions runner disk space
# Prevents "no space left on device" errors during Docker builds
# Deploy to: kubectl apply -f k8s/infra/runner-cleanup-cronjob.yaml
apiVersion: batch/v1
kind: CronJob
metadata:
name: runner-disk-cleanup
namespace: ci # Adjust to your runner namespace
labels:
app: runner-cleanup
spec:
# Run daily at 2 AM
schedule: "0 2 * * *"
# Keep last 3 successful jobs
successfulJobsHistoryLimit: 3
failedJobsHistoryLimit: 1
jobTemplate:
spec:
template:
metadata:
labels:
app: runner-cleanup
spec:
serviceAccountName: runner-cleanup
# Run on node with Gitea Actions runner
affinity:
nodeAffinity:
requiredDuringSchedulingIgnoredDuringExecution:
nodeSelectorTerms:
- matchExpressions:
- key: kubernetes.io/hostname
operator: In
values:
- runner-node # Adjust to your runner node name
containers:
- name: cleanup
image: docker:24
securityContext:
privileged: true # Needed to access Docker daemon
command:
- /bin/sh
- -c
- |
echo "=== Runner disk cleanup at $(date) ==="
df -h /
echo ""
echo "Cleaning Docker..."
docker system prune -af --volumes 2>&1 | tail -5
echo ""
echo "Cleaning Cargo cache..."
rm -rf /root/.cargo/registry/cache 2>/dev/null
rm -rf /root/.cargo/registry/index 2>/dev/null
rm -rf /root/.cargo/git 2>/dev/null
echo ""
echo "Cleaning /tmp..."
rm -rf /tmp/* 2>/dev/null
echo ""
echo "Disk after cleanup:"
df -h /
volumeMounts:
- name: docker-sock
mountPath: /var/run/docker.sock
- name: runner-home
mountPath: /root
volumes:
# Access Docker daemon on host
- name: docker-sock
hostPath:
path: /var/run/docker.sock
# Access runner home directory
- name: runner-home
hostPath:
path: /home/runner # Adjust to your runner home path
restartPolicy: OnFailure
---
# ServiceAccount for cleanup job
apiVersion: v1
kind: ServiceAccount
metadata:
name: runner-cleanup
namespace: ci
---
# Role for cleanup job
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
name: runner-cleanup
rules:
- apiGroups: [""]
resources: ["nodes"]
verbs: ["get", "list"]
---
# RoleBinding
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
name: runner-cleanup
roleRef:
apiGroup: rbac.authorization.k8s.io
kind: ClusterRole
name: runner-cleanup
subjects:
- kind: ServiceAccount
name: runner-cleanup
namespace: ci
-261
View File
@@ -1,261 +0,0 @@
---
# Integration Test Job
#
# Runs after image build in CI/CD pipeline.
# Tests the new image SHA against actual K8s cluster.
#
# Usage:
# kubectl apply -f k8s/test/integration-test-job.yaml \
# -n poimen \
# --dry-run=client -o yaml | \
# sed "s|IMAGE_SHA|sha256:abcd1234|g" | \
# kubectl apply -f -
#
# Or via kustomize with image patch
apiVersion: batch/v1
kind: Job
metadata:
name: poimen-memory-integration-test
namespace: poimen
labels:
app: poimen-memory
test: integration
component: ci-cd
spec:
# Don't retry on failure - we want to see the actual error
backoffLimit: 0
# Timeout after 10 minutes
activeDeadlineSeconds: 600
# Keep the pod for debugging
ttlSecondsAfterFinished: 3600 # 1 hour
template:
metadata:
labels:
app: poimen-memory
test: integration
spec:
serviceAccountName: memory-app
restartPolicy: Never
containers:
# Step 1: Run migrations
- name: migrate
image: forgejo.riotpiao.com/riotpiao-poimen/poimen-memory:IMAGE_SHA
imagePullPolicy: IfNotPresent
command:
- /bin/bash
- -c
- |
set -e
# Copy migrations script from image to working dir
cp /app/migrations/run_migrations.sh /tmp/run_migrations.sh
chmod +x /tmp/run_migrations.sh
# Run migrations
/tmp/run_migrations.sh
echo ""
echo "✓ Migrations complete"
echo "Database ready for tests"
env:
- name: DB_HOST
value: "memory-db-rw.poimen.svc.cluster.local"
- name: DB_PORT
value: "5432"
- name: DB_NAME
value: "memory"
- name: DB_USER
value: "app"
- name: DB_PASSWORD
valueFrom:
secretKeyRef:
name: memory-db-app
key: password
resources:
requests:
memory: "256Mi"
cpu: "100m"
limits:
memory: "512Mi"
cpu: "500m"
# Step 2: Run integration tests
- name: test
image: forgejo.riotpiao.com/riotpiao-poimen/poimen-memory:IMAGE_SHA
imagePullPolicy: IfNotPresent
command:
- /bin/bash
- -c
- |
set -e
echo "=========================================="
echo "Integration Test: Ingest + Embedding"
echo "=========================================="
echo ""
# Start HTTP server
echo "Starting memory-service..."
mem-cli serve --port 8080 &
SERVER_PID=$!
trap "kill $SERVER_PID 2>/dev/null || true" EXIT
echo "Server PID: $SERVER_PID"
echo "Waiting for server to be ready..."
# Wait for /health endpoint
for i in {1..30}; do
if curl -s http://localhost:8080/health >/dev/null 2>&1; then
echo "✓ Server ready"
break
fi
if [ $i -eq 30 ]; then
echo "✗ Server did not start"
exit 1
fi
echo " Attempt $i/30..."
sleep 1
done
echo ""
echo "Running E2E ingest test..."
echo ""
# Send ingest request
INGEST_ID="test-$(date +%s)"
RESPONSE=$(curl -s -X POST http://localhost:8080/memory/ingest \
-H "Content-Type: application/json" \
-H "Authorization: Bearer test-key" \
-d "{
\"project\": \"integration-test\",
\"source\": \"k8s-job-test\",
\"ingest_id\": \"$INGEST_ID\",
\"records\": [
{
\"role\": \"user\",
\"text\": \"Kubernetes [[Docker]] [[Linux]] container platform\",
\"timestamp\": \"2026-09-14T13:00:00Z\",
\"source_position\": 0
},
{
\"role\": \"user\",
\"text\": \"Docker [[Container]] microservices architecture\",
\"timestamp\": \"2026-09-14T13:01:00Z\",
\"source_position\": 1
}
]
}")
# Check response
STATUS=$(echo "$RESPONSE" | jq -r '.status // "error"')
ID=$(echo "$RESPONSE" | jq -r '.ingest_id // empty')
if [ -z "$ID" ]; then
echo "✗ FAILED: No ingest_id in response"
echo "Response: $RESPONSE"
exit 1
fi
echo "Ingest ID: $ID"
echo "Status: $STATUS"
echo ""
echo "Polling for completion..."
# Poll until done
for poll in {1..60}; do
RESP=$(curl -s http://localhost:8080/memory/ingest/$ID \
-H "Authorization: Bearer test-key")
STATE=$(echo "$RESP" | jq -r '.status // "unknown"')
if [ "$STATE" = "done" ]; then
echo "Poll $poll: $STATE ✓"
echo ""
echo "✓ INGEST SUCCESSFUL"
break
elif [ "$STATE" = "failed" ] || [ "$STATE" = "error" ]; then
echo "Poll $poll: $STATE ✗"
echo "Response: $RESP"
echo "✗ INGEST FAILED"
exit 1
fi
echo "Poll $poll: $STATE"
sleep 2
done
echo ""
echo "Testing query endpoint..."
QUERY=$(curl -s "http://localhost:8080/memory/query?project=integration-test&question=what%20is%20docker" \
-H "Authorization: Bearer test-key")
ENTITY_COUNT=$(echo "$QUERY" | jq '.count.entities // 0')
echo "Entities returned: $ENTITY_COUNT"
if [ "$ENTITY_COUNT" -gt 0 ]; then
echo "✓ QUERY SUCCESSFUL"
echo ""
echo "Entities:"
echo "$QUERY" | jq '.entities[].name'
else
echo "⚠ No entities returned (schema issue)"
echo "✗ Query test FAILED"
exit 1
fi
echo ""
echo "=========================================="
echo "✓ ALL TESTS PASSED"
echo "=========================================="
env:
- name: DATABASE_URL
value: "postgresql://[email protected]:5432/memory"
- name: RUST_LOG
value: "info,mem_cli=debug,mem_ingest=debug"
- name: MEM_AUTH_MODE
value: "none"
# Password via secret
- name: PGPASSWORD
valueFrom:
secretKeyRef:
name: memory-db-app
key: password
resources:
requests:
memory: "512Mi"
cpu: "200m"
limits:
memory: "1Gi"
cpu: "1000m"
livenessProbe:
exec:
command:
- /bin/sh
- -c
- curl -s http://localhost:8080/health >/dev/null
initialDelaySeconds: 30
periodSeconds: 10
failureThreshold: 2
---
# ServiceAccount for integration test
apiVersion: v1
kind: ServiceAccount
metadata:
name: memory-app
namespace: poimen
labels:
app: poimen-memory
-127
View File
@@ -1,127 +0,0 @@
#!/bin/bash
# Database Migration Runner
# Used by K8s Job to apply all migrations before integration tests
#
# Environment variables (from K8s):
# DB_HOST - PostgreSQL host
# DB_PORT - PostgreSQL port
# DB_NAME - Database name
# DB_USER - Database user
# DB_PASSWORD - Database password (from Secret)
set -e
DB_HOST="${DB_HOST:-memory-db-rw.poimen.svc.cluster.local}"
DB_PORT="${DB_PORT:-5432}"
DB_NAME="${DB_NAME:-memory}"
DB_USER="${DB_USER:-app}"
if [ -z "$DB_PASSWORD" ]; then
echo "ERROR: DB_PASSWORD not set"
exit 1
fi
echo "=========================================="
echo "Database Migration Runner"
echo "=========================================="
echo ""
echo "Configuration:"
echo " Host: $DB_HOST:$DB_PORT"
echo " Database: $DB_NAME"
echo " User: $DB_USER"
echo ""
# Export for psql
export PGPASSWORD="$DB_PASSWORD"
# Get migration directory (where this script is)
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
MIGRATION_DIR="$SCRIPT_DIR"
echo "Migration directory: $MIGRATION_DIR"
echo ""
# Collect all SQL files
MIGRATIONS=($(ls -1 "$MIGRATION_DIR"/*.sql 2>/dev/null | sort))
if [ ${#MIGRATIONS[@]} -eq 0 ]; then
echo "ERROR: No migration files found in $MIGRATION_DIR"
exit 1
fi
echo "Found ${#MIGRATIONS[@]} migration(s):"
for m in "${MIGRATIONS[@]}"; do
echo " - $(basename $m)"
done
echo ""
# Wait for DB to be ready
echo "Waiting for database to be ready..."
for i in {1..30}; do
if psql -h "$DB_HOST" -p "$DB_PORT" -U "$DB_USER" -d "$DB_NAME" -c "SELECT 1;" >/dev/null 2>&1; then
echo "✓ Database is ready"
break
fi
if [ $i -eq 30 ]; then
echo "✗ Database not ready after 30 attempts"
exit 1
fi
echo " Attempt $i/30..."
sleep 1
done
echo ""
echo "=========================================="
echo "Running Migrations"
echo "=========================================="
echo ""
SUCCESS=0
FAILED=0
for migration in "${MIGRATIONS[@]}"; do
name=$(basename "$migration")
echo -n "$name ... "
if psql -h "$DB_HOST" -p "$DB_PORT" -U "$DB_USER" -d "$DB_NAME" -f "$migration" >/dev/null 2>&1; then
echo "✓"
((SUCCESS++))
else
echo "✗ FAILED"
echo ""
echo "Error output:"
psql -h "$DB_HOST" -p "$DB_PORT" -U "$DB_USER" -d "$DB_NAME" -f "$migration" 2>&1 | sed 's/^/ /'
((FAILED++))
fi
done
echo ""
echo "=========================================="
echo "Migration Summary"
echo "=========================================="
echo " Success: $SUCCESS"
echo " Failed: $FAILED"
echo ""
if [ $FAILED -eq 0 ]; then
echo "✓ All migrations applied successfully"
echo ""
echo "Verifying schema..."
echo ""
# Verify key tables exist
for table in memory_entity memory_edge ingest_jobs; do
if psql -h "$DB_HOST" -p "$DB_PORT" -U "$DB_USER" -d "$DB_NAME" -c "SELECT 1 FROM information_schema.tables WHERE table_name='$table';" 2>&1 | grep -q "1 row"; then
echo " ✓ Table $table exists"
else
echo " ⚠ Table $table not found"
fi
done
exit 0
else
echo "✗ Some migrations failed"
exit 1
fi
-81
View File
@@ -1,81 +0,0 @@
#!/bin/bash
# Master script: Run full production ingest test with logging
#
# Usage:
# ./run_production_test.sh
# ./run_production_test.sh [api-key]
#
# What it does:
# 1. Collect baseline logs
# 2. Run ingest test
# 3. Collect post-test logs
# 4. Analyze for errors
# 5. Display results
set -e
API_KEY="${1:-test-key}"
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
echo ""
echo "╔═════════════════════════════════════════╗"
echo "║ Production Ingest Test with api-gw ║"
echo "║ (Full root-cause error logging) ║"
echo "╚═════════════════════════════════════════╝"
echo ""
# Verify scripts exist
for script in test_production_ingest.sh collect_prod_logs.sh; do
if [ ! -f "$SCRIPT_DIR/$script" ]; then
echo "✗ Missing: $script"
exit 1
fi
done
echo "Step 1: Collecting baseline logs..."
"$SCRIPT_DIR/collect_prod_logs.sh" before
echo ""
echo "Step 2: Running ingest test..."
echo " (Sending records through embedding pipeline to api-gw)"
echo ""
if MEM_API_KEY="$API_KEY" "$SCRIPT_DIR/test_production_ingest.sh"; then
echo ""
echo "✓ Test passed!"
test_status=0
else
echo ""
echo "✗ Test failed!"
test_status=1
fi
echo ""
echo "Step 3: Collecting post-test logs..."
"$SCRIPT_DIR/collect_prod_logs.sh" after
echo ""
echo "Step 4: Analyzing logs for errors..."
echo ""
"$SCRIPT_DIR/collect_prod_logs.sh" analyze
echo ""
echo "════════════════════════════════════════"
if [ $test_status -eq 0 ]; then
echo "✓ INGEST TEST PASSED"
else
echo "✗ INGEST TEST FAILED"
echo ""
echo "Next steps:"
echo " 1. Check logs in prod_logs_*/ directory"
echo " 2. Look for errors in:"
echo " - /memory/ingest endpoint response"
echo " - Embedding service (LLM_ENDPOINT)"
echo " - api-gw gateway logs"
echo " - Database connection"
fi
echo "════════════════════════════════════════"
echo ""
exit $test_status
-197
View File
@@ -1,197 +0,0 @@
#!/bin/bash
# Production real test: Full ingest with api-gw + embedding
# Sends records through the complete pipeline and logs all errors
#
# Usage:
# ./test_prod_ingest_real.sh [--verbose]
set -e
NAMESPACE="poimen"
SERVICE="poimen-memory"
LOCAL_PORT="9990"
TIMESTAMP=$(date +%Y%m%d-%H%M%S)
LOG_FILE="/tmp/ingest_test_${TIMESTAMP}.log"
VERBOSE="${1:-}"
{
echo "=========================================="
echo "Production Ingest Test: $(date)"
echo "=========================================="
echo ""
echo "Namespace: $NAMESPACE"
echo "Service: $SERVICE"
echo "Local Port: $LOCAL_PORT"
echo "Log: $LOG_FILE"
echo ""
# Start port-forward
echo "Starting port-forward..."
kubectl -n "$NAMESPACE" port-forward "svc/$SERVICE" "$LOCAL_PORT:8080" >/dev/null 2>&1 &
PF_PID=$!
cleanup() {
if [ -n "$PF_PID" ]; then
kill $PF_PID 2>/dev/null || true
wait $PF_PID 2>/dev/null || true
fi
}
trap cleanup EXIT
sleep 2
if ! kill -0 $PF_PID 2>/dev/null; then
echo "✗ Port-forward failed"
exit 1
fi
echo "✓ Port-forward running (PID $PF_PID)"
echo ""
# Check health
echo "Checking /health endpoint..."
if ! curl -s "http://localhost:$LOCAL_PORT/health" >/dev/null 2>&1; then
echo "✗ Health check failed"
exit 1
fi
echo "✓ Health check passed"
echo ""
# Prepare ingest request
INGEST_ID="ingest-test-${TIMESTAMP}"
PAYLOAD=$(cat <<'EOFPAYLOAD'
{
"project": "production-real-test",
"source": "integration-test",
"ingest_id": "INGEST_ID_PLACEHOLDER",
"records": [
{
"role": "user",
"text": "Kubernetes [[Docker]] [[Linux]] is an open-source container orchestration platform. It automates many manual processes involved in deploying, managing, and scaling containerized applications.",
"timestamp": "2026-09-14T13:00:00Z",
"source_position": 0
},
{
"role": "user",
"text": "Docker [[Container]] [[Go]] is a containerization platform that makes it easier to build, ship, and run applications. Docker achieves high efficiency through the use of operating system-level virtualization.",
"timestamp": "2026-09-14T13:01:00Z",
"source_position": 1
},
{
"role": "user",
"text": "Go [[Concurrency]] [[Static Typing]] is a programming language designed at Google. It is statically typed, compiled, and known for its simplicity, concurrent programming model, and efficient execution.",
"timestamp": "2026-09-14T13:02:00Z",
"source_position": 2
}
]
}
EOFPAYLOAD
)
# Replace placeholder
PAYLOAD="${PAYLOAD//INGEST_ID_PLACEHOLDER/$INGEST_ID}"
echo "Sending ingest request..."
if [ -n "$VERBOSE" ]; then
echo "Payload:"
echo "$PAYLOAD" | jq . 2>/dev/null || echo "$PAYLOAD"
echo ""
fi
RESPONSE=$(curl -s -w "\n%{http_code}" -X POST \
"http://localhost:$LOCAL_PORT/memory/ingest" \
-H "Content-Type: application/json" \
-H "Authorization: Bearer test-key" \
-d "$PAYLOAD")
HTTP_CODE=$(echo "$RESPONSE" | tail -1)
BODY=$(echo "$RESPONSE" | head -n-1)
echo "HTTP Status: $HTTP_CODE"
if [ "$HTTP_CODE" != "202" ]; then
echo "✗ Unexpected HTTP status"
echo "Response: $BODY"
exit 1
fi
echo "✓ Request accepted"
echo ""
if [ -n "$VERBOSE" ]; then
echo "Response body:"
echo "$BODY" | jq . 2>/dev/null || echo "$BODY"
echo ""
fi
# Extract ID
ID=$(echo "$BODY" | jq -r '.ingest_id // empty' 2>/dev/null)
if [ -z "$ID" ]; then
echo "✗ Missing ingest_id in response"
echo "Response: $BODY"
exit 1
fi
echo "Ingest ID: $ID"
echo ""
# Poll status
echo "Polling job status..."
echo "=========================================="
MAX_POLLS=120 # 10 minutes at 5s intervals
poll_count=0
while [ $poll_count -lt $MAX_POLLS ]; do
poll_count=$((poll_count+1))
STATUS_RESP=$(curl -s "http://localhost:$LOCAL_PORT/memory/ingest/$ID" \
-H "Authorization: Bearer test-key")
STATUS=$(echo "$STATUS_RESP" | jq -r '.status // "unknown"' 2>/dev/null)
printf "[%3d] %-20s" "$poll_count" "$STATUS"
case "$STATUS" in
done)
echo " ✓"
echo "=========================================="
echo ""
echo "✓ SUCCESS: Ingest completed"
if [ -n "$VERBOSE" ]; then
echo ""
echo "Final response:"
echo "$STATUS_RESP" | jq . 2>/dev/null || echo "$STATUS_RESP"
fi
exit 0
;;
failed|error)
echo " ✗"
echo "=========================================="
echo ""
echo "✗ FAILED: Ingest did not complete"
echo ""
echo "Final response:"
echo "$STATUS_RESP" | jq . 2>/dev/null || echo "$STATUS_RESP"
exit 1
;;
processing|queued|pending)
echo ""
sleep 5
;;
*)
echo " (unknown)"
sleep 5
;;
esac
done
echo "=========================================="
echo ""
echo "✗ TIMEOUT: Ingest did not complete after ${MAX_POLLS} polls (${poll_count}m)"
exit 1
} 2>&1 | tee "$LOG_FILE"
echo ""
echo "Full log saved to: $LOG_FILE"
-286
View File
@@ -1,286 +0,0 @@
//! Integration test: Full ingest + embedding flow with api-gw
//!
//! Tests:
//! 1. POST /memory/ingest with sample records
//! 2. Poll /memory/ingest/{id} until done
//! 3. Log root causes of errors
//!
//! Requires:
//! - DATABASE_URL set (postgres)
//! - LLM_ENDPOINT set (for embeddings)
//! - Server running locally or started by test
//!
//! Usage:
//! ```
//! RUST_LOG=debug cargo test --test integration_ingest_with_gw -- --nocapture
//! ```
use std::env;
use std::time::Duration;
use tokio::time::sleep;
use serde_json::json;
#[tokio::test]
#[ignore] // Run manually: cargo test --test integration_ingest_with_gw -- --ignored --nocapture
async fn test_ingest_with_embeddings_and_logging() {
// Initialize tracing with DEBUG level to see all logs
let _ = tracing_subscriber::fmt()
.with_max_level(tracing::Level::DEBUG)
.with_writer(std::io::stderr)
.try_init();
let base_url = env::var("MEM_API_URL").unwrap_or_else(|_| "http://localhost:8080".to_string());
let api_key = env::var("MEM_API_KEY").unwrap_or_else(|_| "test-key".to_string());
let client = reqwest::Client::new();
// Sample ingest payload
let payload = json!({
"project": "test-project",
"records": [
{
"content": "Kubernetes is an open-source container orchestration platform. [[Docker]] [[Go]]",
"source": "wiki/kubernetes"
},
{
"content": "Docker is a containerization platform that makes it easier to build, ship, and run applications. [[Linux]] [[Container]]",
"source": "wiki/docker"
},
{
"content": "Go is a programming language designed at Google. [[Concurrency]] [[Static Typing]]",
"source": "wiki/go"
}
]
});
println!("[TEST] Sending ingest request...");
tracing::info!(
target: "integration_test",
"Ingest payload: {}",
serde_json::to_string_pretty(&payload).unwrap()
);
// POST /memory/ingest
let response = match client
.post(&format!("{}/memory/ingest", base_url))
.header("Authorization", format!("Bearer {}", api_key))
.json(&payload)
.send()
.await
{
Ok(resp) => resp,
Err(e) => {
eprintln!("[ERROR] Failed to send ingest request: {}", e);
tracing::error!(
target: "integration_test",
error = %e,
"Failed to POST /memory/ingest"
);
panic!("Request failed: {}", e);
}
};
let status = response.status();
println!("[TEST] Ingest response status: {}", status);
let body_text = match response.text().await {
Ok(text) => text,
Err(e) => {
tracing::error!(target: "integration_test", error = %e, "Failed to read response body");
panic!("Failed to read response body: {}", e);
}
};
println!("[TEST] Response body:\n{}", body_text);
// Parse response
let resp_json: serde_json::Value = match serde_json::from_str(&body_text) {
Ok(j) => j,
Err(e) => {
tracing::error!(
target: "integration_test",
error = %e,
body = %body_text,
"Failed to parse JSON response"
);
panic!("Failed to parse JSON: {}", e);
}
};
let ingest_id = match resp_json["id"].as_str() {
Some(id) => id.to_string(),
None => {
tracing::error!(
target: "integration_test",
response = %serde_json::to_string_pretty(&resp_json).unwrap(),
"Missing 'id' in response"
);
panic!("Missing 'id' in response: {}", resp_json);
}
};
println!("[TEST] Ingest ID: {}", ingest_id);
tracing::info!(target: "integration_test", ingest_id = %ingest_id, "Ingest queued");
// Poll until complete or timeout
let max_polls = 60; // 10 minutes with 10s intervals
for poll_num in 1..=max_polls {
sleep(Duration::from_secs(10)).await;
println!(
"[TEST] Poll #{}/{}: Checking status of ingest {}",
poll_num, max_polls, ingest_id
);
let status_response = match client
.get(&format!("{}/memory/ingest/{}", base_url, ingest_id))
.header("Authorization", format!("Bearer {}", api_key))
.send()
.await
{
Ok(resp) => resp,
Err(e) => {
tracing::error!(
target: "integration_test",
error = %e,
ingest_id = %ingest_id,
poll = poll_num,
"Failed to fetch status"
);
eprintln!("[ERROR] Failed to fetch status: {}", e);
sleep(Duration::from_secs(5)).await;
continue;
}
};
let status_text = match status_response.text().await {
Ok(text) => text,
Err(e) => {
tracing::error!(
target: "integration_test",
error = %e,
ingest_id = %ingest_id,
"Failed to read status response"
);
eprintln!("[ERROR] Failed to read status: {}", e);
continue;
}
};
let status_json: serde_json::Value = match serde_json::from_str(&status_text) {
Ok(j) => j,
Err(e) => {
tracing::error!(
target: "integration_test",
error = %e,
body = %status_text,
"Failed to parse status JSON"
);
eprintln!("[ERROR] Failed to parse status JSON: {}", e);
continue;
}
};
let status = status_json["status"].as_str().unwrap_or("unknown");
println!(
"[TEST] Poll #{}: status = {}",
poll_num, status
);
tracing::info!(
target: "integration_test",
ingest_id = %ingest_id,
poll = poll_num,
status = %status,
full_response = %serde_json::to_string_pretty(&status_json).unwrap(),
"Status check"
);
match status {
"done" => {
println!("[TEST] ✓ Ingest completed successfully!");
tracing::info!(target: "integration_test", "Ingest completed");
// Extract and log results
if let Some(results) = status_json.get("results") {
println!("[TEST] Results:\n{}", serde_json::to_string_pretty(results).unwrap());
tracing::info!(
target: "integration_test",
results = %serde_json::to_string_pretty(results).unwrap(),
"Ingest results"
);
}
return;
}
"failed" | "error" => {
let error_msg = status_json["error"].as_str().unwrap_or("unknown error");
println!("[TEST] ✗ Ingest FAILED: {}", error_msg);
tracing::error!(
target: "integration_test",
ingest_id = %ingest_id,
error = %error_msg,
full_response = %serde_json::to_string_pretty(&status_json).unwrap(),
"Ingest failed"
);
panic!("Ingest failed: {}", error_msg);
}
"processing" | "queued" => {
// Continue polling
println!("[TEST] Still processing, poll again...");
}
_ => {
println!("[TEST] Unknown status: {}", status);
tracing::warn!(target: "integration_test", status = %status, "Unknown status");
}
}
}
// Timeout
let msg = format!("Ingest did not complete after {} polls (timeout)", max_polls);
println!("[TEST] ✗ {}", msg);
tracing::error!(target: "integration_test", ingest_id = %ingest_id, "Ingest timeout");
panic!("{}", msg);
}
#[tokio::test]
#[ignore]
async fn test_ingest_endpoint_only() {
let _ = tracing_subscriber::fmt()
.with_max_level(tracing::Level::DEBUG)
.try_init();
let base_url = env::var("MEM_API_URL").unwrap_or_else(|_| "http://localhost:8080".to_string());
let api_key = env::var("MEM_API_KEY").unwrap_or_else(|_| "test-key".to_string());
let client = reqwest::Client::new();
let payload = json!({
"project": "test-project",
"records": [
{
"content": "Simple test record",
"source": "test"
}
]
});
println!("[TEST] Testing /memory/ingest endpoint only");
let response = client
.post(&format!("{}/memory/ingest", base_url))
.header("Authorization", format!("Bearer {}", api_key))
.json(&payload)
.send()
.await
.expect("Failed to send request");
println!("[TEST] Status: {}", response.status());
let body = response.text().await.expect("Failed to read body");
println!("[TEST] Response: {}", body);
let json: serde_json::Value = serde_json::from_str(&body).expect("Invalid JSON");
println!("[TEST] Parsed: {}", serde_json::to_string_pretty(&json).unwrap());
assert!(json.get("id").is_some(), "Response should contain 'id'");
}
-221
View File
@@ -1,221 +0,0 @@
//! Unit test: Ingest pipeline with detailed error logging
//!
//! Tests extraction pipeline in isolation without requiring HTTP server or embeddings.
//! Useful for debugging extraction errors.
//!
//! Usage:
//! ```
//! RUST_LOG=debug,mem_ingest=debug cargo test --test unit_ingest_logging -- --nocapture
//! ```
#[cfg(test)]
mod tests {
use mem_ingest::ingest_pipeline::{IngestPipeline, Episode};
use mem_ingest::entity_extractor::WikiLinkFallbackExtractor;
use mem_ingest::fact_extractor::SimpleFactExtractor;
use mem_ingest::contradiction_detector::ContradictionHandler;
use std::sync::Arc;
fn init_logging() {
let _ = tracing_subscriber::fmt()
.with_max_level(tracing::Level::DEBUG)
.with_writer(std::io::stderr)
.try_init();
}
#[tokio::test]
async fn test_wiki_link_extraction() {
init_logging();
println!("\n[TEST] Wiki link extraction with logging\n");
let entity_extractor = Arc::new(WikiLinkFallbackExtractor);
let fact_extractor = Arc::new(SimpleFactExtractor);
let contradiction_detector = Arc::new(ContradictionHandler::default());
let pipeline = IngestPipeline::new(
entity_extractor,
fact_extractor,
contradiction_detector,
);
let episode = Episode {
id: "test-1".to_string(),
project_id: "test-project".to_string(),
text: "Kubernetes [[Docker]] is a [[Container]] orchestration platform. It works with [[Go]] programs."
.to_string(),
wiki_links: vec!["Docker".to_string(), "Container".to_string(), "Go".to_string()],
};
tracing::info!(
target: "test",
episode_id = %episode.id,
wiki_links = ?episode.wiki_links,
"Starting pipeline ingest"
);
match pipeline.ingest(&episode).await {
Ok(result) => {
tracing::info!(
target: "test",
entities = result.entities.len(),
edges = result.edges.len(),
reviews = result.reviews.len(),
"Pipeline succeeded"
);
println!("✓ Extracted {} entities", result.entities.len());
for entity in &result.entities {
println!(" - {} ({}): {}", entity.name, entity.entity_type.as_str(), entity.summary.as_deref().unwrap_or(""));
}
println!("✓ Extracted {} edges", result.edges.len());
for edge in &result.edges {
println!(" - {} --[{}]--> {}", edge.source_entity_id, edge.relation_type, edge.target_entity_id);
}
assert!(result.entities.len() > 0, "Should extract entities");
}
Err(e) => {
tracing::error!(
target: "test",
error = %e,
"Pipeline failed"
);
panic!("Pipeline failed: {}", e);
}
}
}
#[tokio::test]
async fn test_extraction_error_logging() {
init_logging();
println!("\n[TEST] Pipeline error handling with logging\n");
let entity_extractor = Arc::new(WikiLinkFallbackExtractor);
let fact_extractor = Arc::new(SimpleFactExtractor);
let contradiction_detector = Arc::new(ContradictionHandler::default());
let pipeline = IngestPipeline::new(
entity_extractor,
fact_extractor,
contradiction_detector,
);
// Episode with problematic content (empty, or only whitespace)
let episode = Episode {
id: "test-empty".to_string(),
project_id: "test-project".to_string(),
text: "".to_string(),
wiki_links: vec![],
};
tracing::info!(
target: "test",
episode_id = %episode.id,
text_len = episode.text.len(),
"Processing empty episode"
);
match pipeline.ingest(&episode).await {
Ok(result) => {
tracing::info!(
target: "test",
entities = result.entities.len(),
edges = result.edges.len(),
"Empty episode processed (no error expected)"
);
println!("✓ Empty episode handled gracefully");
}
Err(e) => {
tracing::error!(
target: "test",
error = %e,
"Empty episode caused error"
);
// Empty is OK for some extractors
println!("⚠ Empty episode error (may be expected): {}", e);
}
}
}
#[tokio::test]
async fn test_multiple_records_error_accumulation() {
init_logging();
println!("\n[TEST] Processing multiple records and logging errors\n");
let entity_extractor = Arc::new(WikiLinkFallbackExtractor);
let fact_extractor = Arc::new(SimpleFactExtractor);
let contradiction_detector = Arc::new(ContradictionHandler::default());
let pipeline = IngestPipeline::new(
entity_extractor,
fact_extractor,
contradiction_detector,
);
let records = vec![
("Kubernetes [[Docker]] is a container orchestrator", "wiki/k8s"),
("Docker [[Linux]] containers enable microservices", "wiki/docker"),
("", "wiki/empty"),
("Go [[Concurrency]] is powerful for backend services", "wiki/go"),
];
let mut success_count = 0;
let mut error_count = 0;
for (idx, (text, source)) in records.iter().enumerate() {
let episode = Episode {
id: format!("record-{}", idx),
project_id: "test-project".to_string(),
text: text.to_string(),
wiki_links: vec![],
};
tracing::info!(
target: "test",
record_idx = idx,
source = source,
text_len = text.len(),
"Processing record"
);
match pipeline.ingest(&episode).await {
Ok(result) => {
tracing::debug!(
target: "test",
record_idx = idx,
entities = result.entities.len(),
edges = result.edges.len(),
"Record succeeded"
);
println!(" ✓ Record {}: {} entities, {} edges", idx, result.entities.len(), result.edges.len());
success_count += 1;
}
Err(e) => {
tracing::warn!(
target: "test",
record_idx = idx,
error = %e,
source = source,
"Record failed"
);
println!(" ✗ Record {}: {}", idx, e);
error_count += 1;
}
}
}
println!("\nSummary: {} success, {} errors", success_count, error_count);
tracing::info!(
target: "test",
total_records = records.len(),
success = success_count,
errors = error_count,
"Batch processing complete"
);
}
}