Skip to content
Merged
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 README.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,8 @@
| **Build Custom** | 6-12 months, $500K+ engineering | Deploy in minutes |
| **Frameworks** | No observability, no fleet management | Real-time monitoring, scheduling, audit trails |

> **Security**: Trinity received a **Grade A (Excellent)** in an independent web application penetration test by [UnderDefense](https://underdefense.com) (April 2026). All critical and high findings from the initial assessment were fully remediated. See the [attestation letter](docs/security/UnderDefense-Web-Pentest-Attestation-Apr-2026.pdf).

---

## Getting Started — Deploy an Agent in 3 Minutes
Expand Down Expand Up @@ -613,6 +615,7 @@ EMAIL_PROVIDER=console # Use 'resend' or 'smtp' for production
- [Testing Guide](docs/TESTING_GUIDE.md) — Testing approach and standards
- [Contributing Guide](CONTRIBUTING.md) — How to contribute (PRs, code standards)
- [Known Issues](docs/KNOWN_ISSUES.md) — Current limitations and workarounds
- [Security Attestation](docs/security/UnderDefense-Web-Pentest-Attestation-Apr-2026.pdf) — Web pentest by UnderDefense (Apr 2026, Grade A — Excellent)

## Development

Expand Down
1 change: 1 addition & 0 deletions docker/base-image/agent_server/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,7 @@ class ParallelTaskRequest(BaseModel):
max_turns: Optional[int] = None # Maximum agentic turns (--max-turns) for runaway prevention
execution_id: Optional[str] = None # Database execution ID (used for process registry if provided)
resume_session_id: Optional[str] = None # Claude Code session ID for --resume (EXEC-023)
images: Optional[List[Dict[str, str]]] = None # Vision images: [{"media_type": "image/jpeg", "data": "<base64>"}]


class ParallelTaskResponse(BaseModel):
Expand Down
3 changes: 2 additions & 1 deletion docker/base-image/agent_server/routers/chat.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,8 @@ async def execute_task(request: ParallelTaskRequest):
timeout_seconds=request.timeout_seconds or 900, # Default 15 minutes for research tasks
max_turns=request.max_turns,
execution_id=request.execution_id, # Use provided ID for process registry (enables termination)
resume_session_id=request.resume_session_id # Resume previous session (EXEC-023)
resume_session_id=request.resume_session_id, # Resume previous session (EXEC-023)
images=request.images, # Vision images from channel adapters (#562)
)

logger.info(f"[Task] Task {session_id} completed successfully")
Expand Down
57 changes: 47 additions & 10 deletions docker/base-image/agent_server/services/claude_code.py
Original file line number Diff line number Diff line change
Expand Up @@ -140,16 +140,18 @@ async def execute_headless(
timeout_seconds: int = 900,
max_turns: Optional[int] = None,
execution_id: Optional[str] = None,
resume_session_id: Optional[str] = None
resume_session_id: Optional[str] = None,
images: Optional[List[Dict]] = None,
) -> Tuple[str, List[ExecutionLogEntry], ExecutionMetadata, str]:
"""Execute Claude Code in headless mode for parallel tasks.

Args:
resume_session_id: Optional session ID to resume (EXEC-023)
images: Optional list of vision images: [{"media_type": str, "data": base64_str}] (#562)
"""
return await execute_headless_task(
prompt, model, allowed_tools, system_prompt, timeout_seconds,
max_turns, execution_id, resume_session_id
max_turns, execution_id, resume_session_id, images=images,
)


