| """Replaceable controllers. API responses are saved as visible decision records.""" |
| from __future__ import annotations |
| import json |
| import os |
| from importlib.resources import files |
| from .schema import Decision, canonical |
|
|
|
|
| def request_context(state,registry,input_bytes): |
| """Identical deterministic context packing for all controller providers.""" |
| import copy |
| prompt=files('peppa').joinpath('prompts/controller.txt').read_text() |
| prompt+='\nReturn JSON matching this schema:\n'+canonical(Decision.model_json_schema()) |
| view=copy.deepcopy(state) |
| view['context_counts']={k:len(state[k]) for k in ['candidates','measurements','evidence']} |
| tools=[{'name':t.name,'description':t.description,'cost':t.cost,'arguments':t.argument_schema} for t in registry.values()] |
| while True: |
| content=canonical({'state':view,'tools':tools}) |
| if len((prompt+content).encode())<=input_bytes:return prompt,content |
| changed=False |
| for key in ['measurements','candidates','evidence','errors','controller_feedback','tool_messages']: |
| obj=view.get(key,[]) |
| if len(obj)>1: |
| n=max(1,len(obj)//2) |
| view[key]=dict(list(obj.items())[-n:]) if isinstance(obj,dict) else obj[-n:] |
| changed=True;break |
| if not changed:raise ValueError('fixed scientific specification exceeds controller context byte cap') |
|
|
|
|
| class ScriptedController: |
| """Execute a fixed workflow for baselines and deterministic local examples.""" |
| token_reservation=0 |
| def __init__(self, decisions): |
| self.decisions=iter(decisions) |
| def next(self,state,registry): |
| try:d=Decision.model_validate(next(self.decisions)) |
| except StopIteration: |
| d=Decision(tool="stop",arguments={},hypothesis="Workflow complete",evidence_ids=[], |
| decision_summary="All registered workflow steps completed.",expected_observation="",stop=True) |
| return d,{"provider":"scripted","decision":d.model_dump()} |
|
|
|
|
| class APIController: |
| """OpenAI Responses or a local OpenAI-compatible chat endpoint. |
| |
| token_reservation caps the serialized request bytes plus maximum output |
| tokens. Configure model IDs and exact serving revisions in the run manifest. |
| """ |
| def __init__(self, model, provider="responses", base_url=None, |
| api_key_env="OPENAI_API_KEY", max_output_tokens=2500, |
| input_bytes=14000,reasoning_effort=None): |
| self.reasoning_effort=reasoning_effort |
| self.model,self.provider=model,provider |
| self.max_output_tokens=max_output_tokens;self.input_bytes=input_bytes |
| self.token_reservation=input_bytes+max_output_tokens+1024 |
| from openai import OpenAI |
| key=os.environ.get(api_key_env) |
| if not key: |
| raise ValueError(f"set {api_key_env} before using the API controller") |
| self.client=OpenAI(api_key=key,base_url=base_url,max_retries=0,timeout=180) |
| def next(self,state,registry): |
| prompt,content=request_context(state,registry,self.input_bytes) |
| messages=[{"role":"system","content":prompt},{"role":"user","content":content}] |
| if self.provider=="responses": |
| settings={"reasoning":{"effort":self.reasoning_effort}} if self.reasoning_effort else {} |
| out=self.client.responses.create(model=self.model,input=messages, |
| text={"format":{"type":"json_object"}},max_output_tokens=self.max_output_tokens,store=False,**settings) |
| text=out.output_text |
| elif self.provider=="chat": |
| out=self.client.chat.completions.create(model=self.model,messages=messages, |
| response_format={"type":"json_object"},max_tokens=self.max_output_tokens) |
| text=out.choices[0].message.content |
| else:raise ValueError("unknown provider") |
| decision=Decision.model_validate_json(text) |
| usage=out.usage.model_dump() if out.usage else {} |
| return decision,{"provider":self.provider,"requested_model":self.model,"returned_model":out.model, |
| "response_id":out.id,"usage":usage,"prompt":messages,"visible_output":text, |
| "reasoning_effort":self.reasoning_effort,"max_output_tokens":self.max_output_tokens} |
|
|
|
|
| class AnthropicController: |
| """Claude Messages API, sharing the same visible decision schema.""" |
| def __init__(self,model,api_key_env='ANTHROPIC_API_KEY',max_output_tokens=2500,input_bytes=14000,effort=None): |
| self.effort=effort |
| import anthropic |
| self.client=anthropic.Anthropic(api_key=os.environ[api_key_env],max_retries=0,timeout=180) |
| self.model=model;self.max_output_tokens=max_output_tokens;self.input_bytes=input_bytes |
| self.token_reservation=input_bytes+max_output_tokens+1024 |
| def next(self,state,registry): |
| prompt,content=request_context(state,registry,self.input_bytes) |
| settings={'output_config':{'effort':self.effort},'thinking':{'type':'adaptive'}} if self.effort else {} |
| out=self.client.messages.create(model=self.model,max_tokens=self.max_output_tokens,system=prompt,messages=[{'role':'user','content':content}],**settings) |
| text=''.join(b.text for b in out.content if b.type=='text').strip() |
| if text.startswith('```'):text=text.split('\n',1)[1].rsplit('```',1)[0].strip() |
| return Decision.model_validate_json(text),{'provider':'anthropic','requested_model':self.model,'returned_model':out.model,'response_id':out.id,'usage':out.usage.model_dump(),'prompt':{'system':prompt,'user':content},'visible_output':text,'effort':self.effort,'max_output_tokens':self.max_output_tokens} |
|
|