Skip to content

Optimization Module Study Guide

Package: cogniverse_agents (Implementation Layer), cogniverse_runtime (Application Layer), cogniverse_synthetic (Implementation Layer) Module Location: libs/agents/cogniverse_agents/optimizer/, libs/agents/cogniverse_agents/routing/ (training-decision models), libs/runtime/cogniverse_runtime/optimization_cli.py, libs/synthetic/


Package Structure

libs/agents/cogniverse_agents/optimizer/
├── artifact_manager.py             # ArtifactManager: ExperimentMetrics, promote_if_better,
│                                    #   promote_to_canary, promote_canary_to_active, retire_canary,
│                                    #   rollback_to_version, snapshot_active, save/load prompts+demos+blobs
├── dspy_agent_optimizer.py         # DSPyAgentPromptOptimizer + DSPyAgentOptimizerPipeline
│                                    #   (BootstrapFewShot prompt compilation for 4 agent signatures)
├── signature_variants.py           # SignatureVariantRegistry: per-tenant DSPy signature variants
└── strategy_learner.py             # StrategyLearner: pattern + LLM distillation from traces into Mem0

libs/agents/cogniverse_agents/routing/
├── xgboost_meta_models.py          # TrainingDecisionModel, TrainingStrategyModel, FusionBenefitModel
├── profile_performance_optimizer.py  # ProfilePerformanceOptimizer (see docs/modules/routing.md)
└── config.py                       # OnlineEvaluationConfig, AutomationRulesConfig
                                     # (remaining routing/ files — annotation, dspy_relationship_router,
                                     #  orchestration_evaluator, etc. — are documented in routing.md)

libs/runtime/cogniverse_runtime/
├── optimization_cli.py             # CLI for all optimization/maintenance modes:
│                                    # cleanup | triggered | simba | workflow | gateway-thresholds |
│                                    # online-routing-eval | llm-annotate | online-eval | profile | entity-extraction |
│                                    # synthetic |
│                                    # rollback | ab-compare | egress-netpol | monthly-reports
├── quality_monitor_cli.py          # QualityMonitor driver — submits `--mode triggered` Argo Workflows
│                                    #   on quality drops; `--once` forces a distillation pass
└── routers/tenant.py               # POST /admin/tenant/{id}/optimize (on-demand submit) + status/cancel/retry

libs/synthetic/cogniverse_synthetic/
├── service.py                      # SyntheticDataService — orchestrates a generator end-to-end
├── api.py                          # FastAPI router (prefix /synthetic): generate, batch/generate,
│                                    #   optimizers, optimizers/{name}, health
├── registry.py                     # OPTIMIZER_REGISTRY, OptimizerConfig, list_optimizers()
├── schemas.py                      # SyntheticDataRequest/Response, ProfileSelectionExampleSchema,
│                                    #   RoutingExperienceSchema, WorkflowExecutionSchema
├── dspy_modules.py / dspy_signatures.py  # DSPy modules used by generators (e.g. query generation)
├── generators/                     # Optimizer-specific generators
│   ├── base.py                     # Base generator classes
│   ├── profile.py                  # ProfileGenerator: ProfileSelectionAgent training data
│   ├── routing.py                  # RoutingGenerator: routing training data
│   └── workflow.py                 # WorkflowGenerator: workflow-orchestration training data
├── profile_selector.py             # LLM-based profile selection
├── backend_querier.py              # Vespa content sampling
├── approval/                       # Human-in-the-loop approval workflow for synthetic demos
│   ├── confidence_extractor.py     # Extracts confidence signal from generated examples
│   └── feedback_handler.py         # Approve/reject feedback processing
└── utils/                          # Pattern extraction and agent inference

Table of Contents

  1. Module Overview
  2. Architecture
  3. Core Components
  4. Usage Examples
  5. Production Considerations
  6. Testing

Module Overview

Purpose

The Optimization Module provides on-demand, per-agent DSPy prompt/module compilation, gateway threshold calibration, workflow-template learning, XGBoost-based training-decision models, and Mem0-backed strategy distillation. It reads Phoenix telemetry spans as training signal and persists compiled artefacts (prompts, demos, and DSPy module state) via ArtifactManager, which agents reload on each dispatch (the cached gateway agent refreshes on a GATEWAY_ARTIFACT_TTL_S interval).

Key Features

  • DSPy Prompt/Module Compilation: every optimization mode compiles with dspy.teleprompt.BootstrapFewShot (scaled to 8/16/2-round settings once ≥50 training examples are available, 4/8/1-round below that). The one exception is --mode triggered's reflective recompile for an all-failure agent, which compiles with dspy.teleprompt.GEPA (see Triggered Optimization); no mode invokes MIPROv2 or SIMBA. ExperimentMetrics.optimizer is a free-text field intended to record whichever optimizer produced a run, but every serving call site passes "BootstrapFewShot". The --mode simba name is historical (it optimizes QueryEnhancementAgent's DSPy module); it does not run the SIMBA algorithm.
  • Gateway Threshold Tuning: _compute_gateway_thresholds() derives GLiNER/fast-path thresholds from Phoenix cogniverse.gateway spans using a rule-based adjustment (± based on per-branch error rate and mean confidence) plus a p25-percentile-derived gliner_threshold.
  • Online Routing Evaluation: run_online_routing_evaluation scores cogniverse.routing spans (routing outcome + confidence calibration) via OnlineEvaluator and persists the scores as telemetry annotations; driven by automation_rules.online_evaluation (OnlineEvaluationConfig in routing/config.py).
  • Profile Selection / Entity Extraction / SIMBA (query enhancement) Optimization: each reads its own Phoenix span kind through TraceStore.get_all_spans() so no fixed result ceiling can discard older examples, merges tenant-approved synthetic examples even when that span window is empty, projects each approved record onto the exact production DSPy signature, compiles the agent's DSPy module, and saves it as a ("model", <key>) blob via ArtifactManager. Bootstrap candidates are accepted only when every output field exactly matches the reviewed label; otherwise the compiled module retains the labeled example instead of replacing it with teacher-generated content. Current artifacts are contract-checked against the live module before rehydration, and a signature mismatch is treated as missing for scoring rather than overwriting the prompt. The versioned ledger keeps the scored flag, the numeric score used for confirmation when present, and metric_id naming the metric that produced it; training_selection.<optimizer>.confirmation_score_threshold makes confirmation score-aware, and an omitted threshold keeps confirmation presence-based and ignores metric_id. Under a threshold, a score counts only when its metric_id is the optimizer's current one — the ids live beside the metric functions in optimization_cli and are registered in OPTIMIZER_METRIC_IDS. Promotions with no score, and promotions scored under any other metric, are unknown rather than negative evidence: they stay in the unscored bucket, so decay only applies when the known confirmations plus the unscored history still fall below low_confirmation_threshold. SIMBA also enforces the tenant floor from routing config: below min_samples_for_optimization or min_unique_queries, it saves a version with decision insufficient_population and leaves the active artifact unchanged.
  • Monthly Performance Reports: each tenant's bounded-time Phoenix window is read through TraceStore.get_spans(..., columns=("start_time", "end_time", "status_code")), then chunked and spilled to disk so latency and error summaries stay exact without retaining the full span window; a tenant read failure is recorded and makes the command fail instead of being reported as an empty history.
  • Workflow Orchestration Optimization: extracts WorkflowExecution records from cogniverse.orchestration spans via OrchestrationEvaluator, drops demos whose agent_sequence references an agent no longer live in configs/config.json, and replaces templates, execution demos, agent performance profiles, and query patterns through one WorkflowStoreRegistry operation guarded by a renewable per-tenant Redis lease.
  • Training-Decision Meta-Models: TrainingDecisionModel (train/skip gate; wired into QualityMonitor), TrainingStrategyModel (PURE_REAL / HYBRID / SYNTHETIC / SKIP selection), and FusionBenefitModel — all XGBoost classifiers in routing/xgboost_meta_models.py, persisted via the same ArtifactManager.
  • Strategy Distillation: StrategyLearner distills execution traces into reusable Strategy objects (pattern extraction, no LLM; or LLM-based contrastive distillation) and stores them in Vespa memory via Mem0MemoryManager for later retrieval by MemoryAwareMixin.
  • Regression-Reject Promotion Gate: ArtifactManager.promote_if_better only promotes a candidate when its score beats the active baseline (within tolerance); every attempt — win or reject — lands as a typed ExperimentMetrics row.
  • Canary Rollout + Rollback: ArtifactManager's three-slot (active/canary/retired) state machine and the --mode rollback CLI restore previously-snapshotted prompt/demo versions.
  • On-Demand Workflows: the web client (or any client) triggers POST /admin/tenant/{id}/optimize, which submits an Argo Workflow referencing the cogniverse-optimization-runner WorkflowTemplate.
  • Scheduled Workflows: helm-chart CronWorkflows run gateway-thresholds/entity-extraction/simba/ profile/workflow weekly, gateway-thresholds daily, cleanup daily, synthetic weekly, a forced distillation pass daily via quality_monitor_cli --once, and monthly-reports monthly.

Dependencies

Note: Optimizer classes require full module path imports as they are not exported at package level.

# CLI optimization
from cogniverse_runtime.optimization_cli import _compute_gateway_thresholds, GATEWAY_DEFAULT_THRESHOLD

# Synthetic Data Generation (exported at package level)
from cogniverse_synthetic import SyntheticDataService, SyntheticDataRequest, SyntheticDataResponse
from cogniverse_synthetic import OPTIMIZER_REGISTRY

Architecture

1. Optimization Architecture

flowchart TB
    WebClient["<span style='color:#000'>Web client / any client<br/>POST /admin/tenant/{id}/optimize</span>"]
    QM["<span style='color:#000'>QualityMonitor<br/>quality_monitor_cli --once<br/>submits raw Argo Workflow on quality drop</span>"]
    Cron["<span style='color:#000'>Helm CronWorkflows<br/>agent-optimization (weekly) / daily-gateway /<br/>daily-cleanup / synthetic-generation / monthly-reports</span>"]

    WebClient --> OptCLI
    QM --> OptCLI
    Cron --> OptCLI

    OptCLI["<span style='color:#000'>optimization_cli<br/>cogniverse_runtime<br/>15 modes: cleanup, triggered, simba, workflow,<br/>gateway-thresholds, online-routing-eval, llm-annotate, online-eval,<br/>profile, entity-extraction, synthetic, rollback,<br/>ab-compare, egress-netpol, monthly-reports</span>"]

    OptCLI --> GatewayOpt["<span style='color:#000'>Gateway Threshold Optimizer<br/>_compute_gateway_thresholds(spans_df)</span>"]
    OptCLI --> DSPyModes["<span style='color:#000'>DSPy compile modes<br/>profile / entity-extraction / simba / workflow / triggered<br/>all use BootstrapFewShot</span>"]
    OptCLI --> Meta["<span style='color:#000'>XGBoost meta-models<br/>TrainingDecisionModel / TrainingStrategyModel / FusionBenefitModel</span>"]
    OptCLI --> Strategy["<span style='color:#000'>StrategyLearner<br/>(triggered mode only)</span>"]

    GatewayOpt --> AM["<span style='color:#000'>ArtifactManager<br/>Phoenix DatasetStore blobs</span>"]
    DSPyModes --> AM
    Meta --> AM
    Strategy --> Mem0["<span style='color:#000'>Mem0MemoryManager<br/>Vespa memory</span>"]

    style WebClient fill:#90caf9,stroke:#1565c0,color:#000
    style QM fill:#90caf9,stroke:#1565c0,color:#000
    style Cron fill:#90caf9,stroke:#1565c0,color:#000
    style OptCLI fill:#ffcc80,stroke:#ef6c00,color:#000
    style GatewayOpt fill:#ce93d8,stroke:#7b1fa2,color:#000
    style DSPyModes fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Meta fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Strategy fill:#ce93d8,stroke:#7b1fa2,color:#000
    style AM fill:#a5d6a7,stroke:#388e3c,color:#000
    style Mem0 fill:#a5d6a7,stroke:#388e3c,color:#000

2. Gateway Threshold Optimization Architecture

