"""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_assignments import decode, decision_keys 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 wire_value(value, call): """Encode fixed test decisions using the model's required per-alias representation.""" assignments, groups = {}, [] for index, group in enumerate(value['groups']): groups.append({k:v for k,v in group.items() if k not in {'members','variation_sets'}}) variation = {a:i for i,part in enumerate(group.get('variation_sets',[])) for a in part} for alias in group['members']: assignments[alias] = {'family':index,'variation':variation[alias]} if variation else index result = {'groups':groups,'assignments':assignments} if 'decisions' in value: key_by_members = {frozenset(v):k for k,v in call['review_keys'].items()} result['decisions'] = {key_by_members[frozenset(d['source_members'])]: {k:v for k,v in d.items() if k!='source_members'} for d in value['decisions']} return result 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, progress=None): 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))) if progress: progress({'cli_running': True, 'cli_output_bytes': 120, 'last_cli_activity_seconds_ago': 0}) cost = costs[len(calls)-1] if costs else 0.1 return wire_value(copy.deepcopy(answers[invocation['stage']]), invocation), { '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([10,9.9,9.8,9.7,9.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('members' not in c[1]['schema']['properties']['groups']['items']['properties'] for c in calls) assert all(set(c[1]['schema']['properties']['assignments']['required']) == {x['alias'] for x in source['cases']} for c in calls) assert 'proposal_a' in json.loads(calls[2][1]['input'])['review_material'] def test_progress_reports_activity_without_content_or_false_percentage(monkeypatch): source = request(7) family = natural(source, [[c['alias'] for c in source['cases']]]) install_backend(monkeypatch, source, family) updates = [] workflow.generate(source, preflight(source), threading.Event(), '192.168.22.8', updates.append) active = [value for value in updates if value.get('cli_running')] assert len(active) == 5 assert [v['completed_model_passes'] for v in active] == [0,1,2,3,4] assert all(v['maximum_model_passes'] == 5 and v['cli_output_bytes'] == 120 for v in active) assert all(v['heartbeat_at'] > 0 and v['job_remaining_seconds'] <= 1200 for v in active) assert all(v['last_cli_activity_seconds_ago'] == 0 for v in active) assert updates[-1]['cli_running'] is False assert CANARY not in json.dumps(updates) def test_twenty_minute_job_limit_retains_smaller_client_deadlines(): source = request(7) assert source['execution']['max_seconds'] == 1200 source['execution']['max_seconds'] = 900 assert validate_request(source, ['claude'])['execution']['max_seconds'] == 900 source['execution']['max_seconds'] = 1201 with pytest.raises(Problem, match='invalid_timeout'): validate_request(source, ['claude']) def test_cli_estimate_guard_default_and_client_override(): source = request(7) assert source['execution']['max_cost_usd'] == 10 source['execution']['max_cost_usd'] = 5 assert validate_request(source, ['claude'])['execution']['max_cost_usd'] == 5 source['execution']['max_cost_usd'] = 10.1 with pytest.raises(Problem, match='invalid_cost_limit'): validate_request(source, ['claude']) @pytest.mark.parametrize('size', [14,75,363]) def test_every_alias_is_required_in_internal_wire_schema(size): source = request(size) aliases = [c['alias'] for c in source['cases']] value = natural(source, [aliases]) for stage in ('proposal_a','reconciliation','large_family_review','decision_audit'): context = {'original_partition':value} if stage=='decision_audit' else {'natural_partition':value} call = invocation(stage, source, None if stage=='proposal_a' else context) data = public(value) if stage=='proposal_a' else review(value,value) if stage in ('large_family_review','decision_audit') else value wire = wire_value(data, call) schema = call['schema']['properties']['assignments'] assert set(schema['required']) == set(schema['properties']) == set(aliases) assert schema['additionalProperties'] is False assert decode(wire,call,source) == data @pytest.mark.parametrize('mode,stage', [ ('missing','assignment_keys'),('unknown','assignment_keys'),('wrong_index','assignment_index'), ('bool_index','assignment_index'),('unused_family','empty_assigned_group'), ('array_instead','assignment_keys'),('extra_group_field','assignment_group_fields')]) def test_invalid_wire_assignments_fail_without_partial_output_or_content(mode,stage): source = request(14) aliases = [c['alias'] for c in source['cases']] call = invocation('proposal_a',source) wire = wire_value(public(natural(source,[aliases])),call) if mode=='missing': del wire['assignments'][aliases[0]] elif mode=='unknown': wire['assignments'][CANARY]=0 elif mode=='wrong_index': wire['assignments'][aliases[0]]=1 elif mode=='bool_index': wire['assignments'][aliases[0]]=True elif mode=='unused_family': wire['groups'].append(dict(wire['groups'][0])) elif mode=='array_instead': wire['assignments']=[] else: wire['groups'][0]['members']=aliases with pytest.raises(Problem) as raised: decode(wire,call,source) assert raised.value.details['failure_stage']==stage assert CANARY not in json.dumps(raised.value.document()) def test_review_wire_keys_reference_exact_memberships_and_cannot_omit_a_decision(): source=request(14) aliases=[c['alias'] for c in source['cases']] natural_value=natural(source,[aliases[:7],aliases[7:]]) context={'original_partition':natural_value} call=invocation('decision_audit',source,context) keys=decision_keys(context) assert len(keys)==2 and set(call['schema']['properties']['decisions']['required'])==set(keys) wire=wire_value(review(natural_value,natural_value),call) del wire['decisions'][next(iter(keys))] with pytest.raises(Problem,match='incomplete_large_family_review'): decode(wire,call,source) 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) source['execution']['max_cost_usd'] = 5 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'])