Files
2026-09-04 14:58:42 +08:00

1063 lines
44 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Constrained semantic annotation with one OpenAI-compatible request per Skill."""
from __future__ import annotations
import json
import re
from types import SimpleNamespace
from typing import Any, Callable, Protocol
from urllib.error import HTTPError, URLError
from urllib.request import Request, urlopen
from scripts.provider_router import resolve_model_route
from .document import SkillDocument, resolve_annotation_conflicts
from .models import Annotation, AnnotationResult, SemanticPlanResult, Signal
from .semantic_plan import (
PLAN_SCHEMA_VERSION,
SUMMARY_SOURCE_KINDS,
rejection_index,
replacement_literal_delta,
repairable_rewrite_indices,
validate_semantic_plan,
)
class AnnotationError(RuntimeError):
"""Semantic annotation was required but could not be completed."""
class _HttpOpenAICompat:
"""Small OpenAI-compatible client fallback for SDK-free environments."""
def __init__(self, endpoint: str, api_key: str, timeout: float = 90.0):
self._endpoint = endpoint
self._api_key = api_key
self._timeout = timeout
self.chat = SimpleNamespace(completions=SimpleNamespace(create=self._create))
def _create(self, **options: Any) -> Any:
request = Request(
self._endpoint,
data=json.dumps(options, ensure_ascii=False).encode("utf-8"),
headers={
"Authorization": f"Bearer {self._api_key}",
"Content-Type": "application/json",
},
method="POST",
)
try:
with urlopen(request, timeout=self._timeout) as response:
payload = json.loads(response.read().decode("utf-8"))
except HTTPError as exc:
body = exc.read().decode("utf-8", "replace")
error = RuntimeError(f"HTTP {exc.code}: {body[:1000]}")
error.status_code = exc.code
raise error from exc
except URLError as exc:
error = RuntimeError(str(exc.reason))
raise error from exc
if not isinstance(payload, dict):
raise RuntimeError("provider returned a non-object JSON value")
choices = []
for item in payload.get("choices", []):
message = item.get("message", {}) if isinstance(item, dict) else {}
choices.append(
SimpleNamespace(
message=SimpleNamespace(content=message.get("content")),
finish_reason=item.get("finish_reason") if isinstance(item, dict) else None,
)
)
usage = payload.get("usage") if isinstance(payload.get("usage"), dict) else {}
return SimpleNamespace(
choices=choices,
usage=SimpleNamespace(
completion_tokens=usage.get("completion_tokens"),
completion_tokens_details=SimpleNamespace(
reasoning_tokens=(usage.get("completion_tokens_details") or {}).get("reasoning_tokens")
if isinstance(usage.get("completion_tokens_details"), dict) else None
),
),
)
class Annotator(Protocol):
model_id: str
def annotate(
self,
document: SkillDocument,
selected_passes: list[str],
allowed_types: set[str],
) -> dict[str, Any]:
"""Return the raw JSON-compatible annotation payload."""
class SemanticPlanner(Protocol):
model_id: str
def plan(
self,
document: SkillDocument,
signals: dict[str, Signal],
selected_passes: list[str],
) -> dict[str, Any]:
"""Return a source-grounded semantic rewrite plan."""
PASS_ANNOTATION_TYPES = {
# The behavioral profile supports using an LLM to locate task-specific
# directives for this signal. It does not support inferring definitions or
# completion criteria from ordinary prose; those remain static-only.
"contextual_rule_adherence": {"critical_rule"},
"semantic_robustness": {"coreference"},
"uncertainty_calibration": {"uncertainty_rule"},
"ambiguity_handling": {"scope_rule", "decision_criterion"},
"evidence_priority": {"evidence_priority_rule"},
"balanced_presentation": {"viewpoint_side_a", "viewpoint_side_b"},
}
PRIMARY_ANNOTATION_TYPES = {
"contextual_rule_adherence": {"critical_rule"},
"evidence_priority": {"evidence_priority_rule"},
}
SCHEMA_VERSION = "1.0"
_CODE_CONFIG_ASSIGNMENT_RE = re.compile(
r"^\s*(?:export\s+)?(?P<name>[A-Za-z_][A-Za-z0-9_]*)\s*=\s*"
r"(?P<value>(?:[rRuUbBfF]{0,2})?['\"][^'\"]+['\"]|(?:\.?\.?/|/)[^\s#]+)"
)
_CONFIG_NAME_RE = re.compile(
r"(?:PATH|DIR|CACHE|ROOT|HOME|CONFIG|ENDPOINT|HOST|PORT|MODE|OFFLINE|DATABASE|DB)",
re.I,
)
_COMMAND_FLAG_RE = re.compile(r"(?<![\w-])--[A-Za-z0-9][A-Za-z0-9-]*")
_COMMAND_WORKFLOW_RE = re.compile(
r"\b(?:command|cmd|args|subprocess|run|execute|invocation)\b", re.I
)
ANNOTATION_DEFINITIONS = {
"critical_rule": (
"An explicit task-specific directive, prohibition, required action, or "
"conditional instruction. It must tell the agent what to do or not do. "
"Background facts, rationales, benefits, limitations, use cases, and "
"section labels are not critical rules."
),
"evidence_priority_rule": (
"An explicit instruction saying which of two named evidence sources wins "
"when they conflict. A factual statement that one source is authoritative "
"is not enough."
),
"coreference": (
"A pronoun or referential phrase whose exact antecedent also occurs in the "
"input and whose replacement is unambiguous."
),
}
def allowed_annotation_types(
selected_passes: list[str],
signals: dict[str, Signal] | None = None,
) -> set[str]:
allowed: set[str] = set()
for name in selected_passes:
if signals is not None:
signal = signals[name]
if signal.level not in {"low", "medium"}:
continue
allowed.update(PASS_ANNOTATION_TYPES.get(name, set()))
return allowed
def semantic_annotation_needed(
signals: dict[str, Signal],
selected_passes: list[str],
static_annotations: list[Annotation],
) -> bool:
existing = {annotation.type for annotation in static_annotations}
for name in selected_passes:
signal = signals[name]
if signal.level not in {"low", "medium"}:
continue
required = PRIMARY_ANNOTATION_TYPES.get(
name, PASS_ANNOTATION_TYPES.get(name, set())
)
# Coreference and viewpoint annotations are opportunistic. The absence of
# a pronoun or explicit comparison is not evidence that an LLM call is useful.
required = required - {"coreference", "viewpoint_side_a", "viewpoint_side_b"}
if required and not required.intersection(existing):
return True
return False
def _code_configuration_evidence(block: Any) -> list[dict[str, str]]:
"""Expose declarative runtime configuration, not implementation flow."""
evidence: list[dict[str, str]] = []
for line in block.text.splitlines():
match = _CODE_CONFIG_ASSIGNMENT_RE.match(line)
if match is None or _CONFIG_NAME_RE.search(match.group("name")) is None:
continue
quote = line.strip()
if quote:
evidence.append(
{
"block_id": block.id,
"kind": "code_config",
"parent_heading": block.parent_heading,
"text": quote,
}
)
return evidence
def _code_workflow_closure_evidence(block: Any) -> list[dict[str, str]]:
"""Expose a command workflow as one source-grounded semantic unit.
A command flag is rarely an independent instruction: neighbouring flags,
configured resources, and the command construction jointly determine its
behaviour. Supplying the complete code block prevents the planner from
treating one flag as the whole operational rule while still leaving the
original code untouched.
"""
flags = _COMMAND_FLAG_RE.findall(block.text)
if len(set(flags)) < 2 or _COMMAND_WORKFLOW_RE.search(block.text) is None:
return []
return [
{
"block_id": block.id,
"kind": "code_workflow_closure",
"semantic_role": "execution_contract",
"parent_heading": block.parent_heading,
"text": block.text,
}
]
def _overlaps_protected(block_text: str, quote: str, spans: list[tuple[int, int]]) -> bool:
start = block_text.find(quote)
if start < 0:
return True
end = start + len(quote)
return any(start < protected_end and end > protected_start for protected_start, protected_end in spans)
CRITICAL_DIRECTIVE_RE = re.compile(
r"\b(?:must(?:\s+not)?|shall(?:\s+not)?|should(?:\s+not)?|"
r"required|require|ensure|verify|validate|preserve|retain|keep|"
r"always|never|do\s+not|don't|cannot|can't|only|prefer|use|"
r"treat|respect|document|include|exclude|recurse|run|write|report)\b|"
r"必须|不得|禁止|只能|仅可|应当|需要|确保|验证|保留|不要|不可|优先|递归",
re.I,
)
COMPLETION_DIRECTIVE_RE = re.compile(
r"只有.+才(?:算|可以|可|能).*(?:完成|结束)|完成条件\s*[::]|"
r"\bonly\s+.+\s+(?:counts?\s+as|is)\s+(?:complete|done)\b|"
r"\b(?:complete|done|finished|successful|success)\s+(?:when|only\s+if)\b|"
r"\btask\s+is\s+(?:complete|done)\s+(?:when|only\s+if)\b",
re.I,
)
EVIDENCE_PRIORITY_DIRECTIVE_RE = re.compile(
r"以.+为准|.+优先于.+|(?:冲突|不一致)时.+(?:为准|优先)|"
r"\b.+takes?\s+precedence\s+over\b.+|"
r"\b(?:prefer|use)\s+.+\s+(?:over|instead\s+of)\s+.+",
re.I,
)
DEFINITION_DIRECTIVE_RE = re.compile(
r"(?:此处|这里|本任务中).+?(?:是指|指的是|定义为)|"
r"\b[A-Za-z][A-Za-z0-9 _-]{0,40}\s+"
r"(?:means|refers\s+to|is\s+defined\s+as)\b",
re.I,
)
EXPLANATORY_PREFIX_RE = re.compile(
r"^(?:this\s+skill\s+provides|use\s+case\s*:|rationale\s*:|"
r"cause\s*:|benefits?\s*:|limitations?\s*:|"
r"key\s+(?:flags|parameters|features)\s+(?:for|are)\b)",
re.I,
)
def annotation_semantically_valid(
annotation_type: str,
block: Any,
quote: str,
) -> bool:
"""Apply conservative type gates after source-span validation."""
clean = re.sub(r"^(?:\s*[-+*]|\s*\d+[.)])\s+", "", quote.strip())
clean = clean.replace("**", "").replace("__", "")
parent_heading = (getattr(block, "parent_heading", None) or "").strip()
if annotation_type == "critical_rule":
if EXPLANATORY_PREFIX_RE.search(clean):
return False
if re.search(
r"\b(?:overview|introduction|why|benefits?|challenges?|"
r"limitations?|ideal\s+use\s+cases?|rationale)\b",
parent_heading,
re.I,
) and not re.search(
r"\b(?:must(?:\s+not)?|required|always|never|do\s+not|"
r"ensure|only)\b|必须|不得|禁止|只能|确保",
clean,
re.I,
):
return False
return bool(CRITICAL_DIRECTIVE_RE.search(clean))
if annotation_type == "completion_criterion":
return bool(COMPLETION_DIRECTIVE_RE.search(clean))
if annotation_type == "evidence_priority_rule":
return bool(EVIDENCE_PRIORITY_DIRECTIVE_RE.search(clean))
if annotation_type == "definition":
return bool(DEFINITION_DIRECTIVE_RE.search(clean))
# Other types are not currently requested from the LLM unless their signal
# has sufficient multi-item confidence. Source and confidence validation
# still apply to them.
return True
def validate_annotations(
payload: dict[str, Any],
document: SkillDocument,
allowed_types: set[str],
) -> tuple[list[Annotation], int]:
if payload.get("schema_version") != SCHEMA_VERSION:
raise AnnotationError("annotation response has unsupported schema_version")
values = payload.get("annotations")
if not isinstance(values, list):
raise AnnotationError("annotation response requires annotations list")
blocks = document.block_index
accepted: list[Annotation] = []
rejected = 0
for value in values:
if not isinstance(value, dict):
rejected += 1
continue
annotation_type = value.get("type")
block_id = value.get("block_id")
quote = value.get("quote")
confidence = value.get("confidence")
antecedent = value.get("antecedent_quote")
if (
annotation_type not in allowed_types
or not isinstance(block_id, str)
or block_id not in blocks
or not isinstance(quote, str)
or not quote
or isinstance(confidence, bool)
or not isinstance(confidence, (int, float))
or float(confidence) < 0.85
):
rejected += 1
continue
block = blocks[block_id]
if quote not in block.text or _overlaps_protected(block.text, quote, block.protected_spans):
rejected += 1
continue
if not annotation_semantically_valid(annotation_type, block, quote):
rejected += 1
continue
if annotation_type == "coreference":
if (
not isinstance(antecedent, str)
or not antecedent
or antecedent not in document.body
or antecedent == quote
):
rejected += 1
continue
else:
antecedent = None
accepted.append(
Annotation(
annotation_type,
block_id,
quote,
float(confidence),
"llm",
antecedent,
)
)
return resolve_annotation_conflicts(accepted), rejected
def annotate_once(
annotator: Annotator,
document: SkillDocument,
signals: dict[str, Signal],
selected_passes: list[str],
) -> AnnotationResult:
allowed = allowed_annotation_types(selected_passes, signals)
payload = annotator.annotate(document, selected_passes, allowed)
if not isinstance(payload, dict):
raise AnnotationError("annotator returned a non-object payload")
annotations, rejected = validate_annotations(payload, document, allowed)
return AnnotationResult(
annotations=annotations,
used=True,
model=annotator.model_id,
accepted=len(annotations),
rejected=rejected,
)
def _proposed_count(payload: dict[str, Any]) -> int:
rewrites = payload.get("rewrites")
return len(rewrites) if isinstance(rewrites, list) else 0
def _repair_requests(
payload: dict[str, Any],
document: SkillDocument,
rejected_reasons: list[str],
indices: list[int],
) -> list[dict[str, Any]]:
rewrites = payload.get("rewrites")
if not isinstance(rewrites, list):
return []
reasons = {
index: reason
for reason in rejected_reasons
if (index := rejection_index(reason)) is not None
}
requests = []
for index in indices:
if not 0 <= index < len(rewrites) or index not in reasons:
continue
raw = rewrites[index]
if not isinstance(raw, dict):
continue
delta = replacement_literal_delta(raw, document)
requests.append(
{
"rewrite_index": index,
"validation_error": reasons[index],
"literal_delta": delta,
"required_action": {
"restore_each_literal_verbatim": delta["restore_verbatim"],
"remove_each_literal_verbatim": delta["remove_verbatim"],
"replacement_must_change": True,
"source_refs_are_locked": True,
},
"rewrite": raw,
}
)
return requests
def _merge_repair_payload(
initial: dict[str, Any],
repair: dict[str, Any],
allowed_indices: set[int],
) -> tuple[dict[str, Any], set[int], list[str]]:
if repair.get("schema_version") != PLAN_SCHEMA_VERSION:
raise ValueError("repair response has unsupported schema_version")
repairs = repair.get("repairs")
if not isinstance(repairs, list):
raise ValueError("repair response requires repairs list")
raw_rewrites = initial.get("rewrites")
if not isinstance(raw_rewrites, list):
raise ValueError("initial semantic plan requires rewrites list")
merged = dict(initial)
rewrites = list(raw_rewrites)
applied: set[int] = set()
errors: list[str] = []
for position, item in enumerate(repairs):
prefix = f"repair[{position}]"
if not isinstance(item, dict):
errors.append(f"{prefix}: item is not an object")
continue
if set(item) != {"rewrite_index", "replacement"}:
errors.append(
f"{prefix}: only rewrite_index and replacement may be returned"
)
continue
index = item.get("rewrite_index")
replacement = item.get("replacement")
if (
isinstance(index, bool)
or not isinstance(index, int)
or index not in allowed_indices
or not 0 <= index < len(rewrites)
):
errors.append(f"{prefix}: rewrite_index is not repairable")
continue
if index in applied:
errors.append(f"{prefix}: duplicate rewrite_index {index}")
continue
if not isinstance(replacement, str):
errors.append(f"{prefix}: replacement is not text")
continue
original = rewrites[index]
if (
isinstance(original, dict)
and isinstance(original.get("replacement"), str)
and replacement.strip() == original["replacement"].strip()
):
errors.append(
f"{prefix}: replacement is unchanged from the rejected rewrite"
)
continue
updated = dict(original) if isinstance(original, dict) else {}
updated["replacement"] = replacement
rewrites[index] = updated
applied.add(index)
merged["rewrites"] = rewrites
return merged, applied, errors
def plan_once(
planner: SemanticPlanner,
document: SkillDocument,
signals: dict[str, Signal],
selected_passes: list[str],
) -> SemanticPlanResult:
payload = planner.plan(document, signals, selected_passes)
if not isinstance(payload, dict):
raise AnnotationError("semantic planner returned a non-object payload")
try:
initial_units, initial_rejected_reasons = validate_semantic_plan(
payload, document
)
except ValueError as exc:
raise AnnotationError(str(exc)) from exc
units = initial_units
rejected_reasons = initial_rejected_reasons
repairable_indices = repairable_rewrite_indices(initial_rejected_reasons)
repair_attempted = False
repair_proposed = 0
repair_accepted = 0
repair_rejected = 0
repair_rejection_reasons: list[str] = []
repair_error: str | None = None
repair_transport_attempts = 0
repair_request_variant: str | None = None
repair_method = getattr(planner, "repair", None)
if repairable_indices and callable(repair_method):
repair_attempted = True
repair_requests = _repair_requests(
payload, document, initial_rejected_reasons, repairable_indices
)
repair_payload: dict[str, Any] | None = None
try:
repair_payload = repair_method(repair_requests)
if not isinstance(repair_payload, dict):
raise AnnotationError(
"semantic repair planner returned a non-object payload"
)
repairs = repair_payload.get("repairs")
repair_proposed = len(repairs) if isinstance(repairs, list) else 0
merged, repaired_indices, protocol_errors = _merge_repair_payload(
payload, repair_payload, set(repairable_indices)
)
units, rejected_reasons = validate_semantic_plan(merged, document)
final_rejected_indices = {
index
for reason in rejected_reasons
if (index := rejection_index(reason)) is not None
}
repair_accepted = len(
set(repairable_indices) - final_rejected_indices
)
repair_rejected = len(repairable_indices) - repair_accepted
repair_rejection_reasons = protocol_errors + [
reason
for reason in rejected_reasons
if rejection_index(reason) in set(repairable_indices)
]
for index in set(repairable_indices) - repaired_indices:
if index not in final_rejected_indices:
repair_accepted -= 1
repair_rejected += 1
repair_rejection_reasons.append(
f"rewrite[{index}]: repair response did not update replacement"
)
except (AnnotationError, OSError, json.JSONDecodeError, ValueError) as exc:
repair_error = str(exc)
repair_rejected = len(repairable_indices)
repair_rejection_reasons = [
reason
for reason in initial_rejected_reasons
if rejection_index(reason) in set(repairable_indices)
]
units = initial_units
rejected_reasons = initial_rejected_reasons
repair_transport_attempts = int(
getattr(planner, "last_repair_request_attempts", 1)
)
repair_request_variant = getattr(
planner, "last_repair_request_variant", None
)
initial_transport_attempts = int(getattr(planner, "last_request_attempts", 1))
return SemanticPlanResult(
units=units,
used=True,
model=planner.model_id,
transport_attempts=initial_transport_attempts,
request_variant=getattr(planner, "last_request_variant", None),
accepted=len(units),
rejected=len(rejected_reasons),
rejection_reasons=rejected_reasons,
semantic_rounds=1 + int(repair_attempted),
provider_request_count=(
initial_transport_attempts + repair_transport_attempts
),
initial_proposed=_proposed_count(payload),
initial_accepted=len(initial_units),
initial_rejected=len(initial_rejected_reasons),
initial_rejection_reasons=initial_rejected_reasons,
repair_attempted=repair_attempted,
repair_proposed=repair_proposed,
repair_accepted=max(0, repair_accepted),
repair_rejected=repair_rejected,
repair_rejection_reasons=repair_rejection_reasons,
repair_transport_attempts=repair_transport_attempts,
repair_request_variant=repair_request_variant,
repair_error=repair_error,
)
class OpenCodeAnnotator:
"""Constrained annotator using a routed OpenAI-compatible API."""
def __init__(
self,
model_id: str,
progress: Callable[[int, str], None] | None = None,
):
try:
route = resolve_model_route(model_id)
except (RuntimeError, ValueError) as exc:
raise AnnotationError(str(exc)) from exc
assert route is not None
self.model_id = route.reference.value
self.api_model_id = route.reference.model_id
base_url = route.url.removesuffix("/chat/completions").rstrip("/")
try:
from openai import OpenAI
except ImportError:
self._client = _HttpOpenAICompat(f"{base_url}/chat/completions", route.api_key)
else:
self._client = OpenAI(
base_url=base_url,
api_key=route.api_key,
max_retries=0,
timeout=90.0,
)
self.last_request_attempts = 0
self.last_request_variant: str | None = None
self.last_repair_request_attempts = 0
self.last_repair_request_variant: str | None = None
self._progress = progress
@staticmethod
def _empty_content_detail(completion: Any) -> str:
choices = getattr(completion, "choices", None) or []
choice = choices[0] if choices else None
finish_reason = getattr(choice, "finish_reason", None)
usage = getattr(completion, "usage", None)
completion_tokens = getattr(usage, "completion_tokens", None)
details = getattr(usage, "completion_tokens_details", None)
reasoning_tokens = getattr(details, "reasoning_tokens", None)
values = [
f"finish_reason={finish_reason or 'unknown'}",
f"completion_tokens={completion_tokens if completion_tokens is not None else 'unknown'}",
f"reasoning_tokens={reasoning_tokens if reasoning_tokens is not None else 'unknown'}",
]
return ", ".join(values)
def annotate(
self,
document: SkillDocument,
selected_passes: list[str],
allowed_types: set[str],
) -> dict[str, Any]:
blocks = [
{
"block_id": block.id,
"kind": block.kind,
"parent_heading": block.parent_heading,
"text": block.text,
}
for block in document.blocks
if block.kind not in {"code", "code_like", "html", "heading", "table"}
]
request = {
"selected_passes": selected_passes,
"allowed_annotation_types": sorted(allowed_types),
"blocks": blocks,
}
system_prompt = (
"You are a conservative text annotation engine. Identify only exact text "
"already present in the supplied blocks. Do not rewrite, summarize, infer "
"new rules, or return text not found verbatim in a block. Return JSON with "
'schema_version "1.0" and annotations. Each annotation must contain type, '
"block_id, quote, and confidence from 0 to 1. For coreference also return "
"antecedent_quote, copied verbatim from the input. Respect parent_heading "
"when classifying a span. Type definitions: "
+ json.dumps(
{
name: ANNOTATION_DEFINITIONS[name]
for name in sorted(allowed_types)
if name in ANNOTATION_DEFINITIONS
},
ensure_ascii=False,
sort_keys=True,
)
+ ". Return an empty list when evidence is insufficient. Prefer missing "
"a rule over labeling explanatory prose. Return at most 12 annotations."
)
create_options: dict[str, Any] = {
"model": self.api_model_id,
"temperature": 0,
"max_tokens": 1500,
"response_format": {"type": "json_object"},
"messages": [
{"role": "system", "content": system_prompt},
{
"role": "user",
"content": json.dumps(request, ensure_ascii=False, sort_keys=True),
},
],
}
# DeepSeek V4 enables long thinking by default and counts reasoning
# tokens against max_tokens. Annotation is constrained extraction, so
# thinking wastes the entire budget and can leave message.content empty.
# Keep this provider-specific option scoped to model IDs for which the
# API contract defines it; other OpenCode models must not receive it.
if self.api_model_id.startswith("deepseek-v4-"):
create_options["extra_body"] = {"thinking": {"type": "disabled"}}
try:
completion = self._client.chat.completions.create(**create_options)
except Exception as exc:
raise AnnotationError(
f"annotator API request failed ({type(exc).__name__}): {exc}"
) from exc
choices = getattr(completion, "choices", None) or []
if not choices:
raise AnnotationError("annotator returned no choices")
content = choices[0].message.content
if not isinstance(content, str) or not content.strip():
detail = self._empty_content_detail(completion)
raise AnnotationError(f"annotator returned empty content ({detail})")
try:
payload = json.loads(content)
except json.JSONDecodeError as exc:
raise AnnotationError(f"annotator returned invalid JSON: {exc}") from exc
if not isinstance(payload, dict):
raise AnnotationError("annotator returned a non-object JSON value")
return payload
def _semantic_json_request(
self,
create_options: dict[str, Any],
*,
phase: str,
progress_percent: int,
) -> dict[str, Any]:
"""Request one semantic round with provider-compatibility fallbacks."""
plain_json_options = dict(create_options)
plain_json_options.pop("response_format", None)
if self.api_model_id.startswith("deepseek-v4-"):
non_reasoning = dict(plain_json_options)
non_reasoning["reasoning_effort"] = "none"
non_reasoning["max_tokens"] = 8000
template_non_think = dict(plain_json_options)
template_non_think["max_tokens"] = 8000
template_non_think["extra_body"] = {
"chat_template_kwargs": {"thinking": False}
}
reasoning_budget = dict(plain_json_options)
reasoning_budget["max_tokens"] = 12000
variants = (
("deepseek_non_reasoning", non_reasoning),
("deepseek_template_non_think", template_non_think),
("deepseek_reasoning_budget", reasoning_budget),
)
else:
prompt_json_retry = dict(plain_json_options)
prompt_json_retry["max_tokens"] = 6000
variants = (
("json_object", create_options),
("prompt_json", plain_json_options),
("prompt_json_retry", prompt_json_retry),
)
attempts_attr = (
"last_request_attempts"
if phase == "plan"
else "last_repair_request_attempts"
)
variant_attr = (
"last_request_variant"
if phase == "plan"
else "last_repair_request_variant"
)
setattr(self, attempts_attr, 0)
setattr(self, variant_attr, None)
failures: list[str] = []
payload: dict[str, Any] | None = None
label = "semantic plan" if phase == "plan" else "semantic repair"
for attempt, (variant, options) in enumerate(variants, start=1):
setattr(self, attempts_attr, attempt)
progress = getattr(self, "_progress", None)
if progress is not None:
progress(
progress_percent,
f"LLM {label}: attempt {attempt}/{len(variants)} ({variant})",
)
try:
completion = self._client.chat.completions.create(**options)
except Exception as exc:
status = getattr(exc, "status_code", None)
message = str(exc).lower()
retryable = (
status in {408, 409, 429}
or isinstance(status, int) and status >= 500
or "timeout" in type(exc).__name__.lower()
or "connection" in type(exc).__name__.lower()
or (
status == 400
and any(
marker in message
for marker in (
"response_format",
"json_object",
"reasoning_effort",
"chat_template_kwargs",
"unsupported",
"unknown parameter",
)
)
)
)
failures.append(
f"attempt {attempt}/{variant}: "
f"{type(exc).__name__} status={status or 'unknown'}"
)
if not retryable:
detail = "; ".join(failures)
raise AnnotationError(
f"{label} API request failed after {attempt} attempt(s) "
f"for model {self.model_id} ({detail}): {exc}"
) from exc
continue
choices = getattr(completion, "choices", None) or []
if not choices:
failures.append(f"attempt {attempt}/{variant}: no choices")
continue
content = choices[0].message.content
if not isinstance(content, str) or not content.strip():
detail = self._empty_content_detail(completion)
failures.append(
f"attempt {attempt}/{variant}: empty content ({detail})"
)
if progress is not None:
progress(
progress_percent,
f"LLM {label} attempt {attempt}/{len(variants)} exhausted "
"reasoning budget; retrying",
)
continue
candidate = content.strip()
fenced = re.fullmatch(
r"```(?:json)?\s*(\{.*\})\s*```", candidate, re.I | re.S
)
if fenced:
candidate = fenced.group(1)
try:
decoded = json.loads(candidate)
except json.JSONDecodeError as exc:
failures.append(
f"attempt {attempt}/{variant}: invalid JSON ({exc})"
)
continue
if not isinstance(decoded, dict):
failures.append(f"attempt {attempt}/{variant}: non-object JSON")
continue
payload = decoded
setattr(self, variant_attr, variant)
break
if payload is None:
detail = "; ".join(failures)
raise AnnotationError(
f"{label} produced no usable JSON after {len(variants)} "
f"attempt(s) for model {self.model_id} ({detail})"
)
return payload
def plan(
self,
document: SkillDocument,
signals: dict[str, Signal],
selected_passes: list[str],
) -> dict[str, Any]:
# Tables are direct read-only evidence. Code contributes isolated
# configuration assignments and complete command-workflow closures.
blocks: list[dict[str, str]] = []
for block in document.blocks:
if block.kind in SUMMARY_SOURCE_KINDS - {"code"}:
blocks.append(
{
"block_id": block.id,
"kind": block.kind,
"parent_heading": block.parent_heading,
"text": block.text,
}
)
elif block.kind == "code":
blocks.extend(_code_configuration_evidence(block))
blocks.extend(_code_workflow_closure_evidence(block))
has_execution_contract = any(
block.get("semantic_role") == "execution_contract" for block in blocks
)
focus: list[str] = []
robustness = signals["semantic_robustness"]
if robustness.level in {"low", "medium"}:
focus.extend(
[
"canonicalize inconsistent terminology",
"make unambiguous coreferences explicit",
"consolidate duplicated or dispersed requirements",
"classify input, workflow, output, and validation content",
]
)
if signals["contextual_rule_adherence"].level in {"low", "medium"}:
focus.append("surface task-specific directives and prohibitions")
if signals["evidence_priority"].level in {"low", "medium"}:
focus.append("surface existing evidence priority rules")
if signals["ambiguity_handling"].level in {"low", "medium"}:
focus.append("surface existing scope and decision criteria")
if signals["uncertainty_calibration"].level in {"low", "medium"}:
focus.append("surface existing uncertainty boundaries")
if has_execution_contract:
focus.append(
"treat complete execution contracts as candidates before isolated "
"preflight checks or individual command parameters"
)
request = {
"profile_signals": {
name: signal.to_dict() for name, signal in sorted(signals.items())
},
"selected_passes": selected_passes,
"rewrite_focus": focus,
"blocks": blocks,
}
system_prompt = (
"You are a source-grounded Skill rewrite planner. Return JSON only with "
'schema_version "2.0" and a rewrites array of at most 8 items. Each item '
"has kind, target_section, source_refs, replacement, and confidence. "
"Allowed kinds are add_summary and replace_block. add_summary may consolidate "
"existing content under one of: definitions, critical_rules, "
"evidence_priority, scope, decision_criteria, uncertainty_rule, inputs, task, "
"output, validation, completion_criterion. Every source_ref must contain an "
"exact block_id and exact quote copied from that block. Preserve every condition, "
"exception, modality, prohibition, number, path, command, field name, and tool name. "
"Do not add domain knowledge, new requirements, defaults, fallback behavior, tools, "
"or workflow steps. For replace_block cite exactly one complete block and preserve "
"its paragraph/list shape; use it only for terminology consistency, explicit "
"coreference, or clearer equivalent wording. Use confidence >= 0.90 only when the "
"replacement is fully equivalent. add_summary leaves source content in place and "
"must retain all protected literals from its sources. Use add_summary only when it "
"materially reduces execution risk by consolidating duplicated or genuinely "
"dispersed directives, or by surfacing an explicit global contract. Do not summarize "
"an isolated background fact, rationale, example, implementation detail, or single "
"substep. Do not rephrase a concise list item or content already consolidated in a "
"clear summary/list merely to classify it. critical_rules summaries require source "
"text that explicitly prescribes an action, prohibition, condition, or required "
"invariant; never promote a descriptive fact into a rule. task summaries must state "
"the overall objective or a complete workflow, never one branch, pattern, or "
"subsection. Tables are read-only evidence: cite them only for add_summary, and "
"only when they state an explicit directive, prohibition, condition, or required "
"invariant. Code configuration evidence is limited to isolated assignments of "
"runtime paths, locations, modes, endpoints, or similar configuration values; cite "
"it only for add_summary. A code_workflow_closure is a complete command workflow "
"whose flags and configuration jointly constrain one execution step. A closure with "
"semantic_role execution_contract controls a data source, execution mode, target "
"path, or write behaviour. Treat it as one candidate rule and rank it ahead of an "
"isolated preflight check or individual parameter from the same workflow, when that "
"contract materially affects the task result. It is read-only evidence and may be "
"cited only for add_summary. When you use any member of a code_workflow_closure, "
"surface the complete operational contract: retain every relevant command, flag, "
"configured value, condition, and prohibition from that closure. Do not promote one "
"flag or one configured value as an independent rule. Never infer workflow beyond a "
"supplied closure. "
"Before returning an empty rewrites list, you MUST "
"check for an explicit execution-critical contract distributed across two or more "
"source blocks. A configured resource path plus a table or prose parameter that "
"selects a resource location is a resource-parameter binding. An operation "
"restriction plus required command flags is also such a contract. When one is "
"present, you MUST add one concise source-grounded critical_rules or "
"validation summary, even if each source block is individually an implementation "
"detail or workflow substep. For a resource-parameter binding, it MUST be an "
"explicit directive using a verb such as 'set' or 'use': bind the parameter to the "
"declared configured value, retaining both the parameter spelling and value verbatim; "
"do not leave a placeholder unbound or describe the value merely as a possible "
"path. It must preserve every relevant path, command, flag, condition, and "
"prohibition, and may cite table or code evidence. Return an empty "
"rewrites list only when no two source blocks jointly constrain execution. Do not "
"invent resource availability, infer unstated environment facts, or add fallback "
"behavior."
)
create_options: dict[str, Any] = {
"model": self.api_model_id,
"temperature": 0,
"max_tokens": 3000,
"response_format": {"type": "json_object"},
"messages": [
{"role": "system", "content": system_prompt},
{
"role": "user",
"content": json.dumps(request, ensure_ascii=False, sort_keys=True),
},
],
}
return self._semantic_json_request(
create_options, phase="plan", progress_percent=35
)
def repair(self, rejected_rewrites: list[dict[str, Any]]) -> dict[str, Any]:
"""Repair rejected replacements without changing their source or scope."""
system_prompt = (
"You perform mechanical repairs on rejected source-grounded Skill rewrites. "
'Return JSON only with schema_version "2.0" and a repairs array. Each repair '
"must contain exactly rewrite_index and replacement. For every item, obey its "
"required_action literally: every restore_each_literal_verbatim value MUST "
"appear in the new replacement, and every remove_each_literal_verbatim value "
"MUST NOT appear in it. Matching is case-insensitive, but preserve the supplied "
"spelling when practical. The repaired replacement MUST differ from the rejected "
"replacement; never copy it unchanged. Use only the locked source_refs as evidence. "
"If an added literal is absent from those exact quotes, delete that literal or the "
"unsupported clause containing it; never work around an added obligation by merely "
"removing MUST while retaining the same implicit directive. Do not add rewrites or "
"change kind, target_section, source_refs, confidence, or source meaning. Preserve "
"all other inline code, paths, commands, fields, template markers, semantic numbers, "
"obligations, prohibitions, conditions, ordering, exceptions, and scope boundaries. "
"Soft MAY/CAN wording may change only if no source fact or obligation changes. Before "
"returning JSON, self-check each replacement against its restore and remove lists. "
"Return exactly one repair for every supplied rewrite_index."
)
create_options: dict[str, Any] = {
"model": self.api_model_id,
"temperature": 0,
"max_tokens": 2000,
"response_format": {"type": "json_object"},
"messages": [
{"role": "system", "content": system_prompt},
{
"role": "user",
"content": json.dumps(
{"rejected_rewrites": rejected_rewrites},
ensure_ascii=False,
sort_keys=True,
),
},
],
}
return self._semantic_json_request(
create_options, phase="repair", progress_percent=45
)