Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
573 changes: 573 additions & 0 deletions .github/workflows/agentic-tests-mcp.yml

Large diffs are not rendered by default.

134 changes: 132 additions & 2 deletions core/loop_orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,12 @@

from __future__ import annotations

import logging
import time
import uuid
from typing import Any, Callable, Dict, Optional

from .metrics import MetricsEmitter, noop_emitter
from .pal_integration import PALReviewSignal, incorporate_pal_feedback
from .quality_assessment import QualityAssessor
from .termination import (
Expand Down Expand Up @@ -53,16 +56,51 @@ class LoopOrchestrator:

Where skill_invoker is a callable that Claude Code uses to
execute the sc-implement skill and return evidence.

Observability:
The orchestrator supports structured logging and metrics emission:

- Logger: Inject a custom logger via the `logger` parameter.
All log messages include `loop_id` for correlation.
- Metrics: Inject a metrics callback via `metrics_emitter`.
See core.metrics for the MetricsEmitter protocol.

Emitted metrics:
- loop.started.count: Emitted when loop begins
- loop.completed.count: Emitted with termination_reason tag
- loop.duration.seconds: Total loop execution time
- loop.iterations.total.gauge: Number of iterations executed
- loop.quality_score.final.gauge: Final quality score
- loop.errors.count: Skill invocation failures
- loop.iteration.duration.seconds: Per-iteration timing
- loop.iteration.quality_score.gauge: Per-iteration quality
- loop.iteration.quality_delta.gauge: Quality change per iteration

Thread Safety:
This class is NOT thread-safe. Each LoopOrchestrator instance
maintains mutable state (iteration_history, score_history,
all_changed_files) that is modified during run(). Do not share
instances across threads. Create a new orchestrator per task/thread.
"""

def __init__(self, config: Optional[LoopConfig] = None):
def __init__(
self,
config: Optional[LoopConfig] = None,
logger: Optional[logging.Logger] = None,
metrics_emitter: Optional[MetricsEmitter] = None,
):
"""
Initialize the loop orchestrator.

Args:
config: Loop configuration (defaults to LoopConfig())
logger: A logger instance. If not provided, a default will be used.
metrics_emitter: A callable for emitting operational metrics.
"""
self.config = config or LoopConfig()
self.logger = logger or logging.getLogger(__name__)
self.metrics_emitter = metrics_emitter or noop_emitter
self.loop_id = str(uuid.uuid4())[:12]
self.assessor = QualityAssessor(self.config.quality_threshold)
self.iteration_history: list[IterationResult] = []
self.score_history: list[float] = []
Expand Down Expand Up @@ -93,23 +131,46 @@ def run(
LoopResult with final output, assessment, and iteration history
"""
self._start_time = time.monotonic()
self.metrics_emitter("loop.started.count", 1)
self.logger.info(
"Starting agentic loop.",
extra={
"loop_id": self.loop_id,
"max_iterations": self.config.max_iterations,
"quality_threshold": self.config.quality_threshold,
"pal_review_enabled": self.config.pal_review_enabled,
},
)

current_context = initial_context.copy()
termination_reason = TerminationReason.MAX_ITERATIONS
output: dict[str, Any] = {}
assessment = QualityAssessment(overall_score=0.0, passed=False)

for iteration in range(self.config.max_iterations):
iter_start = time.monotonic()
log_context = {"loop_id": self.loop_id, "iteration": iteration}
self.logger.info(
f"Starting iteration {iteration + 1}/{self.config.max_iterations}.",
extra=log_context,
)

# Check timeout
if self._check_timeout():
self.logger.warning("Loop timed out.", extra=log_context)
termination_reason = TerminationReason.TIMEOUT
break

# 1. Execute skill (delegated to Claude Code)
try:
output = skill_invoker(current_context)
except Exception:
self.logger.error(
"Error during skill invocation.", exc_info=True, extra=log_context
)
self.metrics_emitter(
"loop.errors.count", 1, {"reason": "skill_invocation"}
)
termination_reason = TerminationReason.ERROR
self._record_iteration(
iteration=iteration,
Expand All @@ -129,10 +190,26 @@ def run(

# 2. Assess quality
assessment = self.assessor.assess(output)
self.logger.debug(
"Assessment complete.",
extra={
**log_context,
"score": assessment.overall_score,
"passed": assessment.passed,
},
)
self.score_history.append(assessment.overall_score)

# 3. Check if quality threshold met
if assessment.passed:
self.logger.info(
"Quality threshold met.",
extra={
**log_context,
"score": assessment.overall_score,
"threshold": self.config.quality_threshold,
},
)
termination_reason = TerminationReason.QUALITY_MET
self._record_iteration(
iteration=iteration,
Expand All @@ -149,6 +226,7 @@ def run(
self.score_history,
self.config.oscillation_window,
):
self.logger.info("Oscillation detected.", extra=log_context)
termination_reason = TerminationReason.OSCILLATION
self._record_iteration(
iteration=iteration,
Expand All @@ -165,6 +243,7 @@ def run(
self.config.oscillation_window,
self.config.stagnation_threshold,
):
self.logger.info("Stagnation detected.", extra=log_context)
termination_reason = TerminationReason.STAGNATION
self._record_iteration(
iteration=iteration,
Expand All @@ -181,6 +260,7 @@ def run(
self.score_history[-2],
self.config.min_improvement,
):
self.logger.info("Insufficient improvement detected.", extra=log_context)
termination_reason = TerminationReason.INSUFFICIENT_IMPROVEMENT
self._record_iteration(
iteration=iteration,
Expand All @@ -195,6 +275,7 @@ def run(
# 5. Generate PAL review signal (if enabled and not last iteration)
pal_signal = None
if self.config.pal_review_enabled and iteration < self.config.max_iterations - 1:
self.logger.debug("Generating PAL review signal.", extra=log_context)
pal_signal = PALReviewSignal.generate_review_signal(
iteration=iteration,
changed_files=changed_files,
Expand Down Expand Up @@ -223,6 +304,9 @@ def run(
# Generate final signals
if termination_reason == TerminationReason.QUALITY_MET:
# Final validation signal
self.logger.debug(
"Generating final PAL validation signal.", extra={"loop_id": self.loop_id}
)
final_signal = PALReviewSignal.generate_final_validation_signal(
changed_files=self.all_changed_files,
quality_assessment=assessment,
Expand All @@ -237,6 +321,10 @@ def run(
TerminationReason.STAGNATION,
):
# Debug signal for stuck loops
self.logger.debug(
"Generating PAL debug signal for stuck loop.",
extra={"loop_id": self.loop_id, "reason": termination_reason.value},
)
debug_signal = PALReviewSignal.generate_debug_signal(
iteration=len(self.iteration_history) - 1,
termination_reason=termination_reason.value,
Expand All @@ -246,13 +334,34 @@ def run(
if self.iteration_history:
self.iteration_history[-1].pal_review = debug_signal

total_time = time.monotonic() - self._start_time
final_tags = {"termination_reason": termination_reason.value}
self.metrics_emitter("loop.completed.count", 1, final_tags)
self.metrics_emitter("loop.duration.seconds", total_time, final_tags)
self.metrics_emitter(
"loop.iterations.total.gauge", len(self.iteration_history), final_tags
)
self.metrics_emitter(
"loop.quality_score.final.gauge", assessment.overall_score, final_tags
)
self.logger.info(
"Agentic loop finished.",
extra={
"loop_id": self.loop_id,
"termination_reason": termination_reason.value,
"total_iterations": len(self.iteration_history),
"total_time": total_time,
"final_score": assessment.overall_score,
},
)

return LoopResult(
final_output=output,
final_assessment=assessment,
iteration_history=self.iteration_history,
termination_reason=termination_reason,
total_iterations=len(self.iteration_history),
total_time=time.monotonic() - self._start_time,
total_time=total_time,
)

def _check_timeout(self) -> bool:
Expand All @@ -276,6 +385,13 @@ def _record_iteration(
input_quality = self.score_history[-2] if len(self.score_history) >= 2 else 0.0
output_quality = assessment.overall_score

# Emit per-iteration metrics
self.metrics_emitter("loop.iteration.duration.seconds", time_taken)
self.metrics_emitter("loop.iteration.quality_score.gauge", output_quality)
self.metrics_emitter(
"loop.iteration.quality_delta.gauge", output_quality - input_quality
)

self.iteration_history.append(
IterationResult(
iteration=iteration,
Expand All @@ -289,6 +405,20 @@ def _record_iteration(
changed_files=changed_files,
)
)
self.logger.debug(
"Iteration recorded.",
extra={
"loop_id": self.loop_id,
"iteration": iteration,
"input_quality": input_quality,
"output_quality": output_quality,
"time_taken": time_taken,
"success": success,
"termination_reason": termination,
"changed_files_count": len(changed_files),
"pal_signal_generated": pal_signal is not None,
},
)

def _prepare_next_iteration(
self,
Expand Down
Loading
Loading