From 1c21263a375d0e606f33bcb7d939634fbba32213 Mon Sep 17 00:00:00 2001
From: Kirill <106469980+waterflane@users.noreply.github.com>
Date: Mon, 7 Sep 2026 20:16:55 +0300
Subject: [PATCH 01/11] fix(provider): classify failures and stop repeated
provider calls
---
src/contextforge/models/__init__.py | 14 ++
src/contextforge/models/ollama.py | 8 +-
src/contextforge/models/openai_compatible.py | 60 +++++++-
src/contextforge/models/providers.py | 140 ++++++++++++++++++-
tests/test_model_providers.py | 79 ++++++++++-
tests/test_openai_compatible.py | 36 +++--
6 files changed, 313 insertions(+), 24 deletions(-)
diff --git a/src/contextforge/models/__init__.py b/src/contextforge/models/__init__.py
index 4a4d0f1..6d91fa2 100644
--- a/src/contextforge/models/__init__.py
+++ b/src/contextforge/models/__init__.py
@@ -45,11 +45,18 @@
ModelRequest,
ModelResponse,
ModelUsage,
+ ProviderAuthenticationError,
+ ProviderAuthorizationError,
ProviderCancelledError,
ProviderCapabilities,
+ ProviderCircuitOpenError,
ProviderConfiguration,
ProviderConfigurationError,
ProviderDiagnostic,
+ ProviderMissingCredentialError,
+ ProviderModelNotFoundError,
+ ProviderQuotaError,
+ ProviderRateLimitError,
ProviderRequestError,
ProviderRuntime,
ProviderTimeoutError,
@@ -124,10 +131,17 @@
"OpenAICompatibleModelProvider",
"OpenAICompatibleTransport",
"ProviderCancelledError",
+ "ProviderAuthenticationError",
+ "ProviderAuthorizationError",
"ProviderCapabilities",
+ "ProviderCircuitOpenError",
"ProviderConfiguration",
"ProviderConfigurationError",
"ProviderDiagnostic",
+ "ProviderModelNotFoundError",
+ "ProviderMissingCredentialError",
+ "ProviderQuotaError",
+ "ProviderRateLimitError",
"ProviderRequestError",
"ProviderRuntime",
"ProviderTimeoutError",
diff --git a/src/contextforge/models/ollama.py b/src/contextforge/models/ollama.py
index e600460..a4cadf0 100644
--- a/src/contextforge/models/ollama.py
+++ b/src/contextforge/models/ollama.py
@@ -19,6 +19,8 @@
ModelRequest,
ModelResponse,
ModelUsage,
+ ProviderAuthenticationError,
+ ProviderAuthorizationError,
ProviderCapabilities,
ProviderConfiguration,
ProviderConfigurationError,
@@ -378,7 +380,11 @@ def _raise_for_ollama_http_error(status: int, data: bytes) -> None:
raise StructuredOutputSchemaUnsupportedError(
"Ollama rejected the structured output schema"
)
- if status in {400, 401, 403, 404, 422}:
+ if status == 401:
+ raise ProviderAuthenticationError("Ollama rejected authentication (HTTP 401)")
+ if status == 403:
+ raise ProviderAuthorizationError("Ollama rejected authorization (HTTP 403)")
+ if status in {400, 404, 422}:
raise ProviderRequestError(f"Ollama rejected the request (HTTP {status})")
raise ProviderUnavailableError(f"Ollama returned HTTP status {status}")
diff --git a/src/contextforge/models/openai_compatible.py b/src/contextforge/models/openai_compatible.py
index b96a17c..33b940a 100644
--- a/src/contextforge/models/openai_compatible.py
+++ b/src/contextforge/models/openai_compatible.py
@@ -21,10 +21,15 @@
ModelRequest,
ModelResponse,
ModelUsage,
+ ProviderAuthenticationError,
+ ProviderAuthorizationError,
ProviderCancelledError,
ProviderCapabilities,
ProviderConfiguration,
ProviderConfigurationError,
+ ProviderModelNotFoundError,
+ ProviderQuotaError,
+ ProviderRateLimitError,
ProviderRequestError,
ProviderRuntime,
ProviderTimeoutError,
@@ -712,15 +717,21 @@ def _raise_for_status(
if 200 <= status < 300:
return
detail = _safe_error_detail(response.body)
+ error_code = _safe_error_code(response.body)
lowered = "" if detail is None else detail.casefold()
+ classified = "" if error_code is None else error_code.casefold()
suffix = "" if detail is None else f": {detail}"
- if status in {401, 403}:
- raise ProviderRequestError(
- f"OpenAI-compatible authentication failed with HTTP {status}{suffix}"
+ if status == 401:
+ raise ProviderAuthenticationError(
+ "OpenAI-compatible provider rejected authentication (HTTP 401)"
+ )
+ if status == 403:
+ raise ProviderAuthorizationError(
+ "OpenAI-compatible provider rejected authorization (HTTP 403)"
)
if status == 404 and operation == "chat completion":
- raise ProviderRequestError(
- f"model ID {model_id!r} was not found (HTTP 404){suffix}"
+ raise ProviderModelNotFoundError(
+ f"model ID {model_id!r} was not found (HTTP 404)"
)
if (
status in {400, 422}
@@ -773,7 +784,25 @@ def _raise_for_status(
raise ProviderRequestError(
f"OpenAI-compatible server rejected the request (HTTP {status}){suffix}"
)
- if status in {408, 429} or 500 <= status < 600:
+ if status == 408:
+ raise ProviderTimeoutError("OpenAI-compatible provider returned HTTP 408")
+ if status == 429 and any(
+ marker in f"{classified} {lowered}"
+ for marker in (
+ "insufficient_quota",
+ "quota exceeded",
+ "quota_exceeded",
+ "billing hard limit",
+ "billing_hard_limit",
+ "credits exhausted",
+ )
+ ):
+ raise ProviderQuotaError("OpenAI-compatible provider quota is exhausted")
+ if status == 429:
+ raise ProviderRateLimitError(
+ "OpenAI-compatible provider rate limit was reached"
+ )
+ if 500 <= status < 600:
raise ProviderUnavailableError(
f"OpenAI-compatible server failed {operation} with HTTP {status}{suffix}"
)
@@ -802,6 +831,25 @@ def _safe_error_detail(data: bytes) -> str | None:
return None
+def _safe_error_code(data: bytes) -> str | None:
+ """Read only a bounded provider error classifier, never an arbitrary body."""
+
+ try:
+ payload = json.loads(data.decode("utf-8", errors="strict"))
+ except (UnicodeDecodeError, json.JSONDecodeError):
+ return None
+ if not isinstance(payload, dict):
+ return None
+ error = payload.get("error")
+ candidate = error.get("code") if isinstance(error, dict) else payload.get("code")
+ if not isinstance(candidate, str):
+ return None
+ normalized = candidate.strip()
+ if not re.fullmatch(r"[A-Za-z0-9_.-]{1,128}", normalized):
+ return None
+ return normalized
+
+
def _redact_error[ErrorType: (ProviderRequestError, ProviderUnavailableError)](
error: ErrorType, secrets: Sequence[str]
) -> ErrorType:
diff --git a/src/contextforge/models/providers.py b/src/contextforge/models/providers.py
index 9792918..6a7c549 100644
--- a/src/contextforge/models/providers.py
+++ b/src/contextforge/models/providers.py
@@ -40,6 +40,7 @@
MAX_PROVIDER_TIMEOUT_SECONDS = 600.0
MAX_PROVIDER_CONCURRENCY = 8
MAX_PROVIDER_RETRIES = 2
+DEFAULT_PROVIDER_CIRCUIT_FAILURE_THRESHOLD = 3
RETRY_DELAYS_SECONDS = (0.25, 1.0)
MAX_UNTRUSTED_SOURCE_BYTES = 1_000_000
MAX_REQUEST_SOURCE_BYTES = 4_000_000
@@ -225,7 +226,7 @@ def load_credential(
source = os.environ if environment is None else environment
value = source.get(self.credential_env)
if value is None or not value:
- raise ProviderConfigurationError(
+ raise ProviderMissingCredentialError(
f"credential environment variable {self.credential_env!r} is not set"
)
return SecretStr(value)
@@ -809,6 +810,7 @@ class ModelProviderError(RuntimeError):
"""Base provider failure with stable retry classification."""
retry_classification = RetryClassification.NON_RETRYABLE
+ provider_wide = False
def __init__(
self, message: str, *, diagnostic: ProviderDiagnostic | None = None
@@ -818,6 +820,7 @@ def __init__(
self.provider_capability_calls = 0
self.transport_attempts = 1
self.total_provider_http_calls = 1
+ self.circuit_opened = False
super().__init__(message)
def add_http_accounting(
@@ -944,6 +947,52 @@ class ProviderUnavailableError(ModelProviderError):
retry_classification = RetryClassification.RETRYABLE
+class ProviderRequestError(ModelProviderError):
+ """Raised for a non-retryable provider request failure."""
+
+
+class ProviderRateLimitError(ProviderUnavailableError):
+ """Raised for a transient provider rate limit."""
+
+
+class ProviderAuthenticationError(ProviderRequestError):
+ """Raised when provider credentials are absent, expired, or rejected."""
+
+ provider_wide = True
+
+
+class ProviderMissingCredentialError(ProviderAuthenticationError):
+ """Raised when a configured credential reference has no usable value."""
+
+
+class ProviderAuthorizationError(ProviderRequestError):
+ """Raised when valid credentials cannot perform the configured operation."""
+
+ provider_wide = True
+
+
+class ProviderQuotaError(ProviderRequestError):
+ """Raised when provider quota or billing capacity is exhausted."""
+
+ provider_wide = True
+
+
+class ProviderModelNotFoundError(ProviderRequestError):
+ """Raised when the configured model identity is unavailable."""
+
+ provider_wide = True
+
+
+class ProviderCircuitOpenError(ModelProviderError):
+ """Raised before dispatch after the shared provider circuit opens."""
+
+ provider_wide = True
+
+ def __init__(self, message: str) -> None:
+ super().__init__(message)
+ self.circuit_opened = True
+
+
class ProviderCancelledError(ModelProviderError):
"""Raised when explicit or task cancellation stops provider work."""
@@ -951,9 +1000,7 @@ class ProviderCancelledError(ModelProviderError):
class ProviderConfigurationError(ModelProviderError):
"""Raised for a non-retryable local provider configuration failure."""
-
-class ProviderRequestError(ModelProviderError):
- """Raised for a non-retryable provider request failure."""
+ provider_wide = True
class ModelProvider(Protocol):
@@ -1066,6 +1113,7 @@ def __init__(
environment: Mapping[str, str] | None = None,
clock: Callable[[], float] = time.monotonic,
retry_delays: Sequence[float] = RETRY_DELAYS_SECONDS,
+ circuit_failure_threshold: int = DEFAULT_PROVIDER_CIRCUIT_FAILURE_THRESHOLD,
) -> None:
self.configuration = configuration
self._environment = environment
@@ -1081,9 +1129,62 @@ def __init__(
raise ValueError("retry_delays must contain finite non-negative numbers")
if configuration.retry_limit and not self._retry_delays:
raise ValueError("retry_delays must not be empty when retries are enabled")
+ if (
+ type(circuit_failure_threshold) is not int
+ or circuit_failure_threshold < 1
+ or circuit_failure_threshold > 100
+ ):
+ raise ValueError("circuit_failure_threshold must be between 1 and 100")
self._semaphore = asyncio.Semaphore(configuration.concurrency_limit)
+ self._circuit_lock = asyncio.Lock()
+ self._circuit_failure_threshold = circuit_failure_threshold
+ self._consecutive_failure_key: str | None = None
+ self._consecutive_failure_count = 0
+ self._circuit_error: tuple[str, str] | None = None
self._closed = False
+ async def _raise_if_circuit_open(self) -> None:
+ async with self._circuit_lock:
+ if self._circuit_error is None:
+ return
+ code, message = self._circuit_error
+ raise ProviderCircuitOpenError(
+ f"provider circuit is open after {code}: {message}"
+ )
+
+ async def _record_success(self) -> None:
+ async with self._circuit_lock:
+ if self._circuit_error is None:
+ self._consecutive_failure_key = None
+ self._consecutive_failure_count = 0
+
+ async def _record_final_failure(self, error: ModelProviderError) -> None:
+ if isinstance(error, (ProviderCancelledError, ProviderCircuitOpenError)):
+ return
+ code, message = provider_error_details(error)
+ async with self._circuit_lock:
+ if self._circuit_error is not None:
+ return
+ if error.provider_wide:
+ self._circuit_error = (code, message)
+ error.circuit_opened = True
+ return
+ if classify_retry(error) is not RetryClassification.RETRYABLE:
+ self._consecutive_failure_key = None
+ self._consecutive_failure_count = 0
+ return
+ key = (
+ f"{self.configuration.provider_id}:{self.configuration.model_id}:{code}"
+ )
+ if key == self._consecutive_failure_key:
+ self._consecutive_failure_count += 1
+ else:
+ self._consecutive_failure_key = key
+ self._consecutive_failure_count = 1
+ if self._consecutive_failure_count >= self._circuit_failure_threshold:
+ self._circuit_error = (code, message)
+ error.circuit_opened = True
+
async def execute(
self,
request: ModelRequest,
@@ -1095,7 +1196,12 @@ async def execute(
if self._closed:
raise ProviderRequestError("provider is closed")
- credential = self.configuration.load_credential(self._environment)
+ await self._raise_if_circuit_open()
+ try:
+ credential = self.configuration.load_credential(self._environment)
+ except ModelProviderError as exc:
+ await self._record_final_failure(exc)
+ raise
secrets = () if credential is None else (credential.get_secret_value(),)
started = self._clock()
timeout = min(
@@ -1344,6 +1450,7 @@ def counter_data() -> dict[str, Any]:
timeout=timeout,
)
try:
+ await self._raise_if_circuit_open()
transport_attempts += 1
total_http_calls += 1
raw = await _await_bounded(
@@ -1528,6 +1635,7 @@ def counter_data() -> dict[str, Any]:
usage=raw.usage,
)
progress.complete(message="Provider request completed.")
+ await self._record_success()
return ModelResponse(
normalized_json=accepted.normalized_json,
value=accepted.value,
@@ -1912,6 +2020,7 @@ def counter_data() -> dict[str, Any]:
)
else:
progress.fail(message=message)
+ await self._record_final_failure(error)
raise error
async def close(self) -> None:
@@ -1937,6 +2046,20 @@ def provider_error_details(error: BaseException) -> tuple[str, str]:
return "provider_timeout", "provider request timed out"
if isinstance(error, ProviderCancelledError):
return "cancelled", "provider request was cancelled"
+ if isinstance(error, ProviderCircuitOpenError):
+ return "provider_circuit_open", "provider circuit breaker is open"
+ if isinstance(error, ProviderMissingCredentialError):
+ return "missing_credential", "provider credential is not configured"
+ if isinstance(error, ProviderAuthenticationError):
+ return "authentication_failed", "provider authentication failed"
+ if isinstance(error, ProviderAuthorizationError):
+ return "authorization_failed", "provider authorization failed"
+ if isinstance(error, ProviderQuotaError):
+ return "quota_exhausted", "provider quota is exhausted"
+ if isinstance(error, ProviderModelNotFoundError):
+ return "model_not_found", "configured model was not found"
+ if isinstance(error, ProviderRateLimitError):
+ return "rate_limited", "provider rate limit was reached"
if isinstance(error, ContextWindowExceededError):
return (
"context_window_exceeded",
@@ -3544,9 +3667,16 @@ def _redacted_provider_error(
"ModelUsage",
"MissingRequiredFieldIssue",
"ProviderCancelledError",
+ "ProviderAuthenticationError",
+ "ProviderAuthorizationError",
+ "ProviderCircuitOpenError",
"ProviderCapabilities",
"ProviderConfiguration",
"ProviderConfigurationError",
+ "ProviderModelNotFoundError",
+ "ProviderMissingCredentialError",
+ "ProviderQuotaError",
+ "ProviderRateLimitError",
"ProviderDiagnostic",
"ProviderRequestError",
"ProviderRuntime",
diff --git a/tests/test_model_providers.py b/tests/test_model_providers.py
index f346ab0..8e7a02b 100644
--- a/tests/test_model_providers.py
+++ b/tests/test_model_providers.py
@@ -20,9 +20,12 @@
ModelResponse,
ModelUsage,
OllamaModelProvider,
+ ProviderAuthenticationError,
ProviderCancelledError,
+ ProviderCircuitOpenError,
ProviderConfiguration,
ProviderConfigurationError,
+ ProviderMissingCredentialError,
ProviderRequestError,
ProviderTimeoutError,
ProviderTransportResponse,
@@ -576,6 +579,70 @@ async def exercise() -> FakeModelProvider:
assert provider.call_count == 3
+def test_provider_circuit_opens_after_three_matching_transient_failures() -> None:
+ async def exercise() -> FakeModelProvider:
+ provider = FakeModelProvider(
+ _configuration(),
+ scripts=[ProviderUnavailableError("offline")] * 3 + [_valid_json()],
+ )
+ for _ in range(3):
+ with pytest.raises(ProviderUnavailableError):
+ await provider.complete_structured(_request())
+ with pytest.raises(ProviderCircuitOpenError) as captured:
+ await provider.complete_structured(_request())
+ assert captured.value.circuit_opened is True
+ return provider
+
+ provider = asyncio.run(exercise())
+
+ assert provider.call_count == 3
+
+
+def test_provider_success_resets_transient_circuit_sequence() -> None:
+ async def exercise() -> FakeModelProvider:
+ provider = FakeModelProvider(
+ _configuration(),
+ scripts=[
+ ProviderUnavailableError("offline"),
+ ProviderUnavailableError("offline"),
+ _valid_json(),
+ ProviderUnavailableError("offline"),
+ ProviderUnavailableError("offline"),
+ _valid_json(),
+ ],
+ )
+ for _ in range(2):
+ with pytest.raises(ProviderUnavailableError):
+ await provider.complete_structured(_request())
+ await provider.complete_structured(_request())
+ for _ in range(2):
+ with pytest.raises(ProviderUnavailableError):
+ await provider.complete_structured(_request())
+ await provider.complete_structured(_request())
+ return provider
+
+ provider = asyncio.run(exercise())
+
+ assert provider.call_count == 6
+
+
+def test_terminal_provider_failure_opens_circuit_immediately() -> None:
+ async def exercise() -> FakeModelProvider:
+ provider = FakeModelProvider(
+ _configuration(),
+ scripts=[ProviderAuthenticationError("expired"), _valid_json()],
+ )
+ with pytest.raises(ProviderAuthenticationError):
+ await provider.complete_structured(_request())
+ with pytest.raises(ProviderCircuitOpenError):
+ await provider.complete_structured(_request())
+ return provider
+
+ provider = asyncio.run(exercise())
+
+ assert provider.call_count == 1
+
+
def test_environment_credential_loading_and_secret_redaction() -> None:
secret = "sensitive-provider-value"
configuration = ProviderConfiguration(
@@ -623,11 +690,17 @@ def test_missing_environment_credential_reference_is_non_retryable() -> None:
configuration = _configuration(credential_env="MISSING_TEST_TOKEN")
provider = FakeModelProvider(configuration, scripts=[_valid_json()], environment={})
- with pytest.raises(ProviderConfigurationError) as captured:
- asyncio.run(provider.complete_structured(_request()))
+ async def exercise() -> ProviderMissingCredentialError:
+ with pytest.raises(ProviderMissingCredentialError) as captured:
+ await provider.complete_structured(_request())
+ with pytest.raises(ProviderCircuitOpenError):
+ await provider.complete_structured(_request())
+ return captured.value
+
+ error = asyncio.run(exercise())
assert provider.call_count == 0
- assert classify_retry(captured.value) is RetryClassification.NON_RETRYABLE
+ assert classify_retry(error) is RetryClassification.NON_RETRYABLE
def test_secret_values_are_not_persisted_in_contextforge_index(tmp_path: Path) -> None:
diff --git a/tests/test_openai_compatible.py b/tests/test_openai_compatible.py
index 51ad3cd..400a957 100644
--- a/tests/test_openai_compatible.py
+++ b/tests/test_openai_compatible.py
@@ -21,9 +21,14 @@
ModelUsage,
OpenAICompatibleHTTPResponse,
OpenAICompatibleModelProvider,
+ ProviderAuthenticationError,
+ ProviderAuthorizationError,
ProviderCancelledError,
ProviderConfiguration,
ProviderConfigurationError,
+ ProviderModelNotFoundError,
+ ProviderQuotaError,
+ ProviderRateLimitError,
ProviderRequestError,
ProviderTimeoutError,
ProviderUnavailableError,
@@ -301,7 +306,7 @@ async def exercise() -> None:
asyncio.run(exercise())
-def test_http_404_reports_model_error_body_when_safe() -> None:
+def test_http_404_is_a_terminal_typed_model_error() -> None:
call_count = 0
async def transport(
@@ -323,7 +328,7 @@ async def transport(
async def exercise() -> None:
provider = OpenAICompatibleModelProvider(_configuration(), transport=transport)
- with pytest.raises(ProviderRequestError, match="model was unloaded"):
+ with pytest.raises(ProviderRequestError, match="model ID"):
await provider.complete_structured(_request())
asyncio.run(exercise())
@@ -446,7 +451,7 @@ async def exercise() -> ModelResponse:
)
-def test_auth_error_redacts_loaded_credential_and_keeps_safe_body() -> None:
+def test_auth_error_omits_loaded_credential_and_provider_body() -> None:
secret = "credential-that-must-not-leak"
async def transport(
@@ -471,7 +476,6 @@ async def exercise() -> None:
with pytest.raises(ProviderRequestError) as captured:
await provider.complete_structured(_request())
assert secret not in str(captured.value)
- assert "[REDACTED]" in str(captured.value)
assert "HTTP 401" in str(captured.value)
asyncio.run(exercise())
@@ -790,11 +794,13 @@ async def exercise_close() -> None:
@pytest.mark.parametrize(
("status", "operation", "error_type", "message"),
[
- (403, "model diagnostics", ProviderRequestError, "authentication"),
+ (401, "model diagnostics", ProviderAuthenticationError, "authentication"),
+ (403, "model diagnostics", ProviderAuthorizationError, "authorization"),
(404, "model diagnostics", ProviderRequestError, "model diagnostics"),
(418, "chat completion", ProviderRequestError, "chat completion"),
- (408, "chat completion", ProviderUnavailableError, "HTTP 408"),
- (429, "chat completion", ProviderUnavailableError, "HTTP 429"),
+ (404, "chat completion", ProviderModelNotFoundError, "model ID"),
+ (408, "chat completion", ProviderTimeoutError, "HTTP 408"),
+ (429, "chat completion", ProviderRateLimitError, "rate limit"),
(500, "chat completion", ProviderUnavailableError, "HTTP 500"),
(422, "chat completion", ProviderRequestError, "request"),
],
@@ -814,6 +820,18 @@ def test_http_status_classification(
)
+def test_http_quota_is_terminal_and_distinct_from_rate_limit() -> None:
+ response = OpenAICompatibleHTTPResponse(
+ status=429,
+ body=b'{"error":{"code":"insufficient_quota","message":"limit"}}',
+ )
+
+ with pytest.raises(ProviderQuotaError):
+ openai_module._raise_for_status(
+ response, operation="chat completion", model_id="exact/model"
+ )
+
+
@pytest.mark.parametrize(
"body",
[
@@ -884,10 +902,10 @@ async def exercise() -> None:
transport=auth_transport,
environment={"LM_STUDIO_API_KEY": secret},
)
- with pytest.raises(ProviderRequestError) as captured:
+ with pytest.raises(ProviderAuthorizationError) as captured:
await authenticated.list_models()
assert secret not in str(captured.value)
- assert "[REDACTED]" in str(captured.value)
+ assert "HTTP 403" in str(captured.value)
asyncio.run(exercise())
From 2e0088ec8482742bdf982240927c2c8151dac6a1 Mon Sep 17 00:00:00 2001
From: Kirill <106469980+waterflane@users.noreply.github.com>
Date: Mon, 7 Sep 2026 20:17:24 +0300
Subject: [PATCH 02/11] feat(index): add bounded failure policies and JSONL
progress
---
src/contextforge/application.py | 79 +++++--
src/contextforge/cli/intelligence_commands.py | 42 +++-
src/contextforge/cli/progress.py | 10 +-
src/contextforge/intelligence/__init__.py | 6 +
src/contextforge/intelligence/indexer.py | 11 +
src/contextforge/intelligence/models.py | 10 +
src/contextforge/intelligence/semantics.py | 192 +++++++++++++-----
tests/test_cli_intelligence.py | 95 ++++++++-
tests/test_semantic_analysis.py | 114 +++++++++++
9 files changed, 478 insertions(+), 81 deletions(-)
diff --git a/src/contextforge/application.py b/src/contextforge/application.py
index 6ab890b..6a3bc87 100644
--- a/src/contextforge/application.py
+++ b/src/contextforge/application.py
@@ -3,7 +3,6 @@
from __future__ import annotations
import asyncio
-import hashlib
import json
import uuid
from contextlib import suppress
@@ -76,6 +75,7 @@
load_file_semantic_analysis,
load_manifest,
load_repository_overview,
+ normalize_analyzer_identity,
write_manifest,
)
from contextforge.intelligence.models import AnalyzerIdentity
@@ -107,6 +107,10 @@ class MissingIndexError(ApplicationError):
"""Raised when an operation explicitly requires an active index."""
+class IndexSourceChangedError(ApplicationError):
+ """Raised when repository identity changes before atomic publication."""
+
+
class ArtifactReadError(ApplicationError):
"""Raised when a portable handoff cannot be read or validated."""
@@ -197,6 +201,8 @@ async def build_repository_index(
update_only: bool = False,
concurrency: int = 2,
fail_on_error: bool = False,
+ fail_fast: bool = False,
+ max_failures: int | None = None,
force_reanalyze: bool = False,
max_files: int | None = None,
semantic_max_output_tokens: int = 1024,
@@ -205,9 +211,17 @@ async def build_repository_index(
progress: ProgressObserver | None = None,
operation_id: str | None = None,
parent_operation_id: str | None = None,
+ cancellation: asyncio.Event | None = None,
) -> IndexBuildReport:
"""Build/update all index phases while retaining a prior pointer on failure."""
+ if fail_fast and max_failures is not None:
+ raise ValueError("fail_fast and max_failures cannot be used together")
+ if max_failures is not None and (
+ type(max_failures) is not int or max_failures <= 0
+ ):
+ raise ValueError("max_failures must be a positive integer or None")
+ effective_max_failures = 1 if fail_fast else max_failures
reporter = _progress_reporter(
"repository.index.update" if update_only else "repository.index.build",
progress,
@@ -244,12 +258,14 @@ async def build_repository_index(
update_only=update_only,
concurrency=concurrency,
fail_on_error=fail_on_error,
+ max_failures=effective_max_failures,
force_reanalyze=force_reanalyze,
max_files=max_files,
semantic_max_output_tokens=semantic_max_output_tokens,
recover_stale_lock=recover_stale_lock,
confirm_unknown_lock=confirm_unknown_lock,
progress=reporter,
+ cancellation=cancellation,
)
except BaseException as exc:
_report_terminal_exception(reporter, exc)
@@ -264,7 +280,12 @@ async def build_repository_index(
raise
reporter.complete(
message="Repository index build completed.",
- metadata={"partial": report.partial},
+ metadata={
+ "generation_id": report.manifest.generation_id,
+ "snapshot_digest": report.manifest.build.source_snapshot_digest,
+ "index_schema": report.manifest.schema_versions.index_schema_version,
+ "partial": report.partial,
+ },
)
_persist_application_diagnostic(
repository_root,
@@ -285,16 +306,19 @@ async def _build_repository_index(
update_only: bool,
concurrency: int,
fail_on_error: bool,
+ max_failures: int | None,
force_reanalyze: bool,
max_files: int | None,
semantic_max_output_tokens: int,
recover_stale_lock: bool,
confirm_unknown_lock: bool,
progress: ProgressReporter,
+ cancellation: asyncio.Event | None,
) -> IndexBuildReport:
"""Implement index construction under the public progress boundary."""
root = Path(repository_root).expanduser().resolve(strict=True)
+ _raise_if_index_cancelled(cancellation)
initialize_index(root)
previous: IndexManifest | None
try:
@@ -318,7 +342,8 @@ async def _build_repository_index(
phase_weight=scan_end,
activity=ProgressActivity.ACTIVE,
)
- snapshot = scan_repository(root)
+ snapshot = await asyncio.to_thread(scan_repository, root)
+ _raise_if_index_cancelled(cancellation)
progress.report(
"scan",
"Repository scan completed.",
@@ -370,10 +395,12 @@ async def _build_repository_index(
unit_type="files",
activity=ProgressActivity.ACTIVE,
)
- structural = build_structural_index(
+ structural = await asyncio.to_thread(
+ build_structural_index,
snapshot,
lock,
previous_manifest=previous,
+ cancellation=cancellation,
)
progress.report(
"structural_index",
@@ -454,11 +481,13 @@ def observe_semantic(event: ProgressEvent) -> None:
max_files=max_files,
max_output_tokens=semantic_max_output_tokens,
fail_on_error=fail_on_error,
+ max_failures=max_failures,
force_reanalyze=force_reanalyze,
resume=not force_reanalyze,
progress=observe_semantic,
),
previous_manifest=previous,
+ cancellation=cancellation,
)
semantic_event = progress.last_event
if semantic_event is not None:
@@ -529,6 +558,7 @@ def observe_semantic(event: ProgressEvent) -> None:
maps_start, 93.0, phase_prefix="repository_maps"
),
),
+ cancellation=cancellation,
)
map_fallback = any(item.status == "fallback" for item in maps.outcomes)
progress.report(
@@ -575,6 +605,14 @@ def observe_semantic(event: ProgressEvent) -> None:
phase_percent=100,
phase_weight=3 if model_enabled else 15,
)
+ _raise_if_index_cancelled(cancellation)
+ current_snapshot = await asyncio.to_thread(scan_repository, root)
+ if calculate_source_snapshot_digest(
+ current_snapshot
+ ) != calculate_source_snapshot_digest(snapshot):
+ raise IndexSourceChangedError(
+ "repository source identity changed before index publication"
+ )
progress.report(
"validation",
"Validating the active index generation.",
@@ -764,7 +802,8 @@ def _inspect_repository_index(
or analysis.schema_version != SEMANTIC_SCHEMA_VERSION
or (
analysis.record_kind != "deterministic_metadata_interpretation"
- and analysis.semantic_analyzer != expected
+ and normalize_analyzer_identity(analysis.semantic_analyzer)
+ != expected
)
):
stale.add(path)
@@ -1171,13 +1210,8 @@ def _semantic_identity(
return None
return AnalyzerIdentity(
analyzer_id=(GENERIC_SEMANTIC_ANALYZER_ID if generic else SEMANTIC_ANALYZER_ID),
- analyzer_version=_model_dependent_analyzer_version(
- (
- GENERIC_SEMANTIC_ANALYZER_VERSION
- if generic
- else SEMANTIC_ANALYZER_VERSION
- ),
- configuration,
+ analyzer_version=(
+ GENERIC_SEMANTIC_ANALYZER_VERSION if generic else SEMANTIC_ANALYZER_VERSION
),
analysis_prompt_version=SEMANTIC_PROMPT_VERSION,
response_schema_version=SEMANTIC_SCHEMA_VERSION,
@@ -1188,16 +1222,6 @@ def _semantic_identity(
)
-def _model_dependent_analyzer_version(
- analyzer_version: str, configuration: ProviderConfiguration
-) -> str:
- if configuration.provider_id != "openai-compatible":
- return analyzer_version
- canonical = configuration.endpoint.rstrip("/")
- digest = hashlib.sha256(canonical.encode("utf-8")).hexdigest()
- return f"{analyzer_version}+base.{digest}"
-
-
def _manifest_model_identity(
manifest: IndexManifest,
configuration: ProviderConfiguration | None,
@@ -1294,6 +1318,11 @@ def _workflow_snapshot(
return cast(ProjectSnapshot, source)
+def _raise_if_index_cancelled(cancellation: asyncio.Event | None) -> None:
+ if cancellation is not None and cancellation.is_set():
+ raise asyncio.CancelledError
+
+
def _report_terminal_exception(
reporter: ProgressReporter, error: BaseException
) -> None:
@@ -1358,6 +1387,11 @@ def _diagnostic_error_code(error: BaseException | None) -> str | None:
return None
if isinstance(error, ModelProviderError):
return provider_error_details(error)[0]
+ typed_code = getattr(error, "error_code", None)
+ if isinstance(typed_code, str):
+ return typed_code
+ if isinstance(error, IndexSourceChangedError):
+ return "source_identity_changed"
run_record = getattr(error, "run_record", None)
value = getattr(run_record, "failure_code", None)
if isinstance(value, str):
@@ -1423,6 +1457,7 @@ def _reject_json_constant(value: str) -> None:
"ApplicationError",
"ArtifactReadError",
"IndexBuildReport",
+ "IndexSourceChangedError",
"IndexStatusReport",
"MAX_HANDOFF_BYTES",
"MissingIndexError",
diff --git a/src/contextforge/cli/intelligence_commands.py b/src/contextforge/cli/intelligence_commands.py
index 91420a5..2c02c52 100644
--- a/src/contextforge/cli/intelligence_commands.py
+++ b/src/contextforge/cli/intelligence_commands.py
@@ -3,6 +3,7 @@
from __future__ import annotations
import asyncio
+import sys
from contextlib import suppress
from enum import StrEnum
from pathlib import Path
@@ -60,6 +61,8 @@ def _index_operation(
json_repair_attempts: int | None,
max_output_tokens: int | None,
fail_on_error: bool,
+ fail_fast: bool,
+ max_failures: int | None,
force_reanalyze: bool,
max_files: int | None,
local_only: bool,
@@ -67,8 +70,15 @@ def _index_operation(
confirm_unknown_lock: bool,
progress_mode: ProgressMode,
) -> None:
+ if fail_fast and max_failures is not None:
+ _exit_with_error(
+ "--fail-fast and --max-failures cannot be used together", code=2
+ )
provider: ModelProvider | None = None
- progress = CLIProgressRenderer(progress_mode)
+ progress = CLIProgressRenderer(
+ progress_mode,
+ stream=sys.stdout if progress_mode is ProgressMode.JSONL else None,
+ )
try:
project = load_project_configuration(path, config_path=config)
provider_configuration = resolve_provider_configuration(
@@ -100,6 +110,8 @@ def _index_operation(
update_only=update_only,
concurrency=effective_concurrency,
fail_on_error=fail_on_error,
+ fail_fast=fail_fast,
+ max_failures=max_failures,
force_reanalyze=force_reanalyze,
max_files=max_files,
semantic_max_output_tokens=(
@@ -132,7 +144,8 @@ def _index_operation(
with suppress(ModelProviderError):
asyncio.run(provider.close())
- typer.echo(_render_build_summary(report), nl=False)
+ if progress_mode is not ProgressMode.JSONL:
+ typer.echo(_render_build_summary(report), nl=False)
@index_app.command("build")
@@ -205,6 +218,21 @@ def build_index(
help="Keep the prior active generation on any model-analysis failure.",
),
] = False,
+ fail_fast: Annotated[
+ bool,
+ typer.Option(
+ "--fail-fast",
+ help="Stop after the first model-analysis failure.",
+ ),
+ ] = False,
+ max_failures: Annotated[
+ int | None,
+ typer.Option(
+ "--max-failures",
+ min=1,
+ help="Stop after this many model-analysis failures.",
+ ),
+ ] = None,
force_reanalyze: Annotated[
bool,
typer.Option(
@@ -242,7 +270,7 @@ def build_index(
ProgressMode,
typer.Option(
"--progress",
- help="Progress rendering: auto, always when safe, or never.",
+ help="Progress rendering: auto, always, never, or JSONL on stdout.",
case_sensitive=False,
),
] = ProgressMode.AUTO,
@@ -262,6 +290,8 @@ def build_index(
json_repair_attempts=json_repair_attempts,
max_output_tokens=max_output_tokens,
fail_on_error=fail_on_error,
+ fail_fast=fail_fast,
+ max_failures=max_failures,
force_reanalyze=force_reanalyze,
max_files=max_files,
local_only=local_only,
@@ -310,6 +340,8 @@ def update_index(
typer.Option("--max-output-tokens", min=96, max=32_768),
] = None,
fail_on_error: Annotated[bool, typer.Option("--fail-on-error")] = False,
+ fail_fast: Annotated[bool, typer.Option("--fail-fast")] = False,
+ max_failures: Annotated[int | None, typer.Option("--max-failures", min=1)] = None,
force_reanalyze: Annotated[bool, typer.Option("--force-reanalyze")] = False,
max_files: Annotated[int | None, typer.Option("--max-files", min=1)] = None,
local_only: Annotated[bool, typer.Option("--local-only")] = False,
@@ -321,7 +353,7 @@ def update_index(
ProgressMode,
typer.Option(
"--progress",
- help="Progress rendering: auto, always when safe, or never.",
+ help="Progress rendering: auto, always, never, or JSONL on stdout.",
case_sensitive=False,
),
] = ProgressMode.AUTO,
@@ -341,6 +373,8 @@ def update_index(
json_repair_attempts=json_repair_attempts,
max_output_tokens=max_output_tokens,
fail_on_error=fail_on_error,
+ fail_fast=fail_fast,
+ max_failures=max_failures,
force_reanalyze=force_reanalyze,
max_files=max_files,
local_only=local_only,
diff --git a/src/contextforge/cli/progress.py b/src/contextforge/cli/progress.py
index ee7e647..3b53466 100644
--- a/src/contextforge/cli/progress.py
+++ b/src/contextforge/cli/progress.py
@@ -39,6 +39,7 @@ class ProgressMode(StrEnum):
AUTO = "auto"
ALWAYS = "always"
NEVER = "never"
+ JSONL = "jsonl"
class CLIProgressRenderer:
@@ -84,7 +85,7 @@ def __init__(
self._unicode = self._supports_unicode(self._console.encoding)
self._spinner = Spinner("dots" if self._unicode else "line", style="cyan")
self._dynamic = (
- self.mode is not ProgressMode.NEVER
+ self.mode not in {ProgressMode.NEVER, ProgressMode.JSONL}
and self._console.is_terminal
and self._is_interactive(self._stdout)
and self._is_interactive(self._stream)
@@ -124,6 +125,8 @@ def rendering_mode(self) -> str:
if self.mode is ProgressMode.NEVER:
return "disabled"
+ if self.mode is ProgressMode.JSONL:
+ return "jsonl"
return "dynamic" if self._dynamic else "discrete"
def __call__(self, event: ProgressEvent) -> None:
@@ -132,6 +135,11 @@ def __call__(self, event: ProgressEvent) -> None:
with self._state_lock:
if self.mode is ProgressMode.NEVER or self._closed:
return
+ if self.mode is ProgressMode.JSONL:
+ self._stream.write(event.model_dump_json() + "\n")
+ self._stream.flush()
+ self._event = event
+ return
if self._started is None:
self._started = self._clock()
request_key = (event.current_item, event.current_attempt)
diff --git a/src/contextforge/intelligence/__init__.py b/src/contextforge/intelligence/__init__.py
index 6679c01..2007a40 100644
--- a/src/contextforge/intelligence/__init__.py
+++ b/src/contextforge/intelligence/__init__.py
@@ -93,6 +93,7 @@
SchemaVersionMetadata,
SemanticStatus,
calculate_index_statistics,
+ normalize_analyzer_identity,
validate_portable_relative_path,
)
from contextforge.intelligence.polyglot import (
@@ -127,8 +128,10 @@
SEMANTIC_SYSTEM_INSTRUCTIONS,
SemanticAnalysisError,
SemanticAnalysisOptions,
+ SemanticFailureLimitError,
SemanticFileOutcome,
SemanticIndexBuildResult,
+ SemanticProviderCircuitError,
SemanticRoute,
SemanticWorkPlan,
SemanticWorkPlanItem,
@@ -247,9 +250,11 @@
"SchemaVersionMetadata",
"SemanticAnalysisError",
"SemanticAnalysisOptions",
+ "SemanticFailureLimitError",
"SemanticConfidence",
"SemanticFileOutcome",
"SemanticIndexBuildResult",
+ "SemanticProviderCircuitError",
"SemanticRoute",
"SemanticWorkPlan",
"SemanticWorkPlanItem",
@@ -276,6 +281,7 @@
"build_repository_overview",
"calculate_generation_id",
"calculate_index_statistics",
+ "normalize_analyzer_identity",
"calculate_source_snapshot_digest",
"canonical_json_bytes",
"clean_generated_index",
diff --git a/src/contextforge/intelligence/indexer.py b/src/contextforge/intelligence/indexer.py
index ecf58b8..8d862f7 100644
--- a/src/contextforge/intelligence/indexer.py
+++ b/src/contextforge/intelligence/indexer.py
@@ -2,6 +2,7 @@
from __future__ import annotations
+import asyncio
import hashlib
from dataclasses import dataclass
from pathlib import Path
@@ -67,6 +68,7 @@ def build_structural_index(
*,
max_source_bytes: int = DEFAULT_CODEMAP_SOURCE_LIMIT,
previous_manifest: IndexManifest | None = None,
+ cancellation: asyncio.Event | None = None,
) -> StructuralIndexBuildResult:
"""Extract, resolve, and atomically persist facts without semantic analysis."""
@@ -85,6 +87,7 @@ def build_structural_index(
reused: list[str] = []
all_records_valid = previous is not None
for project_file in sorted(snapshot.files, key=lambda item: item.path):
+ _raise_if_cancelled(cancellation)
state = previous_states.get(project_file.path)
code_map = _reuse_code_map(lock, previous, state, project_file)
if code_map is None:
@@ -122,10 +125,12 @@ def build_structural_index(
generation_path=generation,
)
+ _raise_if_cancelled(cancellation)
code_maps = resolve_relationships(tuple(base_maps))
states: list[IndexedFileState] = []
record_digests: list[tuple[str, str]] = []
for code_map in code_maps:
+ _raise_if_cancelled(cancellation)
content = serialize_code_map(code_map)
location = _record_location(code_map.path)
digest = write_index_record(lock, location, content)
@@ -190,6 +195,7 @@ def build_structural_index(
previous.generation_id if previous is not None else None
),
)
+ _raise_if_cancelled(cancellation)
manifest = build_index_manifest(
build=build,
files=states,
@@ -310,6 +316,11 @@ def _validate_build_inputs(
raise ValueError("snapshot root does not match the locked repository")
+def _raise_if_cancelled(cancellation: asyncio.Event | None) -> None:
+ if cancellation is not None and cancellation.is_set():
+ raise asyncio.CancelledError
+
+
def _optional_manifest(lock: IndexWriteLock) -> IndexManifest | None:
try:
return load_manifest(lock.layout.repository_root)
diff --git a/src/contextforge/intelligence/models.py b/src/contextforge/intelligence/models.py
index f758b9e..03a90ab 100644
--- a/src/contextforge/intelligence/models.py
+++ b/src/contextforge/intelligence/models.py
@@ -34,6 +34,7 @@
]
_IDENTIFIER = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:/+@-]{0,127}$")
+_LEGACY_ENDPOINT_SUFFIX = re.compile(r"\+base\.[0-9a-f]{64}$")
class IndexModel(BaseModel):
@@ -339,3 +340,12 @@ def analyzer_identity_key(identity: AnalyzerIdentity) -> tuple[str, ...]:
model.provider_id if model is not None else "",
model.model_id if model is not None else "",
)
+
+
+def normalize_analyzer_identity(identity: AnalyzerIdentity) -> AnalyzerIdentity:
+ """Remove the legacy transport-endpoint suffix from analyzer provenance."""
+
+ version = _LEGACY_ENDPOINT_SUFFIX.sub("", identity.analyzer_version)
+ if version == identity.analyzer_version:
+ return identity
+ return identity.model_copy(update={"analyzer_version": version})
diff --git a/src/contextforge/intelligence/semantics.py b/src/contextforge/intelligence/semantics.py
index 43534b3..fd83d51 100644
--- a/src/contextforge/intelligence/semantics.py
+++ b/src/contextforge/intelligence/semantics.py
@@ -39,6 +39,7 @@
ModelIdentity,
SemanticStatus,
analyzer_identity_key,
+ normalize_analyzer_identity,
)
from contextforge.intelligence.semantic_models import (
SEMANTIC_SCHEMA_VERSION,
@@ -242,6 +243,28 @@ class SemanticAnalysisError(RuntimeError):
"""Raised when semantic analysis cannot safely publish the requested result."""
+class SemanticFailureLimitError(SemanticAnalysisError):
+ """Raised after the configured number of semantic units fail."""
+
+ def __init__(self, failure_count: int, diagnostic: AnalysisDiagnostic) -> None:
+ self.failure_count = failure_count
+ self.error_code = diagnostic.code
+ self.safe_reason = diagnostic.message
+ super().__init__(
+ f"semantic failure limit reached after {failure_count} file(s): "
+ f"{diagnostic.message}"
+ )
+
+
+class SemanticProviderCircuitError(SemanticAnalysisError):
+ """Raised when a provider-wide or repeated transient failure opens the circuit."""
+
+ def __init__(self, error_code: str, safe_reason: str) -> None:
+ self.error_code = error_code
+ self.safe_reason = safe_reason
+ super().__init__(f"provider circuit opened: {safe_reason}")
+
+
class StaleStructuralIndexError(SemanticAnalysisError):
"""Raised when semantics are requested without current deterministic facts."""
@@ -274,6 +297,7 @@ class SemanticAnalysisOptions:
max_chunks_per_file: int = 64
max_requests_per_file: int = 64
max_files: int | None = None
+ max_failures: int | None = None
fail_on_error: bool = False
resume: bool = True
force_reanalyze: bool = False
@@ -307,6 +331,10 @@ def __post_init__(self) -> None:
type(self.max_files) is not int or self.max_files <= 0
):
raise ValueError("max_files must be a positive integer or None")
+ if self.max_failures is not None and (
+ type(self.max_failures) is not int or self.max_failures <= 0
+ ):
+ raise ValueError("max_failures must be a positive integer or None")
if self.max_source_bytes_per_request > self.max_request_bytes:
raise ValueError("source byte limit cannot exceed request byte limit")
if (
@@ -454,11 +482,27 @@ def publish(self) -> None:
self._lifecycle = "published"
self.reporter.complete(message="Semantic generation published atomically.")
- def abort(self, *, cancelled: bool = False) -> None:
+ def abort(
+ self,
+ *,
+ cancelled: bool = False,
+ cancelled_units: int = 0,
+ unstarted_units: int = 0,
+ ) -> None:
+ metadata: dict[str, JsonValue] = {
+ "route_totals": cast(dict[str, JsonValue], self.plan.route_totals),
+ "cancelled_units": cancelled_units,
+ "unstarted_units": unstarted_units,
+ }
if cancelled:
- self.reporter.cancel(message="Semantic analysis cancelled.")
+ self.reporter.cancel(
+ message="Semantic analysis cancelled.", metadata=metadata
+ )
else:
- self.reporter.fail(message="Semantic analysis failed before publication.")
+ self.reporter.fail(
+ message="Semantic analysis failed before publication.",
+ metadata=metadata,
+ )
def _emit(self, message: str, *, current_item: str | None = None) -> None:
total_weight = sum(self._weights.values())
@@ -595,19 +639,17 @@ async def build_semantic_index(
load_file_code_map(snapshot.root, item.path, manifest=structural)
for item in structural.files
)
- provider_id, model_id, base_url_sha256 = _provider_identity(provider)
+ provider_id, model_id = _provider_identity(provider)
rich_analyzer = _semantic_analyzer(
active_options,
provider_id,
model_id,
- base_url_sha256,
analysis_route="rich_model_analysis",
)
generic_analyzer = _semantic_analyzer(
active_options,
provider_id,
model_id,
- base_url_sha256,
analysis_route="generic_model_analysis",
)
options_digest = _analysis_options_digest(active_options)
@@ -616,6 +658,10 @@ async def build_semantic_index(
reusable_manifests = _reuse_manifests(
snapshot.root, structural, previous_manifest=previous_manifest
)
+ identity_migration_needed = any(
+ normalize_analyzer_identity(item) != item
+ for item in structural.semantic_analyzers
+ )
analyses: dict[str, FileSemanticAnalysis] = {}
outcomes: dict[str, SemanticFileOutcome] = {}
@@ -818,7 +864,11 @@ async def build_semantic_index(
plan_item.path, code=diagnostic.code, message=diagnostic.message
)
- if not selected_stale and _manifest_matches_planned_semantics(structural, analyses):
+ if (
+ not selected_stale
+ and not identity_migration_needed
+ and _manifest_matches_planned_semantics(structural, analyses)
+ ):
tracker.publish()
return SemanticIndexBuildResult(
manifest=structural,
@@ -943,41 +993,92 @@ async def analyze_one(
},
)
_emit_status(active_options, project_file.path, "failed")
+ if isinstance(exc, ModelProviderError) and exc.circuit_opened:
+ raise SemanticProviderCircuitError(code, message) from exc
return project_file.path, None, diagnostic, False
task_results: list[
tuple[str, _AnalysisWork | None, AnalysisDiagnostic | None, bool]
] = []
- for offset in range(0, len(selected_stale), active_options.max_concurrency):
- batch = selected_stale[offset : offset + active_options.max_concurrency]
- tasks = [
- asyncio.create_task(
- analyze_one(
- project_file,
- code_map,
- state,
- route,
- expected_analyzer,
- expected_digest,
- )
- )
- for (
+ failure_count = 0
+ pending: set[
+ asyncio.Task[tuple[str, _AnalysisWork | None, AnalysisDiagnostic | None, bool]]
+ ] = set()
+ next_work = 0
+
+ def schedule_available() -> None:
+ nonlocal next_work
+ while len(pending) < active_options.max_concurrency and next_work < len(
+ selected_stale
+ ):
+ (
project_file,
code_map,
state,
route,
expected_analyzer,
expected_digest,
- ) in batch
- ]
- try:
- task_results.extend(await asyncio.gather(*tasks))
- except BaseException:
- for task in tasks:
- task.cancel()
- await asyncio.gather(*tasks, return_exceptions=True)
- tracker.abort(cancelled=True)
- raise
+ ) = selected_stale[next_work]
+ next_work += 1
+ pending.add(
+ asyncio.create_task(
+ analyze_one(
+ project_file,
+ code_map,
+ state,
+ route,
+ expected_analyzer,
+ expected_digest,
+ )
+ )
+ )
+
+ schedule_available()
+ try:
+ while pending:
+ done, pending = await asyncio.wait(
+ pending, return_when=asyncio.FIRST_COMPLETED
+ )
+ limit_diagnostic: AnalysisDiagnostic | None = None
+ for task in done:
+ result = task.result()
+ task_results.append(result)
+ _, work, diagnostic, _ = result
+ if work is not None:
+ continue
+ assert diagnostic is not None
+ failure_count += 1
+ if (
+ active_options.max_failures is not None
+ and failure_count >= active_options.max_failures
+ and limit_diagnostic is None
+ ):
+ limit_diagnostic = diagnostic
+ if limit_diagnostic is not None:
+ cancelled_units = len(pending)
+ for unfinished in pending:
+ unfinished.cancel()
+ await asyncio.gather(*pending, return_exceptions=True)
+ tracker.abort(
+ cancelled_units=cancelled_units,
+ unstarted_units=len(selected_stale) - next_work,
+ )
+ raise SemanticFailureLimitError(failure_count, limit_diagnostic)
+ schedule_available()
+ except BaseException as exc:
+ cancelled_units = len(pending)
+ for task in pending:
+ task.cancel()
+ await asyncio.gather(*pending, return_exceptions=True)
+ if not isinstance(exc, SemanticFailureLimitError):
+ tracker.abort(
+ cancelled=isinstance(
+ exc, (asyncio.CancelledError, ProviderCancelledError)
+ ),
+ cancelled_units=cancelled_units,
+ unstarted_units=len(selected_stale) - next_work,
+ )
+ raise
failures: list[AnalysisDiagnostic] = []
for path, work, diagnostic, resumed in task_results:
@@ -2527,7 +2628,6 @@ def _semantic_analyzer(
options: SemanticAnalysisOptions,
provider_id: str,
model_id: str,
- base_url_sha256: str | None,
*,
analysis_route: Literal["rich_model_analysis", "generic_model_analysis"],
) -> AnalyzerIdentity:
@@ -2543,7 +2643,7 @@ def _semantic_analyzer(
)
return AnalyzerIdentity(
analyzer_id=analyzer_id,
- analyzer_version=_connection_bound_version(analyzer_version, base_url_sha256),
+ analyzer_version=analyzer_version,
analysis_prompt_version=options.prompt_version,
response_schema_version=SEMANTIC_SCHEMA_VERSION,
model_identity=ModelIdentity(
@@ -2553,7 +2653,7 @@ def _semantic_analyzer(
)
-def _provider_identity(provider: ModelProvider) -> tuple[str, str, str | None]:
+def _provider_identity(provider: ModelProvider) -> tuple[str, str]:
provider_id = provider.provider_id
configuration = getattr(provider, "configuration", None)
model_id = getattr(configuration, "model_id", None)
@@ -2561,23 +2661,7 @@ def _provider_identity(provider: ModelProvider) -> tuple[str, str, str | None]:
raise SemanticAnalysisError(
"semantic provider must expose stable provider and model identity"
)
- endpoint = getattr(configuration, "endpoint", None)
- base_url_sha256 = None
- if provider_id == "openai-compatible":
- if not isinstance(endpoint, str):
- raise SemanticAnalysisError(
- "OpenAI-compatible provider must expose a stable base URL identity"
- )
- base_url_sha256 = hashlib.sha256(
- endpoint.rstrip("/").encode("utf-8")
- ).hexdigest()
- return provider_id, model_id, base_url_sha256
-
-
-def _connection_bound_version(version: str, base_url_sha256: str | None) -> str:
- if base_url_sha256 is None:
- return version
- return f"{version}+base.{base_url_sha256}"
+ return provider_id, model_id
def _validate_response_identity(
@@ -2630,7 +2714,7 @@ def _analysis_matches(
and analysis.language == state.language == code_map.language
and analysis.fact_record_sha256 == _required_fact_digest(state)
and analysis.codemap_analyzer == code_map.analyzer
- and analysis.semantic_analyzer == analyzer
+ and normalize_analyzer_identity(analysis.semantic_analyzer) == analyzer
and analysis.analysis_options_digest == options_digest
)
@@ -2663,7 +2747,7 @@ def _find_reusable_analysis(
if analysis.schema_version == SEMANTIC_SCHEMA_VERSION and _analysis_matches(
analysis, state, code_map, analyzer, options_digest
):
- return analysis
+ return analysis.model_copy(update={"semantic_analyzer": analyzer})
return None
@@ -2796,9 +2880,11 @@ def _bounded_error_message(error: BaseException) -> str:
"SEMANTIC_PROMPT_VERSION",
"SEMANTIC_SYSTEM_INSTRUCTIONS",
"SemanticAnalysisError",
+ "SemanticFailureLimitError",
"SemanticAnalysisOptions",
"SemanticFileOutcome",
"SemanticIndexBuildResult",
+ "SemanticProviderCircuitError",
"SemanticRoute",
"SemanticWorkPlan",
"SemanticWorkPlanItem",
diff --git a/tests/test_cli_intelligence.py b/tests/test_cli_intelligence.py
index 47b54c2..427fd70 100644
--- a/tests/test_cli_intelligence.py
+++ b/tests/test_cli_intelligence.py
@@ -1,3 +1,4 @@
+import asyncio
import json
from pathlib import Path
from typing import Any, cast
@@ -6,9 +7,14 @@
import pytest
from typer.testing import CliRunner, Result
+import contextforge.application as application_module
import contextforge.cli.context_commands as context_cli
import contextforge.cli.intelligence_commands as index_cli
-from contextforge.application import render_context_suggestion
+from contextforge.application import (
+ IndexSourceChangedError,
+ build_repository_index,
+ render_context_suggestion,
+)
from contextforge.cli.main import app
from contextforge.discovery import (
CompletenessWarning,
@@ -22,6 +28,7 @@
from contextforge.discovery.renderers import DiscoveryResultFormat
from contextforge.intelligence import IndexManifestNotFoundError, load_manifest
from contextforge.models import FakeModelProvider, ProviderConfiguration
+from contextforge.progress import ProgressEvent, ProgressStatus
runner = CliRunner()
TERMINAL_WIDTH = 140
@@ -230,6 +237,44 @@ def test_index_provider_failure_preserves_previous_active_generation(
assert load_manifest(tmp_path) == previous
+def test_index_rechecks_snapshot_before_atomic_publication(
+ tmp_path: Path, monkeypatch: pytest.MonkeyPatch
+) -> None:
+ source = tmp_path / "app.py"
+ _write(tmp_path, "app.py", "VALUE = 1\n")
+ asyncio.run(
+ build_repository_index(
+ tmp_path,
+ provider=None,
+ provider_configuration=None,
+ )
+ )
+ previous = load_manifest(tmp_path)
+ source.write_text("VALUE = 2\n", encoding="utf-8")
+ original = cast(Any, application_module).build_structural_index
+
+ def mutate_after_structural(*args: Any, **kwargs: Any) -> Any:
+ result = original(*args, **kwargs)
+ source.write_text("VALUE = 3\n", encoding="utf-8")
+ return result
+
+ monkeypatch.setattr(
+ application_module, "build_structural_index", mutate_after_structural
+ )
+
+ with pytest.raises(IndexSourceChangedError):
+ asyncio.run(
+ build_repository_index(
+ tmp_path,
+ provider=None,
+ provider_configuration=None,
+ update_only=True,
+ )
+ )
+
+ assert load_manifest(tmp_path) == previous
+
+
def test_index_cancellation_maps_to_130(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
@@ -281,6 +326,54 @@ def test_progress_never_suppresses_stderr_and_preserves_json_stdout(
assert json.loads(suggested.stdout)["mode"] == "hybrid"
+def test_index_jsonl_progress_is_a_clean_schema_three_stream(tmp_path: Path) -> None:
+ _write(tmp_path, "app.py", "VALUE = 1\n")
+
+ result = _invoke(
+ "--log-level",
+ "quiet",
+ "index",
+ "build",
+ str(tmp_path),
+ "--provider",
+ "none",
+ "--progress",
+ "jsonl",
+ )
+
+ assert result.exit_code == 0, result.output
+ assert result.stderr == ""
+ events = [
+ ProgressEvent.model_validate_json(line) for line in result.stdout.splitlines()
+ ]
+ assert events
+ assert all(event.schema_version == 3 for event in events)
+ assert [event.sequence for event in events] == list(range(len(events)))
+ assert events[-1].status is ProgressStatus.COMPLETED
+ assert events[-1].metadata["generation_id"] == load_manifest(tmp_path).generation_id
+ assert events[-1].metadata["snapshot_digest"]
+ assert events[-1].metadata["index_schema"] == 2
+ assert events[-1].metadata["partial"] is False
+ assert "\x1b[" not in result.stdout
+ assert "Status:" not in result.stdout
+
+
+def test_index_rejects_ambiguous_failure_policy(tmp_path: Path) -> None:
+ result = _invoke(
+ "index",
+ "build",
+ str(tmp_path),
+ "--provider",
+ "none",
+ "--fail-fast",
+ "--max-failures",
+ "2",
+ )
+
+ assert result.exit_code == 2
+ assert "cannot be used together" in _plain(result.stderr)
+
+
def _invoke_focused_suggestion(tmp_path: Path, *arguments: str) -> Result:
_write(tmp_path, "app.py", "def run():\n return 1\n")
return _invoke(
diff --git a/tests/test_semantic_analysis.py b/tests/test_semantic_analysis.py
index bf82411..9921c18 100644
--- a/tests/test_semantic_analysis.py
+++ b/tests/test_semantic_analysis.py
@@ -19,6 +19,7 @@
SemanticAnalysisError,
SemanticAnalysisOptions,
SemanticConfidence,
+ SemanticFailureLimitError,
SemanticIndexBuildResult,
SourceRange,
StaleStructuralIndexError,
@@ -954,6 +955,72 @@ def test_unchanged_analysis_is_reused_without_provider_calls(tmp_path: Path) ->
assert second.generation_path == first.generation_path
+def test_loopback_endpoint_change_does_not_invalidate_semantics(tmp_path: Path) -> None:
+ snapshot = _snapshot_with_facts(tmp_path, {"app.txt": "service notes\n"})
+ first_provider = FakeModelProvider(
+ ProviderConfiguration(
+ provider_id="fake",
+ endpoint="fake://127.0.0.1:1234",
+ model_id="exact/model",
+ retry_limit=0,
+ ),
+ responder=_valid_response,
+ )
+ first = _build_semantics(snapshot, first_provider, run_id="endpoint-first")
+ second_provider = FakeModelProvider(
+ ProviderConfiguration(
+ provider_id="fake",
+ endpoint="fake://127.0.0.1:9999",
+ model_id="exact/model",
+ retry_limit=0,
+ ),
+ responder=_valid_response,
+ )
+
+ second = _build_semantics(snapshot, second_provider, run_id="endpoint-second")
+
+ assert first_provider.call_count == 1
+ assert second_provider.call_count == 0
+ assert second.manifest == first.manifest
+ assert "+base." not in second.analyses[0].semantic_analyzer.analyzer_version
+
+
+def test_legacy_endpoint_identity_is_republished_without_model_calls(
+ tmp_path: Path, monkeypatch: pytest.MonkeyPatch
+) -> None:
+ snapshot = _snapshot_with_facts(tmp_path, {"app.txt": "service notes\n"})
+ original = semantics_module._semantic_analyzer
+
+ def legacy_analyzer(*args: Any, **kwargs: Any) -> Any:
+ identity = original(*args, **kwargs)
+ return identity.model_copy(
+ update={"analyzer_version": identity.analyzer_version + "+base." + "a" * 64}
+ )
+
+ monkeypatch.setattr(semantics_module, "_semantic_analyzer", legacy_analyzer)
+ legacy = _build_semantics(snapshot, _provider(), run_id="legacy-identity")
+ monkeypatch.setattr(semantics_module, "_semantic_analyzer", original)
+ migration_provider = _provider()
+
+ migrated = _build_semantics(snapshot, migration_provider, run_id="migrate-identity")
+ stable_provider = _provider()
+ stable = _build_semantics(snapshot, stable_provider, run_id="stable-identity")
+
+ assert migration_provider.call_count == stable_provider.call_count == 0
+ assert migrated.manifest.generation_id != legacy.manifest.generation_id
+ assert stable.manifest == migrated.manifest
+ assert all(
+ "+base." not in item.analyzer_version
+ for item in migrated.manifest.semantic_analyzers
+ )
+ assert (
+ "+base."
+ not in load_file_semantic_analysis(
+ tmp_path, "app.txt"
+ ).semantic_analyzer.analyzer_version
+ )
+
+
def test_semantic_persistence_is_deterministic_across_repository_roots(
tmp_path: Path,
) -> None:
@@ -1096,6 +1163,53 @@ def test_fail_on_error_keeps_prior_valid_generation_active(tmp_path: Path) -> No
assert load_manifest(tmp_path) == structural
+@pytest.mark.parametrize(
+ ("concurrency", "failure_limit", "expected_calls"),
+ [(1, 1, 1), (1, 2, 2), (2, 1, 2)],
+)
+def test_failure_limit_stops_scheduling_and_keeps_active_generation(
+ tmp_path: Path,
+ concurrency: int,
+ failure_limit: int,
+ expected_calls: int,
+) -> None:
+ snapshot = _snapshot_with_facts(
+ tmp_path, {f"{name}.txt": f"{name}\n" for name in "abcde"}
+ )
+ active = load_manifest(tmp_path)
+ events: list[ProgressEvent] = []
+ scripts: list[Any] = [
+ ProviderRequestError("rejected") for _ in range(failure_limit)
+ ]
+ if concurrency > 1:
+ scripts.insert(1, FakeScript(ProviderRequestError("late"), delay_seconds=1))
+ provider = _provider(concurrency=concurrency, scripts=scripts)
+
+ with (
+ acquire_index_lock(tmp_path, "semantic-failure-limit") as lock,
+ pytest.raises(SemanticFailureLimitError),
+ ):
+ asyncio.run(
+ build_semantic_index(
+ snapshot,
+ lock,
+ provider,
+ options=SemanticAnalysisOptions(
+ max_concurrency=concurrency,
+ max_failures=failure_limit,
+ progress=events.append,
+ ),
+ )
+ )
+
+ assert provider.call_count == expected_calls
+ assert load_manifest(tmp_path) == active
+ assert events[-1].status.value == "failed"
+ assert events[-1].failed_units == failure_limit
+ assert cast(int, events[-1].metadata["unstarted_units"]) > 0
+ assert cast(int, events[-1].metadata["cancelled_units"]) == concurrency - 1
+
+
def test_interrupted_build_resumes_only_validated_checkpoints(tmp_path: Path) -> None:
snapshot = _snapshot_with_facts(tmp_path, {"a.py": "pass\n", "b.py": "pass\n"})
cancellation = asyncio.Event()
From f65ab55bf3216de986531a92e4a47a0493c1ba2c Mon Sep 17 00:00:00 2001
From: Kirill <106469980+waterflane@users.noreply.github.com>
Date: Mon, 7 Sep 2026 20:17:43 +0300
Subject: [PATCH 03/11] feat(bridge): add tracked index jobs
---
.../contextforge-bridge-v2.schema.json | 242 ++++++++++++++
src/contextforge/bridge/models.py | 32 ++
src/contextforge/bridge/protocol.py | 5 +-
src/contextforge/bridge/server.py | 311 +++++++++++++++++-
tests/test_bridge.py | 271 ++++++++++++++-
5 files changed, 846 insertions(+), 15 deletions(-)
create mode 100644 docs/schemas/contextforge-bridge-v2.schema.json
diff --git a/docs/schemas/contextforge-bridge-v2.schema.json b/docs/schemas/contextforge-bridge-v2.schema.json
new file mode 100644
index 0000000..c501cb9
--- /dev/null
+++ b/docs/schemas/contextforge-bridge-v2.schema.json
@@ -0,0 +1,242 @@
+{
+ "$schema": "https://json-schema.org/draft/2020-12/schema",
+ "$id": "https://contextforge.dev/schemas/contextforge-bridge-v2.schema.json",
+ "title": "ContextForge trusted-local bridge protocol v2",
+ "description": "Bridge v1 frames plus opt-in tracked index jobs and correlated progress notifications.",
+ "oneOf": [
+ { "$ref": "contextforge-bridge-v1.schema.json" },
+ { "$ref": "#/$defs/helloV2Request" },
+ { "$ref": "#/$defs/helloV2Success" },
+ { "$ref": "#/$defs/indexRequest" },
+ { "$ref": "#/$defs/indexSuccess" },
+ { "$ref": "#/$defs/indexFailure" },
+ { "$ref": "#/$defs/progressNotification" }
+ ],
+ "$defs": {
+ "id": {
+ "oneOf": [
+ { "type": "string", "minLength": 1, "maxLength": 200 },
+ { "type": "integer" }
+ ]
+ },
+ "sha256": { "type": "string", "pattern": "^[0-9a-f]{64}$" },
+ "timeout": { "type": "integer", "minimum": 1, "maximum": 900000 },
+ "helloV2Request": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["jsonrpc", "id", "method", "params"],
+ "properties": {
+ "jsonrpc": { "const": "2.0" },
+ "id": { "$ref": "#/$defs/id" },
+ "method": { "const": "hello" },
+ "params": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["protocol_version"],
+ "properties": {
+ "protocol_version": { "const": "2.0" },
+ "client_name": { "type": "string", "minLength": 1, "maxLength": 200 },
+ "timeout_ms": { "$ref": "#/$defs/timeout" }
+ }
+ }
+ }
+ },
+ "helloV2Success": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["jsonrpc", "id", "result"],
+ "properties": {
+ "jsonrpc": { "const": "2.0" },
+ "id": { "$ref": "#/$defs/id" },
+ "result": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["protocol_version", "contextforge_version", "supported_protocol_versions", "capabilities", "workspace", "policy"],
+ "properties": {
+ "protocol_version": { "const": "2.0" },
+ "contextforge_version": { "type": "string" },
+ "supported_protocol_versions": { "type": "array", "items": { "type": "string" }, "uniqueItems": true },
+ "capabilities": { "type": "object" },
+ "workspace": { "type": "object" },
+ "policy": { "type": "object" }
+ }
+ }
+ }
+ },
+ "indexRequest": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["jsonrpc", "id", "method", "params"],
+ "properties": {
+ "jsonrpc": { "const": "2.0" },
+ "id": { "$ref": "#/$defs/id" },
+ "method": { "const": "index" },
+ "params": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["action", "expected_snapshot_digest"],
+ "properties": {
+ "action": { "enum": ["build", "update"] },
+ "expected_snapshot_digest": { "$ref": "#/$defs/sha256" },
+ "provider": { "type": "string", "minLength": 1, "maxLength": 128 },
+ "model": { "type": "string", "minLength": 1, "maxLength": 128 },
+ "base_url": { "type": "string", "minLength": 1, "maxLength": 2000 },
+ "concurrency": { "type": "integer", "minimum": 1, "maximum": 8 },
+ "request_timeout": { "type": "number", "minimum": 1, "maximum": 600 },
+ "context_window": { "type": "integer", "minimum": 1024, "maximum": 2000000 },
+ "json_repair_attempts": { "type": "integer", "minimum": 0, "maximum": 10 },
+ "max_output_tokens": { "type": "integer", "minimum": 96, "maximum": 32768 },
+ "fail_on_error": { "type": "boolean" },
+ "fail_fast": { "type": "boolean" },
+ "max_failures": { "type": "integer", "minimum": 1 },
+ "force_reanalyze": { "type": "boolean" },
+ "max_files": { "type": "integer", "minimum": 1 },
+ "local_only": { "type": "boolean" },
+ "recover_stale_lock": { "type": "boolean" },
+ "confirm_unknown_lock": { "type": "boolean" },
+ "timeout_ms": { "$ref": "#/$defs/timeout" }
+ }
+ }
+ }
+ },
+ "indexSuccess": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["jsonrpc", "id", "result"],
+ "properties": {
+ "jsonrpc": { "const": "2.0" },
+ "id": { "$ref": "#/$defs/id" },
+ "result": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["action", "generation_id", "snapshot_digest", "index_schema", "partial", "statistics"],
+ "properties": {
+ "action": { "enum": ["build", "update"] },
+ "generation_id": { "$ref": "#/$defs/sha256" },
+ "snapshot_digest": { "$ref": "#/$defs/sha256" },
+ "index_schema": { "type": "integer", "minimum": 1 },
+ "partial": { "type": "boolean" },
+ "statistics": { "type": "object" }
+ }
+ }
+ }
+ },
+ "indexFailure": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["jsonrpc", "id", "error"],
+ "properties": {
+ "jsonrpc": { "const": "2.0" },
+ "id": { "$ref": "#/$defs/id" },
+ "error": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["code", "message", "data"],
+ "properties": {
+ "code": { "enum": [-32001, -32008, -32009, -32010, -32011, -32012, -32013] },
+ "message": { "type": "string" },
+ "data": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["code", "error_code", "phase", "reason", "retryable", "operation_id"],
+ "properties": {
+ "code": {
+ "enum": [
+ "SOURCE_IDENTITY_CHANGED",
+ "INDEX_BUILD_FAILED",
+ "PROVIDER_FAILURE",
+ "PROVIDER_CONFIGURATION_ERROR",
+ "FAILURE_LIMIT_REACHED",
+ "PROVIDER_CIRCUIT_OPEN",
+ "INDEX_STORAGE_ERROR",
+ "INDEX_LOCKED"
+ ]
+ },
+ "error_code": { "type": "string", "minLength": 1, "maxLength": 128 },
+ "phase": { "type": "string", "minLength": 1, "maxLength": 128 },
+ "reason": { "type": "string", "minLength": 1, "maxLength": 1000 },
+ "retryable": { "type": "boolean" },
+ "operation_id": { "type": "string", "minLength": 1, "maxLength": 128 }
+ }
+ }
+ }
+ }
+ }
+ },
+ "progressNotification": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["jsonrpc", "method", "params"],
+ "properties": {
+ "jsonrpc": { "const": "2.0" },
+ "method": { "const": "$/progress" },
+ "params": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["request_id", "event"],
+ "properties": {
+ "request_id": { "$ref": "#/$defs/id" },
+ "event": {
+ "type": "object",
+ "additionalProperties": false,
+ "required": ["schema_version", "operation_id", "operation_type", "phase_id", "message", "percentage", "status", "sequence", "overall_percent", "phase_label", "phase_percent", "processed_units", "failed_units", "elapsed_seconds"],
+ "properties": {
+ "schema_version": { "const": 3 },
+ "operation_id": { "type": "string" },
+ "operation_type": { "type": "string" },
+ "phase_id": { "type": "string" },
+ "message": { "type": "string" },
+ "completed": { "type": "number" },
+ "total": { "type": ["number", "null"] },
+ "percentage": { "type": "number", "minimum": 0, "maximum": 100 },
+ "status": { "enum": ["running", "completed", "failed", "cancelled"] },
+ "top_level_operation_id": { "type": ["string", "null"] },
+ "parent_operation_id": { "type": ["string", "null"] },
+ "metadata": { "type": "object" },
+ "sequence": { "type": "integer", "minimum": 0 },
+ "indeterminate": { "type": "boolean" },
+ "overall_percent": { "type": "number" },
+ "phase_label": { "type": "string" },
+ "phase_percent": { "type": "number" },
+ "phase_weight": { "type": "number" },
+ "completed_units": { "type": "number" },
+ "total_units": { "type": ["number", "null"] },
+ "unit_type": { "type": "string" },
+ "current_item": { "type": ["string", "null"] },
+ "last_completed_item": { "type": ["string", "null"] },
+ "last_failed_item": { "type": ["string", "null"] },
+ "active_items": { "type": "array", "items": { "type": "string" } },
+ "active_item_count": { "type": "integer" },
+ "reused_units": { "type": "integer" },
+ "skipped_units": { "type": "integer" },
+ "failed_units": { "type": "integer" },
+ "elapsed_seconds": { "type": "number" },
+ "activity": { "enum": ["idle", "active", "waiting"] },
+ "planned_units": { "type": "integer" },
+ "processed_units": { "type": "integer" },
+ "succeeded_units": { "type": "integer" },
+ "fallback_units": { "type": "integer" },
+ "active_units": { "type": "integer" },
+ "current_attempt": { "type": ["integer", "null"] },
+ "max_attempts": { "type": ["integer", "null"] },
+ "lifecycle_state": { "type": "string" },
+ "safe_error_code": { "type": ["string", "null"] },
+ "safe_error_message": { "type": ["string", "null"] },
+ "request_elapsed_seconds": { "type": "number" },
+ "operation_elapsed_seconds": { "type": "number" },
+ "analyzer_kind": { "type": ["string", "null"] },
+ "estimated_input_tokens": { "type": ["integer", "null"] },
+ "output_token_budget": { "type": ["integer", "null"] },
+ "input_truncated": { "type": "boolean" },
+ "configured_context_window": { "type": ["integer", "null"] },
+ "schema_overhead_tokens": { "type": ["integer", "null"] },
+ "safety_margin_tokens": { "type": ["integer", "null"] },
+ "estimated_total_tokens": { "type": ["integer", "null"] }
+ }
+ }
+ }
+ }
+ }
+ }
+ }
+}
diff --git a/src/contextforge/bridge/models.py b/src/contextforge/bridge/models.py
index 1ffea60..9e114d0 100644
--- a/src/contextforge/bridge/models.py
+++ b/src/contextforge/bridge/models.py
@@ -38,6 +38,37 @@ class SnapshotParams(BridgeParams):
pass
+class IndexParams(BridgeParams):
+ """Bridge 2 parameters for one atomic tracked index job."""
+
+ action: Literal["build", "update"]
+ expected_snapshot_digest: Sha256
+ provider: str | None = Field(default=None, min_length=1, max_length=128)
+ model: str | None = Field(default=None, min_length=1, max_length=128)
+ base_url: str | None = Field(default=None, min_length=1, max_length=2_000)
+ concurrency: int | None = Field(default=None, ge=1, le=8, strict=True)
+ request_timeout: float | None = Field(default=None, ge=1, le=600)
+ context_window: int | None = Field(
+ default=None, ge=1_024, le=2_000_000, strict=True
+ )
+ json_repair_attempts: int | None = Field(default=None, ge=0, le=10, strict=True)
+ max_output_tokens: int | None = Field(default=None, ge=96, le=32_768, strict=True)
+ fail_on_error: bool = False
+ fail_fast: bool = False
+ max_failures: int | None = Field(default=None, ge=1, strict=True)
+ force_reanalyze: bool = False
+ max_files: int | None = Field(default=None, ge=1, strict=True)
+ local_only: bool = False
+ recover_stale_lock: bool = False
+ confirm_unknown_lock: bool = False
+
+ @model_validator(mode="after")
+ def validate_failure_policy(self) -> IndexParams:
+ if self.fail_fast and self.max_failures is not None:
+ raise ValueError("fail_fast and max_failures cannot be used together")
+ return self
+
+
class DiscoverParams(BridgeParams):
expected_snapshot_digest: Sha256
task: str = Field(min_length=1, max_length=20_000)
@@ -152,6 +183,7 @@ class ShutdownParams(BridgeParams):
"ExpandParams",
"ExpansionOperation",
"HelloParams",
+ "IndexParams",
"PackageParams",
"ReadParams",
"ShutdownParams",
diff --git a/src/contextforge/bridge/protocol.py b/src/contextforge/bridge/protocol.py
index b629e7b..f49fdfa 100644
--- a/src/contextforge/bridge/protocol.py
+++ b/src/contextforge/bridge/protocol.py
@@ -2,10 +2,11 @@
from typing import Final, Literal
-BRIDGE_PROTOCOL_VERSION: Final[Literal["1.1"]] = "1.1"
-SUPPORTED_BRIDGE_PROTOCOL_VERSIONS: Final[tuple[Literal["1.0", "1.1"], ...]] = (
+BRIDGE_PROTOCOL_VERSION: Final[Literal["2.0"]] = "2.0"
+SUPPORTED_BRIDGE_PROTOCOL_VERSIONS: Final[tuple[Literal["1.0", "1.1", "2.0"], ...]] = (
"1.0",
"1.1",
+ "2.0",
)
__all__ = ["BRIDGE_PROTOCOL_VERSION", "SUPPORTED_BRIDGE_PROTOCOL_VERSIONS"]
diff --git a/src/contextforge/bridge/server.py b/src/contextforge/bridge/server.py
index 9b9b0e2..c6554f3 100644
--- a/src/contextforge/bridge/server.py
+++ b/src/contextforge/bridge/server.py
@@ -9,6 +9,7 @@
import os
import unicodedata
from collections import OrderedDict
+from contextlib import suppress
from dataclasses import dataclass
from pathlib import Path
from typing import Any, BinaryIO, TextIO
@@ -16,7 +17,12 @@
from pydantic import BaseModel, ValidationError
from contextforge._metadata import __version__
-from contextforge.application import inspect_repository_index
+from contextforge.application import (
+ ApplicationError,
+ IndexSourceChangedError,
+ build_repository_index,
+ inspect_repository_index,
+)
from contextforge.context import (
ContextLimitError,
ContextReaderError,
@@ -45,14 +51,33 @@
)
from contextforge.discovery.tools import TOOL_INPUT_MODELS
from contextforge.intelligence import (
+ INDEX_SCHEMA_VERSION,
+ MANIFEST_SCHEMA_VERSION,
+ RECORD_SCHEMA_VERSION,
+ GlobalMapAnalysisError,
+ IndexLockError,
+ IndexStorageError,
+ SemanticAnalysisError,
+ SemanticFailureLimitError,
+ SemanticProviderCircuitError,
calculate_source_snapshot_digest,
canonical_json_bytes,
load_file_code_map,
load_file_semantic_analysis,
load_manifest,
)
+from contextforge.models import (
+ ModelProvider,
+ ModelProviderError,
+ ProviderConfigurationError,
+ RetryClassification,
+ classify_retry,
+ provider_error_details,
+)
+from contextforge.progress import PROGRESS_SCHEMA_VERSION, ProgressEvent
from contextforge.project_config import (
ProjectConfigError,
+ create_model_provider,
load_project_configuration,
resolve_provider_configuration,
)
@@ -65,6 +90,7 @@
ExpandParams,
ExpansionOperation,
HelloParams,
+ IndexParams,
PackageParams,
ReadParams,
ShutdownParams,
@@ -93,11 +119,18 @@
SHUTTING_DOWN = -32005
INCOMPATIBLE_PROTOCOL_VERSION = -32006
PROTOCOL_NEGOTIATION_REQUIRED = -32007
+INDEX_BUILD_FAILED = -32008
+PROVIDER_FAILURE = -32009
+FAILURE_LIMIT_REACHED = -32010
+PROVIDER_CIRCUIT_OPEN = -32011
+INDEX_STORAGE_FAILURE = -32012
+INDEX_LOCKED = -32013
_METHOD_MODELS: dict[str, type[BaseModel]] = {
"hello": HelloParams,
"status": StatusParams,
"snapshot": SnapshotParams,
+ "index": IndexParams,
"discover": DiscoverParams,
"expand": ExpandParams,
"read": ReadParams,
@@ -432,7 +465,7 @@ async def _process_request(
},
)
timeout_ms = getattr(params, "timeout_ms", None)
- operation = self._dispatch(method, params, cancellation)
+ operation = self._dispatch(method, request_id, params, cancellation)
if timeout_ms is None:
result = await operation
else:
@@ -513,7 +546,11 @@ async def _process_request(
self._active.pop(_id_key(request_id), None)
async def _dispatch(
- self, method: str, raw: BaseModel, cancellation: asyncio.Event
+ self,
+ method: str,
+ request_id: str | int,
+ raw: BaseModel,
+ cancellation: asyncio.Event,
) -> dict[str, Any]:
if cancellation.is_set():
raise asyncio.CancelledError
@@ -523,6 +560,16 @@ async def _dispatch(
return await self._status(_require_type(raw, StatusParams), cancellation)
if method == "snapshot":
return await self._snapshot(cancellation)
+ if method == "index":
+ if self._protocol_version != "2.0":
+ raise BridgeFault(
+ METHOD_NOT_FOUND,
+ "METHOD_NOT_FOUND",
+ "The requested method is not supported by this protocol version.",
+ )
+ return await self._index(
+ request_id, _require_type(raw, IndexParams), cancellation
+ )
if method == "discover":
return await self._discover(
_require_type(raw, DiscoverParams), cancellation
@@ -543,6 +590,7 @@ async def _dispatch(
)
def _hello(self) -> dict[str, Any]:
+ bridge_v2 = self._protocol_version == "2.0"
return {
"protocol_version": self._protocol_version or BRIDGE_PROTOCOL_VERSION,
"supported_protocol_versions": list(SUPPORTED_BRIDGE_PROTOCOL_VERSIONS),
@@ -552,6 +600,7 @@ def _hello(self) -> dict[str, Any]:
"hello",
"status",
"snapshot",
+ *(["index"] if bridge_v2 else []),
"discover",
"expand",
"read",
@@ -565,17 +614,42 @@ def _hello(self) -> dict[str, Any]:
"concurrent_requests": True,
"serialized_responses": True,
"max_message_bytes": MAX_JSONRPC_MESSAGE_BYTES,
- "expansion_candidates": self._protocol_version == "1.1",
+ "expansion_candidates": self._protocol_version in {"1.1", "2.0"},
+ "tracked_index_jobs": bridge_v2,
+ "progress_notifications": bridge_v2,
+ "schemas": {
+ "index": {
+ "current": INDEX_SCHEMA_VERSION,
+ "readable": [1, INDEX_SCHEMA_VERSION],
+ },
+ "manifest": {
+ "current": MANIFEST_SCHEMA_VERSION,
+ "readable": [1, MANIFEST_SCHEMA_VERSION],
+ },
+ "record": {
+ "current": RECORD_SCHEMA_VERSION,
+ "readable": [1, RECORD_SCHEMA_VERSION],
+ },
+ "progress": {
+ "current": PROGRESS_SCHEMA_VERSION,
+ "readable": [1, 2, PROGRESS_SCHEMA_VERSION],
+ },
+ "context_package": {"current": 1, "readable": [1]},
+ },
},
"workspace": {
"identity": self.workspace_identity,
},
"policy": {
- "repository_access": "read_only_verified_snapshot",
- "external_data": "disabled",
+ "repository_access": (
+ "verified_snapshot_and_atomic_index_write"
+ if bridge_v2
+ else "read_only_verified_snapshot"
+ ),
+ "external_data": "provider_policy" if bridge_v2 else "disabled",
"portable_paths_only": True,
"source_writes": False,
- "index_mutation": False,
+ "index_mutation": bridge_v2,
"shell": False,
"subprocess_execution": False,
},
@@ -630,7 +704,7 @@ async def _status(
"lock_status": report.lock_status,
},
}
- if self._protocol_version == "1.1":
+ if self._protocol_version in {"1.1", "2.0"}:
result["index"]["coverage"] = await asyncio.to_thread(self._index_coverage)
return result
@@ -650,6 +724,146 @@ async def _snapshot(self, cancellation: asyncio.Event) -> dict[str, Any]:
}
return response
+ async def _index(
+ self,
+ request_id: str | int,
+ params: IndexParams,
+ cancellation: asyncio.Event,
+ ) -> dict[str, Any]:
+ operation_id = (
+ "bridge-index-"
+ + hashlib.sha256(
+ f"{type(request_id).__name__}:{request_id}".encode()
+ ).hexdigest()[:24]
+ )
+ provider: ModelProvider | None = None
+ last_event: ProgressEvent | None = None
+ progress_tail: asyncio.Task[None] | None = None
+
+ def observe(event: ProgressEvent) -> None:
+ nonlocal last_event, progress_tail
+ last_event = event
+ previous = progress_tail
+
+ async def publish() -> None:
+ if previous is not None:
+ await previous
+ await self._require_writer().write(
+ {
+ "jsonrpc": JSONRPC_VERSION,
+ "method": "$/progress",
+ "params": {
+ "request_id": request_id,
+ "event": event.model_dump(mode="json"),
+ },
+ }
+ )
+
+ progress_tail = asyncio.create_task(publish())
+
+ async def flush_progress() -> None:
+ if progress_tail is not None:
+ await progress_tail
+
+ try:
+ snapshot = await asyncio.to_thread(scan_repository, self.workspace)
+ current_digest = calculate_source_snapshot_digest(snapshot)
+ if current_digest != params.expected_snapshot_digest:
+ raise _index_bridge_fault(
+ IndexSourceChangedError(
+ "repository source identity differs from expected snapshot"
+ ),
+ last_event,
+ operation_id,
+ )
+ try:
+ project = load_project_configuration(self.workspace)
+ configuration = resolve_provider_configuration(
+ project,
+ provider=params.provider,
+ model=params.model,
+ base_url=params.base_url,
+ concurrency=params.concurrency,
+ timeout_seconds=params.request_timeout,
+ operation_timeout_seconds=params.request_timeout,
+ context_window=params.context_window,
+ json_repair_attempts=params.json_repair_attempts,
+ local_only=True if params.local_only else None,
+ )
+ if configuration is not None:
+ provider = create_model_provider(configuration)
+ except (ProjectConfigError, ValueError):
+ raise _index_bridge_fault(
+ ProviderConfigurationError(
+ "provider configuration could not be resolved"
+ ),
+ last_event,
+ operation_id,
+ ) from None
+ concurrency = (
+ configuration.concurrency_limit
+ if configuration is not None
+ else (
+ project.models.concurrency_limit
+ if params.concurrency is None
+ else params.concurrency
+ )
+ )
+ report = await build_repository_index(
+ self.workspace,
+ provider=provider,
+ provider_configuration=configuration,
+ update_only=params.action == "update",
+ concurrency=concurrency,
+ fail_on_error=params.fail_on_error,
+ fail_fast=params.fail_fast,
+ max_failures=params.max_failures,
+ force_reanalyze=params.force_reanalyze,
+ max_files=params.max_files,
+ semantic_max_output_tokens=(
+ project.models.semantic_max_output_tokens
+ if params.max_output_tokens is None
+ else params.max_output_tokens
+ ),
+ recover_stale_lock=params.recover_stale_lock,
+ confirm_unknown_lock=params.confirm_unknown_lock,
+ progress=observe,
+ operation_id=operation_id,
+ cancellation=cancellation,
+ )
+ await flush_progress()
+ self._snapshot_digest = report.manifest.build.source_snapshot_digest
+ self._preparations.clear()
+ return {
+ "action": params.action,
+ "generation_id": report.manifest.generation_id,
+ "snapshot_digest": report.manifest.build.source_snapshot_digest,
+ "index_schema": report.manifest.schema_versions.index_schema_version,
+ "partial": report.partial,
+ "statistics": report.manifest.statistics.model_dump(mode="json"),
+ }
+ except asyncio.CancelledError:
+ await flush_progress()
+ raise
+ except BridgeFault:
+ await flush_progress()
+ raise
+ except (
+ ApplicationError,
+ GlobalMapAnalysisError,
+ IndexStorageError,
+ ModelProviderError,
+ ProjectConfigError,
+ SemanticAnalysisError,
+ ValueError,
+ ) as exc:
+ await flush_progress()
+ raise _index_bridge_fault(exc, last_event, operation_id) from None
+ finally:
+ if provider is not None:
+ with suppress(ModelProviderError):
+ await provider.close()
+
async def _discover(
self, params: DiscoverParams, cancellation: asyncio.Event
) -> dict[str, Any]:
@@ -721,7 +935,7 @@ async def _expand(
"made_progress": result.made_progress,
"budget_usage": result.budget_usage.model_dump(mode="json"),
}
- if self._protocol_version == "1.1":
+ if self._protocol_version in {"1.1", "2.0"}:
response["candidates"] = self._register_expansion_candidates(
snapshot, preparation, params.operation, result.data
)
@@ -1246,6 +1460,85 @@ def _validation_details(exc: ValidationError | ValueError) -> list[dict[str, Any
return [{"location": [], "message": str(exc), "type": "value_error"}]
+def _index_bridge_fault(
+ error: BaseException,
+ event: ProgressEvent | None,
+ operation_id: str,
+) -> BridgeFault:
+ error_code = "index_build_failed"
+ reason = "ContextForge could not complete the index operation."
+ retryable = False
+ current: BaseException | None = error
+ seen: set[int] = set()
+ while current is not None and id(current) not in seen:
+ seen.add(id(current))
+ if isinstance(current, ModelProviderError):
+ error_code, reason = provider_error_details(current)
+ retryable = (
+ classify_retry(current) is RetryClassification.RETRYABLE
+ and not current.circuit_opened
+ )
+ break
+ typed_code = getattr(current, "error_code", None)
+ safe_reason = getattr(current, "safe_reason", None)
+ if isinstance(typed_code, str) and isinstance(safe_reason, str):
+ error_code = typed_code
+ reason = safe_reason[:1_000]
+ break
+ current = current.__cause__
+ rpc_code = INDEX_BUILD_FAILED
+ typed_rpc_code = "INDEX_BUILD_FAILED"
+ if isinstance(error, IndexSourceChangedError):
+ rpc_code = SOURCE_IDENTITY_CHANGED
+ typed_rpc_code = "SOURCE_IDENTITY_CHANGED"
+ error_code = "source_identity_changed"
+ reason = "Repository source identity changed before index publication."
+ retryable = True
+ elif isinstance(error, (ProjectConfigError, ProviderConfigurationError)):
+ rpc_code = PROVIDER_FAILURE
+ typed_rpc_code = "PROVIDER_CONFIGURATION_ERROR"
+ error_code = "provider_configuration_error"
+ reason = "Project provider configuration is invalid."
+ elif isinstance(error, SemanticFailureLimitError):
+ rpc_code = FAILURE_LIMIT_REACHED
+ typed_rpc_code = "FAILURE_LIMIT_REACHED"
+ error_code = "failure_limit_reached"
+ reason = "The configured semantic failure limit was reached."
+ elif isinstance(error, SemanticProviderCircuitError):
+ rpc_code = PROVIDER_CIRCUIT_OPEN
+ typed_rpc_code = "PROVIDER_CIRCUIT_OPEN"
+ error_code = "provider_circuit_open"
+ reason = "The provider circuit breaker opened during indexing."
+ retryable = False
+ elif isinstance(error, ModelProviderError):
+ rpc_code = PROVIDER_FAILURE
+ typed_rpc_code = "PROVIDER_FAILURE"
+ elif isinstance(error, IndexLockError):
+ rpc_code = INDEX_LOCKED
+ typed_rpc_code = "INDEX_LOCKED"
+ error_code = "index_lock_unavailable"
+ reason = "Another writer owns the index lock or lock recovery is required."
+ retryable = True
+ elif isinstance(error, IndexStorageError):
+ rpc_code = INDEX_STORAGE_FAILURE
+ typed_rpc_code = "INDEX_STORAGE_ERROR"
+ error_code = "index_storage_error"
+ reason = "ContextForge could not safely access index storage."
+ retryable = True
+ return BridgeFault(
+ rpc_code,
+ typed_rpc_code,
+ "ContextForge index operation failed.",
+ data={
+ "error_code": error_code,
+ "phase": "initialize" if event is None else event.phase_id,
+ "reason": reason,
+ "retryable": retryable,
+ "operation_id": operation_id,
+ },
+ )
+
+
__all__ = [
"BRIDGE_PROTOCOL_VERSION",
"JSONRPC_VERSION",
diff --git a/tests/test_bridge.py b/tests/test_bridge.py
index e3eea8e..9d5ae9e 100644
--- a/tests/test_bridge.py
+++ b/tests/test_bridge.py
@@ -19,9 +19,12 @@
BridgeSelectionItem,
CancelParams,
DiscoverParams,
+ IndexParams,
ReadParams,
)
from contextforge.intelligence import acquire_index_lock, build_structural_index
+from contextforge.models import ProviderAuthenticationError
+from contextforge.progress import ProgressEvent, ProgressStatus
from contextforge.repositories import scan_repository
@@ -42,6 +45,23 @@ def test_bridge_protocol_schema_is_closed_and_matches_v1() -> None:
assert "tool_name" not in expand["properties"]
assert schema["$defs"]["discoverResult"]["additionalProperties"] is False
+ v2 = json.loads(
+ (root / "docs/schemas/contextforge-bridge-v2.schema.json").read_text(
+ encoding="utf-8"
+ )
+ )
+ assert (
+ v2["$defs"]["indexRequest"]["properties"]["params"]["additionalProperties"]
+ is False
+ )
+ progress_event = v2["$defs"]["progressNotification"]["properties"]["params"][
+ "properties"
+ ]["event"]
+ assert progress_event["properties"]["schema_version"] == {"const": 3}
+ assert v2["$defs"]["progressNotification"]["properties"]["method"] == {
+ "const": "$/progress"
+ }
+
def test_bridge_status_reports_structural_index_coverage(tmp_path: Path) -> None:
(tmp_path / "parsed.py").write_text("def run():\n return 1\n", encoding="utf-8")
@@ -121,6 +141,19 @@ def wait(self, count: int) -> list[dict[str, Any]]:
chunks = list(self.chunks[:count])
return [cast(dict[str, Any], json.loads(chunk)) for chunk in chunks]
+ def wait_for_id(self, request_id: str | int) -> dict[str, Any]:
+ deadline = time.monotonic() + 10
+ with self._condition:
+ while True:
+ for chunk in self.chunks:
+ frame = cast(dict[str, Any], json.loads(chunk))
+ if frame.get("id") == request_id:
+ return frame
+ remaining = deadline - time.monotonic()
+ if remaining <= 0:
+ raise AssertionError("timed out waiting for bridge response")
+ self._condition.wait(remaining)
+
class _Harness:
def __init__(
@@ -201,7 +234,11 @@ async def exercise() -> None:
hello = (await harness.response(1))[0]
assert hello["jsonrpc"] == "2.0"
assert hello["result"]["protocol_version"] == "1.0"
- assert hello["result"]["supported_protocol_versions"] == ["1.0", "1.1"]
+ assert hello["result"]["supported_protocol_versions"] == [
+ "1.0",
+ "1.1",
+ "2.0",
+ ]
assert hello["result"]["capabilities"]["model_free_discovery"] is True
assert hello["result"]["policy"]["source_writes"] is False
assert "shell" in hello["result"]["policy"]
@@ -234,13 +271,13 @@ async def exercise() -> None:
assert missing["error"]["data"]["code"] == "INVALID_PARAMS"
harness.input.send(
- _request("incompatible", "hello", {"protocol_version": "2.0"})
+ _request("incompatible", "hello", {"protocol_version": "3.0"})
)
incompatible = (await harness.response(3))[-1]
assert incompatible["error"]["data"] == {
"code": "INCOMPATIBLE_PROTOCOL_VERSION",
- "requested_protocol_version": "2.0",
- "supported_protocol_versions": ["1.0", "1.1"],
+ "requested_protocol_version": "3.0",
+ "supported_protocol_versions": ["1.0", "1.1", "2.0"],
}
harness.input.send(_request("compatible", "hello", {"protocol_version": "1.0"}))
@@ -255,6 +292,223 @@ async def exercise() -> None:
asyncio.run(exercise())
+def test_bridge_v1_does_not_expose_index_mutation(tmp_path: Path) -> None:
+ async def exercise() -> None:
+ harness = _Harness(tmp_path)
+ await harness.start(negotiated=False)
+ harness.input.send(_request("hello", "hello", {"protocol_version": "1.1"}))
+ hello = (await harness.response(1))[-1]
+ assert "index" not in hello["result"]["capabilities"]["methods"]
+ assert hello["result"]["policy"]["index_mutation"] is False
+
+ harness.input.send(
+ _request(
+ "index-v1",
+ "index",
+ {
+ "action": "build",
+ "expected_snapshot_digest": "0" * 64,
+ "provider": "none",
+ },
+ )
+ )
+ rejected = await asyncio.to_thread(harness.output.wait_for_id, "index-v1")
+ assert rejected["error"]["data"]["code"] == "METHOD_NOT_FOUND"
+ await harness.close()
+
+ asyncio.run(exercise())
+
+
+def test_bridge_v2_build_update_and_correlated_progress(tmp_path: Path) -> None:
+ (tmp_path / "app.py").write_text("VALUE = 1\n", encoding="utf-8")
+
+ async def exercise() -> None:
+ harness = _Harness(tmp_path)
+ await harness.start(negotiated=False)
+ harness.input.send(_request("hello", "hello", {"protocol_version": "2.0"}))
+ hello = (await harness.response(1))[-1]["result"]
+ capabilities = hello["capabilities"]
+ assert capabilities["tracked_index_jobs"] is True
+ assert capabilities["progress_notifications"] is True
+ assert capabilities["schemas"] == {
+ "index": {"current": 2, "readable": [1, 2]},
+ "manifest": {"current": 2, "readable": [1, 2]},
+ "record": {"current": 2, "readable": [1, 2]},
+ "progress": {"current": 3, "readable": [1, 2, 3]},
+ "context_package": {"current": 1, "readable": [1]},
+ }
+ assert "index" in capabilities["methods"]
+ digest = await _snapshot(harness, 2)
+
+ for action in ("build", "update"):
+ request_id = f"index-{action}"
+ harness.input.send(
+ _request(
+ request_id,
+ "index",
+ {
+ "action": action,
+ "expected_snapshot_digest": digest,
+ "provider": "none",
+ },
+ )
+ )
+ response = await asyncio.to_thread(harness.output.wait_for_id, request_id)
+ assert response["result"]["action"] == action
+ assert response["result"]["snapshot_digest"] == digest
+ assert response["result"]["index_schema"] == 2
+ assert response["result"]["partial"] is False
+ progress = [
+ frame
+ for frame in await harness.response(len(harness.output.chunks))
+ if frame.get("method") == "$/progress"
+ and frame["params"]["request_id"] == request_id
+ ]
+ assert progress
+ events = [
+ ProgressEvent.model_validate(frame["params"]["event"])
+ for frame in progress
+ ]
+ assert events[-1].status is ProgressStatus.COMPLETED
+ assert (
+ events[-1].metadata["generation_id"]
+ == response["result"]["generation_id"]
+ )
+ await harness.close()
+
+ asyncio.run(exercise())
+
+
+def test_bridge_v2_index_cancellation_is_cooperative(
+ tmp_path: Path, monkeypatch: Any
+) -> None:
+ started = asyncio.Event()
+
+ async def blocked_build(*args: object, **kwargs: object) -> None:
+ del args
+ cancellation = cast(asyncio.Event, kwargs["cancellation"])
+ started.set()
+ while not cancellation.is_set():
+ await asyncio.sleep(0)
+ raise asyncio.CancelledError
+
+ monkeypatch.setattr(bridge_module, "build_repository_index", blocked_build)
+
+ async def exercise() -> None:
+ harness = _Harness(tmp_path)
+ await harness.start(negotiated=False)
+ harness.input.send(_request("hello", "hello", {"protocol_version": "2.0"}))
+ await harness.response(1)
+ digest = await _snapshot(harness, 2)
+ harness.input.send(
+ _request(
+ "slow-index",
+ "index",
+ {
+ "action": "build",
+ "expected_snapshot_digest": digest,
+ "provider": "none",
+ },
+ )
+ )
+ await asyncio.wait_for(started.wait(), timeout=1)
+ harness.input.send(
+ {
+ "jsonrpc": "2.0",
+ "method": "$/cancelRequest",
+ "params": {"id": "slow-index"},
+ }
+ )
+ response = await asyncio.to_thread(harness.output.wait_for_id, "slow-index")
+ assert response["error"]["data"]["code"] == "REQUEST_CANCELLED"
+ await harness.close()
+
+ asyncio.run(exercise())
+
+
+def test_bridge_v2_index_errors_are_typed_and_safe(
+ tmp_path: Path, monkeypatch: Any
+) -> None:
+ async def authentication_failure(*args: object, **kwargs: object) -> None:
+ del args, kwargs
+ raise ProviderAuthenticationError(
+ "Bearer top-secret at http://user:pass@example.invalid"
+ )
+
+ monkeypatch.setattr(bridge_module, "build_repository_index", authentication_failure)
+
+ async def exercise() -> None:
+ harness = _Harness(tmp_path)
+ await harness.start(negotiated=False)
+ harness.input.send(_request("hello", "hello", {"protocol_version": "2.0"}))
+ await harness.response(1)
+ digest = await _snapshot(harness, 2)
+ harness.input.send(
+ _request(
+ "provider-error",
+ "index",
+ {
+ "action": "build",
+ "expected_snapshot_digest": digest,
+ "provider": "none",
+ },
+ )
+ )
+ response = await asyncio.to_thread(harness.output.wait_for_id, "provider-error")
+ data = response["error"]["data"]
+ assert response["error"]["code"] == bridge_module.PROVIDER_FAILURE
+ assert data["code"] == "PROVIDER_FAILURE"
+ assert data["error_code"] == "authentication_failed"
+ assert data["phase"] == "initialize"
+ assert data["reason"] == "provider authentication failed"
+ assert data["retryable"] is False
+ assert data["operation_id"].startswith("bridge-index-")
+ assert "secret" not in json.dumps(response)
+
+ harness.input.send(
+ _request(
+ "snapshot-drift",
+ "index",
+ {
+ "action": "build",
+ "expected_snapshot_digest": "0" * 64,
+ "provider": "none",
+ },
+ )
+ )
+ drift = await asyncio.to_thread(harness.output.wait_for_id, "snapshot-drift")
+ drift_data = drift["error"]["data"]
+ assert drift["error"]["code"] == bridge_module.SOURCE_IDENTITY_CHANGED
+ assert drift_data["code"] == "SOURCE_IDENTITY_CHANGED"
+ assert drift_data["error_code"] == "source_identity_changed"
+ assert drift_data["retryable"] is True
+
+ harness.input.send(
+ _request(
+ "configuration-error",
+ "index",
+ {
+ "action": "build",
+ "expected_snapshot_digest": digest,
+ "provider": "openai-compatible",
+ "model": "exact/model",
+ "base_url": "http://user:top-secret@example.invalid/v1",
+ },
+ )
+ )
+ configuration = await asyncio.to_thread(
+ harness.output.wait_for_id, "configuration-error"
+ )
+ configuration_data = configuration["error"]["data"]
+ assert configuration["error"]["code"] == bridge_module.PROVIDER_FAILURE
+ assert configuration_data["code"] == "PROVIDER_CONFIGURATION_ERROR"
+ assert configuration_data["error_code"] == "provider_configuration_error"
+ assert "top-secret" not in json.dumps(configuration)
+ await harness.close()
+
+ asyncio.run(exercise())
+
+
def test_bridge_rejects_malformed_oversized_unknown_and_invalid_requests(
tmp_path: Path,
) -> None:
@@ -1126,6 +1380,15 @@ async def failure_exercise() -> None:
def test_bridge_parameter_models_reject_noncanonical_values() -> None:
+ with pytest.raises(ValidationError):
+ IndexParams.model_validate(
+ {
+ "action": "build",
+ "expected_snapshot_digest": "0" * 64,
+ "fail_fast": True,
+ "max_failures": 2,
+ }
+ )
with pytest.raises(ValidationError):
DiscoverParams.model_validate_json(
json.dumps(
From f5388938671b44551af5296b57db79edd104a0d4 Mon Sep 17 00:00:00 2001
From: Kirill <106469980+waterflane@users.noreply.github.com>
Date: Mon, 7 Sep 2026 20:18:03 +0300
Subject: [PATCH 04/11] fix(index): stabilize analyzer identity across
endpoints
---
src/contextforge/intelligence/manifest.py | 4 +-
src/contextforge/intelligence/maps.py | 80 ++++++++++++++---------
tests/test_intelligence_manifest.py | 15 ++---
tests/test_repository_maps.py | 31 +++++++++
4 files changed, 88 insertions(+), 42 deletions(-)
diff --git a/src/contextforge/intelligence/manifest.py b/src/contextforge/intelligence/manifest.py
index a3395ce..eaf60ac 100644
--- a/src/contextforge/intelligence/manifest.py
+++ b/src/contextforge/intelligence/manifest.py
@@ -15,6 +15,7 @@
SchemaVersionMetadata,
analyzer_identity_key,
calculate_index_statistics,
+ normalize_analyzer_identity,
validate_portable_relative_path,
)
from contextforge.repositories import ProjectFile, ProjectSnapshot
@@ -183,7 +184,8 @@ def identify_stale_analysis(
and (
invalidate_all
or not _source_matches(indexed, current_by_path[indexed.path])
- or indexed.analyzer != expected_analyzer
+ or normalize_analyzer_identity(indexed.analyzer)
+ != normalize_analyzer_identity(expected_analyzer)
)
)
diff --git a/src/contextforge/intelligence/maps.py b/src/contextforge/intelligence/maps.py
index 752527b..dd8fd74 100644
--- a/src/contextforge/intelligence/maps.py
+++ b/src/contextforge/intelligence/maps.py
@@ -48,6 +48,7 @@
IndexModel,
ModelIdentity,
analyzer_identity_key,
+ normalize_analyzer_identity,
)
from contextforge.intelligence.semantic_models import (
EvidenceReference,
@@ -336,8 +337,8 @@ async def build_repository_maps(
)
overview = build_repository_overview(current, code_maps)
semantic_analyses = _load_available_semantics(snapshot.root, current)
- provider_id, model_id, base_url_sha256 = _provider_identity(provider)
- analyzer = _global_analyzer(active_options, provider_id, model_id, base_url_sha256)
+ provider_id, model_id = _provider_identity(provider)
+ analyzer = _global_analyzer(active_options, provider_id, model_id)
options_digest = _options_digest(active_options)
source_interpretations_digest = _file_interpretations_digest(current)
previous_records = _try_load_global_records(snapshot.root, current)
@@ -362,20 +363,40 @@ async def build_repository_maps(
source_interpretations_digest,
)
):
+ normalized_architecture = old_architecture.model_copy(
+ update={"analyzer": analyzer}
+ )
+ normalized_features = old_features.model_copy(update={"analyzer": analyzer})
+ migrated = (
+ normalized_architecture != old_architecture
+ or normalized_features != old_features
+ )
+ generation_path = lock.layout.generations / current.generation_id
+ manifest = current
+ if migrated:
+ generation_path = _publish_global_records(
+ lock,
+ current,
+ overview,
+ normalized_architecture,
+ normalized_features,
+ analyzer,
+ )
+ manifest = load_manifest(snapshot.root)
return GlobalMapBuildResult(
- manifest=current,
+ manifest=manifest,
overview=overview,
- architecture=old_architecture,
- features=old_features,
+ architecture=normalized_architecture,
+ features=normalized_features,
outcomes=(
GlobalMapOutcome("architecture", "reused", 0),
GlobalMapOutcome("features", "reused", 0),
),
- generation_path=lock.layout.generations / current.generation_id,
+ generation_path=generation_path,
request_count=0,
package_summary_count=0,
group_summary_count=0,
- published=False,
+ published=migrated,
)
required_model_calls = _required_model_calls(code_maps, active_options)
@@ -397,6 +418,10 @@ async def build_repository_maps(
except ProviderCancelledError:
raise
except (ModelProviderError, GlobalMapAnalysisError, ValueError) as exc:
+ if isinstance(exc, ModelProviderError) and exc.circuit_opened:
+ raise GlobalMapAnalysisError(
+ "provider circuit opened during repository map hierarchy"
+ ) from exc
if active_options.fail_on_error:
raise GlobalMapAnalysisError(
"repository map hierarchy failed; index not published"
@@ -477,6 +502,10 @@ async def build_repository_maps(
except ProviderCancelledError:
raise
except (ModelProviderError, GlobalMapAnalysisError, ValueError) as exc:
+ if isinstance(exc, ModelProviderError) and exc.circuit_opened:
+ raise GlobalMapAnalysisError(
+ "provider circuit opened during architecture map analysis"
+ ) from exc
diagnostic = _failure_diagnostic("architecture-map-failed", exc)
outcomes.append(GlobalMapOutcome("architecture", "failed", 1, diagnostic))
@@ -504,6 +533,10 @@ async def build_repository_maps(
except ProviderCancelledError:
raise
except (ModelProviderError, GlobalMapAnalysisError, ValueError) as exc:
+ if isinstance(exc, ModelProviderError) and exc.circuit_opened:
+ raise GlobalMapAnalysisError(
+ "provider circuit opened during feature map analysis"
+ ) from exc
diagnostic = _failure_diagnostic("feature-map-failed", exc)
outcomes.append(GlobalMapOutcome("features", "failed", 1, diagnostic))
@@ -1955,11 +1988,13 @@ def _publish_global_records(
interpretations_digest=interpretations_digest,
previous_generation_id=current.generation_id,
)
- semantic_analyzers = current.semantic_analyzers
+ semantic_analyzers = tuple(
+ normalize_analyzer_identity(item) for item in current.semantic_analyzers
+ )
if analyzer is not None:
semantic_analyzers = tuple(
sorted(
- set((*current.semantic_analyzers, analyzer)),
+ set((*semantic_analyzers, analyzer)),
key=analyzer_identity_key,
)
)
@@ -2116,13 +2151,10 @@ def _global_analyzer(
options: GlobalMapAnalysisOptions,
provider_id: str,
model_id: str,
- base_url_sha256: str | None,
) -> AnalyzerIdentity:
return AnalyzerIdentity(
analyzer_id=GLOBAL_MAP_ANALYZER_ID,
- analyzer_version=_connection_bound_version(
- GLOBAL_MAP_ANALYZER_VERSION, base_url_sha256
- ),
+ analyzer_version=GLOBAL_MAP_ANALYZER_VERSION,
analysis_prompt_version=options.prompt_version,
response_schema_version=GLOBAL_MAP_SCHEMA_VERSION,
model_identity=ModelIdentity(
@@ -2183,12 +2215,12 @@ def _map_inputs_match(
value.source_snapshot_digest == manifest.build.source_snapshot_digest
and value.facts_digest == manifest.build.facts_digest
and value.source_interpretations_digest == source_interpretations_digest
- and value.analyzer == analyzer
+ and normalize_analyzer_identity(value.analyzer) == analyzer
and value.analysis_options_digest == options_digest
)
-def _provider_identity(provider: ModelProvider) -> tuple[str, str, str | None]:
+def _provider_identity(provider: ModelProvider) -> tuple[str, str]:
provider_id = provider.provider_id
configuration = getattr(provider, "configuration", None)
model_id = getattr(configuration, "model_id", None)
@@ -2196,23 +2228,7 @@ def _provider_identity(provider: ModelProvider) -> tuple[str, str, str | None]:
raise GlobalMapAnalysisError(
"global map provider must expose stable provider and model identity"
)
- endpoint = getattr(configuration, "endpoint", None)
- base_url_sha256 = None
- if provider_id == "openai-compatible":
- if not isinstance(endpoint, str):
- raise GlobalMapAnalysisError(
- "OpenAI-compatible provider must expose a stable base URL identity"
- )
- base_url_sha256 = hashlib.sha256(
- endpoint.rstrip("/").encode("utf-8")
- ).hexdigest()
- return provider_id, model_id, base_url_sha256
-
-
-def _connection_bound_version(version: str, base_url_sha256: str | None) -> str:
- if base_url_sha256 is None:
- return version
- return f"{version}+base.{base_url_sha256}"
+ return provider_id, model_id
def _validate_build_inputs(snapshot: ProjectSnapshot, lock: IndexWriteLock) -> None:
diff --git a/tests/test_intelligence_manifest.py b/tests/test_intelligence_manifest.py
index d4d2900..cdfb71e 100644
--- a/tests/test_intelligence_manifest.py
+++ b/tests/test_intelligence_manifest.py
@@ -341,7 +341,7 @@ def test_invalid_build_options_digest_is_rejected() -> None:
)
-def test_openai_compatible_base_url_identity_change_invalidates_model_records() -> None:
+def test_legacy_endpoint_suffixes_are_equivalent_model_identities() -> None:
project_file = _file("app.py", "pass")
first = _analyzer().model_copy(
update={
@@ -359,12 +359,9 @@ def test_openai_compatible_base_url_identity_change_invalidates_model_records()
)
manifest = _manifest((project_file,), analyzer=first)
- assert (
- identify_stale_analysis(
- manifest,
- (project_file,),
- expected_analyzer=changed,
- build_options_digest=_sha("options"),
- )
- == manifest.files
+ assert not identify_stale_analysis(
+ manifest,
+ (project_file,),
+ expected_analyzer=changed,
+ build_options_digest=_sha("options"),
)
diff --git a/tests/test_repository_maps.py b/tests/test_repository_maps.py
index f67bec8..b0c0fb8 100644
--- a/tests/test_repository_maps.py
+++ b/tests/test_repository_maps.py
@@ -7,6 +7,7 @@
import pytest
from pydantic import ValidationError
+import contextforge.intelligence.maps as maps_module
from contextforge.intelligence import (
ArchitectureMap,
GlobalMapAnalysisError,
@@ -281,6 +282,36 @@ def test_one_module_hierarchical_maps_persist_and_reuse_without_source_prompt(
)
+def test_legacy_endpoint_identity_is_republished_for_maps_without_model_calls(
+ tmp_path: Path, monkeypatch: pytest.MonkeyPatch
+) -> None:
+ snapshot = _facts(tmp_path, {"main.py": "def main():\n return 1\n"})
+ original = maps_module._global_analyzer
+
+ def legacy_analyzer(*args: Any, **kwargs: Any) -> Any:
+ identity = original(*args, **kwargs)
+ return identity.model_copy(
+ update={"analyzer_version": identity.analyzer_version + "+base." + "b" * 64}
+ )
+
+ monkeypatch.setattr(maps_module, "_global_analyzer", legacy_analyzer)
+ legacy = _maps(snapshot, _Responder(), run_id="legacy-map-identity")
+ monkeypatch.setattr(maps_module, "_global_analyzer", original)
+ responder = _Responder()
+
+ migrated = _maps(snapshot, responder, run_id="migrate-map-identity")
+
+ assert responder.requests == []
+ assert migrated.published is True
+ assert migrated.manifest.generation_id != legacy.manifest.generation_id
+ assert migrated.architecture is not None
+ assert "+base." not in migrated.architecture.analyzer.analyzer_version
+ assert all(
+ "+base." not in item.analyzer_version
+ for item in migrated.manifest.semantic_analyzers
+ )
+
+
def test_multi_package_hierarchy_entry_adapter_core_and_deterministic_order(
tmp_path: Path,
) -> None:
From 7919fcbaffc61e1d2c4f39019ccb79192ec4182b Mon Sep 17 00:00:00 2001
From: Kirill <106469980+waterflane@users.noreply.github.com>
Date: Mon, 7 Sep 2026 20:18:32 +0300
Subject: [PATCH 05/11] docs(bridge): document host integration contracts
---
CHANGELOG.md | 27 ++++++
README.md | 24 +++--
docs/architecture/model-providers.md | 19 +++-
docs/architecture/progress-reporting.md | 8 ++
docs/decisions/003-host-managed-index-jobs.md | 71 ++++++++++++++
docs/guides/bridge.md | 96 +++++++++++++++----
src/contextforge/cli/main.py | 15 +--
7 files changed, 228 insertions(+), 32 deletions(-)
create mode 100644 docs/decisions/003-host-managed-index-jobs.md
diff --git a/CHANGELOG.md b/CHANGELOG.md
index 16bb097..4d73462 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -8,6 +8,33 @@ and Python distribution versions follow PEP 440.
## [Unreleased]
+### Added
+
+- Added independent `--fail-fast` and `--max-failures N` index policies, a
+ bounded semantic scheduler, and a job-scoped provider circuit breaker.
+- Added `--progress jsonl` for clean, flushed `ProgressEvent` schema 3 streams
+ from `index build` and `index update`.
+- Added opt-in Bridge 2.0 tracked `build`/`update` jobs, correlated
+ `$/progress` notifications, cooperative cancellation, schema capabilities,
+ and a normative Bridge 2 JSON Schema. Bridge 1.0/1.1 remain supported.
+
+### Changed
+
+- Classified authentication, authorization, missing credential, quota, rate
+ limit, model, configuration, timeout, and service failures so terminal
+ provider-wide failures are not retried per file.
+- Made semantic and repository-map analyzer identity depend on provider/model
+ and analysis contracts rather than an OpenAI-compatible endpoint. Legacy
+ `+base.