flowchart TB
    Phoenix["<span style='color:#000'>Phoenix Spans<br/>cogniverse.gateway routing telemetry</span>"]

    Phoenix --> Compute["<span style='color:#000'>_compute_gateway_thresholds(spans_df)<br/>• Read output.value.complexity / .confidence (via read_span_io)<br/>• Rule-based ± adjustment from error rate + mean confidence<br/>• p25-percentile-derived gliner_threshold<br/>• GATEWAY_DEFAULT_THRESHOLD = 0.4</span>"]

    Compute --> Thresholds["<span style='color:#000'>Optimized Thresholds<br/>fast_path_confidence_threshold, gliner_threshold<br/>saved via ArtifactManager.save_blob('config', 'gateway_thresholds')</span>"]

    Thresholds --> Gateway["<span style='color:#000'>GatewayAgent<br/>am.load_blob('config', 'gateway_thresholds') at first build,<br/>refreshed every GATEWAY_ARTIFACT_TTL_S (5 min)</span>"]

    style Phoenix fill:#90caf9,stroke:#1565c0,color:#000
    style Compute fill:#ffcc80,stroke:#ef6c00,color:#000
    style Thresholds fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Gateway fill:#a5d6a7,stroke:#388e3c,color:#000

3. Profile Selection Optimization Architecture

flowchart TB
    Source["<span style='color:#000'>Tenant ground truth<br/>ArtifactManager blob kind config, key profile_selection_ground_truth<br/>• Uploaded through the admin profile_selection_ground_truth route<br/>• One query + expected_videos per row</span>"]

    Source --> Derive["<span style='color:#000'>derive_profile_labels<br/>• SearchService.search per query × usable profile, top_k 10<br/>• Result matches an expected video when its title basename (schema document_mapping.title, extension stripped) equals the expected id<br/>• Label = the one profile whose results match every expected video<br/>• Unrecovered, untitled-result and tied queries reported under label_exclusions</span>"]

    Derive --> RunOpt["<span style='color:#000'>run_profile_optimization tenant_id, lookback_hours<br/>• Build dspy.Example trainset from derived labels<br/>• Holdout = tail of the derived labels<br/>• Guard on single-label collapse</span>"]

    Synthetic["<span style='color:#000'>Approved synthetic demos<br/>• _load_approved_synthetic_data 'profile'<br/>• Merged into trainset</span>"] --> RunOpt

    RunOpt --> Compile["<span style='color:#000'>BootstrapFewShot teleprompter<br/>• Compile ProfileSelectionModule<br/>• Save via ArtifactManager.save_blob 'model','profile_selection'</span>"]

    Compile --> Reload["<span style='color:#000'>Next agent dispatch<br/>• ProfileSelectionAgent loads via am.load_blob 'model','profile_selection'<br/>• dspy_module.load_state applied to live module</span>"]

    style Source fill:#90caf9,stroke:#1565c0,color:#000
    style Derive fill:#90caf9,stroke:#1565c0,color:#000
    style Synthetic fill:#90caf9,stroke:#1565c0,color:#000
    style RunOpt fill:#ffcc80,stroke:#ef6c00,color:#000
    style Compile fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Reload fill:#a5d6a7,stroke:#388e3c,color:#000

Core Components

1. Gateway Threshold Optimizer (optimization_cli.py)

On-demand gateway confidence threshold tuning using Phoenix span data.

Key Functions:

GATEWAY_DEFAULT_THRESHOLD = 0.4

def _compute_gateway_thresholds(spans_df) -> dict:
    """
    Pure function: calibrate gateway thresholds from a spans DataFrame.

    Branches:
      1. simple_error_rate > 0.2         -> threshold = min(current + 0.1, 0.95)
      2. complex_error_rate < 0.05 and mean_confidence > 0.8
                                          -> threshold = max(current - 0.05, 0.5)
      3. otherwise                        -> threshold unchanged

    gliner_threshold = round(max(0.15, min(p25_confidence * 0.8, 0.5)), 3)

    Returns {"status": "no_data", ...} or
            {"status": "ready", "thresholds": {...}, "spans_found": N}
    """

On-Demand Submission:

# Web client or any client submits via runtime API:
# POST /admin/tenant/{tenant_id}/optimize
# Body: {"mode": "gateway-thresholds", "lookback_hours": 48}
# Synthetic training data for chosen optimizer types, queued for review:
# Body: {"mode": "synthetic", "optimizers": ["profile", "routing"]}
# Returns: {workflow_name, namespace, mode, status_url}

# Check run status:
# GET /admin/tenant/{tenant_id}/optimize/runs/{workflow_name}
# Returns: {phase, started_at, finished_at, message}

# List the tenant's runs (manual + CronWorkflow-spawned), newest first:
# GET /admin/tenant/{tenant_id}/optimize/runs?limit=20
# Returns: {runs: [{workflow_name, mode, trigger, phase,
#                   started_at, finished_at}, ...]}

File: libs/runtime/cogniverse_runtime/optimization_cli.py


2. ProfileSelectionAgent Optimization

DSPy module optimization for ProfileSelectionAgent — the per-query classifier that picks the backend profile and emits modality/complexity/intent in a cogniverse.profile_selection span.

CLI mode: --mode profile

Key function:

async def run_profile_optimization(
    tenant_id: str,
    lookback_hours: float = 24.0,
) -> dict:
    """
    Optimize ProfileSelectionAgent's DSPy module:

    1. Load the tenant's versioned ground-truth blob
       ("config", "profile_selection_ground_truth") before any scoring. If the
       blob is absent, return status "failed" with the missing-label reason
       and retryable=False. If the artifact store is unreachable or the payload
       is malformed, return status "failed" with the store-unavailable reason,
       retryable=True, and a chained cause. Both failures exit the CLI nonzero.
    2. Derive (query, available_profiles) -> selected_profile labels from those
       rows with derive_profile_labels: the tenant's SearchService runs every
       query against each profile from
       tenant_usable_profile_names(ConfigManager, tenant_id) at top_k=10. Each
       candidate's schema and type are read from the tenant's profile catalog
       (ConfigUtils.backend_profiles: shipped and system profiles with the
       tenant's stored profiles on top), the catalog search resolves it from,
       so a shipped profile the tenant serves but never stored is resolved. A
       result matches an expected video when the basename of its title,
       extension stripped, equals the expected id; the title field is the one
       the profile's schema names under document_mapping.title (video_title,
       audio_title, document_title, image_title, title, chunk_name). Candidate
       profiles are first restricted to the shipped backend types that can
       serve the row's expected media type, so video rows only score video
       profiles. The label is the single profile whose results match all of the
       row's expected_videos. Queries with no serving profile emit the named
       `no_profile_serves_media_type` exclusion. A failed retrieval is attempted
       at most three times, waiting `PROFILE_SELECTION_RETRIEVAL_BACKOFF_SECONDS`
       before the second attempt and twice that before the third; an exhausted
       comparison excludes the entire query with `incomplete_comparison` and
       per-profile failure context, and a run whose `incomplete_comparison`
       share exceeds `PROFILE_SELECTION_MAX_INCOMPLETE_SHARE` stops before
       scoring. No recovered video, untitled
       results, or profile ties are also excluded and reported under
       label_exclusions. cogniverse.profile_selection spans are counted as
       spans_found and not read.
    3. Merge in approved synthetic demos for optimizer type "profile". The
       consumer projection supplies the signature-required string confidence
       sentinel and preserves the two exact input fields.
    4. Compile ProfileSelectionModule via BootstrapFewShot (scaled: 8/16/2-round
       once >= 50 examples, else 4/8/1-round).
    5. Save compiled module as artifact ("model", "profile_selection").

    The agent reloads the artifact on each dispatch via am.load_blob("model", "profile_selection").

    Success and compile payloads also include `distinct_queries` and `holdout_queries`, so row counts and held-out query-key counts stay separate.

    Returns (every shape carries
    "label_exclusions": {"count": int, "queries": list[str],
    "incomplete_comparison_rate": float}):
      - {"status": "success", "spans_found": int, "served_examples": int,
         "approved_examples": int, "served_scoreable_examples": int,
         "training_examples": int, "holdout_examples": int,
         "holdout_source": "derived_labels", "selection": {...},
         "baseline_score": float, "current_score": float | None,
         "candidate_score": float | None,
         "decision": "promote" | "keep" | "rollback" | "reject",
         "version": int, "consumed_example_ids": list[str],
         "labels_by_profile": dict[str, int],
         "exclusions_by_reason": dict[str, int]}
      - {"status": "profile_labels_degenerate", "spans_found": int,
         "served_examples": int, "approved_examples": int,
         "served_scoreable_examples": int, "label_exclusions": {...},
         "labels_by_profile": dict[str, int],
         "exclusions_by_reason": dict[str, int]}
      - {"status": "no_data", "spans_found": int, "examples": 0}
      - {"status": "failed", "reason": "retrieval_exclusion_rate_exceeded", "retryable": True,
         "error": str, "spans_found": int, "max_incomplete_share": float}
      - {"status": "skipped", "reason": "profile_selection_ground_truth_missing", "retryable": False,
         "error": str}
      - {"status": "failed", "reason": "profile_selection_ground_truth_store_unavailable", "retryable": True,
         "error": str, "cause": {"type": str, "message": str}}
      - {"status": "failed", "reason": "profile_selection_ground_truth_invalid", "retryable": False,
         "error": str, "cause": {"type": str, "message": str}}
      - {"status": "no_eval_material", "spans_found": int,
         "served_scoreable_examples": int, "training_examples": int,
         "holdout_examples": 0, "holdout_source": "derived_labels",
         "selection": {...}}
      - {"status": "insufficient_population", "spans_found": int, "examples": int,
         "distinct_queries": int, "min_samples": int, "min_unique_queries": int,
         "version": int}
      - {"status": "failed", "error": str, "selection": {...}}
    """

Profile-selection ground truth is tenant-owned state in the artifact store, and the profile-selection optimizer returns profile_labels_degenerate only when the derived labels collapse to a single label value while recording labels_by_profile, dominant_label_share, and exclusions_by_reason.

Training: Always BootstrapFewShot; the 50-example threshold only changes its max_bootstrapped_demos/max_labeled_demos/max_rounds settings, it does not switch optimizers.

File: libs/runtime/cogniverse_runtime/optimization_cli.py


3. DSPyAgentPromptOptimizer

DSPy prompt optimizer for the 4 core agent-orchestration signatures — query analysis, agent routing, summary generation, detailed report — always compiled with BootstrapFewShot.

File: libs/agents/cogniverse_agents/optimizer/dspy_agent_optimizer.py


4. SIMBA (Query Enhancement) Optimization

Despite the mode name, this compiles QueryEnhancementAgent's DSPy module with BootstrapFewShot, not the SIMBA algorithm.

CLI mode: --mode simba

Key function: run_simba_optimization(tenant_id, lookback_hours=24.0) — reads cogniverse.query_enhancement spans into served-call records (input.value, input.source_text, input.grounding_context → enhanced_query, expansion_terms, synonyms, context_additions, confidence), appends approved synthetic demos for "query_enhancement", and splits the records with _split_served_holdout(..., holdout_eligible_predicate=_served_span_is_holdout_eligible) (deterministic ~25% tail holdout over served-scoreable spans; approved synthetic rows stay train-only, and unknown example_id families raise). Trainable records — a non-identity enhancement with at least one expansion term — compile the module via _create_teleprompter(len(trainset), metric=_query_enhancement_metric); records the metric cannot score (no source text and no grounding context) are ordered after every scoreable record so the bootstrap walk reaches them only when the scoreable records did not fill the demos, and their count is reported as unscoreable_examples. The compiled candidate is then cut to fit the student before it is scored: _student_demo_budget reads the optimization endpoint's served window (max_model_len, else its declared context_window) less its max_tokens, with served_message_counter counting requests by the tokenizer and chat template the endpoint serves, and _bound_candidate_demos keeps each predictor's longest prefix of demonstrations whose request fits that allowance with every distinct served input of the trainset and holdout, under both the chat adapter the run scores with and the LenientJSONAdapter the runtime serves with. The cut program is what is scored, saved and served; a window that cannot be read fails the run with nothing persisted, as does an endpoint that cannot count a request (TokenCountUnavailableError). Every holdout record is an evaluation probe. _query_enhancement_quality scores a module's own output for the probe inputs (1.0 when the enhanced query differs from the query, has expansion terms and, given a grounding context, names one of its entities; else 0.0). _query_enhancement_scores runs a module once per distinct (query, source_text, grounding_context) and counts that score once per holdout record carrying it, so a call served hundreds of times costs one LM request and keeps its weight in the mean. The base module, the persisted artifact and the compiled candidate are scored on the same holdout and _select_simba_artifact decides: promote persists the candidate (it beats the served module by the tenant's optimization_improvement_threshold), keep leaves the artifact, rollback persists the base state over an artifact that scores below base, reject persists nothing. A persisted artifact whose state is the base module's (_current_score) takes the baseline score instead of being scored again: a second scoring of the same module measures only the LM's noise, and a lower draw would roll back to that same state and discard the candidate. The SIMBA, profile and entity-extraction runs all score the served artifact this way. Without a holdout the run returns no_eval_material and persists nothing. The artifact key is ("model", "simba_query_enhancement"); QueryEnhancementAgent reloads it via am.load_blob("model", "simba_query_enhancement").