Expand Down Expand Up @@ -1016,7 +1018,8 @@ async def execute_headless_task(
timeout_seconds: int = 900,
max_turns: Optional[int] = None,
execution_id: Optional[str] = None,
resume_session_id: Optional[str] = None
resume_session_id: Optional[str] = None,
images: Optional[List[Dict]] = None,
) -> tuple[str, List[ExecutionLogEntry], ExecutionMetadata, str]:
"""
Execute Claude Code in headless mode for parallel task execution.
Expand Down Expand Up @@ -1101,6 +1104,12 @@ async def execute_headless_task(
cmd.extend(["--disallowedTools", ",".join(disallowed_tools)])
logger.info(f"[Headless Task] Guardrails disallow tools: {disallowed_tools}")

# #562: when images are present, use stream-json stdin format so images
# are delivered as proper vision content blocks, not base64 text strings.
if images:
cmd.extend(["--input-format", "stream-json"])
logger.info(f"[Headless Task] {len(images)} image(s) — switching to stream-json input")

# Add system prompt if specified
if system_prompt:
cmd.extend(["--append-system-prompt", system_prompt])
Expand Down Expand Up @@ -1153,10 +1162,6 @@ async def execute_headless_task(
"pgid": process_pgid,
})

# Write prompt to stdin and close it
process.stdin.write(prompt)
process.stdin.close()

# Issue #285: Event to signal auth failure detected in stderr
# When set, stdout loop should stop and process should be killed
auth_abort_event = threading.Event()
Expand Down Expand Up @@ -1263,14 +1268,46 @@ def _run_stdout():
pass

def read_subprocess_output_with_timeout():
"""Runs in thread pool. Waits for subprocess with bounded timeout,
then drains reader threads (killing process-group stragglers if
they hold pipes open — Issue #407)."""
"""Runs in thread pool. Writes stdin, starts reader threads, and
waits for subprocess with bounded timeout, then drains reader
threads (killing process-group stragglers if they hold pipes
open — Issue #407).

Stdin is written here (not in the async coroutine) so that:
1. Large payloads (e.g. base64 images) do not block the event loop.
2. Reader threads are active before the write, preventing pipe-
buffer deadlock if claude writes stdout before stdin is closed.
"""
stderr_thread = threading.Thread(target=read_stderr, daemon=True)
stdout_thread = threading.Thread(target=_run_stdout, daemon=True)
stderr_thread.start()
stdout_thread.start()

# Build and write stdin payload. For vision tasks use stream-json
# format so images arrive as proper content blocks (#562).
if images:
content_blocks: List[Dict] = [
{
"type": "image",
"source": {
"type": "base64",
"media_type": img["media_type"],
"data": img["data"],
},
}
for img in images
]
content_blocks.append({"type": "text", "text": prompt})
stdin_payload = (
json.dumps({"type": "user", "message": {"role": "user", "content": content_blocks}})
+ "\n"
)
else:
stdin_payload = prompt

process.stdin.write(stdin_payload)
process.stdin.close()

# Bounded wait on the subprocess itself. If claude hangs, we
# never wedge the executor thread for more than timeout_seconds.
try:
Expand Down
3 changes: 2 additions & 1 deletion docker/base-image/agent_server/services/gemini_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -506,7 +506,8 @@ async def execute_headless(
timeout_seconds: int = 900,
max_turns: Optional[int] = None,
execution_id: Optional[str] = None,
resume_session_id: Optional[str] = None
resume_session_id: Optional[str] = None,
images: Optional[List[Dict]] = None,
) -> Tuple[str, List[ExecutionLogEntry], ExecutionMetadata, str]:
"""
Execute Gemini CLI in headless mode for parallel tasks.
Expand Down
3 changes: 2 additions & 1 deletion docker/base-image/agent_server/services/runtime_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,8 @@ async def execute_headless(
timeout_seconds: int = 900,
max_turns: Optional[int] = None,
execution_id: Optional[str] = None,
resume_session_id: Optional[str] = None
resume_session_id: Optional[str] = None,
images: Optional[List[Dict]] = None,
) -> Tuple[str, List[ExecutionLogEntry], ExecutionMetadata, str]:
"""
Execute a stateless task in headless mode (no conversation context).
Expand Down
1 change: 1 addition & 0 deletions docs/memory/feature-flows/task-execution-service.md
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,7 @@ async def execute_task(
parent_activity_id: Optional[str] = None, # Issue #95: CHAT_START parent linkage
extra_activity_details: Optional[dict] = None, # Issue #95: merged into CHAT_START details
slot_already_held: bool = False, # Issue #95: async path pre-acquires slot upfront
images: Optional[list] = None, # #562: vision content blocks for channel images
) -> TaskExecutionResult:
```

Expand Down
39 changes: 38 additions & 1 deletion docs/memory/feature-flows/telegram-integration.md
Original file line number Diff line number Diff line change
Expand Up @@ -808,7 +808,7 @@ the Slack pattern).

**Chat injection format** (every channel, not just Telegram):
- Successful file: `[File uploaded by {uploader}]: {filename} ({size}) saved to {dest_path}`
- Successful image: `[File uploaded by {uploader}]: {filename} ({size}) — image attached inline` followed by `![{name}](data:{mime};base64,…)`
- Successful image: `[File uploaded by {uploader}]: {filename} ({size}) — image provided for visual analysis` (image delivered as vision content block — see #562)
- Workspace write failure: `[File upload failed]: {filename} — {reason}`

`{uploader}` is the verified email when `adapter.resolve_verified_email`
Expand Down Expand Up @@ -837,8 +837,44 @@ human-readable provenance for any file in its workspace.
`dest_path`, `storage` (`container_file` or `inline_base64`), and
`uploader` so the audit trail captures who uploaded what to which agent.

### Phase 3: Vision Delivery via stream-json (#562)

**Problem**: Images embedded as `data:` URIs in the text prompt were opaque strings — the Claude Code CLI never forwarded them to the Claude API as real image content blocks, so agents could not see the images.

**Fix**: `_handle_file_uploads()` now returns image MIME/base64 dicts as a 4th element of its tuple (`image_data: list`). For image MIME types, the file is still written to the container workspace (read tool fallback), but it is also collected for vision delivery. The description changes from `"image attached inline"` to `"image provided for visual analysis"` (no `data:` URI appended).

`message_router.py` passes `images=image_data or None` into `task_execution_service.execute_task()`. The service adds it to the `/api/task` HTTP payload. The agent-side `chat.py` router passes it to `runtime.execute_headless()`. In `claude_code.py`, when `images` is truthy, two changes take effect:

1. `--input-format stream-json` is appended to the `claude` CLI command.
2. `stdin` payload is a JSON message instead of plain text:
```json
{"type": "user", "message": {"role": "user", "content": [
{"type": "image", "source": {"type": "base64", "media_type": "image/jpeg", "data": "<b64>"}},
{"type": "text", "text": "<prompt>"}
]}}
```

The `GeminiRuntime` and `AgentRuntime` ABC gained an `images` parameter (ignored by Gemini today) to prevent `TypeError` on any image task.

**Stdin write ordering**: stdout/stderr reader threads are started before `process.stdin.write()` to prevent pipe deadlock when the image payload is large. The write itself runs inside the executor thread (not the async event loop) to avoid blocking.

**Key files**:
- `src/backend/adapters/message_router.py` — 4-tuple return, `images or None` passthrough
- `src/backend/services/task_execution_service.py` — `images` param, forwarded in payload
- `docker/base-image/agent_server/models.py` — `ParallelTaskRequest.images` field
- `docker/base-image/agent_server/routers/chat.py` — passes `images` to `execute_headless()`
- `docker/base-image/agent_server/services/claude_code.py` — stream-json stdin + flag
- `docker/base-image/agent_server/services/runtime_adapter.py` — ABC signature updated
- `docker/base-image/agent_server/services/gemini_runtime.py` — GeminiRuntime signature updated

### Tests

17 unit tests in `tests/unit/test_channel_image_vision.py` (#562):
- `TestParallelTaskRequestImages` (4 tests): model field defaults and acceptance
- `TestStreamJsonPayload` (6 tests): payload construction for 0/1/N images
- `TestCmdContainsInputFormat` (4 tests): `--input-format stream-json` added iff images present
- `TestHandleFileUploadsImageReturn` (3 tests): 4-tuple return, image vs non-image MIME

27 unit tests in `tests/unit/test_file_upload.py`:
- `TestTelegramFileExtraction` (4 tests): Photo/document extraction
- `TestTelegramFileDownload` (3 tests): Bot API two-step download
Expand Down Expand Up @@ -884,3 +920,4 @@ human-readable provenance for any file in its workspace.
| 2026-04-16 | #354 Phase 1: Telegram file upload support. `_extract_files()` and `download_file()` in adapter. Post-download size/MIME validation in router. python-magic dependency. 11 unit tests. |
| 2026-04-18 | #318: Voice transcription via Gemini. `process_voice()` in telegram_media.py, voice processing hook in message_router.py. Limits: 5 min duration, 10MB size. 22 unit tests. |
| 2026-04-25 | #487 Phase 2: workspace delivery hardened. New `_sanitize_filename` helper (NFKC + basename + safe-chars + 200-char truncation + collision dedup). Chat injection format `[File uploaded by {uploader}]: {name} ({size}) saved to {path}`. All-writes-failed now replies via channel and aborts execution. Audit entries include `uploader`. 16 new tests (27 total in `test_file_upload.py`). |
| 2026-04-28 | #562: Vision delivery fixed. Replaced broken base64 data URI text embedding with proper `--input-format stream-json` vision content blocks delivered via Claude CLI stdin. `_handle_file_uploads` returns 4-tuple with `image_data`. `GeminiRuntime` and `AgentRuntime` ABC updated to accept `images` param. 17 new tests in `test_channel_image_vision.py`. |
Binary file not shown.
34 changes: 20 additions & 14 deletions src/backend/adapters/message_router.py
Original file line number Diff line number Diff line change
Expand Up @@ -465,8 +465,9 @@ async def _handle_message_inner(self, adapter: ChannelAdapter, message: Normaliz
# Issue #487: workspace-write failures abort execution and surface a
# channel-native error so the user knows the upload didn't land.
upload_dir = None # Track for cleanup
image_data: list = []
if message.files:
file_descriptions, upload_dir, all_writes_failed = await self._handle_file_uploads(
file_descriptions, upload_dir, all_writes_failed, image_data = await self._handle_file_uploads(
adapter, message, agent_name, container, session_id,
verified_email=verified_email,
)
Expand Down Expand Up @@ -504,10 +505,9 @@ async def _handle_message_inner(self, adapter: ChannelAdapter, message: Normaliz
# Configurable via settings_service (default: WebSearch, WebFetch)
public_allowed_tools = _get_channel_allowed_tools()

# If user uploaded non-image files, agent needs Read to access them.
# Images are excluded: Claude Code crashes when reading PNGs with
# --allowedTools (API returns 400 "Could not process image" and the
# process exits without flushing stdout, hanging the pipe reader).
# Non-image files land at upload_dir in the container; agent needs Read
# to access them. Images are delivered as vision content blocks via
# --input-format stream-json (#562), so no Read is needed for them.
_IMAGE_MIMES = {"image/png", "image/jpeg", "image/gif", "image/webp", "image/svg+xml"}
has_readable_files = message.files and any(
f.mimetype not in _IMAGE_MIMES for f in message.files
Expand All @@ -525,6 +525,7 @@ async def _handle_message_inner(self, adapter: ChannelAdapter, message: Normaliz
source_user_email=source_email,
timeout_seconds=None, # Uses agent's configured timeout (TIMEOUT-001)
allowed_tools=public_allowed_tools,
images=image_data or None,
)

if result.status == "failed":
Expand Down Expand Up @@ -697,15 +698,16 @@ async def _handle_file_uploads(
) -> tuple:
"""
Download files via adapter and either:
- Images: embed as base64 data URI in the prompt (Claude vision)
- Images: collect as base64 vision objects for stream-json delivery (#562)
- Other files: copy into per-session dir in agent container

Returns (descriptions, upload_dir, all_writes_failed):
Returns (descriptions, upload_dir, all_writes_failed, image_data):
- descriptions: list of context strings for prompt injection
- upload_dir: container path to clean up after execution, or None
- all_writes_failed: True iff at least one file attempted a workspace
write but every such attempt failed; the caller should reply with
an explicit error and skip agent execution (Issue #487 AC6).
- image_data: list of {"media_type": str, "data": base64_str} for vision
"""
import base64

Expand All @@ -721,6 +723,7 @@ async def _handle_file_uploads(
safe_session_id = re.sub(r"[^a-zA-Z0-9_-]", "_", session_id)
upload_dir = f"{UPLOAD_BASE}/{safe_session_id}"
descriptions = []
image_data: list = []
dir_created = False
total_image_bytes = 0
used_names: set = set()
Expand Down Expand Up @@ -814,22 +817,25 @@ async def _handle_file_uploads(
size_str = _format_file_size(actual_size)

if is_image:
# Check total inline image budget
# Check total image budget
if total_image_bytes + len(data) > MAX_TOTAL_IMAGE_SIZE:
logger.warning(f"[ROUTER] Skipping {safe_name}: total image budget exceeded")
descriptions.append(f"{safe_name} ({size_str}) — skipped (total image size limit reached)")
continue

# Image embedding is the "write" for vision-mode files.
# #562: collect as vision content block for stream-json delivery.
# Passing images as proper API content blocks (via --input-format
# stream-json) instead of base64 data URIs in text, which are
# invisible to the model — Claude Code passes stdin as plain text.
write_attempted += 1
total_image_bytes += len(data)
b64 = base64.b64encode(data).decode()
image_data.append({"media_type": actual_mime, "data": b64})
descriptions.append(
f"[File uploaded by {uploader}]: {safe_name} ({size_str}) — image attached inline\n"
f"![{safe_name}](data:{actual_mime};base64,{b64})"
f"[File uploaded by {uploader}]: {safe_name} ({size_str}) — image provided for visual analysis"
)
write_succeeded += 1
logger.info(f"[ROUTER] Embedded {safe_name} ({size_str}) as base64 for {agent_name}")
logger.info(f"[ROUTER] Queued {safe_name} ({size_str}) as vision block for {agent_name}")

# Audit log for image upload
await platform_audit_service.log(
Expand All @@ -842,7 +848,7 @@ async def _handle_file_uploads(
"filename": safe_name,
"size_bytes": actual_size,
"mime_type": actual_mime,
"storage": "inline_base64",
"storage": "stream_json_vision",
"sender_id": message.sender_id,
"channel_id": message.channel_id,
"uploader": uploader,
Expand Down Expand Up @@ -915,7 +921,7 @@ async def _handle_file_uploads(
descriptions.append(f"({len(files) - MAX_FILES} more file(s) skipped — max {MAX_FILES} per message)")

all_writes_failed = write_attempted > 0 and write_succeeded == 0
return descriptions, upload_dir if dir_created else None, all_writes_failed
return descriptions, upload_dir if dir_created else None, all_writes_failed, image_data


# Singleton instance
Expand Down
Loading
Loading