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¶
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 withdspy.teleprompt.GEPA(see Triggered Optimization); no mode invokes MIPROv2 or SIMBA.ExperimentMetrics.optimizeris a free-text field intended to record whichever optimizer produced a run, but every serving call site passes"BootstrapFewShot". The--mode simbaname is historical (it optimizesQueryEnhancementAgent's DSPy module); it does not run the SIMBA algorithm. - Gateway Threshold Tuning:
_compute_gateway_thresholds()derives GLiNER/fast-path thresholds from Phoenixcogniverse.gatewayspans using a rule-based adjustment (± based on per-branch error rate and mean confidence) plus a p25-percentile-derivedgliner_threshold. - Online Routing Evaluation:
run_online_routing_evaluationscorescogniverse.routingspans (routing outcome + confidence calibration) viaOnlineEvaluatorand persists the scores as telemetry annotations; driven byautomation_rules.online_evaluation(OnlineEvaluationConfiginrouting/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 viaArtifactManager. 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 thescoredflag, the numericscoreused for confirmation when present, andmetric_idnaming the metric that produced it;training_selection.<optimizer>.confirmation_score_thresholdmakes confirmation score-aware, and an omitted threshold keeps confirmation presence-based and ignoresmetric_id. Under a threshold, a score counts only when itsmetric_idis the optimizer's current one — the ids live beside the metric functions inoptimization_cliand are registered inOPTIMIZER_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 belowlow_confirmation_threshold. SIMBA also enforces the tenant floor from routing config: belowmin_samples_for_optimizationormin_unique_queries, it saves a version with decisioninsufficient_populationand 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
WorkflowExecutionrecords fromcogniverse.orchestrationspans viaOrchestrationEvaluator, drops demos whoseagent_sequencereferences an agent no longer live inconfigs/config.json, and replaces templates, execution demos, agent performance profiles, and query patterns through oneWorkflowStoreRegistryoperation guarded by a renewable per-tenant Redis lease. - Training-Decision Meta-Models:
TrainingDecisionModel(train/skip gate; wired intoQualityMonitor),TrainingStrategyModel(PURE_REAL / HYBRID / SYNTHETIC / SKIP selection), andFusionBenefitModel— all XGBoost classifiers inrouting/xgboost_meta_models.py, persisted via the sameArtifactManager. - Strategy Distillation:
StrategyLearnerdistills execution traces into reusableStrategyobjects (pattern extraction, no LLM; or LLM-based contrastive distillation) and stores them in Vespa memory viaMem0MemoryManagerfor later retrieval byMemoryAwareMixin. - Regression-Reject Promotion Gate:
ArtifactManager.promote_if_betteronly promotes a candidate when its score beats the active baseline (withintolerance); every attempt — win or reject — lands as a typedExperimentMetricsrow. - Canary Rollout + Rollback:
ArtifactManager's three-slot (active/canary/retired) state machine and the--mode rollbackCLI 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 thecogniverse-optimization-runnerWorkflowTemplate. - Scheduled Workflows: helm-chart
CronWorkflows rungateway-thresholds/entity-extraction/simba/profile/workflowweekly,gateway-thresholdsdaily,cleanupdaily,syntheticweekly, a forced distillation pass daily viaquality_monitor_cli --once, andmonthly-reportsmonthly.
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:
- Reads the live
agentsblock fromconfigs/config.json(notAgentRegistry, which starts empty in the optimization CLI's own pod) to build a set of currently-enabled agent names. - Drops any execution whose
agent_sequencereferences an agent not in that live set (stale demos from renamed/removed agents can't be replayed). - Replaces the surviving executions, agent performance profiles, query-type patterns, and derived templates via
WorkflowStoreRegistry.get(name="telemetry")— the same storeWorkflowIntelligencereads 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.WorkflowIntelligencevalidates 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:
- Builds a labeled
dspy.Exampleset from the high-scoring rows (agent-specific field mapping forsearch/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) whenenable_reflective_recompileis on and, after the same ~25% tail split is applied to the low-scoring rows, the resulting GEPA trainset has at leastmin_reflective_failuresrows; 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 thanmin_reflective_failures(the threshold checks the split trainset, not the raw failing-row count). - Compiles the module the runtime actually serves (
_served_module:SearchOptimizationModule,SummarizationModule,ReportGenerationModule) withBootstrapFewShoton the trainset, scoped insidedspy.context(lm=optimizer.lm)(initialize_language_modelonly setsoptimizer.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). - Scores the compiled candidate against the currently-active baseline (
_holdout_scoresvia_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 withreason: "baseline_not_reconstructable"rather than scoring against a stock module. - Publishes the compiled module's whole
dump_state()via_serve_compiled_promptsonly if the candidate wins by at least the tenant'soptimization_improvement_threshold— the call routes throughArtifactManager.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 withpromoted=False,--mode rollbackrestores a prior version, and the result reports the outcome under"served"(served_agent,version,active,promoted, plusbaseline_score/candidate_scorewhen eval material was available or areasonwhen 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 viasnapshot_activewhensnapshot_before_promote=True, the default). Withserve_versioned=Truethe 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 inextra_metrics["served_version"]. With the defaultserve_versioned=Falsethe win overwrites the un-versioned active artefacts only. - Rejected otherwise: prompts/demos are not saved; the run is recorded with
promoted=Falseand arejection_reason.--mode triggeredwires the tenant'soptimization_improvement_thresholdconfig knob intomin_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'sschemas_deployedlist andschema_count; top-level summary{org_count, tenant_count, schema_count}.performance-YYYYMM.json— per-tenant Phoenix span count, latencymean / p50 / p95, anderror_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 exceptcleanup,egress-netpol, andmonthly-reports -
--lookback-hours: hours of span history to analyze (default 24.0, accepts fractions) -
--log-retention-days/--memory-retention-days: cleanup mode; overrideLOG_RETENTION_DAYS/MEMORY_RETENTION_DAYS(7 / 30). Cleanup requires existing dedicatedLOG_DIRandTEMP_DIRroots;TEMP_RETENTION_DAYSdefaults to 1. Unsafe roots raiseCleanupRootErrorbefore 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):
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:
-
Automatic (Scheduled): CronWorkflows run
optimization_cliper the schedule table above -
Automatic (Quality-triggered):
QualityMonitorsubmits--mode triggeredon detected degradation -
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:
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:
Cause: DSPy Examples store all fields as strings internally.
Fix: Convert to strings:
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¶
- Define all output fields upfront - Check what your metrics access
- Use consistent field names - Match your DSPy signature output fields
- Validate before optimization - Catch missing fields early
- Use string values - DSPy converts everything to strings
- Document required fields - Keep a reference list for your team
File References¶
libs/agents/cogniverse_agents/optimizer/dspy_agent_optimizer.py- Training data loadingtests/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_atas 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_promptsusesDatasetStore.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 viareplace_dataset, same as prompts). Saving the empty set clears the store:save_demonstrations([])is expressed asclear_demonstrations, which deletes the dataset (an empty frame cannot be persisted), so a laterload_demonstrationsreturnsNone.TelemetryWorkflowStorerelies on this to roll an empty-prior learning corpus back on a failed multi-step write.ArtifactManagerdataset readers treatDatasetNotFoundErroras the absent-artifact signal and let other storeValueErrors propagate, so a real store failure is not mistaken for "no optimized prompts", "no demonstrations", or a missing versioned snapshot.
Concurrency contract:
replace_datasetserializes same-name writers only within one sharedDatasetStoreinstance on one event loop. If a torn delete/create cannot restore the old frame, it raisesDatasetReplaceRestoreFailedError. 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 typedExperimentMetricsrows (one per run viasave_experiment; read the latest withload_latest_experiment) -("model", <key>)blobs — compiled DSPy module state forprofile_selection,entity_extraction,simba_query_enhancement(triggered mode publishes compiled instructions as versioned prompts instead of a module-state blob); each version ledger storesconsumed_example_ids,decision,scored,score,metric_id,base_score,candidate_score, andcreated_at; rows written before a field existed omit it and read back asNone. 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 withqueryplus normalizedexpected_videos; the quality-monitor Deployment seeds it once from its mounted JSON file, then reads it throughArtifactManagerfor golden evals -("config", "profile_selection_ground_truth")blob — tenant-uploaded profile-selection labels, validated as rows withqueryplus normalizedexpected_videos-("config", "entity_extraction_ground_truth")blob — tenant-uploaded entity-extraction labels, validated as rows withqueryplus normalizedentities[{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 |
Related Documentation¶
- 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:
-
Review the DSPy
BootstrapFewShotteleprompter docs; MIPROv2/SIMBA are not wired into any mode today (GEPA is wired only for--mode triggered's all-failure reflective recompile) -
Extend
_create_teleprompterif a data-size-tiered optimizer switch (e.g. MIPROv2 above a higher example count) becomes worth the added compile cost -
Monitor optimization metrics in production via
ExperimentMetricsrows (load_latest_experiment) -
Wire
TrainingStrategyModel/FusionBenefitModelinto a production caller (currently exercised only in tests) -
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