File: libs/runtime/cogniverse_runtime/optimization_cli.py


5. Entity Extraction Optimization

CLI mode: --mode entity-extraction

run_entity_extraction_optimization(tenant_id, lookback_hours=24.0) loads the tenant-owned ("config", "entity_extraction_ground_truth") blob, loads approved synthetic rows for "entity_extraction", and queries cogniverse.entity_extraction spans only to report spans_found, served_examples, and served_scoreable_examples. The truth rows define the holdout: _split_holdout(truth_records, scoreable_predicate=_entity_extraction_is_scoreable, holdout_eligible_predicate=_ground_truth_is_holdout_eligible) splits the tenant's ground-truth rows by distinct casefold query, approved synthetic rows stay in train only, unknown example_id families raise, and distinct_queries / holdout_queries report the truth rows only. Each truth or approved row becomes dspy.Example(query, entities=[EntityMention(text, type), ...]) via _entity_extraction_example, which raises ValueError naming the record when an entity's type is outside EntityType. The metric (entity_extraction.pair_set_f1.v1) is F1 over the (casefold stripped text, type) pairs of the typed mentions (_entity_extraction_pair_set); an empty recorded label set raises. Before the bootstrap walk, _sample_entity_self_consistency draws the teacher SELF_CONSISTENCY_SAMPLES (3) times per training record through create_sampling_dspy_lm(..., temperature=SELF_CONSISTENCY_TEMPERATURE, max_tokens=SELF_CONSISTENCY_MAX_TOKENS) — an LM with the cache off and no seed, so each draw is an independent request. A sampled draw occasionally falls into a repetition loop that runs to the token limit and never closes its JSON: on the served teacher (Qwen3-14B-AWQ) about 1 in 100 draws of a query it reasons about word by word, while the longest answer that parsed over 3,160 measured draws was 470 tokens. SELF_CONSISTENCY_MAX_TOKENS (1024) ends such a loop at half the teacher's completion budget. A draw whose answer does not parse (AdapterParseError) is asked once more from a new seed (SELF_CONSISTENCY_DRAW_ATTEMPTS, 2); a draw that fails both raises DrawNotParsed. vLLM's repetition_penalty is not sent: at 1.05 it removed the loops but lowered the draws' pair-F1 against the ground truth on the looping query from 0.756 to 0.606. Per (casefold text, type) mention, agreement is the fraction of draws carrying it; a mention below 1.0 carries needs_review. review_row builds the approval-queue payload: data is the query with the unanimous mentions only, metadata.self_consistency carries samples and one {text, type, agreement, needs_review} record per mention. Rows with at least one flagged mention are queued as a PENDING_REVIEW ApprovalBatch through ApprovalStorageImpl, with confidence the mean agreement; the approval queue tab renders one agreement line per mention. A record with no unanimous mention is queued the same way, carrying an empty training example and every mention flagged, and its query is listed under NO_UNANIMOUS_KEY (no_unanimous_examples). A record whose draws did not all complete — a draw that never parsed, or a failed request — is a non-vote: it contributes no row, its cause is recorded through BootstrapErrorLog.record_cause (so it counts in the bootstrap's errors and error_causes), and it is listed under NON_VOTES_KEY (non_votes) as {query, cause}. Every requested record is therefore either sampled or a non-vote. The run reports the pass under self_consistency: samples, temperature, max_tokens, examples_requested, examples_sampled, retries (draws asked again), non_votes, rows_needing_review, batch_id, rows_queued and no_unanimous_examples. Approved rows re-enter training the same way every approved synthetic row does — through approved_synthetic_data-{tenant}; ground truth is never rewritten. validate_approved_training_values refuses an entity whose type is outside ENTITY_TYPES (cogniverse_foundation.common.entity_types), so a human correction naming an unknown type is refused at approval time rather than at the next optimization.

Bootstrap uses BootstrapMetricRecorder with _entity_bootstrap_threshold(...), which keeps the bar at ENTITY_BOOTSTRAP_METRIC_THRESHOLD (1.0) and never below the served module's holdout score. The recorder appends each attempt as a JSONL row under ~/.cache/cogniverse/bootstrap_attempts.jsonl. The entity floor is min_samples_for_optimization: 30 and min_unique_queries: 15. Ground-truth loading renders the shared contract in ground_truth_blob.py: nothing uploaded is {"status": "skipped", "reason": "entity_extraction_ground_truth_missing"}, a store outage is {"status": "failed", ... "retryable": True}, and an unusable payload is {"status": "failed", ... "retryable": False}. The artifact key is ("model", "entity_extraction"); EntityExtractionAgent reloads it via am.load_blob("model", "entity_extraction").

Returns: - {"status": "success", "spans_found": int, "served_examples": int, "served_scoreable_examples": int, "distinct_queries": int, "holdout_queries": int, "label_rows": int, "truth_rows": int, "approved_rows": int, "training_examples": int, "holdout_examples": int, "holdout_source": "ground_truth", "bootstrap": {"trainset": int, "max_bootstrapped_demos": int, "max_labeled_demos": int, "max_rounds": int, "metric_threshold": float, "attempts": int, "errors": int, "error_causes": list[str], "examples_walked": int, "accepted": int, "bootstrapped_demos": int, "labeled_demos": int, "metric_values": list[float]} | None, "baseline_score": float, "current_score": float | None, "candidate_score": float | None, "decision": "promote" | "keep" | "rollback" | "reject", "version": int, "consumed_example_ids": list[str]} - {"status": "skipped", "reason": "entity_extraction_ground_truth_missing", "retryable": False, "error": str} - {"status": "no_data", "spans_found": int, "served_examples": int, "served_scoreable_examples": int, "label_rows": int, "truth_rows": int, "approved_rows": int, "examples": 0} - {"status": "insufficient_population", "spans_found": int, "examples": int, "distinct_queries": int, "min_samples": int, "min_unique_queries": int, "version": int} - {"status": "failed", "error": str}

File: libs/runtime/cogniverse_runtime/optimization_cli.py


6. Workflow Orchestration Optimization

CLI mode: --mode workflow

