hermes: reconcile whole-suite proposals and cap tasks at five cases
This commit is contained in:
parent
edb83b78d9
commit
7af4f4f206
@ -21,7 +21,7 @@
|
||||
"name": {
|
||||
"type": "string",
|
||||
"minLength": 1,
|
||||
"maxLength": 80
|
||||
"maxLength": 64
|
||||
},
|
||||
"description": {
|
||||
"type": "string",
|
||||
@ -30,6 +30,7 @@
|
||||
},
|
||||
"members": {
|
||||
"type": "array",
|
||||
"maxItems": 5,
|
||||
"minItems": 1,
|
||||
"items": {
|
||||
"type": "string"
|
||||
|
||||
150
docs/hermes_suite_multipass.md
Normal file
150
docs/hermes_suite_multipass.md
Normal file
@ -0,0 +1,150 @@
|
||||
# Multi-pass implementation grouping
|
||||
|
||||
The suite planner keeps the same HTTPS endpoint, request fields, authentication,
|
||||
provider permissions, generalized-data approval, and asynchronous job lifecycle.
|
||||
The required result object remains `result.groups`, with `name`, `description`,
|
||||
and `members`. Final groups now contain **one to five** cases and names are at
|
||||
most **64 characters**, unique after whitespace and case normalization.
|
||||
|
||||
## Revisions and process
|
||||
|
||||
- Configuration: `suite-v6-20260929` (HTTP compatibility identifier).
|
||||
- Policy: `implementation-five-v1-20260929`.
|
||||
- Prompt: `implementation-proximity-multipass-v3-20260929`.
|
||||
- Execution: `suite-multipass-v1-20260929`.
|
||||
|
||||
A single server-side job performs:
|
||||
|
||||
1. Proposal A on the complete suite, in stable alias order, without a size cap.
|
||||
2. Proposal B in a reproducible SHA-256 alias order, in a separate fresh CLI
|
||||
invocation. It sees no proposal A or earlier conversation. Both use exactly
|
||||
the same implementation-effort system instructions and output schema.
|
||||
3. Whole-suite reconciliation using both memberships and all original fields.
|
||||
Success criteria remain primary evidence. There is no majority voting or
|
||||
transitive closure of pairwise similarity. Natural families are still uncapped.
|
||||
4. If any reconciled family exceeds five, review every such original family using
|
||||
all case fields, distinguishing different machinery from cheap variations.
|
||||
5. One additional bounded audit of those decisions, including descendants that
|
||||
were split into smaller groups. It may reverse an unjustified split or refine
|
||||
a broad family, but cannot cross the reconciled family boundaries. There is
|
||||
no agreement loop.
|
||||
|
||||
The final two passes record explicit keep/split rationales and exact short
|
||||
source-field quotations. The code verifies that evidence belongs to an assigned
|
||||
alias and is present verbatim in that source field. This checks support provenance,
|
||||
not the correctness of a model's engineering interpretation. Unknown implementation
|
||||
details remain concise uncertainty, not invented equipment or procedures.
|
||||
|
||||
Each invocation retains six CLI turns for structured output. Model review passes
|
||||
and CLI turns are separate counters. All calls use the originally selected provider
|
||||
and pinned model. The model remains `claude-opus-4-8[1m]`, with canonical runtime
|
||||
identity checked as `claude-opus-4-8`, firstParty, medium effort, reported 1M context
|
||||
and 64K output ceiling. No new provider fallback or tools are enabled.
|
||||
|
||||
## Deterministic task sizing and names
|
||||
|
||||
For each remaining coherent large family, `k = ceil(n / 5)` and
|
||||
`base, remainder = divmod(n, k)`. The first remainder parts have base+1 cases;
|
||||
the rest have base. Thus 6 becomes 3+3, 7 becomes 4+3, 11 becomes 4+4+3, and
|
||||
14 becomes 5+5+4. No leftover singleton is created by a five-at-a-time slice.
|
||||
|
||||
The semantic review supplies internal variation sets. A deterministic packer uses
|
||||
those hints and stable complete-record hashes/aliases to fill the calculated
|
||||
sizes. Equivalent content remains distinct by alias. Hint order, source order,
|
||||
and member order do not cause gratuitous reshuffling. Different model decisions
|
||||
can still change membership between fresh jobs; determinism is conditional on the
|
||||
same natural families and variation hints, not a claim that inference is deterministic.
|
||||
|
||||
Natural names must be concise, meaningful, distinct, and at most 56 characters,
|
||||
leaving room for numbering within the 64-character final limit. Collisions or
|
||||
invalid names fail validation; the service never truncates them or adds numbers
|
||||
to unrelated families. Capacity-only parts have the exact same base name:
|
||||
`Reset recovery (1/3)`, `Reset recovery (2/3)`, `Reset recovery (3/3)`.
|
||||
Their descriptions identify a work-size division and use only the reviewed common
|
||||
work applicable to all members. No other part's specific objectives are copied.
|
||||
|
||||
Every intermediate partition and final result is checked for exact alias coverage.
|
||||
Final size, names, balanced part sizes, and campaign/suite ownership are checked
|
||||
programmatically. Incomplete or invalid work is never returned as completed.
|
||||
|
||||
## Review artifact and operational metadata
|
||||
|
||||
`GET /suite-planning/v1/jobs/<id>/result` returns an optional top-level
|
||||
**`review_summary`**, alongside the unchanged **`result`** object. It contains:
|
||||
|
||||
- `proposal_disagreements`: affected aliases, differing pair count, and both full
|
||||
membership partitions in compressed form. Group names/numbers are not compared.
|
||||
- `reconciled_families`: the first natural partition with concise rationale,
|
||||
common work, exact source-field support, and uncertainty.
|
||||
- `large_family_review` and `decision_audit`: explicit decisions for each original
|
||||
oversized membership set, including semantic-split versus coherent-keep reasons.
|
||||
- `natural_families`: the final conceptual partition before task sizing.
|
||||
- `capacity_divisions`: base names, calculated sizes, final part names, common
|
||||
implementation rationale, and remaining uncertainty.
|
||||
- `unresolved_uncertainties` and `counts`: separate natural-family/final-task counts
|
||||
and natural/final singleton counts. Size compliance is not semantic-quality evidence.
|
||||
|
||||
Review content lives only in the authorized in-memory result cache, with the same
|
||||
one-hour TTL and restart loss as results. It is absent from routine logs, status
|
||||
responses, and SQLite job metadata. It contains source-derived material and exact
|
||||
quotations, so it is not an export-safe ClickUp artifact. No new hierarchy level
|
||||
or client request fields are introduced.
|
||||
|
||||
Status/result metadata includes `policy_revision`, `execution_revision`,
|
||||
`prompt_revision`, `model_pass_count`, `natural_family_count`, `final_task_count`,
|
||||
`singleton_statistics`, aggregate usage/cost, and safe per-invocation `passes`.
|
||||
Each pass records its stage, actual model, system/schema hashes, ordering hash,
|
||||
capacity check, remaining budgets, usage, timing, and CLI diagnostics. The top-level
|
||||
`cli_diagnostics` concerns the last invocation; `turns` is the sum of reported CLI
|
||||
turns, while `model_pass_count` counts complete model review invocations.
|
||||
`execution_progress.current_pass` reports the current stage during a running job.
|
||||
|
||||
## Shared bounds and failure behavior
|
||||
|
||||
The original 900-second maximum job deadline and USD 5 CLI estimated-cost allowance
|
||||
apply to **all calls combined**, including routing time. Each invocation receives
|
||||
only the remaining time and estimated-cost allowance. Available usage/cost is
|
||||
accumulated after every call; unknown cost accounting stops further hosted calls.
|
||||
No partial partition substitutes for unfinished passes. The CLI's estimated-cost
|
||||
ceiling can overshoot within one generation, as before; this is not a guaranteed
|
||||
subscription billing limit. An observed overrun fails the job rather than allowing
|
||||
further calls or reporting successful completion.
|
||||
|
||||
Request and complete response envelopes stay bounded at 1 MiB. Every expanded pass
|
||||
is checked before dispatch, including original source, proposals/reviews, system
|
||||
instructions, and schema. The same conservative byte bound and six possible CLI
|
||||
outputs are reserved against context. Exact tokenization and output reservation
|
||||
remain unverified estimates; actual incomplete generation fails explicitly.
|
||||
Preflight checks initial and minimum reconciliation capacity; it identifies that
|
||||
actual later-pass capacity will be checked before each invocation. It does not
|
||||
claim to know model-generated proposal sizes in advance.
|
||||
|
||||
The existing 8K-context local route cannot admit the new reconciliation schema and
|
||||
instructions. Local-only requests fail with capacity errors and never start hosted
|
||||
jobs. The separate `/local-model` endpoint and its limits are unchanged. No bounds
|
||||
were increased, and no provider was switched after an error.
|
||||
|
||||
Failures include existing CLI diagnostics and a fixed `review_pass`, completed-pass
|
||||
count, and available aggregate usage. New explicit codes include `pass_capacity`,
|
||||
`pass_request_too_large`, `job_time_budget_exhausted`, `job_cost_budget_exhausted`,
|
||||
`budget_accounting_unavailable`, and bounded review/assignment validation errors.
|
||||
Raw provider messages, source text, and generated explanations are not error details.
|
||||
|
||||
The laptop submits one request with its existing scoped token, then polls the same
|
||||
job. Retry a timed-out submission using the same key/body; an intentional fresh
|
||||
attempt needs a new Idempotency-Key. Include policy/prompt/execution revisions in
|
||||
local cache provenance. No real suite is submitted by deployment checks.
|
||||
|
||||
## Verification and rollback
|
||||
|
||||
Synthetic unit tests cover the balanced-size examples, semantic subdivisions and
|
||||
a bounded reversal, differing proposals, duplicate text, name normalization and
|
||||
collisions, source support, deadline/cost exhaustion, cancellation, privacy,
|
||||
authorization, idempotency, and 14/75/363-case capacity and sizing. Native CLI
|
||||
loopback and authenticated LAN/provider acceptance measurements are recorded after
|
||||
deployment. Mocked outputs alone do not establish model quality or throughput.
|
||||
|
||||
Deploy or roll back through Git and Flux. Revert only the multi-pass commits to
|
||||
restore the preceding single-pass policy while preserving the earlier CLI diagnostic
|
||||
fix. Wait for active jobs before a worker restart because completed results are
|
||||
held in memory. No Vault credential or ingress/routing changes are required.
|
||||
@ -4,6 +4,19 @@ This optional API groups one complete campaign/suite into implementation familie
|
||||
The existing `/local-model` API, GPU allocation, serving model, and limits are unchanged.
|
||||
Deployment uses Flux; the application and workbook stay on the laptop.
|
||||
|
||||
## Current multi-pass policy
|
||||
|
||||
New jobs use [the multi-pass implementation policy](hermes_suite_multipass.md):
|
||||
two independent full-suite proposals, reconciliation, up to two bounded semantic
|
||||
reviews, and a programmatic five-case final cap. Names are at most 64 characters.
|
||||
The endpoint and request fields are unchanged. Optional top-level `review_summary`
|
||||
is returned only with the authorized result. Configuration remains `suite-v6-20260929`;
|
||||
policy, prompt, and execution revisions identify the changed grouping behavior.
|
||||
The prior single-pass acceptance results below are historical and do not measure
|
||||
the new workflow. The new process shares the existing 900-second / USD 5 estimated
|
||||
job limits across all invocations. The separate local-model endpoint is unchanged;
|
||||
the 8K local route cannot admit the new multi-pass reconciliation schema.
|
||||
|
||||
## CLI failure diagnostics update
|
||||
|
||||
Execution revision `claude-diagnostics-turns-v1-20260929` adds content-free
|
||||
|
||||
@ -8,7 +8,8 @@ import json, threading, tempfile, pathlib, subprocess, sys, time
|
||||
from http.server import HTTPServer, BaseHTTPRequestHandler
|
||||
sys.path.insert(0, '/opt/planner')
|
||||
from suite_backends import claude_environment, claude_command
|
||||
from suite_contract import SYSTEM, prompt, validate_result
|
||||
from suite_contract import SYSTEM, prompt, validate_partition
|
||||
from suite_policy import DISCOVERY_SCHEMA
|
||||
from suite_synthetic import fixture
|
||||
seen = []
|
||||
request = fixture(14)[0]
|
||||
@ -55,6 +56,7 @@ for limit in [3, 6]:
|
||||
env['ANTHROPIC_BASE_URL'] = 'http://127.0.0.1:' + str(server.server_port)
|
||||
cmd = claude_command('claude-opus-4-8', 5)
|
||||
cmd[cmd.index('--max-turns') + 1] = str(limit)
|
||||
cmd[cmd.index('--json-schema') + 1] = json.dumps(DISCOVERY_SCHEMA)
|
||||
started = time.monotonic()
|
||||
p = subprocess.run(cmd, input=source, capture_output=True, text=True, env=env, cwd=directory, timeout=45)
|
||||
events = [json.loads(line) for line in p.stdout.splitlines()]
|
||||
@ -64,7 +66,7 @@ for limit in [3, 6]:
|
||||
assert not final.get('structured_output')
|
||||
else:
|
||||
assert p.returncode == 0 and final.get('subtype') == 'success'
|
||||
validate_result(final['structured_output'], request)
|
||||
validate_partition(final['structured_output'], request)
|
||||
assert all(item['whole_input'] and item['system_preserved'] for item in seen)
|
||||
print(json.dumps({'max_turns': limit, 'exit_code': p.returncode, 'wall_seconds': round(time.monotonic()-started, 3),
|
||||
'final': {k: final.get(k) for k in ('type', 'subtype', 'is_error', 'num_turns', 'stop_reason', 'usage')},
|
||||
|
||||
130
scripts/ops/hermes_suite_multipass_transport_probe.py
Executable file
130
scripts/ops/hermes_suite_multipass_transport_probe.py
Executable file
@ -0,0 +1,130 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Run synthetic multi-pass transport checks against a loopback provider only.
|
||||
|
||||
Execute inside the planner using Python stdin. No hosted inference is performed.
|
||||
This checks full-field transmission and fresh CLI sessions, not grouping quality.
|
||||
"""
|
||||
import json
|
||||
from http.server import BaseHTTPRequestHandler, HTTPServer
|
||||
import sys
|
||||
import threading
|
||||
|
||||
sys.path.insert(0, '/opt/planner')
|
||||
import suite_backends
|
||||
from suite_contract import validate_request, validate_result
|
||||
from suite_multipass import generate, preflight_workflow
|
||||
from suite_synthetic import fixture
|
||||
|
||||
active = {}
|
||||
seen = []
|
||||
|
||||
|
||||
def strings(value):
|
||||
"""Yield nested strings for byte-exact prompt transmission assertions."""
|
||||
if isinstance(value, str):
|
||||
yield value
|
||||
elif isinstance(value, list):
|
||||
for item in value:
|
||||
yield from strings(item)
|
||||
elif isinstance(value, dict):
|
||||
for item in value.values():
|
||||
yield from strings(item)
|
||||
|
||||
|
||||
def response_value():
|
||||
"""Create schema-valid mock replies from independent synthetic expectations."""
|
||||
source = json.loads(active['input'])['suite']
|
||||
expected = fixture(len(source['cases']))[1]
|
||||
buckets = {}
|
||||
records = {c['alias']: c for c in source['cases']}
|
||||
for alias, family in expected.items():
|
||||
buckets.setdefault(family, []).append(alias)
|
||||
groups = []
|
||||
for name, members in buckets.items():
|
||||
group = {'name': name, 'description': 'Mock transport test machinery.', 'members': members}
|
||||
if active['stage'] not in ('proposal_a', 'proposal_b'):
|
||||
group.update(common_work='Mock transport test machinery.', rationale='Mock only; no semantic inference.',
|
||||
uncertainty='', variation_sets=[members], evidence=[{
|
||||
'alias': members[0], 'field': 'success_criteria',
|
||||
'quote': records[members[0]]['success_criteria'][:100]}])
|
||||
groups.append(group)
|
||||
value = {'groups': groups}
|
||||
if active['stage'] in ('large_family_review', 'decision_audit'):
|
||||
value['decisions'] = [{'source_members': g['members'], 'decision': 'keep',
|
||||
'rationale': 'Mock transport check of a coherent implementation.',
|
||||
'evidence': g['evidence']} for g in groups if len(g['members']) > 5]
|
||||
return value
|
||||
|
||||
|
||||
class Provider(BaseHTTPRequestHandler):
|
||||
"""Return bounded Anthropic-compatible SSE with a structured tool result."""
|
||||
|
||||
def log_message(self, *args):
|
||||
"""Never record headers or synthetic bodies."""
|
||||
|
||||
def do_POST(self):
|
||||
"""Check full input before emitting a fixed synthetic structured result."""
|
||||
raw = self.rfile.read(int(self.headers['Content-Length']))
|
||||
body = json.loads(raw)
|
||||
complete = active['input'] in list(strings(body.get('messages', [])))
|
||||
system = any(active['system'] in s for s in strings(body.get('system', [])))
|
||||
seen.append({'stage': active['stage'], 'bytes': len(raw), 'complete_input': complete,
|
||||
'complete_system': system, 'max_tokens': body.get('max_tokens')})
|
||||
assert complete and system
|
||||
value = response_value()
|
||||
block = {'type': 'tool_use', 'id': 'mock-output', 'name': 'StructuredOutput', 'input': {}}
|
||||
message = {'id': 'mock', 'type': 'message', 'role': 'assistant', 'model': 'claude-opus-4-8',
|
||||
'content': [], 'stop_reason': None, 'stop_sequence': None,
|
||||
'usage': {'input_tokens': 10, 'output_tokens': 0}}
|
||||
events = [
|
||||
('message_start', {'type': 'message_start', 'message': message}),
|
||||
('content_block_start', {'type': 'content_block_start', 'index': 0, 'content_block': block}),
|
||||
('content_block_delta', {'type': 'content_block_delta', 'index': 0,
|
||||
'delta': {'type': 'input_json_delta', 'partial_json': json.dumps(value)}}),
|
||||
('content_block_stop', {'type': 'content_block_stop', 'index': 0}),
|
||||
('message_delta', {'type': 'message_delta', 'delta': {'stop_reason': 'tool_use', 'stop_sequence': None},
|
||||
'usage': {'output_tokens': 5}}),
|
||||
('message_stop', {'type': 'message_stop'})]
|
||||
data = ''.join('event: ' + kind + '\ndata: ' + json.dumps(event) + '\n\n' for kind, event in events).encode()
|
||||
self.send_response(200)
|
||||
self.send_header('Content-Type', 'text/event-stream')
|
||||
self.send_header('Content-Length', str(len(data)))
|
||||
self.end_headers()
|
||||
self.wfile.write(data)
|
||||
|
||||
|
||||
def main():
|
||||
"""Run complete suites through the installed binary with fake loopback auth."""
|
||||
server = HTTPServer(('127.0.0.1', 0), Provider)
|
||||
threading.Thread(target=server.serve_forever, daemon=True).start()
|
||||
original_env = suite_backends.claude_environment
|
||||
original_generate = suite_backends.claude_generate
|
||||
|
||||
def environment(directory, token):
|
||||
env = original_env(directory, 'synthetic-unused-oauth-token')
|
||||
env['ANTHROPIC_BASE_URL'] = 'http://127.0.0.1:' + str(server.server_port)
|
||||
return env
|
||||
|
||||
def backend(request, cancel, *, invocation):
|
||||
active.clear()
|
||||
active.update(invocation)
|
||||
return original_generate(request, cancel, invocation=invocation)
|
||||
|
||||
suite_backends.claude_environment = environment
|
||||
suite_backends.claude_generate = backend
|
||||
try:
|
||||
for size in (14, 75, 363):
|
||||
seen.clear()
|
||||
request = fixture(size)[0]
|
||||
request['routing'] = {'allow_external': True, 'allowed_external_providers': ['claude']}
|
||||
request = validate_request(request, ['claude'])
|
||||
result, metadata = generate(request, preflight_workflow(request), threading.Event(), '192.168.22.8')
|
||||
validate_result(result, request)
|
||||
print(json.dumps({'size': size, 'model_passes': metadata['model_pass_count'],
|
||||
'final_tasks': len(result['groups']), 'requests': seen}), flush=True)
|
||||
finally:
|
||||
server.shutdown()
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
main()
|
||||
@ -61,6 +61,9 @@ configMapGenerator:
|
||||
- suite_jobs.py=scripts/suite_jobs.py
|
||||
- suite_backends.py=scripts/suite_backends.py
|
||||
- suite_cli_diagnostics.py=scripts/suite_cli_diagnostics.py
|
||||
- suite_policy.py=scripts/suite_policy.py
|
||||
- suite_sizing.py=scripts/suite_sizing.py
|
||||
- suite_multipass.py=scripts/suite_multipass.py
|
||||
- suite_contract.py=scripts/suite_contract.py
|
||||
- suite_synthetic.py=scripts/suite_synthetic.py
|
||||
options:
|
||||
|
||||
@ -16,6 +16,7 @@ from suite_contract import (EXECUTION_REVISION, MAX_BODY, MAX_CASES, MAX_RESULT,
|
||||
PROMPT_SHA256, REVISION, TIMEOUT, Problem, encoded,
|
||||
preflight, validate_request)
|
||||
from suite_jobs import Jobs
|
||||
from suite_policy import POLICY_REVISION
|
||||
from suite_synthetic import allowed_synthetic, fixture
|
||||
|
||||
|
||||
@ -119,10 +120,14 @@ class Handler(BaseHTTPRequestHandler):
|
||||
if method == "GET" and self.path == "/healthz":
|
||||
return self.send(200, {"status": "ready", "configuration_revision": REVISION,
|
||||
"execution_revision": EXECUTION_REVISION,
|
||||
"policy_revision": POLICY_REVISION,
|
||||
"prompt_revision": PROMPT_REVISION, "prompt_sha256": PROMPT_SHA256})
|
||||
if method == "GET" and self.path == "/v1/capabilities":
|
||||
return self.send(200, {"configuration_revision": REVISION, "models": MODELS,
|
||||
"execution_revision": EXECUTION_REVISION,
|
||||
"policy_revision": POLICY_REVISION, "max_final_group_cases": 5,
|
||||
"max_final_name_characters": 64, "model_passes": {"minimum": 3, "maximum": 5},
|
||||
"review_summary_location": "GET /v1/jobs/<id>/result: top-level review_summary",
|
||||
"prompt_revision": PROMPT_REVISION, "prompt_sha256": PROMPT_SHA256,
|
||||
"allowed_external_providers": providers,
|
||||
"external_scope": "generalized_claude_and_exact_synthetic_fixtures"
|
||||
|
||||
@ -64,16 +64,17 @@ def switchyard_decision(provider):
|
||||
return expected
|
||||
|
||||
|
||||
def local_generate(request, cancel, client_ip):
|
||||
def local_generate(request, cancel, client_ip, *, invocation=None):
|
||||
"""Reuse the unchanged RTX API and its pinned-model, local-only safeguards."""
|
||||
headers = {"Authorization": "Bearer " + Path("/vault/secrets/local-token").read_text().strip(),
|
||||
"X-Forwarded-For": client_ip}
|
||||
schema = json.loads(encoded(SCHEMA))
|
||||
schema = json.loads(encoded(invocation["schema"] if invocation else SCHEMA))
|
||||
# Constrain the local decoder to aliases, excluding copied case descriptions.
|
||||
schema["properties"]["groups"]["items"]["properties"]["members"]["items"]["enum"] = [
|
||||
case["alias"] for case in request["cases"]]
|
||||
response = post(LOCAL + "/api/generate", {
|
||||
"model": MODELS["local"]["model"], "prompt": SYSTEM + "\n" + prompt(request),
|
||||
"model": MODELS["local"]["model"],
|
||||
"prompt": invocation["system"] + "\n" + invocation["input"] if invocation else SYSTEM + "\n" + prompt(request),
|
||||
"stream": False, "format": schema,
|
||||
"options": {"num_predict": 2048, "temperature": 0, "seed": 0}},
|
||||
request["execution"]["max_seconds"], headers)
|
||||
@ -203,7 +204,7 @@ def parse_claude(raw, expected_model, **process_info):
|
||||
"turns": diagnostics["turns"], "cli_diagnostics": diagnostics}
|
||||
|
||||
|
||||
def claude_generate(request, cancel):
|
||||
def claude_generate(request, cancel, *, invocation=None):
|
||||
"""Run one fresh job in tmpfs; input, output, configuration and caches expire together."""
|
||||
token = Path("/vault/secrets/claude-token").read_text().strip()
|
||||
if not token:
|
||||
@ -211,13 +212,17 @@ def claude_generate(request, cancel):
|
||||
model = MODELS["claude"]["model"]
|
||||
with tempfile.TemporaryDirectory(prefix="suite-", dir="/jobs") as directory:
|
||||
root = Path(directory)
|
||||
(root / "input").write_text(prompt(request))
|
||||
(root / "input").write_text(invocation["input"] if invocation else prompt(request))
|
||||
command = claude_command(model, request["execution"]["max_cost_usd"])
|
||||
if invocation:
|
||||
command[command.index("--system-prompt") + 1] = invocation["system"]
|
||||
command[command.index("--json-schema") + 1] = encoded(invocation["schema"]).decode()
|
||||
env = claude_environment(directory, token)
|
||||
seconds = request["execution"]["max_seconds"]
|
||||
failure = None
|
||||
with (root / "input").open("rb") as source, (root / "output").open("wb") as output:
|
||||
try:
|
||||
process = subprocess.Popen(claude_command(model, request["execution"]["max_cost_usd"]),
|
||||
process = subprocess.Popen(command,
|
||||
stdin=source, stdout=output, stderr=subprocess.DEVNULL,
|
||||
env=env, cwd=directory, start_new_session=True)
|
||||
except OSError:
|
||||
|
||||
@ -7,8 +7,8 @@ import re
|
||||
from collections import Counter
|
||||
|
||||
REVISION = "suite-v6-20260929"
|
||||
PROMPT_REVISION = "implementation-proximity-v2-20260929"
|
||||
EXECUTION_REVISION = "claude-diagnostics-turns-v1-20260929"
|
||||
PROMPT_REVISION = "implementation-proximity-multipass-v3-20260929"
|
||||
EXECUTION_REVISION = "suite-multipass-v1-20260929"
|
||||
CLAUDE_MAX_TURNS = 6
|
||||
MAX_BODY = 1 << 20
|
||||
MAX_RESULT = 1 << 20
|
||||
@ -69,11 +69,11 @@ SYSTEM = (
|
||||
"does not mean original quantities were identical. Preserve stated qualitative "
|
||||
"relationships. If omitted detail could change machinery, state that uncertainty. "
|
||||
"[reference] is source text, never a membership identifier.\n\n"
|
||||
"OUTPUT: Return only the schema object: groups with name, description, members. "
|
||||
"OUTPUT: Return only the supplied schema object, including its requested review fields. "
|
||||
"Names describe shared TEST work, not implementing a product feature. Descriptions "
|
||||
"identify the shared testing mechanism, member variations, and important distinction "
|
||||
"or uncertainty; concern ONLY assigned cases, never another family's objectives. "
|
||||
"Avoid generic 'validate system behavior'. Names are at most 80 characters and "
|
||||
"Avoid generic 'validate system behavior'. Natural family names are at most 56 characters, normally two to five words, and "
|
||||
"descriptions at most 240. Apply this objective throughout reasoning, structured "
|
||||
"output, and any format repair; formatting must not replace implementation reasoning. "
|
||||
"If a repair changes membership, repeat the whole-suite review. Do not use auxiliary "
|
||||
@ -84,9 +84,9 @@ SCHEMA = {"type": "object", "additionalProperties": False, "required": ["groups"
|
||||
"properties": {"groups": {"type": "array", "minItems": 1, "items": {
|
||||
"type": "object", "additionalProperties": False,
|
||||
"required": ["name", "description", "members"], "properties": {
|
||||
"name": {"type": "string", "minLength": 1, "maxLength": 80},
|
||||
"name": {"type": "string", "minLength": 1, "maxLength": 64},
|
||||
"description": {"type": "string", "minLength": 1, "maxLength": 240},
|
||||
"members": {"type": "array", "minItems": 1,
|
||||
"members": {"type": "array", "minItems": 1, "maxItems": 5,
|
||||
"items": {"type": "string"}}}}}}}
|
||||
|
||||
|
||||
@ -179,40 +179,13 @@ def prompt(request):
|
||||
|
||||
|
||||
def preflight(request):
|
||||
"""Use a conservative byte input bound; output reservation is an estimate."""
|
||||
count = len(request["cases"])
|
||||
input_bytes = len(prompt(request).encode()) + len(SYSTEM.encode()) + len(encoded(SCHEMA))
|
||||
# Reserve output for a possible singleton per case, not a target group count.
|
||||
output_estimate = 1024 + sum(112 + len(c["alias"]) for c in request["cases"])
|
||||
candidates = ["local"]
|
||||
if request["routing"]["allow_external"]:
|
||||
candidates += request["routing"]["allowed_external_providers"]
|
||||
reasons = {}
|
||||
for provider in candidates:
|
||||
model = MODELS[provider]
|
||||
if not model["enabled"]:
|
||||
reasons[provider] = "unverified_output_capacity"
|
||||
continue
|
||||
reserve = output_estimate + (8192 if provider == "claude" else 0)
|
||||
schema_extra = len(encoded([c["alias"] for c in request["cases"]])) + 16 if provider == "local" else 0
|
||||
total_input = input_bytes + schema_extra
|
||||
if reserve > model["output"] or total_input + model["overhead"] + model["output"] * model.get("max_turns", 1) > model["context"]:
|
||||
reasons[provider] = "capacity"
|
||||
continue
|
||||
return {"provider": provider, **model, "configuration_revision": REVISION,
|
||||
"execution_revision": EXECUTION_REVISION,
|
||||
"prompt_revision": PROMPT_REVISION, "prompt_sha256": PROMPT_SHA256,
|
||||
"input_bytes": total_input, "input_token_count": None,
|
||||
"input_token_bound": total_input + model["overhead"],
|
||||
"input_count_method": "UTF-8 byte upper bound plus reserved harness overhead; not a tokenizer",
|
||||
"output_reservation_tokens": reserve, "output_reservation_verified": False,
|
||||
"case_count": count, "source_sha256": digest(request["cases"])}
|
||||
raise Problem("capacity_or_unsupported_backend", 422, candidates=reasons,
|
||||
input_bytes=input_bytes, output_estimate=output_estimate)
|
||||
"""Select one permitted provider that can admit the multi-pass policy."""
|
||||
from suite_multipass import preflight_workflow
|
||||
return preflight_workflow(request)
|
||||
|
||||
|
||||
def validate_result(result, request):
|
||||
"""Validate shape and exact alias coverage independently of model claims."""
|
||||
def validate_partition(result, request, *, name_limit=64, group_limit=None, unique_names=True):
|
||||
"""Validate exact whole-suite membership; natural discovery may be uncapped."""
|
||||
if not isinstance(result, dict) or set(result) != {"groups"}:
|
||||
raise Problem("invalid_json_result", 502)
|
||||
groups = result["groups"]
|
||||
@ -223,18 +196,25 @@ def validate_result(result, request):
|
||||
for group in groups:
|
||||
if not isinstance(group, dict) or set(group) != {"name", "description", "members"}:
|
||||
raise Problem("invalid_json_result", 502)
|
||||
for field, limit in (("name", 80), ("description", 240)):
|
||||
if type(group[field]) is not str or not 1 <= len(group[field].strip()) <= limit:
|
||||
for field, limit in (("name", name_limit), ("description", 240)):
|
||||
if type(group[field]) is not str or not group[field].strip() or len(group[field]) > limit:
|
||||
raise Problem("invalid_json_result", 502)
|
||||
if group["name"].casefold().strip() in names:
|
||||
if unique_names and " ".join(group["name"].split()).casefold() in names:
|
||||
raise Problem("duplicate_family_name", 502)
|
||||
names.add(group["name"].casefold().strip())
|
||||
names.add(" ".join(group["name"].split()).casefold())
|
||||
members = group["members"]
|
||||
if not isinstance(members, list) or not members or any(type(a) is not str for a in members):
|
||||
raise Problem("invalid_case_assignments", 502)
|
||||
if group_limit is not None and len(members) > group_limit:
|
||||
raise Problem("group_size_limit", 502)
|
||||
found.extend(members)
|
||||
if Counter(found) != Counter(c["alias"] for c in request["cases"]):
|
||||
raise Problem("invalid_case_assignments", 502)
|
||||
if len(encoded(result)) > MAX_RESULT:
|
||||
raise Problem("response_too_large", 502)
|
||||
return result
|
||||
|
||||
|
||||
def validate_result(result, request):
|
||||
"""Enforce the final one-to-five-case policy independently of the model."""
|
||||
return validate_partition(result, request, name_limit=64, group_limit=5)
|
||||
|
||||
@ -8,8 +8,10 @@ import time
|
||||
import uuid
|
||||
|
||||
import suite_backends
|
||||
import suite_multipass
|
||||
from suite_contract import (EXECUTION_REVISION, PROMPT_REVISION, PROMPT_SHA256, Problem, RESULT_TTL,
|
||||
REVISION, digest, validate_result)
|
||||
MAX_RESULT, REVISION, digest, encoded, validate_result)
|
||||
from suite_policy import POLICY_REVISION
|
||||
|
||||
IDEMPOTENCY_TTL = 7 * 86400
|
||||
|
||||
@ -38,7 +40,8 @@ class Jobs:
|
||||
|
||||
def _prune(self):
|
||||
now = time.time()
|
||||
for job_id, (expires, _) in list(self.results.items()):
|
||||
for job_id, retained in list(self.results.items()):
|
||||
expires = retained[0]
|
||||
if expires < now:
|
||||
del self.results[job_id]
|
||||
self.db.execute("DELETE FROM jobs WHERE created < ?", (now - IDEMPOTENCY_TTL,))
|
||||
@ -57,6 +60,8 @@ class Jobs:
|
||||
if job_id not in self.results:
|
||||
raise Problem("result_expired_or_worker_restarted", 410)
|
||||
document["result"] = self.results[job_id][1]
|
||||
if len(self.results[job_id]) > 2:
|
||||
document["review_summary"] = self.results[job_id][2]
|
||||
return document
|
||||
|
||||
def submit(self, owner, key, request, selected, client_ip, launch=True):
|
||||
@ -78,6 +83,7 @@ class Jobs:
|
||||
document = {"job_id": job_id, "status": "accepted",
|
||||
"configuration_revision": REVISION, "routing": request["routing"],
|
||||
"execution_revision": EXECUTION_REVISION,
|
||||
"policy_revision": POLICY_REVISION,
|
||||
"prompt_revision": PROMPT_REVISION, "prompt_sha256": PROMPT_SHA256,
|
||||
"selection": selected, "attempted_destinations": [],
|
||||
"compaction": None, "truncation": None, "usage": None,
|
||||
@ -104,7 +110,7 @@ class Jobs:
|
||||
return document
|
||||
|
||||
def run(self, job_id, owner, request, selected, client_ip):
|
||||
"""Authorize a fixed Switchyard decision, then make one inference attempt."""
|
||||
"""Authorize a fixed route, then run one bounded multi-pass suite job."""
|
||||
started = time.monotonic()
|
||||
provider = selected["provider"]
|
||||
document = self.get(job_id, owner)
|
||||
@ -131,10 +137,12 @@ class Jobs:
|
||||
self._save(job_id, document)
|
||||
print(json.dumps({"job_id": job_id, "event": "inference_attempt",
|
||||
"provider": provider, "model": selected["model"]}), flush=True)
|
||||
if provider == "local":
|
||||
result, metadata = suite_backends.local_generate(effective, event, client_ip)
|
||||
else:
|
||||
result, metadata = suite_backends.claude_generate(effective, event)
|
||||
def progress(value):
|
||||
document["execution_progress"] = value
|
||||
with self.lock:
|
||||
self._save(job_id, document)
|
||||
result, metadata = suite_multipass.generate(effective, selected, event, client_ip, progress)
|
||||
review = metadata.pop("review_summary")
|
||||
document.update(metadata)
|
||||
try:
|
||||
validate_result(result, request)
|
||||
@ -144,16 +152,21 @@ class Jobs:
|
||||
if event.is_set():
|
||||
raise Problem("cancelled", 409)
|
||||
document.update(metadata, status="completed")
|
||||
if len(encoded({**document, "result": result, "review_summary": review})) > MAX_RESULT:
|
||||
raise Problem("response_too_large", 502)
|
||||
with self.lock:
|
||||
if len(self.results) >= 128:
|
||||
del self.results[min(self.results, key=lambda key: self.results[key][0])]
|
||||
self.results[job_id] = (time.time() + RESULT_TTL, result)
|
||||
self.results[job_id] = (time.time() + RESULT_TTL, result, review)
|
||||
except Problem as exc:
|
||||
if "final_event_seen" in exc.details:
|
||||
document.update(cli_diagnostics=exc.details,
|
||||
usage=exc.details.get("usage"), turns=exc.details.get("turns"),
|
||||
duration_api_ms=exc.details.get("duration_api_ms"),
|
||||
cost_usd_estimate=exc.details.get("cost_usd_estimate"))
|
||||
if "aggregate_usage" in exc.details:
|
||||
document.update(usage=exc.details["aggregate_usage"],
|
||||
cost_usd_estimate=exc.details.get("aggregate_cost_usd_estimate"))
|
||||
document.update(status="cancelled" if exc.code == "cancelled" else "failed",
|
||||
**exc.document())
|
||||
except Exception:
|
||||
|
||||
211
services/hermes/scripts/suite_multipass.py
Normal file
211
services/hermes/scripts/suite_multipass.py
Normal file
@ -0,0 +1,211 @@
|
||||
"""Bounded whole-suite orchestration with one shared deadline and cost allowance."""
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import math
|
||||
import time
|
||||
|
||||
import suite_backends
|
||||
from suite_contract import (EXECUTION_REVISION, MAX_BODY, MAX_RESULT, MODELS, PROMPT_REVISION,
|
||||
PROMPT_SHA256, REVISION, Problem, digest, encoded, validate_partition)
|
||||
from suite_policy import (BASE_NAME_LIMIT, MAX_GROUP, POLICY_REVISION, invocation,
|
||||
validate_natural)
|
||||
from suite_sizing import cap_families, review_summary
|
||||
|
||||
MAX_PASSES = 5
|
||||
|
||||
|
||||
def ordered_request(request, alternative=False):
|
||||
"""Use stable independent orders while retaining every complete case record."""
|
||||
def key(case):
|
||||
alias = case["alias"]
|
||||
return hashlib.sha256((POLICY_REVISION + ":order-b:" + alias).encode()).hexdigest() if alternative else alias
|
||||
cases = sorted(request["cases"], key=key)
|
||||
if alternative and len(cases) > 1 and [c["alias"] for c in cases] == sorted(c["alias"] for c in cases):
|
||||
cases.reverse()
|
||||
return {**request, "cases": cases}
|
||||
|
||||
|
||||
def capacity(call, provider, count, family_count=None):
|
||||
"""Check the actual complete pass, including proposal/review input and schema."""
|
||||
model = MODELS[provider]
|
||||
input_bytes = len(call["input"].encode()) + len(call["system"].encode()) + len(encoded(call["schema"]))
|
||||
if input_bytes > MAX_BODY:
|
||||
raise Problem("pass_request_too_large", 413, review_pass=call["stage"], input_bytes=input_bytes)
|
||||
if family_count is None:
|
||||
estimate = 1024 + 128 * count
|
||||
else:
|
||||
estimate = 1024 + 48 * count + 384 * family_count
|
||||
reserve = estimate + (8192 if provider == "claude" else 0)
|
||||
bound = input_bytes + model["overhead"]
|
||||
if reserve > model["output"] or bound + model["output"] * model.get("max_turns", 1) > model["context"]:
|
||||
raise Problem("pass_capacity", 422, review_pass=call["stage"], input_bytes=input_bytes,
|
||||
output_reservation_tokens=reserve)
|
||||
return {"input_bytes": input_bytes, "input_token_bound": bound, "input_token_count": None,
|
||||
"output_reservation_tokens": reserve, "output_reservation_verified": False,
|
||||
"input_count_method": "Complete UTF-8 input/system/schema byte bound plus harness overhead; not a tokenizer"}
|
||||
|
||||
|
||||
def preflight_workflow(request):
|
||||
"""Choose once, respecting permissions; recheck actual expanded inputs at every pass."""
|
||||
count = len(request["cases"])
|
||||
candidates = ["local"]
|
||||
if request["routing"]["allow_external"]:
|
||||
candidates += request["routing"]["allowed_external_providers"]
|
||||
reasons = {}
|
||||
for provider in candidates:
|
||||
model = MODELS[provider]
|
||||
if not model["enabled"]:
|
||||
reasons[provider] = "unverified_output_capacity"
|
||||
continue
|
||||
try:
|
||||
initial = capacity(invocation("proposal_a", request), provider, count)
|
||||
# Even a minimum-size reconciliation must fit; its real proposals and
|
||||
# every subsequent request get checked again before any model call.
|
||||
capacity(invocation("reconciliation", request), provider, count, 1)
|
||||
except Problem as exc:
|
||||
reasons[provider] = exc.code
|
||||
continue
|
||||
return {"provider": provider, **model, **initial, "configuration_revision": REVISION,
|
||||
"prompt_revision": PROMPT_REVISION, "prompt_sha256": PROMPT_SHA256,
|
||||
"execution_revision": EXECUTION_REVISION, "policy_revision": POLICY_REVISION,
|
||||
"case_count": count, "source_sha256": digest(request["cases"]),
|
||||
"minimum_model_passes": 3, "maximum_model_passes": MAX_PASSES,
|
||||
"max_final_group_cases": MAX_GROUP, "max_final_name_characters": 64,
|
||||
"later_pass_capacity_verified": False, "later_pass_checks": "before_each_invocation"}
|
||||
raise Problem("capacity_or_unsupported_backend", 422, candidates=reasons)
|
||||
|
||||
|
||||
def sum_usage(records):
|
||||
"""Aggregate measured counts only; keep an unknown counter unknown."""
|
||||
values = [r.get("usage") for r in records]
|
||||
if not values or any(not isinstance(v, dict) for v in values):
|
||||
return None
|
||||
keys = set().union(*(v.keys() for v in values))
|
||||
result = {}
|
||||
for key in keys:
|
||||
counts = [v.get(key) for v in values]
|
||||
if all(type(n) in (int, float) for n in counts):
|
||||
result[key] = sum(counts)
|
||||
elif all(isinstance(n, dict) for n in counts):
|
||||
result[key] = sum_usage([{"usage": n} for n in counts])
|
||||
else:
|
||||
result[key] = None
|
||||
return result
|
||||
|
||||
|
||||
class Workflow:
|
||||
"""Keep content in memory and allow only one pinned provider across model passes."""
|
||||
|
||||
def __init__(self, request, selected, cancel, client_ip, progress=None):
|
||||
self.request, self.provider, self.cancel = request, selected["provider"], cancel
|
||||
self.client_ip, self.progress = client_ip, progress or (lambda _: None)
|
||||
self.deadline = time.monotonic() + request["execution"]["max_seconds"]
|
||||
self.cost_limit, self.spent = request["execution"]["max_cost_usd"], 0.0
|
||||
self.records = []
|
||||
self.last_metadata = {}
|
||||
self.stage = "proposal_a"
|
||||
|
||||
def checkpoint(self):
|
||||
"""Fail closed on cancellation or exhaustion of the shared job budgets."""
|
||||
if self.cancel.is_set():
|
||||
raise Problem("cancelled", 409)
|
||||
if time.monotonic() >= self.deadline:
|
||||
raise Problem("job_time_budget_exhausted", 504)
|
||||
if self.spent > self.cost_limit:
|
||||
raise Problem("job_cost_budget_exhausted", 502)
|
||||
|
||||
def call(self, stage, source, context=None, family_count=None):
|
||||
"""Run one fresh complete-input invocation, retaining metadata before validation."""
|
||||
self.stage = stage
|
||||
self.checkpoint()
|
||||
if len(self.records) >= MAX_PASSES:
|
||||
raise Problem("model_pass_limit", 502)
|
||||
call = invocation(stage, source, context)
|
||||
limits = capacity(call, self.provider, len(source["cases"]), family_count)
|
||||
remaining_seconds = self.deadline - time.monotonic()
|
||||
remaining_cost = self.cost_limit - self.spent
|
||||
if remaining_cost <= 0:
|
||||
raise Problem("job_cost_budget_exhausted", 502)
|
||||
effective = {**source, "execution": {**source["execution"],
|
||||
"max_seconds": remaining_seconds, "max_cost_usd": remaining_cost}}
|
||||
self.progress({"current_pass": stage, "completed_model_passes": len(self.records), "passes": self.records})
|
||||
started = time.monotonic()
|
||||
if self.provider == "claude":
|
||||
value, metadata = suite_backends.claude_generate(effective, self.cancel, invocation=call)
|
||||
elif self.provider == "local":
|
||||
value, metadata = suite_backends.local_generate(effective, self.cancel, self.client_ip, invocation=call)
|
||||
else:
|
||||
raise Problem("unsupported_backend", 422)
|
||||
cost = metadata.get("cost_usd_estimate") if self.provider == "claude" else 0.0
|
||||
record = {"stage": stage, "provider": self.provider, "model": MODELS[self.provider]["model"],
|
||||
"wall_seconds": round(time.monotonic() - started, 3), **limits,
|
||||
"system_sha256": hashlib.sha256(call["system"].encode()).hexdigest(),
|
||||
"schema_sha256": digest(call["schema"]),
|
||||
"case_order_sha256": digest([c["alias"] for c in source["cases"]]),
|
||||
"allocated_seconds": remaining_seconds, "allocated_cost_usd": remaining_cost,
|
||||
"usage": metadata.get("usage"), "cli_turns": metadata.get("turns"),
|
||||
"duration_api_ms": metadata.get("duration_api_ms"),
|
||||
"cost_usd_estimate": cost, "cli_diagnostics": metadata.get("cli_diagnostics")}
|
||||
self.records.append(record)
|
||||
self.last_metadata = metadata
|
||||
if type(cost) not in (int, float) or not math.isfinite(cost) or cost < 0:
|
||||
raise Problem("budget_accounting_unavailable", 502)
|
||||
self.spent += cost
|
||||
self.checkpoint()
|
||||
self.progress({"current_pass": stage, "completed_model_passes": len(self.records), "passes": self.records})
|
||||
return value
|
||||
|
||||
def execute(self):
|
||||
"""Perform independent discovery, reconciliation, and bounded semantic reviews."""
|
||||
source = ordered_request(self.request)
|
||||
a = self.call("proposal_a", source)
|
||||
validate_partition(a, source, name_limit=BASE_NAME_LIMIT, unique_names=False)
|
||||
b = self.call("proposal_b", ordered_request(source, True))
|
||||
validate_partition(b, source, name_limit=BASE_NAME_LIMIT, unique_names=False)
|
||||
count = max(len(a["groups"]), len(b["groups"]))
|
||||
reconciled = self.call("reconciliation", source, {"proposal_a": a, "proposal_b": b}, count)
|
||||
validate_natural(reconciled, source)
|
||||
large = [g for g in reconciled["groups"] if len(g["members"]) > MAX_GROUP]
|
||||
reviewed = audited = None
|
||||
if large:
|
||||
reviewed = self.call("large_family_review", source,
|
||||
{"natural_partition": reconciled, "oversized_families": [g["members"] for g in large]},
|
||||
len(reconciled["groups"]))
|
||||
validate_natural(reviewed, source, reconciled["groups"])
|
||||
audited = self.call("decision_audit", source,
|
||||
{"original_partition": reconciled, "reviewed_partition": reviewed,
|
||||
"oversized_families": [g["members"] for g in large]},
|
||||
len(reviewed["groups"]))
|
||||
validate_natural(audited, source, reconciled["groups"])
|
||||
natural = audited or reconciled
|
||||
final, divisions = cap_families(natural, source)
|
||||
review = review_summary(a, b, reconciled, reviewed, audited, final, divisions)
|
||||
self.checkpoint()
|
||||
metadata = {**self.last_metadata, "passes": self.records,
|
||||
"model_pass_count": len(self.records), "usage": sum_usage(self.records),
|
||||
"cli_diagnostics_scope": "last_model_pass",
|
||||
"duration_api_ms": sum(r["duration_api_ms"] for r in self.records) if all(type(r["duration_api_ms"]) in (int, float) for r in self.records) else None,
|
||||
"cost_usd_estimate": round(self.spent, 8),
|
||||
"turns": sum(r["cli_turns"] for r in self.records) if all(type(r["cli_turns"]) is int for r in self.records) else None,
|
||||
"natural_family_count": len(natural["groups"]), "final_task_count": len(final["groups"]),
|
||||
"singleton_statistics": {k: v for k, v in review["counts"].items() if "singleton" in k},
|
||||
"policy_revision": POLICY_REVISION, "review_summary": review}
|
||||
if len(encoded({"result": final, **metadata})) > MAX_RESULT - 16384:
|
||||
raise Problem("response_too_large", 502)
|
||||
return final, metadata
|
||||
|
||||
|
||||
def generate(request, selected, cancel, client_ip, progress=None):
|
||||
"""Never expose partial partitions as completion or broaden a failed route."""
|
||||
workflow = Workflow(request, selected, cancel, client_ip, progress)
|
||||
try:
|
||||
return workflow.execute()
|
||||
except Problem as exc:
|
||||
details = {"failure_stage": "multi_pass_orchestration", **exc.details, "review_pass": workflow.stage,
|
||||
"completed_model_passes": len(workflow.records), "passes": workflow.records,
|
||||
"aggregate_usage": sum_usage(workflow.records + ([{"usage": exc.details["usage"]}] if isinstance(exc.details.get("usage"), dict) else [])),
|
||||
"aggregate_cost_usd_estimate": None if exc.code == "budget_accounting_unavailable" else round(workflow.spent, 8)}
|
||||
if exc.details.get("cost_usd_estimate") is not None:
|
||||
details["aggregate_cost_usd_estimate"] += exc.details["cost_usd_estimate"]
|
||||
raise Problem(exc.code, exc.status, **details) from None
|
||||
201
services/hermes/scripts/suite_policy.py
Normal file
201
services/hermes/scripts/suite_policy.py
Normal file
@ -0,0 +1,201 @@
|
||||
"""Internal multi-pass schemas, prompts, and validation; no client fields added."""
|
||||
from __future__ import annotations
|
||||
|
||||
import copy
|
||||
import re
|
||||
from collections import Counter
|
||||
|
||||
from suite_contract import FIELDS, Problem, SCHEMA, SYSTEM, encoded, validate_partition
|
||||
|
||||
POLICY_REVISION = "implementation-five-v1-20260929"
|
||||
MAX_GROUP = 5
|
||||
BASE_NAME_LIMIT = 56 # Leaves room for ' (80/80)' at the 400-case input limit.
|
||||
STAGES = ("proposal_a", "proposal_b", "reconciliation", "large_family_review", "decision_audit")
|
||||
|
||||
DISCOVERY = (
|
||||
"Discover natural implementation families across this COMPLETE suite, without any "
|
||||
"task-size cap. Do not split a coherent family for a numerical quota. Names normally "
|
||||
"use two to five words describing test machinery; never repeat the hierarchy. "
|
||||
"Use distinct meaningful mechanism names, not numbering to hide collisions. "
|
||||
"Use StructuredOutput directly when available; do not first emit a prose draft."
|
||||
)
|
||||
RECONCILE = (
|
||||
"Independently reassess two complete candidate partitions using ALL original cases. "
|
||||
"Compare memberships, never proposal names or group numbers. Resolve disagreements "
|
||||
"against success_criteria and shared test machinery. Do not vote, count majority "
|
||||
"co-memberships, or take transitive connected components. Agreement is not proof. "
|
||||
"Return one coherent NATURAL partition, with no size cap. Review cross-family merges "
|
||||
"and unjustified splits. Record concise engineering rationales and unresolved "
|
||||
"uncertainty, not private reasoning traces. Evidence must quote a short exact substring "
|
||||
"from a named source field of an assigned alias. Prefer success_criteria when present. "
|
||||
"common_work must describe ONLY machinery applicable to EVERY assigned member, "
|
||||
"without enumerating objectives/variations that apply to only some members. "
|
||||
"variation_sets partition each family's aliases into closely related variations; "
|
||||
"these are ordering hints, not additional task levels or semantic families. "
|
||||
"Keep identical case content together in variation_sets. Never impose five-member "
|
||||
"limits on natural discovery. Use unique mechanism names, normally two to five words "
|
||||
"and at most 56 characters, without numbered suffixes. Rewrite long/colliding names "
|
||||
"meaningfully; do not blindly truncate or add arbitrary identifiers. "
|
||||
"Use StructuredOutput directly when available; no prose draft."
|
||||
)
|
||||
REVIEW = (
|
||||
"Review EVERY listed oversized family against its original source case fields, "
|
||||
"especially success_criteria. Decide whether different control/measurement/evidence "
|
||||
"mechanisms warrant separate natural families with meaningful distinct names. "
|
||||
"Do not split by source order, quota, nominal/fault labels alone, or cheap assertion "
|
||||
"variations. If coherent, explicitly explain the common machinery and why differences "
|
||||
"are inexpensive variations. Preserve that coherent natural family even if it has "
|
||||
"more than five members: deterministic work sizing happens later. Return a complete "
|
||||
"natural partition of the whole suite and one explicit decision per reviewed family. "
|
||||
"Decisions refer to the exact reviewed source_members, not names. Splits must partition "
|
||||
"that family's members; do not cross existing family boundaries during this review. "
|
||||
"Retain and improve concise source-field support and uncertainty."
|
||||
)
|
||||
AUDIT = (
|
||||
"Perform ONE bounded additional review of the previous large-family decisions. "
|
||||
"Using all original fields, challenge unjustified semantic splits, unjustified broad "
|
||||
"merges, and assertions that would need different machinery. You may correct the "
|
||||
"decision by merging/splitting only within each ORIGINAL oversized family. Re-examine "
|
||||
"all its descendants even if each is now small. Return the complete final NATURAL "
|
||||
"partition and final decisions for every original oversized family. There will be "
|
||||
"no loop seeking agreement; explicitly retain unresolved uncertainty. Preserve "
|
||||
"coherent large families for deterministic capacity division afterwards."
|
||||
)
|
||||
|
||||
|
||||
def text_schema(limit):
|
||||
"""Return a bounded string schema for internal model output."""
|
||||
return {"type": "string", "minLength": 1, "maxLength": limit}
|
||||
|
||||
|
||||
ALIAS_LIST = {"type": "array", "minItems": 1, "items": text_schema(53)}
|
||||
EVIDENCE = {"type": "array", "minItems": 1, "maxItems": 3, "items": {
|
||||
"type": "object", "additionalProperties": False,
|
||||
"required": ["alias", "field", "quote"], "properties": {
|
||||
"alias": text_schema(53), "field": {"type": "string", "enum": sorted(FIELDS)},
|
||||
"quote": text_schema(160)}}}
|
||||
DISCOVERY_SCHEMA = copy.deepcopy(SCHEMA)
|
||||
DISCOVERY_SCHEMA["properties"]["groups"]["items"]["properties"]["members"].pop("maxItems", None)
|
||||
DISCOVERY_SCHEMA["properties"]["groups"]["items"]["properties"]["name"]["maxLength"] = BASE_NAME_LIMIT
|
||||
NATURAL_SCHEMA = copy.deepcopy(DISCOVERY_SCHEMA)
|
||||
GROUP_PROPERTIES = NATURAL_SCHEMA["properties"]["groups"]["items"]["properties"]
|
||||
GROUP_PROPERTIES["name"]["maxLength"] = BASE_NAME_LIMIT
|
||||
GROUP_PROPERTIES.update({
|
||||
"common_work": text_schema(160), "rationale": text_schema(320),
|
||||
"uncertainty": {"type": "string", "maxLength": 240}, "evidence": EVIDENCE,
|
||||
"variation_sets": {"type": "array", "minItems": 1, "items": ALIAS_LIST},
|
||||
})
|
||||
NATURAL_SCHEMA["properties"]["groups"]["items"]["required"] = list(GROUP_PROPERTIES)
|
||||
REVIEW_SCHEMA = copy.deepcopy(NATURAL_SCHEMA)
|
||||
REVIEW_SCHEMA["required"].append("decisions")
|
||||
REVIEW_SCHEMA["properties"]["decisions"] = {"type": "array", "minItems": 1, "items": {
|
||||
"type": "object", "additionalProperties": False,
|
||||
"required": ["source_members", "decision", "rationale", "evidence"], "properties": {
|
||||
"source_members": ALIAS_LIST, "decision": {"type": "string", "enum": ["keep", "split"]},
|
||||
"rationale": text_schema(320), "evidence": EVIDENCE}}}
|
||||
|
||||
|
||||
def invocation(stage, request, context=None):
|
||||
"""Build one full-input pass; independent discovery has no proposal context."""
|
||||
if stage not in STAGES:
|
||||
raise ValueError("unknown internal stage")
|
||||
discovery = stage in STAGES[:2]
|
||||
if discovery and context:
|
||||
raise ValueError("discovery must be independent")
|
||||
schema = DISCOVERY_SCHEMA if discovery else NATURAL_SCHEMA if stage == "reconciliation" else REVIEW_SCHEMA
|
||||
instruction = DISCOVERY if discovery else RECONCILE
|
||||
if stage == "large_family_review":
|
||||
instruction += "\n" + REVIEW
|
||||
if stage == "decision_audit":
|
||||
instruction += "\n" + REVIEW + "\n" + AUDIT
|
||||
source = {k: request[k] for k in ("campaign", "suite", "cases")}
|
||||
return {"stage": stage, "system": SYSTEM + "\n\nPASS INSTRUCTIONS:\n" + instruction,
|
||||
"schema": schema, "input": encoded({"suite": source, "review_material": context or {}}).decode()}
|
||||
|
||||
|
||||
def projection(value):
|
||||
"""Extract only the unchanged public group shape from a natural partition."""
|
||||
return {"groups": [{k: group[k] for k in ("name", "description", "members")}
|
||||
for group in value["groups"]]}
|
||||
|
||||
|
||||
def bounded_string(value, limit, *, empty=False):
|
||||
"""Reject invalid internal strings without echoing their contents."""
|
||||
if type(value) is not str or len(value) > limit or (not empty and not value.strip()):
|
||||
raise Problem("invalid_review_output", 502)
|
||||
|
||||
|
||||
def evidence(value, members, by_alias):
|
||||
"""Require short exact source-field support within the explained membership."""
|
||||
if not isinstance(value, list) or not 1 <= len(value) <= 3:
|
||||
raise Problem("invalid_review_evidence", 502)
|
||||
for item in value:
|
||||
if not isinstance(item, dict) or set(item) != {"alias", "field", "quote"}:
|
||||
raise Problem("invalid_review_evidence", 502)
|
||||
alias, field, quote = item["alias"], item["field"], item["quote"]
|
||||
if type(alias) is not str or alias not in members or type(field) is not str or field not in FIELDS:
|
||||
raise Problem("invalid_review_evidence", 502)
|
||||
bounded_string(quote, 160)
|
||||
source = by_alias[alias].get(field)
|
||||
if not isinstance(source, str) or quote not in source:
|
||||
raise Problem("invalid_review_evidence", 502)
|
||||
|
||||
|
||||
def validate_natural(value, request, originals=None):
|
||||
"""Validate complete natural assignments and bounded explanations, without a cap."""
|
||||
expected_keys = {"groups", "decisions"} if originals is not None else {"groups"}
|
||||
if not isinstance(value, dict) or set(value) != expected_keys or not isinstance(value["groups"], list):
|
||||
raise Problem("invalid_review_output", 502)
|
||||
required = set(GROUP_PROPERTIES)
|
||||
for group in value["groups"]:
|
||||
if not isinstance(group, dict) or set(group) != required:
|
||||
raise Problem("invalid_review_output", 502)
|
||||
validate_partition(projection(value), request, name_limit=BASE_NAME_LIMIT)
|
||||
by_alias = {case["alias"]: case for case in request["cases"]}
|
||||
for group in value["groups"]:
|
||||
members = set(group["members"])
|
||||
for field, limit in (("common_work", 160), ("rationale", 320), ("uncertainty", 240)):
|
||||
bounded_string(group[field], limit, empty=field == "uncertainty")
|
||||
if re.search(r"\(\s*\d+\s*/\s*\d+\s*\)$", group["name"]):
|
||||
raise Problem("invalid_natural_family_name", 502)
|
||||
evidence(group["evidence"], members, by_alias)
|
||||
variations = group["variation_sets"]
|
||||
if not isinstance(variations, list) or not variations or any(
|
||||
not isinstance(part, list) or not part or any(type(a) is not str for a in part)
|
||||
for part in variations):
|
||||
raise Problem("invalid_variation_assignments", 502)
|
||||
if Counter(a for part in variations for a in part) != Counter(group["members"]):
|
||||
raise Problem("invalid_variation_assignments", 502)
|
||||
if originals is not None:
|
||||
validate_decisions(value, originals, by_alias)
|
||||
return value
|
||||
|
||||
|
||||
def validate_decisions(value, originals, by_alias):
|
||||
"""Audit every original oversized family and forbid cross-family review drift."""
|
||||
decisions = value["decisions"]
|
||||
large = {frozenset(g["members"]): g for g in originals if len(g["members"]) > MAX_GROUP}
|
||||
if not isinstance(decisions, list) or len(decisions) != len(large):
|
||||
raise Problem("incomplete_large_family_review", 502)
|
||||
seen = set()
|
||||
for decision in decisions:
|
||||
if not isinstance(decision, dict) or set(decision) != {"source_members", "decision", "rationale", "evidence"}:
|
||||
raise Problem("invalid_review_output", 502)
|
||||
members = decision["source_members"]
|
||||
if not isinstance(members, list) or any(type(a) is not str for a in members):
|
||||
raise Problem("invalid_review_output", 502)
|
||||
key = frozenset(members)
|
||||
if key not in large or key in seen or len(key) != len(members):
|
||||
raise Problem("incomplete_large_family_review", 502)
|
||||
seen.add(key)
|
||||
descendants = [g for g in value["groups"] if set(g["members"]) <= key]
|
||||
split = len(descendants) > 1
|
||||
if not descendants or decision["decision"] != ("split" if split else "keep"):
|
||||
raise Problem("invalid_review_decision", 502)
|
||||
bounded_string(decision["rationale"], 320)
|
||||
evidence(decision["evidence"], key, by_alias)
|
||||
for group in value["groups"]:
|
||||
parents = [g for g in originals if set(group["members"]) <= set(g["members"])]
|
||||
if len(parents) != 1 or (len(parents[0]["members"]) <= MAX_GROUP and
|
||||
set(group["members"]) != set(parents[0]["members"])):
|
||||
raise Problem("review_crossed_family_boundary", 502)
|
||||
114
services/hermes/scripts/suite_sizing.py
Normal file
114
services/hermes/scripts/suite_sizing.py
Normal file
@ -0,0 +1,114 @@
|
||||
"""Deterministic balanced work sizing after semantic family review."""
|
||||
from __future__ import annotations
|
||||
|
||||
import math
|
||||
from collections import Counter
|
||||
|
||||
from suite_contract import Problem, digest, validate_result
|
||||
from suite_policy import MAX_GROUP, POLICY_REVISION, projection
|
||||
|
||||
|
||||
def balanced_sizes(count):
|
||||
"""Compute the minimum number of bounded parts without remainder singletons."""
|
||||
if type(count) is not int or count < 1:
|
||||
raise ValueError("positive case count required")
|
||||
parts = math.ceil(count / MAX_GROUP)
|
||||
base, remainder = divmod(count, parts)
|
||||
return [base + 1] * remainder + [base] * (parts - remainder)
|
||||
|
||||
|
||||
def ordered_parts(group, by_alias):
|
||||
"""Pack variation hints into fixed balanced sizes; equivalent text stays stable."""
|
||||
def content(alias):
|
||||
return digest({k: v for k, v in by_alias[alias].items() if k != "alias"})
|
||||
# Canonicalize every hint, not its model-generated order. Identical records use
|
||||
# the same hint owner even if the model placed their aliases in different hints.
|
||||
hints = [sorted(part, key=lambda a: (content(a), a)) for part in group["variation_sets"]]
|
||||
hints.sort(key=lambda part: tuple(sorted((content(a), a) for a in part)))
|
||||
owner = {}
|
||||
for index, part in enumerate(hints):
|
||||
for alias in part:
|
||||
owner.setdefault(content(alias), index)
|
||||
buckets = {}
|
||||
for alias in group["members"]:
|
||||
buckets.setdefault(owner[content(alias)], []).append(alias)
|
||||
blocks = [sorted(part, key=lambda a: (content(a), a)) for part in buckets.values()]
|
||||
blocks.sort(key=lambda part: (-len(part), tuple((content(a), a) for a in part)))
|
||||
sizes = balanced_sizes(len(group["members"]))
|
||||
parts = [[] for _ in sizes]
|
||||
for block in blocks:
|
||||
left = list(block)
|
||||
while left:
|
||||
room = [size - len(part) for size, part in zip(sizes, parts)]
|
||||
fits = [i for i, free in enumerate(room) if free >= len(left)]
|
||||
index = min(fits, key=lambda i: (room[i], i)) if fits else max(range(len(room)), key=lambda i: (room[i], -i))
|
||||
count = min(room[index], len(left))
|
||||
if count <= 0:
|
||||
raise Problem("invalid_capacity_partition", 502)
|
||||
parts[index].extend(left[:count])
|
||||
left = left[count:]
|
||||
parts = [sorted(part) for part in parts]
|
||||
parts.sort(key=lambda part: (-len(part), tuple(part)))
|
||||
if [len(part) for part in parts] != sizes or Counter(a for p in parts for a in p) != Counter(group["members"]):
|
||||
raise Problem("invalid_capacity_partition", 502)
|
||||
return parts
|
||||
|
||||
|
||||
def cap_families(natural, request):
|
||||
"""Create only final tasks, retaining conceptual-family sizing explanations."""
|
||||
by_alias = {case["alias"]: case for case in request["cases"]}
|
||||
final, divisions = [], []
|
||||
for group in sorted(natural["groups"], key=lambda g: (" ".join(g["name"].split()).casefold(), sorted(g["members"]))):
|
||||
base_name = " ".join(group["name"].split())
|
||||
if len(group["members"]) <= MAX_GROUP:
|
||||
final.append({"name": base_name, "description": group["description"], "members": sorted(group["members"])})
|
||||
continue
|
||||
parts = ordered_parts(group, by_alias)
|
||||
names = []
|
||||
for index, members in enumerate(parts, 1):
|
||||
name = f"{base_name} ({index}/{len(parts)})"
|
||||
names.append(name)
|
||||
final.append({"name": name,
|
||||
"description": f"Work-size part {index}/{len(parts)} of one family. {group['common_work']}",
|
||||
"members": members})
|
||||
divisions.append({"family_name": base_name, "natural_case_count": len(group["members"]),
|
||||
"part_sizes": [len(part) for part in parts], "part_names": names,
|
||||
"common_work": group["common_work"], "rationale": group["rationale"],
|
||||
"uncertainty": group["uncertainty"]})
|
||||
result = validate_result({"groups": final}, request)
|
||||
# Recheck final names and memberships against each declared capacity division.
|
||||
by_name = {group["name"]: group for group in final}
|
||||
for division in divisions:
|
||||
if [len(by_name[name]["members"]) for name in division["part_names"]] != balanced_sizes(division["natural_case_count"]):
|
||||
raise Problem("invalid_capacity_partition", 502)
|
||||
return result, divisions
|
||||
|
||||
|
||||
def disagreements(a, b):
|
||||
"""Compare co-membership sets, storing compressed partitions rather than O(n^2) pairs."""
|
||||
def memberships(value):
|
||||
return {alias: frozenset(g["members"]) for g in value["groups"] for alias in g["members"]}
|
||||
first, second = memberships(a), memberships(b)
|
||||
changed = sorted(alias for alias in first if first[alias] != second[alias])
|
||||
pair_count = sum(len(first[alias] ^ second[alias]) for alias in first) // 2
|
||||
return {"aliases": changed, "pair_count": pair_count,
|
||||
"proposal_a": sorted(sorted(g["members"]) for g in a["groups"]),
|
||||
"proposal_b": sorted(sorted(g["members"]) for g in b["groups"])}
|
||||
|
||||
|
||||
def review_summary(a, b, reconciled, reviewed, audited, final, divisions):
|
||||
"""Build an authorized result-only explanation, never routine operational metadata."""
|
||||
natural = audited or reconciled
|
||||
return {"policy_revision": POLICY_REVISION,
|
||||
"proposal_disagreements": disagreements(a, b),
|
||||
"reconciled_families": reconciled["groups"],
|
||||
"large_family_review": reviewed.get("decisions", []) if reviewed else [],
|
||||
"decision_audit": audited.get("decisions", []) if audited else [],
|
||||
"natural_families": natural["groups"], "capacity_divisions": divisions,
|
||||
"unresolved_uncertainties": [{"family_name": g["name"], "members": sorted(g["members"]),
|
||||
"uncertainty": g["uncertainty"]}
|
||||
for g in natural["groups"] if g["uncertainty"].strip()],
|
||||
"counts": {"natural_families": len(natural["groups"]), "final_tasks": len(final["groups"]),
|
||||
"natural_singletons": sum(len(g["members"]) == 1 for g in natural["groups"]),
|
||||
"final_singletons": sum(len(g["members"]) == 1 for g in final["groups"])},
|
||||
"sizing_is_not_semantic_evidence": True}
|
||||
@ -36,7 +36,7 @@ spec:
|
||||
app: hermes-suite-planner
|
||||
annotations:
|
||||
fluentbit.io/exclude: "true"
|
||||
ai.bstein.dev/config-rev: suite-v6-prompt-v2-diagnostics-v1-20260929
|
||||
ai.bstein.dev/config-rev: suite-v6-multipass-cap5-v1-20260929
|
||||
vault.hashicorp.com/agent-inject: "true"
|
||||
vault.hashicorp.com/agent-pre-populate-only: "true"
|
||||
vault.hashicorp.com/agent-init-first: "true"
|
||||
|
||||
@ -9,6 +9,7 @@ import pytest
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "services/hermes/scripts"))
|
||||
import suite_backends
|
||||
import suite_multipass
|
||||
from suite_cli_diagnostics import snapshot
|
||||
from suite_contract import CLAUDE_MAX_TURNS, MODELS, Problem, preflight, validate_request
|
||||
from suite_jobs import Jobs
|
||||
@ -140,7 +141,8 @@ def test_validation_failure_retains_usage_without_accepting_content(tmp_path, mo
|
||||
else:
|
||||
result["groups"][0]["members"].append("CASE-1")
|
||||
monkeypatch.setattr(suite_backends, "switchyard_decision", lambda *_: None)
|
||||
monkeypatch.setattr(suite_backends, "claude_generate", lambda *_: (result, metadata))
|
||||
metadata["review_summary"] = {}
|
||||
monkeypatch.setattr(suite_multipass, "generate", lambda *_: (result, metadata))
|
||||
jobs = Jobs(tmp_path / "jobs.sqlite")
|
||||
job, _ = jobs.submit("owner", "synthetic-validation", value, selected, "192.168.22.8", launch=False)
|
||||
jobs.run(job["job_id"], "owner", value, selected, "192.168.22.8")
|
||||
@ -204,7 +206,7 @@ def test_failed_cli_usage_survives_job_cleanup_without_content(tmp_path, monkeyp
|
||||
value, jobs = request(), Jobs(tmp_path / "jobs.sqlite")
|
||||
selected = preflight(value)
|
||||
monkeypatch.setattr(suite_backends, "switchyard_decision", lambda *_: None)
|
||||
def fail(*_):
|
||||
def fail(*_, **kwargs):
|
||||
suite_backends.parse_claude(envelope(subtype="error_max_turns", is_error=True,
|
||||
structured_output=None), MODEL, exit_code=1)
|
||||
monkeypatch.setattr(suite_backends, "claude_generate", fail)
|
||||
|
||||
339
testing/tests/test_suite_multipass.py
Normal file
339
testing/tests/test_suite_multipass.py
Normal file
@ -0,0 +1,339 @@
|
||||
"""Whole-suite orchestration, semantic review boundaries, and hard task sizing."""
|
||||
import copy
|
||||
import json
|
||||
from pathlib import Path
|
||||
import sys
|
||||
import threading
|
||||
|
||||
import pytest
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[2] / 'services/hermes/scripts'))
|
||||
import suite_backends
|
||||
import suite_multipass as workflow
|
||||
from suite_contract import MODELS, Problem, encoded, preflight, validate_partition, validate_request, validate_result
|
||||
from suite_jobs import Jobs
|
||||
from suite_policy import invocation, validate_natural
|
||||
from suite_sizing import balanced_sizes, cap_families, disagreements
|
||||
from suite_synthetic import fixture
|
||||
|
||||
CANARY = 'SYNTHETIC_CONTENT_NOT_FOR_DATABASE_OR_ROUTINE_LOGS'
|
||||
|
||||
|
||||
def request(count):
|
||||
"""Build fully synthetic coherent cases with separate aliases."""
|
||||
value = {'campaign': 'SYNTHETIC', 'suite': 'COHERENT', 'cases': [
|
||||
{'alias': f'CASE-{i:04d}', 'description': CANARY + ': inspect parser response.',
|
||||
'success_criteria': 'Call the JSON parser and assert returned fields against expected values.',
|
||||
'preconditions': 'Initialize an in-memory parser fixture with deterministic input.',
|
||||
'case_type': 'nominal' if i % 2 else 'fault injection'} for i in range(count)],
|
||||
'routing': {'allow_external': True, 'allowed_external_providers': ['claude']}}
|
||||
return validate_request(value, ['claude'])
|
||||
|
||||
|
||||
def natural(source, parts, names=None):
|
||||
"""Supply schema-valid mock engineering decisions and exact source support."""
|
||||
by_alias = {c['alias']: c for c in source['cases']}
|
||||
names = names or [f'Parser machinery {i}' for i in range(len(parts))]
|
||||
groups = []
|
||||
for name, members in zip(names, parts):
|
||||
groups.append({'name': name, 'description': 'Implement the shared parser fixture and returned-field assertions.',
|
||||
'members': list(members), 'common_work': 'Implement the shared parser fixture and returned-field assertions.',
|
||||
'rationale': CANARY + ': observations already exist; new expected values are inexpensive.',
|
||||
'uncertainty': '', 'variation_sets': [list(members)],
|
||||
'evidence': [{'alias': members[0], 'field': 'success_criteria',
|
||||
'quote': by_alias[members[0]]['success_criteria'][:100]}]})
|
||||
return {'groups': groups}
|
||||
|
||||
|
||||
def public(value):
|
||||
return {'groups': [{k: g[k] for k in ('name', 'description', 'members')} for g in value['groups']]}
|
||||
|
||||
|
||||
def review(value, original):
|
||||
"""Attach one explicit keep/split decision for each original oversized family."""
|
||||
value = copy.deepcopy(value)
|
||||
value['decisions'] = []
|
||||
for old in original['groups']:
|
||||
if len(old['members']) > 5:
|
||||
descendants = [g for g in value['groups'] if set(g['members']) <= set(old['members'])]
|
||||
value['decisions'].append({'source_members': old['members'],
|
||||
'decision': 'split' if len(descendants) > 1 else 'keep',
|
||||
'rationale': 'Different machinery warrants splitting.' if len(descendants) > 1 else
|
||||
'The same fixture and returned observations support all assertion variations.',
|
||||
'evidence': old['evidence']})
|
||||
return value
|
||||
|
||||
|
||||
def install_backend(monkeypatch, source, reconciled, *, a=None, b=None, reviewed=None, audited=None, costs=None):
|
||||
"""Capture complete invocations while substituting content-free usage metadata."""
|
||||
calls = []
|
||||
reviewed = review(reviewed or reconciled, reconciled)
|
||||
audited = review(audited or reviewed, reconciled)
|
||||
answers = {'proposal_a': a or public(reconciled), 'proposal_b': b or public(reconciled),
|
||||
'reconciliation': reconciled, 'large_family_review': reviewed, 'decision_audit': audited}
|
||||
def backend(value, cancel, *, invocation):
|
||||
payload = json.loads(invocation['input'])
|
||||
assert {c['alias']: c for c in payload['suite']['cases']} == {c['alias']: c for c in source['cases']}
|
||||
assert payload['suite']['campaign'] == source['campaign']
|
||||
assert payload['suite']['suite'] == source['suite']
|
||||
calls.append((copy.deepcopy(value), copy.deepcopy(invocation)))
|
||||
cost = costs[len(calls)-1] if costs else 0.1
|
||||
return copy.deepcopy(answers[invocation['stage']]), {
|
||||
'model': MODELS['claude']['model'], 'usage': {'input_tokens': 100, 'output_tokens': 40},
|
||||
'cost_usd_estimate': cost, 'turns': 2, 'duration_api_ms': 10,
|
||||
'compaction': False, 'truncation': False, 'cli_diagnostics': {'exit_code': 0}}
|
||||
monkeypatch.setattr(suite_backends, 'claude_generate', backend)
|
||||
return calls
|
||||
|
||||
|
||||
@pytest.mark.parametrize('count,sizes', [(6,[3,3]), (7,[4,3]), (11,[4,4,3]), (14,[5,5,4])])
|
||||
def test_coherent_families_are_balanced_not_semantically_fragmented(monkeypatch, count, sizes):
|
||||
source = request(count)
|
||||
aliases = [c['alias'] for c in source['cases']]
|
||||
family = natural(source, [aliases], ['Reset recovery'])
|
||||
calls = install_backend(monkeypatch, source, family)
|
||||
result, metadata = workflow.generate(source, preflight(source), threading.Event(), '192.168.22.8')
|
||||
assert len(calls) == 5 and [len(g['members']) for g in result['groups']] == sizes
|
||||
assert metadata['natural_family_count'] == 1 and metadata['final_task_count'] == len(sizes)
|
||||
assert metadata['model_pass_count'] == 5 and metadata['turns'] == 10
|
||||
assert metadata['singleton_statistics'] == {'natural_singletons': 0, 'final_singletons': 0}
|
||||
assert [g['name'] for g in result['groups']] == [f'Reset recovery ({i}/{len(sizes)})' for i in range(1,len(sizes)+1)]
|
||||
assert all('Work-size part' in g['description'] for g in result['groups'])
|
||||
assert metadata['review_summary']['decision_audit'][0]['decision'] == 'keep'
|
||||
assert [call[0]['execution']['max_cost_usd'] for call in calls] == pytest.approx([5,4.9,4.8,4.7,4.6])
|
||||
assert all(calls[i+1][0]['execution']['max_seconds'] < calls[i][0]['execution']['max_seconds'] for i in range(4))
|
||||
validate_result(result, source)
|
||||
|
||||
|
||||
def test_independent_orders_and_content_are_reproducible(monkeypatch):
|
||||
source = request(14)
|
||||
family = natural(source, [[c['alias'] for c in source['cases']]])
|
||||
calls = install_backend(monkeypatch, source, family)
|
||||
workflow.generate(source, preflight(source), threading.Event(), '192.168.22.8')
|
||||
a, b = (json.loads(calls[i][1]['input']) for i in (0,1))
|
||||
assert a['review_material'] == b['review_material'] == {}
|
||||
assert calls[0][1]['system'] == calls[1][1]['system']
|
||||
assert a['suite']['cases'] != b['suite']['cases']
|
||||
assert b['suite']['cases'] == workflow.ordered_request(source, True)['cases']
|
||||
assert all('maxItems' not in c[1]['schema']['properties']['groups']['items']['properties']['members'] for c in calls)
|
||||
assert 'proposal_a' in json.loads(calls[2][1]['input'])['review_material']
|
||||
|
||||
|
||||
def test_real_semantic_subdivisions_precede_work_sizing(monkeypatch):
|
||||
source = request(11)
|
||||
aliases = [c['alias'] for c in source['cases']]
|
||||
for case in source['cases'][6:]:
|
||||
case['success_criteria'] = 'Capture reset-line waveform with an oscilloscope and measure transition duration.'
|
||||
case['preconditions'] = 'Configure pulse source and digital capture fixture.'
|
||||
original = natural(source, [aliases], ['Response observations'])
|
||||
split = natural(source, [aliases[:6],aliases[6:]], ['JSON response assertions','Waveform timing capture'])
|
||||
calls = install_backend(monkeypatch, source, original, reviewed=split, audited=split)
|
||||
result, metadata = workflow.generate(source, preflight(source), threading.Event(), '192.168.22.8')
|
||||
assert len(calls) == 5 and metadata['natural_family_count'] == 2
|
||||
assert sorted(len(g['members']) for g in result['groups']) == [3,3,5]
|
||||
assert metadata['review_summary']['decision_audit'][0]['decision'] == 'split'
|
||||
assert [d['family_name'] for d in metadata['review_summary']['capacity_divisions']] == ['JSON response assertions']
|
||||
|
||||
|
||||
def test_bounded_audit_can_reverse_an_unjustified_semantic_split(monkeypatch):
|
||||
source = request(7)
|
||||
aliases = [c['alias'] for c in source['cases']]
|
||||
original = natural(source,[aliases],['Parser assertions'])
|
||||
split = natural(source,[aliases[:3],aliases[3:]],['Nominal parser checks','Rejection parser checks'])
|
||||
calls = install_backend(monkeypatch,source,original,reviewed=split,audited=original)
|
||||
result, metadata = workflow.generate(source,preflight(source),threading.Event(),'192.168.22.8')
|
||||
assert len(calls) == 5 and [len(g['members']) for g in result['groups']] == [4,3]
|
||||
assert metadata['review_summary']['large_family_review'][0]['decision'] == 'split'
|
||||
assert metadata['review_summary']['decision_audit'][0]['decision'] == 'keep'
|
||||
|
||||
|
||||
def test_disagreement_is_not_resolved_by_transitive_closure(monkeypatch):
|
||||
source = request(3)
|
||||
aliases = [c['alias'] for c in source['cases']]
|
||||
a = public(natural(source,[aliases[:2],aliases[2:]],['Parser checks','Parser checks']))
|
||||
b = public(natural(source,[aliases[:1],aliases[1:]],['Other labels','Other labels']))
|
||||
final = natural(source,[aliases[:2],aliases[2:]],['Frame parser checks','Report parser checks'])
|
||||
final['groups'][1]['uncertainty'] = 'The observation interface remains unspecified.'
|
||||
calls = install_backend(monkeypatch,source,final,a=a,b=b)
|
||||
result, metadata = workflow.generate(source,preflight(source),threading.Event(),'192.168.22.8')
|
||||
assert len(calls) == 3 and len(result['groups']) == 2
|
||||
assert metadata['review_summary']['proposal_disagreements']['pair_count'] == 2
|
||||
assert metadata['review_summary']['unresolved_uncertainties']
|
||||
assert set(result['groups'][0]['members']) != set(aliases)
|
||||
|
||||
|
||||
@pytest.mark.parametrize('count', [6,7,11,14,400])
|
||||
def test_identical_text_distinct_aliases_stable_under_reordering(count):
|
||||
source = request(count)
|
||||
for case in source['cases']:
|
||||
case['case_type'] = 'nominal'
|
||||
aliases = [c['alias'] for c in source['cases']]
|
||||
family = natural(source,[aliases],['X'*56])
|
||||
first, _ = cap_families(family, source)
|
||||
source['cases'].reverse()
|
||||
family['groups'][0]['members'].reverse()
|
||||
family['groups'][0]['variation_sets'][0].reverse()
|
||||
second, _ = cap_families(family, source)
|
||||
assert first == second
|
||||
assert max(len(g['name']) for g in first['groups']) <= 64
|
||||
assert sorted(a for g in first['groups'] for a in g['members']) == sorted(aliases)
|
||||
sizes = [len(g['members']) for g in first['groups']]
|
||||
assert sizes == balanced_sizes(count) and max(sizes)-min(sizes) <= 1
|
||||
|
||||
|
||||
@pytest.mark.parametrize('mutation,code', [('collision','duplicate_family_name'), ('long','invalid_json_result'),
|
||||
('invented','invalid_case_assignments'), ('duplicate','invalid_case_assignments'), ('overcap','group_size_limit')])
|
||||
def test_final_validation_is_independent_of_model_claims(mutation,code):
|
||||
source = request(6)
|
||||
aliases = [c['alias'] for c in source['cases']]
|
||||
value = public(natural(source,[aliases[:3],aliases[3:]],['Parser fixtures','Report fixtures']))
|
||||
if mutation == 'collision':
|
||||
value['groups'][1]['name'] = ' PARSER fixtures '
|
||||
elif mutation == 'long': value['groups'][0]['name'] = 'x'*65
|
||||
elif mutation == 'invented': value['groups'][0]['members'][0] = 'CASE-NOT-SUPPLIED'
|
||||
elif mutation == 'duplicate': value['groups'][0]['members'].append(aliases[0])
|
||||
else: value = public(natural(source,[aliases]))
|
||||
with pytest.raises(Problem,match=code): validate_result(value,source)
|
||||
|
||||
|
||||
def test_review_requires_every_oversized_family_and_real_source_support():
|
||||
source = request(14)
|
||||
aliases = [c['alias'] for c in source['cases']]
|
||||
original = natural(source,[aliases[:7],aliases[7:]])
|
||||
value = review(original,original)
|
||||
validate_natural(value,source,original['groups'])
|
||||
value['decisions'].pop()
|
||||
with pytest.raises(Problem,match='incomplete_large_family_review'):
|
||||
validate_natural(value,source,original['groups'])
|
||||
value = review(original,original)
|
||||
value['groups'][0]['evidence'][0]['quote'] = 'Invented equipment not in the source'
|
||||
with pytest.raises(Problem,match='invalid_review_evidence'):
|
||||
validate_natural(value,source,original['groups'])
|
||||
|
||||
|
||||
@pytest.mark.parametrize('size', [14,75,363])
|
||||
def test_existing_suite_sizes_have_separate_natural_and_task_counts(monkeypatch,size):
|
||||
source, expected = fixture(size)
|
||||
source['routing'] = {'allow_external':True,'allowed_external_providers':['claude']}
|
||||
source = validate_request(source,['claude'])
|
||||
buckets = {}
|
||||
for alias,family in expected.items(): buckets.setdefault(family,[]).append(alias)
|
||||
families = natural(source,list(buckets.values()),list(buckets))
|
||||
install_backend(monkeypatch,source,families)
|
||||
result, metadata = workflow.generate(source,preflight(source),threading.Event(),'192.168.22.8')
|
||||
assert metadata['natural_family_count'] == 9
|
||||
assert metadata['final_task_count'] == sum(len(balanced_sizes(len(p))) for p in buckets.values())
|
||||
assert all(1 <= len(g['members']) <= 5 for g in result['groups'])
|
||||
if size == 363: assert metadata['final_task_count'] > 7*metadata['natural_family_count']
|
||||
assert len(encoded({'result':result,**metadata})) < 1<<20
|
||||
|
||||
|
||||
def test_whole_job_cost_budget_and_failure_no_partial_answer(monkeypatch):
|
||||
source = request(14)
|
||||
family = natural(source,[[c['alias'] for c in source['cases']]])
|
||||
calls = install_backend(monkeypatch,source,family,costs=[1,2,3])
|
||||
with pytest.raises(Problem,match='job_cost_budget_exhausted') as raised:
|
||||
workflow.generate(source,preflight(source),threading.Event(),'192.168.22.8')
|
||||
assert [c[0]['execution']['max_cost_usd'] for c in calls] == [5,4,2]
|
||||
assert raised.value.details['completed_model_passes'] == 3
|
||||
assert CANARY not in json.dumps(raised.value.document())
|
||||
|
||||
|
||||
def test_expanded_reconciliation_capacity_checked_before_launch(monkeypatch):
|
||||
source = request(14)
|
||||
family = natural(source,[[c['alias'] for c in source['cases']]])
|
||||
calls = install_backend(monkeypatch,source,family)
|
||||
selected = preflight(source)
|
||||
original = workflow.capacity
|
||||
def capacity(call,*args):
|
||||
if call['stage'] == 'reconciliation':
|
||||
call['input'] += 'x'*(1<<20)
|
||||
return original(call,*args)
|
||||
monkeypatch.setattr(workflow,'capacity',capacity)
|
||||
with pytest.raises(Problem,match='pass_request_too_large'):
|
||||
workflow.generate(source,selected,threading.Event(),'192.168.22.8')
|
||||
assert len(calls) == 2
|
||||
|
||||
|
||||
def test_local_only_fails_capacity_without_hosted_calls(monkeypatch):
|
||||
source = request(14)
|
||||
source['routing'] = {'allow_external':False,'allowed_external_providers':[]}
|
||||
monkeypatch.setattr(suite_backends,'claude_generate',lambda *a,**k: pytest.fail('hosted call'))
|
||||
with pytest.raises(Problem,match='capacity_or_unsupported_backend') as raised:
|
||||
preflight(source)
|
||||
assert set(raised.value.details['candidates']) == {'local'}
|
||||
|
||||
|
||||
def test_review_is_authorized_result_only_and_never_persisted(tmp_path,monkeypatch,capsys):
|
||||
source = request(7)
|
||||
family = natural(source,[[c['alias'] for c in source['cases']]])
|
||||
install_backend(monkeypatch,source,family)
|
||||
monkeypatch.setattr(suite_backends,'switchyard_decision',lambda *_:None)
|
||||
jobs = Jobs(tmp_path/'jobs.sqlite')
|
||||
selected=preflight(source)
|
||||
job,_ = jobs.submit('owner','multi-pass-key',source,selected,'192.168.22.8',launch=False)
|
||||
jobs.run(job['job_id'],'owner',source,selected,'192.168.22.8')
|
||||
assert jobs.get(job['job_id'],'owner')['status'] == 'completed'
|
||||
assert 'review_summary' not in jobs.get(job['job_id'],'owner')
|
||||
assert CANARY in json.dumps(jobs.get(job['job_id'],'owner',result=True)['review_summary'])
|
||||
with pytest.raises(Problem,match='job_not_found'): jobs.get(job['job_id'],'other',result=True)
|
||||
assert CANARY not in jobs.db.execute('select document from jobs').fetchone()[0]
|
||||
assert CANARY not in capsys.readouterr().out
|
||||
replay,new = jobs.submit('owner','multi-pass-key',source,selected,'192.168.22.8',launch=False)
|
||||
assert not new and replay['job_id'] == job['job_id']
|
||||
jobs.results[job['job_id']] = (0,{}, {})
|
||||
with pytest.raises(Problem,match='result_expired_or_worker_restarted'):
|
||||
jobs.get(job['job_id'],'owner',result=True)
|
||||
|
||||
|
||||
def test_shared_deadline_stops_after_prior_calls(monkeypatch):
|
||||
source = request(7)
|
||||
source['execution']['max_seconds'] = 60
|
||||
family = natural(source,[[c['alias'] for c in source['cases']]])
|
||||
calls = install_backend(monkeypatch,source,family)
|
||||
clock = [100.0]
|
||||
backend = suite_backends.claude_generate
|
||||
def advancing(*args,**kwargs):
|
||||
answer = backend(*args,**kwargs)
|
||||
clock[0] += 31
|
||||
return answer
|
||||
monkeypatch.setattr(workflow.time,'monotonic',lambda:clock[0])
|
||||
monkeypatch.setattr(suite_backends,'claude_generate',advancing)
|
||||
with pytest.raises(Problem,match='job_time_budget_exhausted'):
|
||||
workflow.generate(source,preflight(source),threading.Event(),'192.168.22.8')
|
||||
assert len(calls) == 2
|
||||
assert [c[0]['execution']['max_seconds'] for c in calls] == [60,29]
|
||||
|
||||
|
||||
def test_missing_cost_measurement_stops_future_paid_calls(monkeypatch):
|
||||
source = request(7)
|
||||
family = natural(source,[[c['alias'] for c in source['cases']]])
|
||||
calls = install_backend(monkeypatch,source,family,costs=[None])
|
||||
with pytest.raises(Problem,match='budget_accounting_unavailable'):
|
||||
workflow.generate(source,preflight(source),threading.Event(),'192.168.22.8')
|
||||
assert len(calls) == 1
|
||||
|
||||
|
||||
def test_cancellation_stops_between_model_passes(monkeypatch):
|
||||
source=request(7)
|
||||
family=natural(source,[[c['alias'] for c in source['cases']]])
|
||||
calls=install_backend(monkeypatch,source,family)
|
||||
cancel=threading.Event()
|
||||
backend=suite_backends.claude_generate
|
||||
def cancelling(*args,**kwargs):
|
||||
answer=backend(*args,**kwargs)
|
||||
cancel.set()
|
||||
return answer
|
||||
monkeypatch.setattr(suite_backends,'claude_generate',cancelling)
|
||||
with pytest.raises(Problem,match='cancelled'):
|
||||
workflow.generate(source,preflight(source),cancel,'192.168.22.8')
|
||||
assert len(calls) == 1
|
||||
|
||||
|
||||
def test_review_cannot_cross_a_natural_boundary():
|
||||
source=request(14)
|
||||
aliases=[c['alias'] for c in source['cases']]
|
||||
original=natural(source,[aliases[:7],aliases[7:]])
|
||||
mixed=natural(source,[aliases[::2],aliases[1::2]])
|
||||
value=review(mixed,original)
|
||||
with pytest.raises(Problem): validate_natural(value,source,original['groups'])
|
||||
@ -11,7 +11,7 @@ sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "services/hermes/sc
|
||||
import suite_api
|
||||
import suite_backends
|
||||
from suite_contract import (MODELS, PROMPT_REVISION, PROMPT_SHA256, SYSTEM, Problem,
|
||||
preflight, prompt, validate_request, validate_result)
|
||||
preflight, prompt, validate_partition, validate_request, validate_result)
|
||||
from suite_jobs import Jobs
|
||||
from suite_synthetic import fixture, score
|
||||
|
||||
@ -40,11 +40,11 @@ def test_full_suite_capacity_and_coverage(size):
|
||||
# Every field and record, including duplicate text, survives preparation exactly.
|
||||
assert json.loads(prompt(request))["cases"] == request["cases"]
|
||||
assert len({case["alias"] for case in request["cases"]}) == size
|
||||
result = validate_result(expected_result(size), request)
|
||||
result = validate_partition(expected_result(size), request)
|
||||
assert score(result, fixture(size)[1])["pair_recall"] == 1
|
||||
result["groups"][0]["members"].append("[reference]")
|
||||
with pytest.raises(Problem, match="invalid_case_assignments"):
|
||||
validate_result(result, request)
|
||||
validate_partition(result, request)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("policy", [
|
||||
@ -175,13 +175,16 @@ def test_no_fallback_after_local_failure(tmp_path, monkeypatch, capsys):
|
||||
request = {"campaign": "SYNTHETIC", "suite": "TINY", "cases": [
|
||||
{"alias": "CASE-1", "description": "Read parser status"}]}
|
||||
request = validate_request(request, [])
|
||||
selected = preflight(request)
|
||||
import suite_multipass
|
||||
# Isolate runtime fail-closed behavior from the new workflow admission limit.
|
||||
selected = {"provider": "local", **MODELS["local"]}
|
||||
monkeypatch.setattr(suite_multipass, "capacity", lambda *a: {})
|
||||
calls = []
|
||||
monkeypatch.setattr(suite_backends, "switchyard_decision", lambda provider: calls.append(provider))
|
||||
def fail(*args):
|
||||
def fail(*args, **kwargs):
|
||||
raise Problem("backend_unavailable", 503)
|
||||
monkeypatch.setattr(suite_backends, "local_generate", fail)
|
||||
monkeypatch.setattr(suite_backends, "claude_generate", lambda *args: pytest.fail("external launch"))
|
||||
monkeypatch.setattr(suite_backends, "claude_generate", lambda *args, **kwargs: pytest.fail("external launch"))
|
||||
jobs = Jobs(tmp_path / "jobs.sqlite")
|
||||
document, _ = jobs.submit("owner", "local-failure", request, selected, "192.168.22.8", launch=False)
|
||||
jobs.run(document["job_id"], "owner", request, selected, "192.168.22.8")
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user