Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/fuzz.yml
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,9 @@ jobs:
- name: Fuzz orchestration engine
run: python fuzz/fuzz_orchestration.py -max_total_time=${FUZZ_SECONDS} -artifact_prefix=crash- fuzz/corpus/orchestration

- name: Fuzz image placement catalog
run: python fuzz/fuzz_image_catalog.py -max_total_time=${FUZZ_SECONDS} -artifact_prefix=crash- fuzz/corpus/image_catalog

- name: Upload crash artifacts
if: failure()
uses: actions/upload-artifact@330a01c490aca151604b8cf639adc76d48f6c5d4 # actions/upload-artifact@v5
Expand Down
31 changes: 31 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
# Changelog

All notable changes to this project are documented in this file.

The format follows [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
and this project uses [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [Unreleased]

### Added

- Chat messages accept OpenAI `text` + `image_url` content parts. The gateway
records a 3NF `image_content_catalog` (`image_payload` / `image_placement` /
`image_recognition_event`) so an invoice PNG stays next to
`Please pay invoice 1042`. Raw base64 is hashed, not stored. Next action:
send the figure as `data:image/png;base64,...` or `https://...` and read
`orchestration.image_content_catalog` to find it.

### References

- Faysse, M., Sibille, H., Wu, T., Omrani, B., Viaud, G., Hudelot, C., &
Colombo, P. (2024). *ColPali: Efficient document retrieval with vision
language models* (arXiv:2407.01449). arXiv.
https://doi.org/10.48550/arXiv.2407.01449
- Xu, Y., Li, M., Cui, L., Huang, S., Wei, F., & Zhou, M. (2020). LayoutLM:
Pre-training of text and layout for document image understanding. In
*Proceedings of the 26th ACM SIGKDD International Conference on Knowledge
Discovery & Data Mining* (pp. 1192–1200). Association for Computing
Machinery. https://doi.org/10.1145/3394486.3403172
- Masinter, L. (1998). *The "data" URL scheme* (RFC 2397). Internet
Engineering Task Force. https://doi.org/10.17487/RFC2397
2 changes: 1 addition & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ A stdlib-Python lab implementing a single OpenAI-compatible API that routes, del

### Modules (`contextual_orchestrator/`)

- `orchestrator.py` — the domain heart: `ModelAgent`, `WorkflowStep`, `OrchestrationPolicy`, `ModelClient`, `TaskOrchestrator`, secret/PII redaction, budget enforcement, spend analytics, and the commercial-readiness report generators behind `/api/v1/*`. Domain code stays here until a second implementation forces extraction (see `docs/code_conventions.md`).
- `orchestrator.py` — the domain heart: `ModelAgent`, `WorkflowStep`, `OrchestrationPolicy`, `ModelClient`, `TaskOrchestrator`, `collect_image_catalog` (3NF figure placement), secret/PII redaction, budget enforcement, spend analytics, and the commercial-readiness report generators behind `/api/v1/*`. Domain code stays here until a second implementation forces extraction (see `docs/code_conventions.md`).
- `server.py` — HTTP delivery adapter and `SecurityConfig`; all request validation lives here.
- `admin.py` — static HTML/CSS/JS for the `/admin` operator console (stays inline while the product is dependency-free).
- `credentials.py` / `kv_config.py` — the KV seam: `get_credential`/`register_credential` over pluggable backends (`InMemoryCredentialBackend` default; pgcrypto-encrypted `PostgresCredentialBackend`, selected via `CONTEXTUAL_ORCHESTRATOR_KV_BACKEND`).
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -256,6 +256,7 @@ python tests/test_admin_contract.py
python tests/test_conventions.py
python tests/test_api_contract.py
python tests/test_security_hardening.py
python tests/test_image_placement_catalog.py
python tests/test_repository_security_metadata.py
python tests/test_product_planning_contract.py
python tests/test_plugin_driven_artifacts.py
Expand Down
1 change: 1 addition & 0 deletions conductor/product.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ Provide one API and one domain model:
- TRINITY: make thinker, worker, and verifier roles visible in the trace.
- Conductor: show natural-language subtasks and access lists as first-class audit objects.
- Enterprise operations: treat provider exclusion, locale bundles, and replayable workflow evidence as product surfaces.
- Image placement: keep invoice and email figures searchable at the text they sat next to (ColPali / LayoutLM).

## Non-Goals

Expand Down
1 change: 1 addition & 0 deletions conductor/tracks.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,3 +4,4 @@
|---|---|---|
| 001-paper-grounded-orchestrator | active | Implement the source-backed orchestration contract with TDD, DDD, and CDD |
| 002-enterprise-design-foundation | active | Add paper-grounded screen design, user stories, REST API, code/DB conventions, and i18n |
| 005-image-placement-catalog | active | Keep invoice/email figures searchable at their source offset (ColPali / LayoutLM). 3NF payload + placement + later recognition events. |
129 changes: 123 additions & 6 deletions contextual_orchestrator/orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@

from collections import Counter, deque, OrderedDict
from contextvars import ContextVar
import base64
import binascii
import copy
from dataclasses import dataclass, replace
from functools import wraps
Expand All @@ -29,7 +31,7 @@
from .credentials import NotConfigured, get_credential


ChatMessage = dict[str, str]
ChatMessage = dict[str, Any]

class BudgetExceededError(RuntimeError):
"""Raised when an operator-configured spend budget is already exhausted."""
Expand All @@ -39,6 +41,106 @@ def __init__(self, message: str, detail: dict[str, Any] | None = None) -> None:
self.detail = detail or {}


def flatten_message_text(content: Any) -> str:
"""Return concatenated text parts from a chat ``content`` value.

OpenAI vision callers send a list of ``text`` and ``image_url`` parts.
Routing and adjacent-text anchors need the words that sat next to the
figure, not the base64 payload.
"""
if isinstance(content, str):
return content
if not isinstance(content, list):
return ""
texts: list[str] = []
for part in content:
if isinstance(part, dict) and isinstance(part.get("text"), str):
texts.append(part["text"])
return " ".join(texts)


def _parse_image_source(url: str) -> tuple[str, str, int, str] | None:
"""Return ``(payload_digest, mime_type, byte_length, source_kind)`` or None."""
if url.startswith("data:"):
header, separator, payload = url.partition(",")
if not separator or ";base64" not in header.lower():
return None
media = header[5:].split(";", 1)[0].strip().lower()
if not media.startswith("image/"):
return None
try:
raw = base64.b64decode(payload, validate=True)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

validate=True rejects RFC 2397 whitespace. Combined with url.startswith("data:") (case-sensitive), a wrapped DATA:IMAGE/PNG;BASE64,... invoice is accepted by HTTP and then omitted from the catalog. Strip whitespace, casefold the scheme, then decode. Fixed on #684.

except (ValueError, binascii.Error):
return None
if not raw:
return None
return hashlib.sha256(raw).hexdigest(), media, len(raw), "inline_data_uri"
parsed = urlparse(url)
if parsed.scheme != "https" or not parsed.hostname:
return None
return hashlib.sha256(url.encode("utf-8")).hexdigest(), "image/remote", 0, "remote_https"


def collect_image_catalog(messages: list[Any]) -> dict[str, Any]:
"""Build a 3NF image catalog that keeps each figure at its source offset.

``image_payload`` is identity by digest so the same invoice PNG on a
reminder thread is one payload with two ``image_placement`` rows.
``image_recognition_event`` stays empty until a later vision/OCR pass
(temporal modeling: tags are not attributes of the bytes).
"""
payloads: dict[str, dict[str, Any]] = {}
placements: list[dict[str, Any]] = []
if not isinstance(messages, list):
return {
"image_payloads": [],
"image_placements": [],
"image_recognition_events": [],
}
for message_index, message in enumerate(messages):
if not isinstance(message, dict):
continue
content = message.get("content")
adjacent_text = flatten_message_text(content)
if not isinstance(content, list):
continue
for part_index, part in enumerate(content):
if not isinstance(part, dict) or part.get("type") != "image_url":
continue
image_url = part.get("image_url")
if isinstance(image_url, str):
url = image_url
elif isinstance(image_url, dict):
url = image_url.get("url")
else:
continue
if not isinstance(url, str) or not url.strip():
continue
parsed = _parse_image_source(url.strip())
if parsed is None:
continue
payload_digest, mime_type, byte_length, source_kind = parsed
payloads[payload_digest] = {
"payload_digest": payload_digest,
"mime_type": mime_type,
"byte_length": byte_length,
}
placements.append(
{
"payload_digest": payload_digest,
"message_index": message_index,
"part_index": part_index,
"source_kind": source_kind,
"adjacent_text": adjacent_text,
Comment on lines +128 to +134

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

SQL/ERD image_placement.placement_id has no JSON counterpart. Add a stable image_placement_{message_index}_{part_index} so a later writer can round-trip. Fixed on #684.

}
)
return {
"image_payloads": list(payloads.values()),
"image_placements": placements,
"image_recognition_events": [],
}


def estimate_tokens(text: str) -> int:
"""Rough token estimate (~4 chars/token). ponytail: heuristic, not a real tokenizer.

Expand Down Expand Up @@ -492,9 +594,13 @@ def _provider_url(self, agent: ModelAgent, path: str) -> str:
return f"{agent.base_url.rstrip('/')}{path}"

def _mock(self, agent: ModelAgent, messages: list[ChatMessage]) -> str:
last = next((m["content"] for m in reversed(messages) if m.get("role") == "user"), "")
last = next(
(flatten_message_text(m.get("content", "")) for m in reversed(messages) if m.get("role") == "user"),
"",
)
role = "worker"
system = messages[0]["content"] if messages and messages[0].get("role") == "system" else ""
system_content = messages[0].get("content", "") if messages and messages[0].get("role") == "system" else ""
system = flatten_message_text(system_content)
match = re.search(r"Role: ([a-z]+)", system)
if match:
role = match.group(1)
Expand Down Expand Up @@ -924,8 +1030,11 @@ def complete(self, messages: list[ChatMessage], mode: str = "auto") -> dict[str,
def _dispatch(self, messages: list[ChatMessage], mode: str) -> dict[str, Any]:
text = self._latest_user_text(messages)
if mode == "route" or (mode == "auto" and not self._needs_workflow(text)):
return self.route_once(messages)
return self.conduct(messages)
result = self.route_once(messages)
else:
result = self.conduct(messages)
result["image_content_catalog"] = collect_image_catalog(messages)
return result

def would_route(self, messages: list[ChatMessage], mode: str = "auto") -> bool:
"""True when this request takes the single-worker route path (vs the conduct workflow)."""
Expand Down Expand Up @@ -993,6 +1102,8 @@ def run(self, messages: list[ChatMessage], mode: str = "auto", workflow_run_id:
"trace": result["trace"],
"policy_snapshot": self.policy.as_dict(),
"verification": result.get("verification"),
"image_content_catalog": result.get("image_content_catalog")
or collect_image_catalog(messages),
}
self._workflow_runs[record["workflow_run_id"]] = record
self._run_order.appendleft(record["workflow_run_id"])
Expand Down Expand Up @@ -1600,7 +1711,10 @@ def _needs_workflow(self, text: str) -> bool:
return hits >= self.policy.conduct_hint_threshold or len(text) > 700

def _latest_user_text(self, messages: list[ChatMessage]) -> str:
return next((m.get("content", "") for m in reversed(messages) if m.get("role") == "user"), "") # pragma: no cover
return next(
(flatten_message_text(m.get("content", "")) for m in reversed(messages) if m.get("role") == "user"),
"",
)

def _model_judge_verification(self, task: str, fallback: dict[str, Any]) -> dict[str, Any]:
"""Ask a model to judge the verifier report (fixes term-matching false negatives).
Expand Down Expand Up @@ -8500,6 +8614,9 @@ def chat_completion_response(
}
if include_trace:
orchestration["trace"] = redact_value(result["trace"])
catalog = result.get("image_content_catalog")
if catalog:
orchestration["image_content_catalog"] = catalog

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This attaches adjacent_text without credential redaction. An api_key= next to the pay line leaves /v1/chat/completions even when include_orchestration_trace is false. Redact credential shapes only — keep invoice 1042 and ap@acme.com. Fixed on #684.

return {
"id": f"chatcmpl-{int(time.time() * 1000)}",
"object": "chat.completion",
Expand Down
46 changes: 43 additions & 3 deletions contextual_orchestrator/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -183,16 +183,56 @@ def _validate_mode(mode: Any) -> str:
return mode


def _validate_messages(messages: Any) -> list[dict[str, str]]:
def _validate_content_parts(content: list[Any]) -> list[dict[str, Any]]:
"""Accept OpenAI text + image_url parts; keep the figure at its offset."""
if not content:
raise RequestError(400, "invalid_message_content", "content parts must be a non-empty array")
cleaned: list[dict[str, Any]] = []
for part in content:
if not isinstance(part, dict):
raise RequestError(400, "invalid_message_content", "each content part must be an object")
part_type = part.get("type")
if part_type == "text":
text = part.get("text")
if not isinstance(text, str) or not text.strip():
raise RequestError(400, "invalid_message_content", "text content part requires non-empty text")
cleaned.append({"type": "text", "text": text})
continue
if part_type == "image_url":
image_url = part.get("image_url")
if isinstance(image_url, str):
url = image_url
elif isinstance(image_url, dict):
url = image_url.get("url")
else:
raise RequestError(400, "invalid_message_content", "image_url content part requires a url")
if not isinstance(url, str) or not url.strip():
raise RequestError(400, "invalid_message_content", "image_url content part requires a non-empty url")
url = url.strip()
if url.startswith("javascript:") or url.startswith("data:text/"):
raise RequestError(400, "invalid_message_content", "image_url must be https or data:image")
if not (url.startswith("https://") or url.lower().startswith("data:image/")):

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

https:// is case-sensitive here while data:image/ is not. DATA:... is accepted then dropped by the parser; HTTPS:// is 400. Use the same casefold gate in both places. Fixed on #684.

raise RequestError(400, "invalid_message_content", "image_url must be https or data:image")
cleaned.append({"type": "image_url", "image_url": {"url": url}})
continue
raise RequestError(400, "invalid_message_content", "content part type must be text or image_url")
return cleaned


def _validate_messages(messages: Any) -> list[dict[str, Any]]:
if not isinstance(messages, list) or not messages:
raise RequestError(400, "invalid_message", "messages must be a non-empty array")
validated: list[dict[str, str]] = []
validated: list[dict[str, Any]] = []
for message in messages:
if not isinstance(message, dict):
raise RequestError(400, "invalid_message", "each message must be an object")
role = message.get("role")
content = message.get("content")
if not isinstance(role, str) or role not in ALLOWED_MESSAGE_ROLES or not isinstance(content, str):
if not isinstance(role, str) or role not in ALLOWED_MESSAGE_ROLES:
raise RequestError(400, "invalid_message", "message role or content is invalid")
if isinstance(content, list):
content = _validate_content_parts(content)
elif not isinstance(content, str):
raise RequestError(400, "invalid_message", "message role or content is invalid")
validated.append({"role": role, "content": content})
return validated
Expand Down
49 changes: 41 additions & 8 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,15 @@

## Sources Read

- Sakana AI launch article, "Sakana Fugu: One Model to Command Them All" (June 22, 2026): https://sakana.ai/fugu-release/
- Sakana Fugu Technical Report: https://github.com/SakanaAI/fugu/blob/main/Fugu_technical_report.pdf
- TRINITY: An Evolved LLM Coordinator: https://arxiv.org/abs/2512.04695
- Learning to Orchestrate Agents in Natural Language with the Conductor: https://arxiv.org/abs/2512.04388
APA 7th citations (titles retained for paper-contract search):

- Sakana AI. (2026, June 22). *Sakana Fugu: One model to command them all*. https://sakana.ai/fugu-release/
- Tang, Y., Cetin, E., Xu, J., Sun, Q., Nielsen, S., Richard, V., Goda, H., Tymchenko, I., Nguyen, N., Lee, H., Ashiga, M., Kotyan, S., Kuroki, S., & Clanuwat, T. (2026). *Sakana Fugu technical report* (arXiv:2606.21228). arXiv. https://doi.org/10.48550/arXiv.2606.21228
- Xu, J., Sun, Q., Schwendeman, P., Nielsen, S., Cetin, E., & Tang, Y. (2026). *TRINITY: An evolved LLM coordinator* (arXiv:2512.04695). arXiv. https://doi.org/10.48550/arXiv.2512.04695
- Nielsen, S., Cetin, E., Schwendeman, P., Sun, Q., Xu, J., & Tang, Y. (2026). *Learning to orchestrate agents in natural language with the Conductor* (arXiv:2512.04388). arXiv. https://doi.org/10.48550/arXiv.2512.04388
- Faysse, M., Sibille, H., Wu, T., Omrani, B., Viaud, G., Hudelot, C., & Colombo, P. (2024). *ColPali: Efficient document retrieval with vision language models* (arXiv:2407.01449). arXiv. https://doi.org/10.48550/arXiv.2407.01449
- Xu, Y., Li, M., Cui, L., Huang, S., Wei, F., & Zhou, M. (2020). LayoutLM: Pre-training of text and layout for document image understanding. In *Proceedings of the 26th ACM SIGKDD International Conference on Knowledge Discovery & Data Mining* (pp. 1192–1200). Association for Computing Machinery. https://doi.org/10.1145/3394486.3403172
- Masinter, L. (1998). *The "data" URL scheme* (RFC 2397). Internet Engineering Task Force. https://doi.org/10.17487/RFC2397

## What The Architecture Is

Expand All @@ -31,12 +36,40 @@ The Fugu report combines these ideas into production constraints:

This repository implements the interface and control plane, not the trained coordinator.

- `contextual_orchestrator.orchestrator.Agent`: one configured worker model.
- `Orchestrator.route_once`: the low-latency routing path.
- `Orchestrator.conduct`: the workflow path with planner, worker, verifier, and synthesizer steps.
- `contextual_orchestrator.orchestrator.ModelAgent`: one configured worker model.
- `TaskOrchestrator.route_once`: the low-latency routing path.
- `TaskOrchestrator.conduct`: the workflow path with planner, worker, verifier, and synthesizer steps.
- `WorkflowStep.access`: Conductor-style visibility control.
- `collect_image_catalog`: 3NF image payload / placement / recognition-event split so a figure stays next to its pay line.

```mermaid
erDiagram
IMAGE_PAYLOAD ||--o{ IMAGE_PLACEMENT : appears_on
IMAGE_PAYLOAD ||--o{ IMAGE_RECOGNITION_EVENT : recognized_as
WORKFLOW_RUN ||--o{ IMAGE_PLACEMENT : contains
IMAGE_PAYLOAD {
text payload_digest PK
text mime_type
int byte_length
}
IMAGE_PLACEMENT {
text placement_id PK
text payload_digest FK
int message_index
int part_index
text source_kind
text adjacent_text
}
IMAGE_RECOGNITION_EVENT {
text recognition_event_id PK
text payload_digest FK
text recognized_text
text object_tags
timestamptz observed_at
}
```
- `ModelClient`: OpenAI-compatible HTTP client, with `mock://` for local checks.
- `contextual_orchestrator.server`: small `/v1/chat/completions` HTTP server.
- `contextual_orchestrator.server`: small `/v1/chat/completions` HTTP server. Buyer next action: send OpenAI `text` + `image_url` parts; read `orchestration.image_content_catalog` to find the figure.

The deliberate simplification is the policy. The paper systems learn routing and topology from rewards; this lab uses deterministic keyword scoring so the repo runs without training data, GPUs, or vendor credentials.

Expand Down
Loading
Loading