run_workflow_optimization(tenant_id, lookback_hours=24.0) reads cogniverse.orchestration spans, feeds them through OrchestrationEvaluator.evaluate_orchestration_spans (backed by WorkflowIntelligence) to extract WorkflowExecution records, then:

  1. Reads the live agents block from configs/config.json (not AgentRegistry, which starts empty in the optimization CLI's own pod) to build a set of currently-enabled agent names.
  2. Drops any execution whose agent_sequence references an agent not in that live set (stale demos from renamed/removed agents can't be replayed).
  3. Replaces the surviving executions, agent performance profiles, query-type patterns, and derived templates via WorkflowStoreRegistry.get(name="telemetry") — the same store WorkflowIntelligence reads at orchestrator startup. The store acquires a renewable per-tenant Redis lease before reading Phoenix, then owns the complete replacement and its compensation. A second optimizer replica waits outside the storage boundary. If any Phoenix write fails, the lease holder restores the exact prior values on all four channels before releasing the lease and propagating the failure. Serving loads acquire that same lease and return one typed four-channel snapshot. WorkflowIntelligence validates and stages the complete snapshot before replacing its in-memory executions, profiles, learned patterns, and templates in one publication step, so a reload removes obsolete entries, never duplicates history, and retains the prior complete snapshot when a Phoenix read fails. Redis is required; there is no process-local fallback.

Raises RuntimeError if the live-agents set is empty (refuses to guess whether every demo is stale or none are). Returns {"status": "success"|"no_data", "spans_found", "workflows_extracted", "execution_demos_saved", "agent_profiles_saved"}.

File: libs/runtime/cogniverse_runtime/optimization_cli.py


7. Online Routing Evaluation

CLI mode: --mode online-routing-eval

run_online_routing_evaluation(tenant_id, lookback_hours=24.0) reads AutomationRulesConfig.online_evaluation (OnlineEvaluationConfig — enabled, sampling_rate, evaluators, persist_scores, score_annotation_name) from config; if disabled, returns {"status": "disabled"} immediately. Otherwise reads cogniverse.routing spans and scores each with OnlineEvaluator (routing outcome + confidence calibration), persisting scores as telemetry annotations for drift detection. Returns {"status": "success"|"no_data", "spans_found", "scores_persisted", "statistics"}.

File: libs/runtime/cogniverse_runtime/optimization_cli.py


8. Triggered Optimization (quality-monitor driven)

CLI mode: --mode triggered

Invoked by QualityMonitor (not an on-demand submission — triggered is excluded from _MANUAL_OPTIMIZE_MODES) when golden/live evaluation detects degradation. QualityMonitor builds and submits its own raw Argo Workflow manifest (not the shared cogniverse-optimization-runner WorkflowTemplate) running:

uv run python -m cogniverse_runtime.optimization_cli \
  --mode triggered --tenant-id <tid> \
  --agents <comma-separated agent names> \
  --trigger-dataset <phoenix dataset name>

run_triggered_optimization loads the named Phoenix dataset (flattening input/output dict columns), splits each agent's rows into low_scoring/high_scoring by category, and for each agent in --agents:

  1. Builds a labeled dspy.Example set from the high-scoring rows (agent-specific field mapping for search / summary / report) and splits it deterministically into a trainset and a ~25% tail holdout (_split_train_holdout). An agent whose rows are all failures (empty positives trainset) falls back to a reflective recompile (see below) when enable_reflective_recompile is on and, after the same ~25% tail split is applied to the low-scoring rows, the resulting GEPA trainset has at least min_reflective_failures rows; otherwise it returns {"status": "skipped", "reason": "no_positive_examples", "negative_examples": N} — or "insufficient_failures_to_reflect" when reflection is on but that post-split trainset is smaller than min_reflective_failures (the threshold checks the split trainset, not the raw failing-row count).
  2. Compiles the module the runtime actually serves (_served_module: SearchOptimizationModule, SummarizationModule, ReportGenerationModule) with BootstrapFewShot on the trainset, scoped inside dspy.context(lm=optimizer.lm) (initialize_language_model only sets optimizer.lm; the compile reads the LM from the task-local binding). Training examples carry the served signature's own input and output fields (_served_example / _served_inputs).
  3. Scores the compiled candidate against the currently-active baseline (_holdout_scores via _probe_score): held-out positives contribute token-F1 of the served output field (_EVAL_FIELD: enhanced_query / summary / executive_summary) against the labeled output, and the low-scoring rows become known-bad probes (_negative_probes) that reward NOT reproducing the recorded failing output. The baseline is the same served module with the active compiled state loaded into it (_active_compiled_payload). An agent whose active prompts are not a compiled module state has no reconstructable baseline, so the run fails with reason: "baseline_not_reconstructable" rather than scoring against a stock module.
  4. Publishes the compiled module's whole dump_state() via _serve_compiled_prompts only if the candidate wins by at least the tenant's optimization_improvement_threshold — the call routes through ArtifactManager.promote_if_better(serve_versioned=True) (versioned save → canary → active). The published prompts dict has the single reserved key __dspy_module__ (cogniverse_core.agents.base.COMPILED_MODULE_PROMPT_KEY) holding {"dspy_version", "module", "state"} as JSON, so the instructions AND the learned demonstrations that produced the winning score serve together: the per-request overlay loads it into its per-call copy of the served module (load_compiled_module_state), which refuses a state produced for a different module class. The payload is loaded into a fresh served module before promotion, so a state the runtime cannot load fails the run instead of every request. A loser is recorded in the experiments ledger with promoted=False, --mode rollback restores a prior version, and the result reports the outcome under "served" (served_agent, version, active, promoted, plus baseline_score/candidate_score when eval material was available or a reason when it wasn't).

Reflective recompile (all-failure agents). When an agent's rows are all failures there are no positives to bootstrap from, so BootstrapFewShot cannot run. With enable_reflective_recompile (on by default for the three servable agents search / summary / report), _reflect_or_skip first splits the failing rows into a GEPA trainset and a held-out negatives slice (_split_train_holdout over the failing rows — GEPA never sees the held-out slice, ~75%/25%); only when that trainset has at least min_reflective_failures rows does _optimize_agent recompile with dspy.GEPA (_reflect_or_skip → _reflective_compile) — the threshold is checked against the post-split trainset size, not the raw failing-row count, so it takes more raw failures than min_reflective_failures to clear it. Each GEPA training example carries the recorded failing output as a _bad_output attribute, and a GEPA feedback metric (gold, pred, trace=None, pred_name=None, pred_trace=None) -> ScoreWithFeedback (_reflective_metric) rewards a candidate for not reproducing that failing output (1 - token_f1(pred, _bad_output)) while its feedback string names the failing output and what a good one must avoid. Only gold and pred are required, which is the arity DSPy's own Evaluate and bootstrap tracing call a metric with. GEPA's reflection LM (the resolved optimization LM) reads the failing rollouts plus that feedback and proposes improved instructions, capped at reflective_max_metric_calls metric calls (_build_gepa is the injectable GEPA seam). The GEPA candidate STILL flows through the same _score_and_serve → promote_if_better(serve_versioned=True) gate, scored on the held-out negatives — a reflective candidate that does not beat baseline + optimization_improvement_threshold is rejected and the base prompt is left byte-unchanged, so an all-failure agent that cannot be improved keeps its base prompt rather than a worse one. A successful run returns {"status": "success", "reflective": True, ..., "served": {...}}; too few failing rows to reflect returns {"status": "skipped", "reason": "insufficient_failures_to_reflect", "negative_examples": N}.

After the per-agent loop, it also runs strategy distillation: builds (or reuses) a Mem0MemoryManager for the tenant (requires SystemConfig.backend_url/backend_port, an api_base on the resolved LLM endpoint, and a configured denseon inference-service URL — raises ValueError if any are missing) and calls StrategyLearner(memory_manager, tenant_id, llm_config).learn_from_trigger_dataset(trigger_df). Distillation failure is caught and logged as non-fatal (results["strategies_distilled"] = 0).

File: libs/runtime/cogniverse_runtime/optimization_cli.py


9. XGBoost Training-Decision Meta-Models

Three independent XGBoost classifiers, each backed by ArtifactManager for persistence, all requiring a non-empty tenant_id.

File: libs/agents/cogniverse_agents/routing/xgboost_meta_models.py

Class Purpose Wired into
TrainingDecisionModel Binary train/skip gate + expected-improvement estimate from ModelingContext features (real/synthetic sample counts, success rate, confidence, days since last training, etc.) cogniverse_evaluation.quality_monitor.QualityMonitor._apply_training_decision_model — confirms or overrides the monitor's own train/skip verdict
TrainingStrategyModel Multi-class selection among TrainingStrategy.{PURE_REAL, HYBRID, SYNTHETIC, SKIP} Not called by any production pipeline today; exercised directly in tests/routing/unit/test_xgboost_meta_models.py
FusionBenefitModel Regression estimate of whether multi-signal fusion helps for a given context Not called by any production pipeline today; exercised directly in tests

Both models fall back to a hand-tuned heuristic when is_trained is False (before enough historical data exists to fit XGBoost). TrainingStrategyModel._fallback_strategy's actual thresholds:

# Fallback heuristic (untrained model only — select_strategy() uses the
# fitted XGBoost classifier once train() has run on >= 20 samples):
if context.real_sample_count < 10:
    return SYNTHETIC if context.synthetic_sample_count >= 50 else SKIP
if context.real_sample_count >= 100:
    return PURE_REAL
if context.real_sample_count >= 30:
    return HYBRID if context.synthetic_sample_count >= 50 else PURE_REAL
return SKIP  # 10 <= real_sample_count < 30

TrainingDecisionModel.train() requires ≥10 historical (ModelingContext, outcome) pairs; TrainingStrategyModel.train() requires ≥20. Both persist via save_to_telemetry() / load_from_telemetry() on the injected ArtifactManager.


10. StrategyLearner

Distills execution traces into reusable Strategy records, stored in Vespa memory via Mem0MemoryManager (type="strategy" metadata) for retrieval by MemoryAwareMixin.

File: libs/agents/cogniverse_agents/optimizer/strategy_learner.py

Two distillation paths: 1. Pattern extraction (_extract_patterns) — statistical analysis of which profiles/strategies/ parameters scored best for different query types; requires ≥MIN_TRACES_FOR_PATTERN (5) traces; no LLM call. 2. LLM distillation (_distill_with_llm) — contrastive analysis of high- vs low-scoring trace pairs via DSPy to surface workflow-level insights.

learn_from_trigger_dataset(trigger_df) is the entry point, called from run_triggered_optimization (mode triggered) after per-agent DSPy compilation. get_strategies_for_agent / rank_strategies_with_decay support retrieval-side confidence decay; format_strategies_for_context renders retrieved strategies as prompt context (used by MemoryAwareMixin). _store_strategy deduplicates near-identical strategies (DEDUP_SIMILARITY_THRESHOLD = 0.9) by bumping confirmation_count on the existing record instead of inserting a duplicate.


11. Signature Variants

Per-tenant named-variant registry for DSPy signatures. Each agent has at least a "default" variant; tenants opt into variants like "with_jurisdiction" via TenantConfig.metadata['signature_variants'][agent_type]. The artefact manager keys prompts / demos / experiments on (tenant_id, agent_type, variant_id) so each variant has its own compiled artefacts.

File: libs/agents/cogniverse_agents/optimizer/signature_variants.py

Key methods:

from cogniverse_agents.optimizer.signature_variants import (
    SignatureVariantRegistry,
    variant_qualified_agent_key,
    DEFAULT_VARIANT_ID,  # "default"
)

reg = SignatureVariantRegistry()

# Register a variant. Idempotent for identical (agent_type, variant_id, description);
# raises ValueError when replace=False and the variant exists with a different definition.
reg.register("legal_qa", "with_jurisdiction", description="adds jurisdiction input")

# Tenant lookup. Falls back to "default" when:
#   * TenantConfig is None / has no metadata dict
#   * signature_variants key is missing or not a dict
#   * the requested variant id is not registered (logged at WARNING — operators want typo signal)
variant_id = reg.selected_for_tenant(tenant_config, "legal_qa")

# Artefact dataset key. Default variant returns the bare agent_type so pre-variant
# artefacts keep working; non-default variants get a ``::variant=<safe_id>`` suffix.
key = variant_qualified_agent_key("legal_qa", variant_id)
# -> "legal_qa"                              when variant_id == "default"
# -> "legal_qa::variant=with_jurisdiction"   otherwise

The registry is intentionally schemaless — variants only track which ids are valid for an agent; the agent owns the actual DSPy signature class lookup table.


12. Canary FSM

ArtifactManager maintains a three-slot state — active, canary, retired — per (tenant_id, agent_type) and persists it via Phoenix DatasetStore under the config blob key. Routing decisions are stable per request (sha1 of request_seed, bucket % 100).

File: libs/agents/cogniverse_agents/optimizer/artifact_manager.py

from cogniverse_agents.optimizer.artifact_manager import ArtifactManager

am = ArtifactManager(telemetry_provider=phoenix_provider, tenant_id="acme:production")

# Promote a versioned artefact to canary at a traffic percentage.
# Raises ValueError if traffic_pct not in [1, 100]. Replaces any existing canary
# (the prior canary moves to retired with reason="superseded_by_new_canary").
state = await am.promote_to_canary("search_agent", version=7, traffic_pct=10)

# Promote canary to active. Prior active goes to retired with
# reason="superseded_by_canary_promotion". Restores prompts + demos at the
# un-versioned dataset name agents read at __init__. Raises ValueError if no canary set.
state = await am.promote_canary_to_active("search_agent")

# Drop the current canary back to retired without promoting. No-op when no canary set.
# Common reason values: "metric_regression", "admin_retire", "superseded_by_new_canary".
state = await am.retire_canary("search_agent", reason="metric_regression")

State shape:

{
  "active":  {"version": 6, "promoted_at": "..."},
  "canary":  {"version": 7, "promoted_at": "...", "traffic_pct": 10},
  "retired": [{"version": 5, "retired_at": "...", "reason": "..."}, ...],
}

Serving blobs

ArtifactManager.save_blob(kind, key, content) publishes a serving revision as a single immutable row. Revisions rotate through three datasets, dspy-{kind}-{tenant}-{key}--r0, --r1 and --r2, each row carrying its blob_revision. A publication only adds: it writes revision n into slot n % 3, which holds revision n - 3 and which no reader resolves, so the committed revision and its predecessor stay readable throughout and a failed or killed publication leaves both in place. The slot holding revision n - 2 is pruned afterwards, and only when a re-read shows it still carries a revision that far behind. load_blob reads all three slots and returns the greatest revision, or None when none exists. A store failure propagates and names the logical blob, not the slot whose read failed.

Phoenix has no compare-and-set, so two overlapping publications are last-write-wins within one slot. No publication deletes the slot it writes, and none deletes a slot a re-read showed carrying a newer revision, so readers resolve a complete revision throughout.

Serving revisions and optimizer candidates are separate namespaces. save_blob_versioned records candidates and their evaluation ledger under -v{n}; only activate_version copies a candidate into the serving blob, so a rejected candidate never reaches serving.

13. Regression-Reject Promotion Gate (promote_if_better)

ArtifactManager.promote_if_better is the guarded alternative to unconditionally overwriting active prompts/demos: it compares a candidate against the currently-active baseline on a held-out score and only promotes when the candidate wins.

File: libs/agents/cogniverse_agents/optimizer/artifact_manager.py

metrics = await am.promote_if_better(
    agent_type="search_agent",
    candidate_prompts=new_prompts,
    candidate_demos=new_demos,          # or None
    baseline_score=0.72,
    candidate_score=0.75,
    tolerance=0.0,                      # allowed regression band (0 = strict win)
    min_improvement=0.05,               # required margin over baseline
    serve_versioned=True,               # winner goes versioned → canary → active
    optimizer="BootstrapFewShot",
    train_examples=64,
)
  • Promoted when candidate_score >= baseline_score + min_improvement - tolerance: prompts/demos are saved and become active (current active is snapshotted first via snapshot_active when snapshot_before_promote=True, the default). With serve_versioned=True the win lands through the canary state machine (save_prompts_versioned → promote_to_canary → promote_canary_to_active) so both read seams — init-time un-versioned datasets and the request-time overlay keyed off the state blob — serve the new version; the version number is recorded in extra_metrics["served_version"]. With the default serve_versioned=False the win overwrites the un-versioned active artefacts only.
  • Rejected otherwise: prompts/demos are not saved; the run is recorded with promoted=False and a rejection_reason. --mode triggered wires the tenant's optimization_improvement_threshold config knob into min_improvement.

Either outcome lands as a typed ExperimentMetrics row via save_experiment in the per-agent experiments dataset, so the promotion ledger is observable end-to-end — rejected runs stay visible with their scores instead of being silently discarded.


14. Rollback CLI

Restore active artefacts to a previously snapshotted version. Wraps ArtifactManager.rollback_to_version and snapshots the current active first so the rollback is itself reversible.

A rollback carrying --prompts-version also moves the artefact state machine that load_for_request reads: the superseded active and any canary are retired with reason rollback, active becomes the rollback target and canary is cleared, all under the same compensation scope as the content writes. Every request seed therefore serves the rollback target on the next dispatch. A demos-only rollback leaves the state alone — the served identity is the prompts version.

File: libs/runtime/cogniverse_runtime/optimization_cli.py::run_rollback

# Roll back search_agent's active prompts to v2
uv run python -m cogniverse_runtime.optimization_cli \
  --mode rollback \
  --tenant-id acme:production \
  --agent search_agent \
  --prompts-version 2

# Demos rollback is independent; supply either or both
uv run python -m cogniverse_runtime.optimization_cli \
  --mode rollback \
  --tenant-id acme:production \
  --agent search_agent \
  --prompts-version 2 \
  --demos-version 3

Required: --tenant-id, --agent, plus at least one of --prompts-version / --demos-version. The Phoenix provider is built from TELEMETRY_HTTP_ENDPOINT / TELEMETRY_OTLP_ENDPOINT, the variables the chart sets on the pod; either one unset fails the run.

Returns {summary: ..., backup_versions: {prompts: int?, demos: int?}} — pass those versions to a follow-up --mode rollback to undo.


15. A/B Comparison (--mode ab-compare)

Runs RLMABRunner (cogniverse_agents.inference.ab_harness) over a Phoenix dataset of (query, context) rows, comparing a with-RLM and without-RLM arm per row, and emits a rlm.ab_compare Phoenix span per row with the harness's to_telemetry_dict() as attributes for the web client's RLM A/B view to aggregate.

uv run python -m cogniverse_runtime.optimization_cli \
  --mode ab-compare \
  --tenant-id acme:production \
  --queries-dataset golden_eval_v1 \
  --judge-substring "Paris"

--judge-substring (optional) enables a deterministic substring-match judge (1.0 if the substring appears in the answer, else 0.0) — described in code as "the minimum viable judge for getting a judge_delta populated in CI"; a real eval-time judge should be wired by the caller for production use. --rlm-max-iterations (default 10) and --rlm-max-llm-calls (default 30) cap the RLM arm's cost per row.

Returns {status, rows_compared, avg_latency_delta_ms, avg_tokens_delta, avg_judge_delta, rlm_fallback_rate, ab_ids}.

File: libs/runtime/cogniverse_runtime/optimization_cli.py::run_ab_compare


16. Egress NetworkPolicy Generation (--mode egress-netpol)

Not a training/compilation mode — a code-generation utility that reads agent egress policy YAMLs from configs/agent_policies/ and emits Kubernetes NetworkPolicy manifests so the cluster's CNI plugin (Cilium/Calico/etc.) enforces per-agent egress at the kernel level, independent of in-process HTTP enforcement.

uv run python -m cogniverse_runtime.optimization_cli \
  --mode egress-netpol \
  --policy-dir configs/agent_policies/ \
  --output-dir charts/cogniverse/templates/networkpolicies/ \
  --service-map vespa=cogniverse/vespa-service:8080 \
  --service-map llm=cogniverse/llm-service:11434

Two emit modes: - Per-agent (default): one NetworkPolicy per agent, selecting on app=<pod_app_label>, cogniverse-agent=<agent> — for topologies where each agent runs in its own Deployment. - Unified-runtime (--unified-pod-selector app.kubernetes.io/component=runtime): a single runtime-egress-netpol.yaml whose egress rules are the de-duplicated union of every agent's allowed destinations — needed for this project's default shared-runtime-pod topology, where per-agent L4 enforcement is impossible.

--helm-conditional wraps each emitted YAML in {{- if <expr> }} ... {{- end }} so a chart values flag can toggle it.

File: libs/runtime/cogniverse_runtime/optimization_cli.py::run_egress_netpol


17. --mode cleanup (memory + logs + temp + config vacuum)

Daily-cleanup workflow body (per-tenant when --tenant-id is set, global sweep when omitted). The CLI resolves cleanup environment values once and passes explicit data to run_cleanup. Both roots are validated before any work; unset, empty, missing, broad system/home directories, and roots containing the interpreter or checkout raise CleanupRootError naming the parameter. The chart sets LOG_DIR=/logs and TEMP_DIR=/tmp/cogniverse-cleanup; its scratch emptyDir is also TMPDIR.

File: libs/runtime/cogniverse_runtime/optimization_cli.py::run_cleanup

Section Source Knob
memory_cleanup Mem0MemoryManager.cleanup_with_schema(build_default_registry(), PinService(mm, registry).pinned_target_ids(tenant)) per tenant whose agent_memories schema is deployed (backend.schema_exists); the manager is initialised with auto_create_schema=False, so the sweep never deploys a schema. The pin read raises on a store outage, so the tenant is reported failed and nothing is deleted; a tenant whose memory partition schema is not deployed raises MemoryPartitionMissingError rather than sweeping against reads that answer "no memories" per-kind TTLs in KnowledgeRegistry
log_cleanup _prune_aged_files(LOG_DIR, older_than_days=log_retention_days) LOG_DIR required existing directory; --log-retention-days overrides LOG_RETENTION_DAYS (7)
temp_cleanup _prune_aged_files(TEMP_DIR, older_than_days=TEMP_RETENTION_DAYS) TEMP_DIR required existing directory; TEMP_RETENTION_DAYS (1)
config_vacuum VespaConfigStore.prune_all_configs(keep=CONFIG_KEEP_VERSIONS) CONFIG_KEEP_VERSIONS env (default 10)

Both roots must exist before cleanup starts. _vacuum_config_metadata reports a skip when the backing store is not VespaConfigStore. File reports include the resolved root, scanned, deleted, and errors; config vacuum reports dropped.

# Global sweep (every org / every tenant)
LOG_DIR=/logs TEMP_DIR=/tmp/cogniverse-cleanup \
  uv run python -m cogniverse_runtime.optimization_cli --mode cleanup --log-retention-days 7

# Per-tenant sweep
LOG_DIR=/logs TEMP_DIR=/tmp/cogniverse-cleanup \
  uv run python -m cogniverse_runtime.optimization_cli \
  --mode cleanup --tenant-id acme:production --log-retention-days 7

Result dict shape: {log_retention_days, memory_retention_days, memory_cleanup: {tid: entry}, memory_cleanup_summary: {completed, skipped, failed}, tenants_processed, log_cleanup, temp_cleanup, config_vacuum}. Each per-tenant entry is one of {"status": "completed", "deleted_by_kind": {kind: count}}, {"status": "skipped", "reason": "no memory schema deployed"}, or {"status": "failed", "error": "<ExcType>: <message>"}. A failed entry makes the CLI exit 1 (_run_failed), so a per-tenant outage fails the workflow instead of reading as a completed run.


18. --mode synthetic (Synthetic Data Generation)

run_synthetic_generation(tenant_id, optimizer_types=None, options=None) generates training data for one or more optimizer types (default: every type in APPROVED_TRAINING_AGENT_BY_OPTIMIZER) via SyntheticDataService. options is a SyntheticRunOptions (cogniverse_runtime.optimization_options, the CLI's --options JSON): count (default 50), vespa_sample_size, strategy, max_profiles and human_review (default true). submit_synthetic_outcome turns each type's non-empty output into an ApprovalBatch and hands it to HumanApprovalAgent.submit_for_review over ApprovalStorageImpl: items at or above ApprovalConfig.confidence_threshold are approved into the tenant's approved synthetic dataset at once, the rest stay pending review; without human review every item is approved. Each type's result carries batch_id, examples_generated, auto_approved, pending_review, avg_confidence, schema_name, selected_profiles, profile_selection_reasoning and generation_time_ms, which GET /admin/tenant/{tenant_id}/optimize/runs/{name}/synthetic reads back from the run.

The result retains one entry per requested optimizer. Its aggregate status is failed when any entry is failed or error, including a partial run where another optimizer succeeded. Otherwise it is success when at least one entry succeeded, including success plus no_data, and is no_data when every entry is no_data. The process exits 1 for an aggregate or nested failure and exits 0 for success or no_data, so scheduled workflows cannot treat partial failure as success.

Generators that wrap DSPy modules run under llm_config.primary, bound for the duration of each generate call with dspy.context(lm=...) — task-local, so concurrent runs never share or overwrite one another's LM.

uv run python -m cogniverse_runtime.optimization_cli \
  --mode synthetic --tenant-id acme:production --agents profile,routing
# NOTE: --agents is reused as the optimizer-types list for this mode

SyntheticDataService is also reachable directly as a REST API (mounted at /synthetic, not /admin/tenant):

Method + path Purpose
POST /synthetic/generate Generate synthetic data for one optimizer type
POST /synthetic/batch/generate Generate for multiple optimizer types in one call
GET /synthetic/optimizers List registered optimizer types (OPTIMIZER_REGISTRY)
GET /synthetic/optimizers/{optimizer_name} Schema + config for one optimizer type
GET /synthetic/health Service health check

File: libs/runtime/cogniverse_runtime/optimization_cli.py::run_synthetic_generation, libs/synthetic/cogniverse_synthetic/api.py


18a. --mode routing / --mode unified (module optimization)

run_module_optimization(module, tenant_id, lookback_hours, options) runs a module's steps in one pod (optimization_options.MODULE_STEPS): routing runs gateway-thresholds, entity-extraction and profile; unified runs those, then workflow. options is a ModuleRunOptions: max_iterations caps each DSPy compile's bootstrap rounds (_create_teleprompter(max_rounds=...)) and the workflow optimizer's 50-span evaluation batches (run_workflow_optimization(max_batches=...), also what --mode workflow --options sets); use_synthetic_data=false keeps approved synthetic examples out of the entity-extraction and profile compiles; and dataset_name (only without synthetic data) names a telemetry dataset of the tenant (<name>-<tenant>, query/expected_videos rows) that is the profile step's ground truth in place of the uploaded one. The result is {status, module, failed_steps, results: {step: result}}; any failed step fails it.

File: libs/runtime/cogniverse_runtime/optimization_cli.py::run_module_optimization


19. --mode monthly-reports

Generates usage + performance JSON for the prior period (default 30 days).

File: libs/runtime/cogniverse_runtime/optimization_cli.py::run_monthly_reports

Writes two files into --reports-output-dir (default ./reports):

  • usage-YYYYMM.json — per-org tenant counts (organization_metadata + tenant_metadata), each tenant's schemas_deployed list and schema_count; top-level summary {org_count, tenant_count, schema_count}.
  • performance-YYYYMM.json — per-tenant Phoenix span count, latency mean / p50 / p95, and error_rate (status_code != OK) over the lookback window. Empty-data tenants record {span_count: 0, latency_*: null, error_rate: 0.0}; Phoenix query failures record {"error": "phoenix query failed: ..."} and continue. Latency chunk merges run in 32-file passes before the final exact percentile scan, so open descriptors stay bounded without changing p50/p95.
uv run python -m cogniverse_runtime.optimization_cli \
  --mode monthly-reports \
  --reports-output-dir /reports \
  --lookback-hours 720

--tenant-id is not required (the workflow sweeps every tenant the metadata schemas know about). CronWorkflow {fullname}-monthly-reports (chart, schedule 0 5 1 * *, 1st of month 5 AM UTC) runs this followed by an mc step (minio.mcImage) that uploads to the configured MinIO bucket under reports/ (argo.optimization.monthlyReports.uploadPrefix).

Returns {period, generated_at, output_dir, files_written: [usage_path, perf_path], summary: {org_count, tenant_count, perf_tenants_with_data}}.


Usage Examples

Example 1: On-Demand Gateway Threshold Optimization

# Submit gateway-threshold optimization via the runtime API (web client or CLI):
curl -X POST http://localhost:8000/admin/tenant/acme:production/optimize \
  -H "Content-Type: application/json" \
  -d '{"mode": "gateway-thresholds"}'
# Returns: {"workflow_name": "opt-gateway-acme-...", "namespace": "cogniverse",
#           "mode": "gateway-thresholds", "status_url": "/admin/tenant/acme:production/optimize/runs/opt-gateway-acme-..."}

# Check status:
curl http://localhost:8000/admin/tenant/acme:production/optimize/runs/opt-gateway-acme-...
# Returns: {"phase": "Succeeded", "started_at": "...", "finished_at": "...", "message": ""}

# Terminate a running Workflow:
curl -X POST http://localhost:8000/admin/tenant/acme:production/optimize/runs/opt-gateway-acme-.../cancel
# Returns the run as Argo holds it right after the terminate, e.g. {"phase": "Running", ...};
# once Argo finishes it, status and the run list report it as {"phase": "Cancelled", ...}

# Retry a failed Workflow (restarts only the failed nodes, reuses successful ones):
curl -X POST http://localhost:8000/admin/tenant/acme:production/optimize/runs/opt-gateway-acme-.../retry
# Returns: {"phase": "Running", ...}

A finished run that was shut down (its Workflow's spec.shutdown is set) reports phase Cancelled in the status and run-list endpoints: Argo itself ends a run cancelled while it waited on the tenant mutex Succeeded (its step Skipped) and one cancelled while running Failed. A cancelled run is not retried; start a new one.

The status, cancel, and retry endpoints scope each action to the path tenant. Before returning status or issuing the terminate/retry, the runtime reads the Workflow's tenant-id argument (falling back to the cogniverse.ai/tenant label) and compares its canonical form to the path tenant. A Workflow owned by another tenant — or carrying no tenant tag — returns 404 Workflow not found, so one tenant cannot read or disrupt another tenant's optimization run by name.

The runtime does not inline the container spec into each submitted Workflow. Instead it references a chart-installed WorkflowTemplate (cogniverse-optimization-runner) via spec.workflowTemplateRef: the scheduled CronWorkflows (agent-optimization weekly, daily-gateway daily) use the same template, so the image/env/resource/mutex spec lives in one place (charts/cogniverse/templates/optimization-workflow-template.yaml).

The WorkflowTemplate declares a per-tenant mutex so multiple submits for the same tenant serialise (prevents repeated on-demand submits from stacking pods); different tenants optimize independently.

On-demand runs take the modes in _MANUAL_OPTIMIZE_MODES (libs/runtime/cogniverse_runtime/routers/tenant.py): gateway-thresholds, simba, workflow, profile, entity-extraction, llm-annotate and synthetic (with the optimizer types to generate data for, passed as --agents). triggered and cleanup aren't meant for interactive use and are CLI/cron-only.

# The Argo Workflow runs optimization_cli internally:
# uv run python -m cogniverse_runtime.optimization_cli --mode gateway-thresholds --tenant-id acme:production

# To call the threshold computation directly in tests:
from cogniverse_runtime.optimization_cli import _compute_gateway_thresholds, GATEWAY_DEFAULT_THRESHOLD

thresholds = _compute_gateway_thresholds(spans_df)
# Returns: {"status": "ready", "spans_found": N, "thresholds": {"fast_path_confidence_threshold": 0.5, "gliner_threshold": 0.4, "analysis": {...}}}
print(f"Default threshold: {GATEWAY_DEFAULT_THRESHOLD}")  # 0.4

Example 2: Profile Selection Optimization

# Trigger via the runtime CLI (or Argo Workflow --mode profile):
# uv run python -m cogniverse_runtime.optimization_cli --mode profile --tenant-id acme:production

# Or trigger on-demand via the admin API:
import requests

response = requests.post(
    "http://localhost:8000/admin/tenant/acme:production/optimize",
    json={"mode": "profile"}
)
result = response.json()
print(f"Workflow: {result['workflow_name']}")
print(f"Status URL: {result['status_url']}")

Production Considerations

1. Performance Optimization

Optimization on Demand:

# optimization_cli runs modes in isolation; no long-running optimizer process
# Triggered via Argo Workflow or direct API call

Asynchronous Optimization:

# Argo Workflow runs CLI modes in background without blocking runtime
# POST /admin/tenant/{tenant_id}/optimize → submits Argo Workflow → returns immediately


2. Data Quality and Safety

Label Derivation:

# Profile optimization trains only on queries that exactly one usable profile
# recovers in full, matched on result title basenames (derive_profile_labels);
# unrecovered, untitled-result and tied queries are reported under
# label_exclusions and never trained on.

Synthetic Data Control (TrainingStrategyModel, routing/xgboost_meta_models.py):

# select_strategy() uses the trained XGBoost classifier once fit(); before
# that (or on any untrained instance) it falls back to a fixed heuristic:
strategy = training_strategy_model.select_strategy(context)

# Fallback heuristic thresholds (see xgboost_meta_models.py _fallback_strategy):
#   real_sample_count < 10   -> SYNTHETIC if synthetic_sample_count >= 50 else SKIP
#   real_sample_count >= 100 -> PURE_REAL
#   30 <= real < 100         -> HYBRID if synthetic_sample_count >= 50 else PURE_REAL
#   10 <= real < 30          -> SKIP


3. Multi-Tenant Isolation

Tenant-Specific Optimization:

# Each tenant has isolated optimization state via telemetry
# Optimization is submitted per-tenant via the runtime API:
# POST /admin/tenant/tenant_a/optimize  {"mode": "gateway-thresholds"}
# POST /admin/tenant/tenant_b/optimize  {"mode": "profile"}

# Optimization CLI uses tenant-scoped telemetry provider for artifact isolation
# Artifacts saved as ("model", "profile_selection") scoped to tenant_id


4. Monitoring and Observability

Workflow Status Monitoring:

# Check status of a submitted optimization workflow
curl http://localhost:8000/admin/tenant/acme:production/optimize/runs/<workflow_name>
# {"phase": "Running", "started_at": "2026-04-24T03:00:00Z", "finished_at": null, "message": ""}

# List all optimization workflows
argo list -n cogniverse --selector workflow-type=optimization

Performance Degradation Detection via QualityMonitor:

# QualityMonitor (cogniverse_evaluation) checks agent scores and submits its own
# Argo Workflow (--mode triggered) automatically via quality_monitor_cli.
# --llm-model is required (must match evaluators.llm_judge.model in config);
# --runtime-url defaults to http://localhost:28000. TELEMETRY_OTLP_ENDPOINT and
# TELEMETRY_HTTP_ENDPOINT (or --phoenix-url) must name Phoenix; the CLI exits 2
# when either is missing.
TELEMETRY_OTLP_ENDPOINT=localhost:4317 \
uv run python -m cogniverse_runtime.quality_monitor_cli \
  --tenant-id default \
  --runtime-url http://localhost:28000 \
  --phoenix-url http://localhost:26006 \
  --llm-model google/gemma-4-e4b-it

# --once forces a single distillation-only pass (bypasses the quality-drop
# threshold check) — this is what the daily scheduled-distillation
# CronWorkflow uses.


5. Production Deployment

Optimization is on-demand, not a long-running service.

# Trigger via the web client or direct API call:
curl -X POST http://localhost:8000/admin/tenant/acme:production/optimize \
  -d '{"mode": "gateway-thresholds"}'

# Or run directly via CLI (e.g., from Argo Workflow template):
uv run python -m cogniverse_runtime.optimization_cli \
  --mode gateway-thresholds --tenant-id acme:production

6. Production Deployment Infrastructure

The optimization module includes production-ready deployment infrastructure with CLI tools and Argo Workflows for batch and scheduled optimization.

CLI: cogniverse_runtime.optimization_cli

Command-line interface for per-agent optimization (14 modes total):

# Optimize query enhancement (mode named "simba"; runs BootstrapFewShot)
uv run python -m cogniverse_runtime.optimization_cli \
  --mode simba \
  --tenant-id default

# Optimize gateway thresholds
uv run python -m cogniverse_runtime.optimization_cli \
  --mode gateway-thresholds \
  --tenant-id default

# Optimize entity extraction
uv run python -m cogniverse_runtime.optimization_cli \
  --mode entity-extraction \
  --tenant-id default

# Optimize workflow orchestration
uv run python -m cogniverse_runtime.optimization_cli \
  --mode workflow \
  --tenant-id default

# Optimize search profile selection
uv run python -m cogniverse_runtime.optimization_cli \
  --mode profile \
  --tenant-id default

# Clean up old optimization logs
LOG_DIR=/logs TEMP_DIR=/tmp/cogniverse-cleanup \
  uv run python -m cogniverse_runtime.optimization_cli \
  --mode cleanup \
  --log-retention-days 7

Available Options (subset — see build_parser() for the full set):

  • --mode: cleanup | triggered | simba | workflow | gateway-thresholds | online-routing-eval | llm-annotate | online-eval | profile | entity-extraction | synthetic | rollback | ab-compare | egress-netpol | monthly-reports

  • --tenant-id: required for every mode except cleanup, egress-netpol, and monthly-reports

  • --lookback-hours: hours of span history to analyze (default 24.0, accepts fractions)

  • --log-retention-days / --memory-retention-days: cleanup mode; override LOG_RETENTION_DAYS / MEMORY_RETENTION_DAYS (7 / 30). Cleanup requires existing dedicated LOG_DIR and TEMP_DIR roots; TEMP_RETENTION_DAYS defaults to 1. Unsafe roots raise CleanupRootError before cleanup.

DSPy Optimizer Selection (actual behavior): The simba, profile, and entity-extraction modes use _create_teleprompter(trainset_size, teacher_settings=None, metric=..., metric_threshold=None), which always returns dspy.teleprompt.BootstrapFewShot — scaled by a single threshold, not a multi-tier optimizer selection:

  • < 50 examples → BootstrapFewShot(metric=metric, metric_threshold=metric_threshold, max_bootstrapped_demos=4, max_labeled_demos=8, max_rounds=1, max_errors=5)
  • ≥ 50 examples → BootstrapFewShot(metric=metric, metric_threshold=metric_threshold, max_bootstrapped_demos=8, max_labeled_demos=16, max_rounds=2, max_errors=10)

metric defaults to _approved_example_exact_metric; simba passes _query_enhancement_metric, profile _profile_selection_metric, entity-extraction a BootstrapMetricRecorder over _entity_extraction_quality with metric_threshold=_entity_bootstrap_threshold(...).

All three pass teacher_settings={"lm": teacher_lm_or_raise(llm_config)}, so the bootstrap teacher runs on the centralized llm_config.teacher endpoint. Both shipped configs give it request_timeout 210 s: the teacher's Modal app scales to zero, so the first call after an idle window waits out the engine start (up to 135 s measured on the Modal chat apps), and the Modal proxy's 300 s upstream timeout stays above it. That helper probes the endpoint and then builds the LM through create_budgeted_dspy_lm, so the teacher's prompt is bounded by the window it actually serves: max_model_len read from the endpoint, or the endpoint's declared context_window when the listing carries none, minus the endpoint's max_tokens reservation. Bootstrapping appends a demonstration per accepted trace, so requests over the allowance shed whole demonstrations, oldest first, and one that cannot fit without them raises PromptBudgetExceededError naming the window, the reservation, the input size and the number dropped. The student LM (llm_config.resolve("optimization")) is bound with with dspy.context(lm=...) around the teleprompter.compile(...) call that reads it — task-local, so a mode runs no matter which async task in the process already owns DSPy's ambient binding, and concurrent runs never share or overwrite one another's LM. triggered (the search/summary/report agents) does not go through _create_teleprompter — _optimize_agent builds BootstrapFewShot directly from DSPyAgentPromptOptimizer.optimization_settings (a fixed configuration — max_bootstrapped_demos=8, max_labeled_demos=16, max_rounds=3, max_errors=10 — not scaled by trainset size), with the teacher threaded through _optimize_agent(teacher_endpoint=llm_config.resolve_teacher()) → initialize_language_model(teacher_endpoint_config=...).

No mode in this CLI selects MIPROv2 or SIMBA. GEPA is selected only for --mode triggered's reflective recompile of an all-failure agent (not by data size), never for the positives-trained path.

Argo Workflows Integration

On-demand submission (web client or any client):

curl -X POST http://localhost:8000/admin/tenant/acme:production/optimize \
  -d '{"mode": "profile"}'

submits a Workflow via spec.workflowTemplateRef against the chart-installed cogniverse-optimization-runner WorkflowTemplate (charts/cogniverse/templates/optimization-workflow-template.yaml).

Scheduled CronWorkflows (charts/cogniverse/templates/optimization-workflows.yaml, enabled/scheduled via values.yaml's argo.optimization.*):

CronWorkflow name Schedule (default) What it runs
{fullname}-agent-optimization 0 3 * * 0 (Sunday 3 AM UTC) Step 1 (parallel): gateway-thresholds, entity-extraction, simba (168h lookback), profile-ground-truth (the profile-ground-truth-check mode, printing present or absent). Step 2: profile (48h lookback), run only when the check printed present, so a tenant with no uploaded ground truth has the step omitted and recorded as skipped. Step 3: workflow. Step 4: rolling-restart the runtime Deployment to pick up new artifacts.
{fullname}-daily-gateway 0 4 * * * (daily 4 AM UTC) gateway-thresholds only — warm runtime pods pick up the recalibration via the dispatcher's gateway reload interval, no restart
{fullname}-daily-cleanup 0 4 * * * (daily 4 AM UTC) cleanup
{fullname}-synthetic-generation 0 1 * * 6 (Saturday 1 AM UTC) synthetic
{fullname}-scheduled-distillation daily quality_monitor_cli --once (forces a distillation pass even when quality is stable, so learning doesn't stall during long healthy periods)
{fullname}-monthly-reports 0 5 1 * * (1st of month 5 AM UTC) monthly-reports, then uploads output to MinIO via mc (minio.mcImage)
# View a schedule
kubectl get cronworkflow cogniverse-agent-optimization -n cogniverse

# Check last run
argo list -n cogniverse --selector workflows.argoproj.io/cron-workflow=cogniverse-agent-optimization --limit 1

# Trigger manually
argo submit --from cronwf/cogniverse-agent-optimization -n cogniverse

# Suspend/resume
argo cron suspend cogniverse-daily-gateway -n cogniverse
argo cron resume cogniverse-daily-gateway -n cogniverse

The e2e stack guard annotates suspended CronWorkflows with cogniverse.io/e2e-suspended=<uuid>. If a session stops before teardown, run uv run python -m tests.e2e.cron_guard restore-stale to clear those stale annotations and re-enable the affected CronWorkflows.

Quality-triggered optimization: QualityMonitor submits its own ad-hoc Workflow (not the shared WorkflowTemplate) running --mode triggered whenever golden/live evaluation detects a quality drop for one or more agents — see Triggered Optimization above.

Web Client Integration

The optimization infrastructure integrates with the web client:

Optimization Runs View:

  • Submit on-demand runs in any mode the runtime accepts (GET /admin/tenant/optimize-modes)

  • Monitor workflow progress (phase, started/finished timestamps)

  • Cancel or retry a run

Execution Modes:

  1. Automatic (Scheduled): CronWorkflows run optimization_cli per the schedule table above

  2. Automatic (Quality-triggered): QualityMonitor submits --mode triggered on detected degradation

  3. Manual (on demand): the web client calls POST /admin/tenant/{id}/optimize → submits an Argo Workflow on demand

See Web Client for full UI documentation.


7. Error Handling and Recovery

Workflow Failure Handling:

# Check failed workflow details
argo get <workflow-name> -n cogniverse -o json | jq '.status.message'

# Retry a failed optimization (only the failed nodes re-run)
curl -X POST http://localhost:8000/admin/tenant/acme:production/optimize/runs/<workflow-name>/retry

Artifact Persistence:

Optimization artifacts are persisted to the telemetry store via ArtifactManager using Phoenix DatasetStore. The profile, entity-extraction, and simba compile modes (and DSPyAgentPromptOptimizer) save compiled modules that agents reload on each dispatch via am.load_blob(...); triggered mode instead publishes versioned prompts served by the per-request overlay (_serve_compiled_prompts).

Loaded DSPy module blobs must match the live module's signature contract exactly: each predictor's saved signature.instructions and ordered signature.fields prefix/description pairs must match before load_state() runs. When the contract differs, the agent keeps its in-memory defaults, records artifact_load_status = "signature_mismatch", and skips the reload.


Testing

Test Files

Area Location
CLI argument parser, batch-mode branches (gateway-thresholds/workflow/entity-extraction/simba no-data paths, synthetic-data merge, _create_teleprompter tiering) tests/runtime/unit/test_batch_optimization_modes.py
_compute_gateway_thresholds — tight per-field assertions across all 3 calibration branches tests/runtime/unit/test_batch_optimization_modes.py::TestComputeGatewayThresholdsAlgorithm
_monthly_report_percentiles bounded merge fan-in + exact percentile preservation tests/runtime/unit/test_optimization_cli_monthly_reports.py
run_rollback / ArtifactManager.rollback_to_version round-trip tests/runtime/integration/test_optimization_cli_rollback.py
run_cleanup (memory/log/temp/config vacuum) tests/runtime/integration/test_optimization_cli_cleanup.py
run_monthly_reports tests/runtime/integration/test_optimization_cli_monthly_reports.py
run_triggered_optimization + strategy distillation wiring tests/runtime/integration/test_optimization_cli_triggered.py
Reflective GEPA recompile — config defaults, 5-arg feedback metric, GEPA trainset/gate wiring for all-failure agents tests/runtime/unit/test_reflective_recompile.py
Reflective GEPA recompile against a real Phoenix + real LM (_optimize_agent end-to-end) tests/runtime/integration/test_optimization_cli_reflective.py
run_ab_compare / RLMABRunner tests/runtime/integration/test_optimization_cli_ab_compare.py
ArtifactManager — canary FSM, promote_if_better tests/agents/integration/test_artifact_manager_experiments.py, tests/agents/integration/test_artifact_manager_variants.py
SignatureVariantRegistry tests/agents/unit/test_signature_variants.py, tests/runtime/integration/test_signature_variant_admin_consumption.py, tests/e2e/test_signature_variants_e2e.py
StrategyLearner tests/agents/unit/test_strategy_learner.py, tests/memory/integration/test_strategy_learner_integration.py
XGBoost meta-models (TrainingDecisionModel, TrainingStrategyModel, FusionBenefitModel) tests/routing/unit/test_xgboost_meta_models.py, tests/evaluation/integration/test_xgboost_quality_monitor.py
DSPyAgentPromptOptimizer / DSPyAgentOptimizerPipeline tests/agents/unit/test_dspy_optimization_integration.py
Profile selection artifact save/load round-trip tests/agents/unit/test_profile_selection_agent.py, tests/e2e/test_optimizer_persistence_e2e.py
End-to-end batch optimization tests/e2e/test_batch_optimization_e2e.py
Synthetic data service/generators tests/synthetic/unit/test_profile_generator.py, tests/synthetic/integration/test_profile_synthetic_service.py

Test Scenarios

1. Gateway Threshold Computation:

def test_compute_gateway_thresholds():
    """Test threshold derivation from Phoenix spans"""
    from cogniverse_runtime.optimization_cli import _compute_gateway_thresholds, GATEWAY_DEFAULT_THRESHOLD
    import json
    import pandas as pd

    spans_df = pd.DataFrame({
        "attributes.output.value": [
            json.dumps({"complexity": "simple", "confidence": c})
            for c in [0.3, 0.5, 0.6, 0.7, 0.8, 0.9]
        ],
    })

    result = _compute_gateway_thresholds(spans_df)
    assert result["status"] == "ready"
    assert "fast_path_confidence_threshold" in result["thresholds"]
    assert GATEWAY_DEFAULT_THRESHOLD == 0.4


Coverage:

  • Unit tests: XGBoost meta-models, gateway threshold computation, CLI argument parsing, teleprompter tiering

  • Integration tests: rollback round-trip, cleanup vacuum, monthly reports, triggered optimization + distillation, A/B comparison, artifact manager canary FSM

  • End-to-end tests: batch optimization, optimizer artifact persistence, signature variants


DSPy Training Data Requirements

Overview

DSPy optimizers require properly formatted training examples with all expected output fields defined. Missing fields will cause AttributeError: 'Example' object has no attribute 'field_name' during metric evaluation.

Training Data Format

Each DSPy Example must include: 1. Input fields (marked with .with_inputs()) 2. All output fields that metrics will access

Example Structure:

import dspy

example = dspy.Example(
    # Input fields
    query="user query here",
    context="optional context",

    # Output fields (ALL fields that metrics check must be present)
    primary_intent="search",
    confidence=0.9,
    recommended_agent="video_search"
).with_inputs("query", "context")  # Specify which fields are inputs

Query Analysis Training Data

Required Output Fields:

  • primary_intent: Main intent category

  • complexity_level: "simple" | "complex"

  • needs_video_search: "true" | "false"

  • needs_text_search: "true" | "false"

  • multimodal_query: "true" | "false"

  • temporal_pattern: Temporal info or "none"

Example:

training_data = [
    dspy.Example(
        query="Show me videos of robots from yesterday",
        context="",
        # All output fields required for metrics
        primary_intent="video_search",
        complexity_level="simple",
        needs_video_search="true",
        needs_text_search="false",
        multimodal_query="false",
        temporal_pattern="yesterday",
    ).with_inputs("query", "context"),

    dspy.Example(
        query="Compare research papers on deep learning",
        context="academic",
        primary_intent="analysis",
        complexity_level="complex",
        needs_video_search="false",
        needs_text_search="true",
        multimodal_query="false",
        temporal_pattern="none",
    ).with_inputs("query", "context"),
]

Agent Routing Training Data

Required Output Fields:

  • recommended_workflow: Workflow type

  • primary_agent: Main agent to use

  • routing_confidence: Confidence score (0.0-1.0 as string)

Example:

training_data = [
    dspy.Example(
        query="Show me videos",
        analysis_result="simple search",
        available_agents=["video_search"],
        # All output fields
        recommended_workflow="raw_results",
        primary_agent="video_search",
        routing_confidence="0.9",
    ).with_inputs("query", "analysis_result", "available_agents"),

    dspy.Example(
        query="Analyze data trends",
        analysis_result="complex analysis",
        available_agents=["detailed_report"],
        recommended_workflow="detailed_report",
        primary_agent="detailed_report",
        routing_confidence="0.85",
    ).with_inputs("query", "analysis_result", "available_agents"),
]

Common Errors

Missing Output Fields

Error:

AttributeError: 'Example' object has no attribute 'primary_intent'

Cause: Metric function accesses example.primary_intent but Example doesn't have that field.

Fix: Add the missing field to all training examples:

example = dspy.Example(
    query="...",
    primary_intent="search",  # ← Add missing field
    # ... other fields
).with_inputs("query")

Incorrect Field Types

Error:

TypeError: expected str, got bool

Cause: DSPy Examples store all fields as strings internally.

Fix: Convert to strings:

# ❌ Bad
needs_video_search=True

# ✅ Good
needs_video_search="true"

Validation

Before running optimization, validate your training data:

def validate_training_data(examples, required_fields):
    """Validate all examples have required output fields."""
    for i, ex in enumerate(examples):
        for field in required_fields:
            if not hasattr(ex, field):
                raise ValueError(
                    f"Example {i} missing required field '{field}'"
                )
    print(f"✅ All {len(examples)} examples valid")

# Usage
required = ["primary_intent", "complexity_level", "needs_video_search"]
validate_training_data(training_data, required)

Best Practices

  1. Define all output fields upfront - Check what your metrics access
  2. Use consistent field names - Match your DSPy signature output fields
  3. Validate before optimization - Catch missing fields early
  4. Use string values - DSPy converts everything to strings
  5. Document required fields - Keep a reference list for your team

File References

  • libs/agents/cogniverse_agents/optimizer/dspy_agent_optimizer.py - Training data loading
  • tests/agents/unit/test_dspy_optimization_integration.py - Example tests with proper format

DSPy Agent Optimization Infrastructure

The libs/agents/cogniverse_agents/optimizer/ package provides infrastructure for running DSPy optimization with configurable model providers (Modal GPU, local Ollama, cloud APIs).

Location: libs/agents/cogniverse_agents/optimizer/

flowchart TB
    subgraph "Configuration"
        Config["<span style='color:#000'>SystemConfig</span>"]
    end

    subgraph "Orchestration"
        Orch["<span style='color:#000'>DSPyAgentOptimizerPipeline</span>"]
        Client["<span style='color:#000'>DSPyAgentPromptOptimizer</span>"]
    end

    subgraph "LLM Factory"
        Factory["<span style='color:#000'>create_dspy_lm()</span>"]
        Modal["<span style='color:#000'>Modal (OpenAI-compatible)</span>"]
        Local["<span style='color:#000'>Ollama / Local</span>"]
        API["<span style='color:#000'>Anthropic / OpenAI</span>"]
    end

    subgraph "DSPy Optimization"
        AgentOpt["<span style='color:#000'>DSPyAgentPromptOptimizer</span>"]
        Boot["<span style='color:#000'>BootstrapFewShot</span>"]
    end

    subgraph "Output"
        Artifacts["<span style='color:#000'>Prompt Artifacts</span>"]
    end

    Config --> Orch
    Orch --> Client
    Client --> Factory

    Factory --> Modal
    Factory --> Local
    Factory --> API

    Modal --> AgentOpt
    Local --> AgentOpt
    API --> AgentOpt

    AgentOpt --> Boot
    Boot --> Artifacts

    style Config fill:#a5d6a7,stroke:#388e3c,color:#000
    style Orch fill:#ffcc80,stroke:#ef6c00,color:#000
    style Client fill:#ffcc80,stroke:#ef6c00,color:#000
    style Factory fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Modal fill:#90caf9,stroke:#1565c0,color:#000
    style Local fill:#90caf9,stroke:#1565c0,color:#000
    style API fill:#90caf9,stroke:#1565c0,color:#000
    style AgentOpt fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Boot fill:#ffcc80,stroke:#ef6c00,color:#000
    style Artifacts fill:#a5d6a7,stroke:#388e3c,color:#000

LLM Factory

Location: libs/foundation/cogniverse_foundation/config/llm_factory.py

All DSPy LM creation goes through the centralized create_dspy_lm() factory. Provider is encoded in the model string (LiteLLM convention). The Cogniverse chart always emits openai/<bare-model> for every in-cluster backend (vLLM or Ollama); the actual destination is selected by api_base. External SaaS providers use their own litellm prefix (anthropic/, together/, etc.).

from cogniverse_foundation.config.unified_config import LLMEndpointConfig
from cogniverse_foundation.config.llm_factory import create_dspy_lm

# Local Ollama (modern Ollama exposes /v1/chat/completions; use openai/ prefix)
local_lm = create_dspy_lm(LLMEndpointConfig(
    model="openai/llama3:8b",
    api_base="http://localhost:11434/v1",
))

# Modal (OpenAI-compatible endpoint)
modal_lm = create_dspy_lm(LLMEndpointConfig(
    model="openai/HuggingFaceTB/SmolLM3-3B",
    api_base="https://username--general-inference-service-serve.modal.run",
))

# Cloud API (Anthropic)
teacher_lm = create_dspy_lm(LLMEndpointConfig(
    model="anthropic/claude-3-5-sonnet-20241022",
    api_key="sk-ant-...",
))

# Use with scoped context (never global dspy.settings.configure)
import dspy
with dspy.context(lm=local_lm):
    result = module(query="Route this query...")

LLM Configuration Structure

LLM configuration is centralized in the llm_config section of config.json:

{
  "llm_config": {
    "primary": {
      "model": "openai/google/gemma-4-e4b-it",
      "api_base": "http://localhost:11434/v1",
      "api_key": "placeholder-no-auth-needed"
    },
    "teacher": {
      "model": "anthropic/claude-3-5-sonnet-20241022",
      "api_key": "sk-ant-..."
    },
    "overrides": {
      "orchestrator_agent": {
        "model": "openai/qwen3:8b",
        "api_base": "http://localhost:11434/v1",
        "api_key": "placeholder-no-auth-needed"
      }
    }
  }
}

All agents and optimizers resolve their LLM config from this section via LLMConfig.from_dict() and create_dspy_lm().

DSPyAgentPromptOptimizer

Location: libs/agents/cogniverse_agents/optimizer/dspy_agent_optimizer.py

Optimizes prompts for multiple agent types using DSPy signatures.

from cogniverse_agents.optimizer.dspy_agent_optimizer import (
    DSPyAgentPromptOptimizer,
    DSPyAgentOptimizerPipeline,
)

# Initialize optimizer
optimizer = DSPyAgentPromptOptimizer(config={
    "optimization": {
        "max_bootstrapped_demos": 8,
        "max_labeled_demos": 16,
        "max_rounds": 3,
        "stop_at_score": 0.95
    }
})

# Initialize language model via centralized LLM config
from cogniverse_foundation.config.unified_config import LLMEndpointConfig

endpoint_config = LLMEndpointConfig(
    model="openai/llama3:8b",
    api_base="http://localhost:11434/v1",
)
optimizer.initialize_language_model(endpoint_config=endpoint_config)

# Create and run pipeline
pipeline = DSPyAgentOptimizerPipeline(optimizer)
pipeline.initialize_modules()
pipeline.load_training_data()
compiled = pipeline.optimize_module("query_analysis", pipeline.training_data["query_analysis"])

initialize_language_model also accepts an optional teacher_endpoint_config (typically llm_config.resolve_teacher()). When given, it builds a separate self.teacher_lm and sets optimization_settings["teacher_settings"] = {"lm": self.teacher_lm}, so the bootstrap teacher runs on that LM instead of the student teaching itself.

Optimizable Modules

Module DSPy Signature Purpose
query_analysis QueryAnalysisSignature Intent detection, complexity analysis
agent_routing AgentRoutingSignature Agent selection, workflow recommendation
summary_generation SummaryGenerationSignature Summary quality optimization
detailed_report DetailedReportSignature Report structure optimization

All four are compiled with dspy.teleprompt.BootstrapFewShot — DSPyAgentOptimizerPipeline.optimize_module builds the BootstrapFewShot config from optimizer.optimization_settings (including teacher_settings, populated when initialize_language_model was given a teacher_endpoint_config) and runs the compile scoped inside dspy.context(lm=optimizer.lm) when a real (non-mock) LM is configured. optimize_all_modules holds out 20% of each module's examples as a validation split and passes it as optimize_module's validation_examples; the compiled module is scored against the un-compiled baseline on that split (per-module metric, a raising prediction scores 0) and is discarded in favour of the baseline when it underperforms — the overfit gate for small training sets.

Live Routing Optimization

The A2A entry point is the GLiNER-based GatewayAgent, whose routing thresholds are tuned by the gateway-thresholds optimization mode in optimization_cli.py (persists the gateway_thresholds artifact via run_gateway_thresholds_optimization, loaded back with ArtifactManager.load_blob("config", "gateway_thresholds")). Agent prompts (query_analysis / summary / detailed_report) are optimized separately by DSPyAgentPromptOptimizer.

CLI Usage

# Optimize gateway routing thresholds
uv run python -m cogniverse_runtime.optimization_cli --mode gateway-thresholds --tenant-id default

Output Artifacts

After optimization, artifacts are persisted to the telemetry store via ArtifactManager using Phoenix DatasetStore, in datasets named dspy-{kind}-{tenant_id}-{key} (kind one word) by artifact_dataset_name in cogniverse_foundation.telemetry.providers.base; is_artifact_dataset(name, tenant_id) reads the name back, and the runtime's evaluation dataset listing leaves these out:

  • Every artifact dataset writes metadata.created_at as a timezone-aware UTC ISO-8601 timestamp, including stable and versioned prompts, demonstrations, and blobs.
  • dspy-prompts-{tenant_id}-{agent_type} — Optimized system prompts for an agent (last-write-wins: save_prompts uses DatasetStore.replace_dataset, so each save replaces the prior rather than appending a version)
  • dspy-demos-{tenant_id}-{agent_type} — Few-shot demonstration examples (last-write-wins via replace_dataset, same as prompts). Saving the empty set clears the store: save_demonstrations([]) is expressed as clear_demonstrations, which deletes the dataset (an empty frame cannot be persisted), so a later load_demonstrations returns None. TelemetryWorkflowStore relies on this to roll an empty-prior learning corpus back on a failed multi-step write.
  • ArtifactManager dataset readers treat DatasetNotFoundError as the absent-artifact signal and let other store ValueErrors propagate, so a real store failure is not mistaken for "no optimized prompts", "no demonstrations", or a missing versioned snapshot.

Concurrency contract: replace_dataset serializes same-name writers only within one shared DatasetStore instance on one event loop. If a torn delete/create cannot restore the old frame, it raises DatasetReplaceRestoreFailedError. The protection does not extend across threads, processes, or replicas, and a hard kill between delete and create can still lose the dataset because Phoenix has no atomic swap primitive. - dspy-experiments-{tenant_id}-{agent_type} — Optimization run metrics as typed ExperimentMetrics rows (one per run via save_experiment; read the latest with load_latest_experiment) - ("model", <key>) blobs — compiled DSPy module state for profile_selection, entity_extraction, simba_query_enhancement (triggered mode publishes compiled instructions as versioned prompts instead of a module-state blob); each version ledger stores consumed_example_ids, decision, scored, score, metric_id, base_score, candidate_score, and created_at; rows written before a field existed omit it and read back as None. A version that records a score must name its metric. - ("config", "gateway_thresholds") blob — calibrated gateway thresholds - ("config", "golden_set_ground_truth") blob — tenant-uploaded retrieval golden set, validated as rows with query plus normalized expected_videos; the quality-monitor Deployment seeds it once from its mounted JSON file, then reads it through ArtifactManager for golden evals - ("config", "profile_selection_ground_truth") blob — tenant-uploaded profile-selection labels, validated as rows with query plus normalized expected_videos - ("config", "entity_extraction_ground_truth") blob — tenant-uploaded entity-extraction labels, validated as rows with query plus normalized entities[{text, type}]

Stored prompt artifact structure (retrieved from DatasetStore):

{
  "system_prompt": "Optimized instructions...",
  "few_shot_examples": [
    {
      "conversation_history": "",
      "user_query": "Show me tutorials on Python",
      "routing_decision": {"search_modality": "video", "generation_type": "raw_results"}
    }
  ],
  "model_config": {
    "student_model": "google/gemma-3-1b-it",
    "temperature": 0.1,
    "max_tokens": 100
  }
}

File References

File Purpose
optimizer/dspy_agent_optimizer.py Multi-agent prompt optimization
optimizer/artifact_manager.py Artifact persistence, canary FSM, regression-reject promotion, via Phoenix DatasetStore
optimizer/signature_variants.py Per-tenant DSPy signature variant registry
optimizer/strategy_learner.py Trace-to-strategy distillation into Mem0
routing/xgboost_meta_models.py Training-decision / strategy-selection / fusion-benefit meta-models

  • Routing Module Study Guide: docs/modules/routing.md - Tiered routing strategies, ProfilePerformanceOptimizer
  • Agents Module Study Guide: docs/modules/agents.md - OrchestratorAgent integration, full agent roster
  • Telemetry Module Study Guide: docs/modules/telemetry.md - Phoenix span collection
  • Evaluation Module Study Guide: docs/modules/evaluation.md - QualityMonitor, RoutingEvaluator

Next Steps:

  1. Review the DSPy BootstrapFewShot teleprompter docs; MIPROv2/SIMBA are not wired into any mode today (GEPA is wired only for --mode triggered's all-failure reflective recompile)

  2. Extend _create_teleprompter if a data-size-tiered optimizer switch (e.g. MIPROv2 above a higher example count) becomes worth the added compile cost

  3. Monitor optimization metrics in production via ExperimentMetrics rows (load_latest_experiment)

  4. Wire TrainingStrategyModel / FusionBenefitModel into a production caller (currently exercised only in tests)

  5. Test synthetic data generation for profile selection cold start via POST /synthetic/generate


File References:

  • libs/agents/cogniverse_agents/routing/xgboost_meta_models.py - XGBoost meta-models for training decisions

  • libs/agents/cogniverse_agents/optimizer/dspy_agent_optimizer.py - Multi-agent DSPy prompt optimization

  • libs/agents/cogniverse_agents/optimizer/artifact_manager.py - Artifact persistence, canary FSM, promotion gate

  • libs/agents/cogniverse_agents/optimizer/strategy_learner.py - Trace distillation into reusable strategies

  • libs/runtime/cogniverse_runtime/optimization_cli.py - CLI entry point for all 14 optimization/maintenance modes

  • libs/runtime/cogniverse_runtime/quality_monitor_cli.py - Quality-triggered optimization driver

  • libs/runtime/cogniverse_runtime/routers/tenant.py - On-demand optimization API endpoints

  • libs/synthetic/cogniverse_synthetic/ - Synthetic training data generation service + REST API