diff --git a/.claude/skills/forge.md b/.claude/skills/forge.md index 1c6cf0b6..c469f3cc 100644 --- a/.claude/skills/forge.md +++ b/.claude/skills/forge.md @@ -1200,7 +1200,7 @@ when OTel tracing is enabled (OTel v1 / Phase 4 / #105). Both use | `AuditScheduleModify` | `schedule_modify` | Schedule mutated at runtime | | `EventAuthVerify` | `auth_verify` | Inbound request authenticated (`provider`, `user_id`, `org_id`, `token_kind`; `email` when the identity carries one). **Channel invoker:** for a channel-originated request the transport credential is the loopback token (`provider:internal`/`user_id:forge-internal`, recorded truthfully) and the human sender is stamped as `channel`/`channel_user`/`channel_email` from the `X-Forge-Channel*` headers — honored only for the runtime-internal identity (same trust gate as `applyChannelOnBehalfOf`). Slack/Teams resolve `channel_email`; Telegram (numeric id) & WhatsApp (msisdn) carry `channel_user` only | | `EventAuthFail` | `auth_fail` | Inbound request rejected (`reason`, `token_kind`) | -| `AuditInputMediaRejected` | `input_media_rejected` | Inbound `file` parts the runtime can't forward to the model → rejected 4xx, not silently dropped (#255). Fields: `dropped` (`["file:"]`), `count`, `reason` (`model_not_vision_capable` \| `unsupported_media_type`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image` source blocks / OpenAI `image_url` data URLs. Docs/video still rejected (Phase 3+) | +| `AuditInputMediaRejected` | `input_media_rejected` | Inbound `file` parts the runtime can't forward to the model → rejected 4xx, not silently dropped (#255). Fields: `dropped` (`["file:"]`), `count`, `reason` (`model_not_vision_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `image_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image` source blocks / OpenAI `image_url` data URLs. Docs/video still rejected (Phase 3+). **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP (`CheckImageLimits`, header-only decode defuses bombs); ≤20 images/msg; ≤4 concurrent media requests (excess shed with 429/unavailable) | | `EventMCPServerStarted` | `mcp_server_started` | MCP server handshake succeeded | | `EventMCPServerFailed` | `mcp_server_failed` | MCP server dial / handshake failed | | `EventMCPServerDegraded` | `mcp_server_degraded` | MCP server in soft-fail | diff --git a/docs/core-concepts/runtime-engine.md b/docs/core-concepts/runtime-engine.md index 648ffe83..820dbcd3 100644 --- a/docs/core-concepts/runtime-engine.md +++ b/docs/core-concepts/runtime-engine.md @@ -39,6 +39,15 @@ An image `file` part (`image/png`, `image/jpeg`, `image/gif`, `image/webp`) is f Media the model can't consume is **rejected loudly, never silently dropped** (the `checkInboundMedia` ingest gate): an image on a text-only model, or a document/video part (not yet supported), returns a 4xx and emits the `input_media_rejected` audit event. Note the image **bytes** themselves are not text-scannable, so guardrail/intent scanning still applies only to the text/data projection; this is an accepted limitation. +**DoS bounds.** Because inline images raise the inbound-body cap to 32 MiB (both transports), the gate also enforces per-image and per-message limits, and a concurrency semaphore bounds how many media-bearing requests run at once — a flat body cap alone is not media DoS protection: + +| Bound | Limit | On breach | +|-------|-------|-----------| +| Per-image bytes | 5 MiB (`MaxImagePartBytes`) | 4xx `image_limit_exceeded` | +| Decoded dimensions | 100 000 px per side **and** 50 MP total (`MaxImagePixels`, png/jpeg/gif via header-only `DecodeConfig`; webp bounded by bytes) | 4xx `image_limit_exceeded` (defuses decompression bombs; the per-side bound also keeps the pixel product from overflowing int64) | +| Images per message | 20 (`MaxImagePartsPerMessage`) | 4xx `too_many_image_parts` | +| Concurrent media requests | 4 (`maxConcurrentMediaRequests`) | `429`/unavailable — request is shed, not queued | + **Note for guardrail pattern authors:** parts join with **newlines** (matching what the model sees). A pattern intended to match content that may span a part boundary should use `\s+` rather than a literal space — a payload split across two text parts joins as `…end\nstart…`. ### Session Recovery Deduplication diff --git a/docs/security/audit-logging.md b/docs/security/audit-logging.md index 26dfd4dc..b664c8d7 100644 --- a/docs/security/audit-logging.md +++ b/docs/security/audit-logging.md @@ -32,7 +32,7 @@ All runtime security events are emitted as structured NDJSON to stderr with corr | `mcp_auth_resolved` | The parked call's consent arrived (or the wait was canceled) and it resumed (#330). Carries `server`, `subject`, `wait_ms`. Emitted **once**, attributed to the parked invocation (#366). | | `mcp_auth_timeout` | No consent within the window; the parked MCP call fails `no_token` (#330). Carries `server`, `subject`, `wait_ms`, `decision`. | | `auth_fail` | Inbound request rejected (with `reason`, `token_kind`). No `task_id` (none is ever created), but carries `workflow_execution_id` when the request had the execution header — so a rejected request is still attributable to its workflow run (#278). | -| `input_media_rejected` | An inbound `tasks/send`/`sendSubscribe` message carried `file` parts the runtime cannot forward to the model, so the request was rejected with a 4xx (JSON-RPC invalid-params / HTTP 400) instead of silently dropping the attachment and returning a plausible answer that ignored it (#255). Carries `fields.dropped` (`["file:", …]`), `fields.count`, and `fields.reason` — `model_not_vision_capable` (an image sent to a text-only model) or `unsupported_media_type` (documents/video — not yet accepted). **Images** (`png`/`jpeg`/`gif`/`webp`) sent to a **vision-capable model** are NOT rejected: they are projected into the model request as inline vision input and do not emit this event. | +| `input_media_rejected` | An inbound `tasks/send`/`sendSubscribe` message carried `file` parts the runtime cannot forward to the model, so the request was rejected with a 4xx (JSON-RPC invalid-params / HTTP 400) instead of silently dropping the attachment and returning a plausible answer that ignored it (#255). Carries `fields.dropped` (`["file:", …]`), `fields.count`, and `fields.reason` — `model_not_vision_capable` (image sent to a text-only model), `unsupported_media_type` (documents/video — not yet accepted), `too_many_image_parts` (over the per-message image limit), or `image_limit_exceeded` (an image over the per-image byte or pixel bound). **Images** (`png`/`jpeg`/`gif`/`webp`) within those bounds, sent to a **vision-capable model**, are NOT rejected: they are projected into the model request as inline vision input and do not emit this event. A separate load-shedding path returns an unavailable/`429` (no audit event) when the runner is already at its concurrent-media-request cap. | | `agent_card_published` | Agent Card finalized at startup or hot-reload (with `name`, `version`, `protocol_version`, `url`, `skill_count`, `capabilities`, `security_schemes`, `card_size_bytes`, `card_sha256`). See [Agent Card reference](../reference/a2a-agent-card.md). | | `policy_loaded` | One per non-empty policy layer at startup (system / user / workspace). Carries `fields.layer`, `source` (file path), deny-list size counts, and max bounds. See [Platform Policy](platform-policy.md). | | `policy_violation_at_build_time` | One per violation when `forge.yaml` conflicts with any policy layer. Agent refuses to start. Carries `fields.violation_kind` / `offending_value` / `forge_yaml_field` plus `layer` + `source` identifying the enforcing file. See [Platform Policy](platform-policy.md). | diff --git a/forge-cli/internal/surface/knowledge/forge.md b/forge-cli/internal/surface/knowledge/forge.md index 1c6cf0b6..c469f3cc 100644 --- a/forge-cli/internal/surface/knowledge/forge.md +++ b/forge-cli/internal/surface/knowledge/forge.md @@ -1200,7 +1200,7 @@ when OTel tracing is enabled (OTel v1 / Phase 4 / #105). Both use | `AuditScheduleModify` | `schedule_modify` | Schedule mutated at runtime | | `EventAuthVerify` | `auth_verify` | Inbound request authenticated (`provider`, `user_id`, `org_id`, `token_kind`; `email` when the identity carries one). **Channel invoker:** for a channel-originated request the transport credential is the loopback token (`provider:internal`/`user_id:forge-internal`, recorded truthfully) and the human sender is stamped as `channel`/`channel_user`/`channel_email` from the `X-Forge-Channel*` headers — honored only for the runtime-internal identity (same trust gate as `applyChannelOnBehalfOf`). Slack/Teams resolve `channel_email`; Telegram (numeric id) & WhatsApp (msisdn) carry `channel_user` only | | `EventAuthFail` | `auth_fail` | Inbound request rejected (`reason`, `token_kind`) | -| `AuditInputMediaRejected` | `input_media_rejected` | Inbound `file` parts the runtime can't forward to the model → rejected 4xx, not silently dropped (#255). Fields: `dropped` (`["file:"]`), `count`, `reason` (`model_not_vision_capable` \| `unsupported_media_type`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image` source blocks / OpenAI `image_url` data URLs. Docs/video still rejected (Phase 3+) | +| `AuditInputMediaRejected` | `input_media_rejected` | Inbound `file` parts the runtime can't forward to the model → rejected 4xx, not silently dropped (#255). Fields: `dropped` (`["file:"]`), `count`, `reason` (`model_not_vision_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `image_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image` source blocks / OpenAI `image_url` data URLs. Docs/video still rejected (Phase 3+). **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP (`CheckImageLimits`, header-only decode defuses bombs); ≤20 images/msg; ≤4 concurrent media requests (excess shed with 429/unavailable) | | `EventMCPServerStarted` | `mcp_server_started` | MCP server handshake succeeded | | `EventMCPServerFailed` | `mcp_server_failed` | MCP server dial / handshake failed | | `EventMCPServerDegraded` | `mcp_server_degraded` | MCP server in soft-fail | diff --git a/forge-cli/runtime/media_gate_test.go b/forge-cli/runtime/media_gate_test.go index 4170c7fa..194d3fcb 100644 --- a/forge-cli/runtime/media_gate_test.go +++ b/forge-cli/runtime/media_gate_test.go @@ -1,7 +1,10 @@ package runtime import ( + "bytes" "context" + "encoding/binary" + "hash/crc32" "strings" "testing" @@ -14,12 +17,28 @@ func runnerWithModel(model string) *Runner { return &Runner{modelConfig: &coreruntime.ModelConfig{Client: llm.ClientConfig{Model: model}}} } -func imagePartMsg(mime string) a2a.Message { +// validPNG builds a minimal decodable PNG (signature + IHDR) with the given +// dimensions — enough for image.DecodeConfig, without a real pixel buffer. +func validPNG(w, h uint32) []byte { + var buf bytes.Buffer + buf.Write([]byte{0x89, 'P', 'N', 'G', 0x0d, 0x0a, 0x1a, 0x0a}) + ihdr := make([]byte, 0, 13) + ihdr = binary.BigEndian.AppendUint32(ihdr, w) + ihdr = binary.BigEndian.AppendUint32(ihdr, h) + ihdr = append(ihdr, 8, 2, 0, 0, 0) + _ = binary.Write(&buf, binary.BigEndian, uint32(len(ihdr))) + buf.WriteString("IHDR") + buf.Write(ihdr) + _ = binary.Write(&buf, binary.BigEndian, crc32.ChecksumIEEE(append([]byte("IHDR"), ihdr...))) + return buf.Bytes() +} + +func fileMsg(mime string, data []byte) a2a.Message { return a2a.Message{ Role: a2a.MessageRoleUser, Parts: []a2a.Part{ a2a.NewTextPart("look"), - a2a.NewFilePart(a2a.FileContent{MimeType: mime, Bytes: []byte{1, 2, 3}}), + a2a.NewFilePart(a2a.FileContent{MimeType: mime, Bytes: data}), }, } } @@ -31,14 +50,14 @@ func TestCheckInboundMedia_Phase2(t *testing.T) { t.Run("image accepted on vision model", func(t *testing.T) { r := runnerWithModel("gpt-4o") - if got := r.checkInboundMedia(ctx, imagePartMsg("image/png"), nil); got != "" { + if got := r.checkInboundMedia(ctx, fileMsg("image/png", validPNG(64, 64)), nil); got != "" { t.Errorf("image on a vision model must be accepted, got reject: %q", got) } }) t.Run("image rejected on non-vision model", func(t *testing.T) { r := runnerWithModel("gpt-3.5-turbo") - got := r.checkInboundMedia(ctx, imagePartMsg("image/png"), nil) + got := r.checkInboundMedia(ctx, fileMsg("image/png", validPNG(64, 64)), nil) if got == "" { t.Fatal("image on a non-vision model must be rejected") } @@ -49,12 +68,37 @@ func TestCheckInboundMedia_Phase2(t *testing.T) { t.Run("document rejected even on vision model", func(t *testing.T) { r := runnerWithModel("gpt-4o") - got := r.checkInboundMedia(ctx, imagePartMsg("application/pdf"), nil) + got := r.checkInboundMedia(ctx, fileMsg("application/pdf", []byte("%PDF-1.7")), nil) + if got == "" || !strings.Contains(got, "application/pdf") { + t.Errorf("a document part must still be rejected in Phase 2; got %q", got) + } + }) + + t.Run("oversized image rejected", func(t *testing.T) { + r := runnerWithModel("gpt-4o") + got := r.checkInboundMedia(ctx, fileMsg("image/png", make([]byte, coreruntime.MaxImagePartBytes+1)), nil) + if got == "" || !strings.Contains(got, "limit") { + t.Errorf("oversized image must be rejected with a size reason; got %q", got) + } + }) + + t.Run("decompression-bomb dimensions rejected", func(t *testing.T) { + r := runnerWithModel("gpt-4o") + got := r.checkInboundMedia(ctx, fileMsg("image/png", validPNG(100_000, 100_000)), nil) if got == "" { - t.Fatal("a document part must still be rejected in Phase 2") + t.Error("a gigapixel image must be rejected") } - if !strings.Contains(got, "application/pdf") { - t.Errorf("reject message should name the pdf mime; got %q", got) + }) + + t.Run("too many image parts rejected", func(t *testing.T) { + r := runnerWithModel("gpt-4o") + parts := []a2a.Part{a2a.NewTextPart("many")} + for range coreruntime.MaxImagePartsPerMessage + 1 { + parts = append(parts, a2a.NewFilePart(a2a.FileContent{MimeType: "image/png", Bytes: validPNG(16, 16)})) + } + got := r.checkInboundMedia(ctx, a2a.Message{Role: a2a.MessageRoleUser, Parts: parts}, nil) + if got == "" || !strings.Contains(got, "image") { + t.Errorf("more than %d images must be rejected; got %q", coreruntime.MaxImagePartsPerMessage, got) } }) @@ -68,8 +112,45 @@ func TestCheckInboundMedia_Phase2(t *testing.T) { t.Run("nil modelConfig treats as non-vision", func(t *testing.T) { r := &Runner{} - if got := r.checkInboundMedia(ctx, imagePartMsg("image/png"), nil); got == "" { + if got := r.checkInboundMedia(ctx, fileMsg("image/png", validPNG(16, 16)), nil); got == "" { t.Error("with no resolved model, image input must be rejected (fail closed)") } }) } + +// TestAcquireMediaSlot verifies the concurrency semaphore bounds media requests +// and never gates non-media requests (#255 DoS control). +func TestAcquireMediaSlot(t *testing.T) { + r := &Runner{mediaSem: make(chan struct{}, 2)} + + // Non-media requests are never bounded. + for range 5 { + if _, ok := r.acquireMediaSlot(false); !ok { + t.Fatal("non-media request must never be shed") + } + } + + // Media requests fill the 2 slots, then the 3rd is shed. + rel1, ok1 := r.acquireMediaSlot(true) + rel2, ok2 := r.acquireMediaSlot(true) + if !ok1 || !ok2 { + t.Fatal("first two media slots must be granted") + } + if _, ok3 := r.acquireMediaSlot(true); ok3 { + t.Fatal("third media slot must be shed when at capacity") + } + // Releasing one frees a slot. + rel1() + rel3, ok := r.acquireMediaSlot(true) + if !ok { + t.Fatal("a slot must be available after release") + } + rel2() + rel3() + + // A nil semaphore disables the bound (never sheds). + rNil := &Runner{} + if _, ok := rNil.acquireMediaSlot(true); !ok { + t.Fatal("nil mediaSem must disable the bound") + } +} diff --git a/forge-cli/runtime/runner.go b/forge-cli/runtime/runner.go index a3feb969..c1eb4bf6 100644 --- a/forge-cli/runtime/runner.go +++ b/forge-cli/runtime/runner.go @@ -203,6 +203,7 @@ type Runner struct { taskStore *a2a.TaskStore // shared task store, populated once srv is built; read by defer hook when it fires platformCommandGuard *coreruntime.PlatformCommandGuard // #238 (ASI02) operator-authored command deny, applied to every tool call; empty when no layer declares denied_command_patterns killed atomic.Bool // kill switch: set by the admin/kill handler; when true, tasks/send + tasks/sendSubscribe refuse new work (in-flight work is cancelled via cancelRegistry.CancelAll, then the platform scales the workload to zero) + mediaSem chan struct{} // #255 DoS control: bounds concurrent media-bearing requests so raising the body cap can't be turned into a memory-exhaustion vector; nil disables the bound } // NewRunner creates a Runner from the given config. @@ -225,6 +226,7 @@ func NewRunner(cfg RunnerConfig) (*Runner, error) { logger: logger, cancelRegistry: coreruntime.NewCancellationRegistry(), seqRegistry: coreruntime.NewSequenceRegistry(), + mediaSem: make(chan struct{}, maxConcurrentMediaRequests), }, nil } @@ -1685,6 +1687,12 @@ func (r *Runner) registerHandlers(srv *server.Server, executor coreruntime.Agent r.logger.Warn("tasks/send rejected: unsupported media", map[string]any{"task_id": params.ID, "reason": reason}) return a2a.NewErrorResponse(id, a2a.ErrCodeInvalidParams, reason) } + // Bound concurrent media-bearing requests (#255 DoS control). + release, ok := r.acquireMediaSlot(len(params.Message.FileParts()) > 0) + if !ok { + return a2a.NewErrorResponse(id, a2a.ErrCodeUnavailable, "server is busy processing media requests; retry shortly") + } + defer release() r.logger.Info("tasks/send", map[string]any{"task_id": params.ID}) // Delegate to executeTask so JSON-RPC and REST share the same // audit + accumulator + invocation_complete wiring (issue #87 / @@ -1743,6 +1751,14 @@ func (r *Runner) registerHandlers(srv *server.Server, executor coreruntime.Agent a2a.NewErrorResponse(id, a2a.ErrCodeInvalidParams, reason)) return } + // Bound concurrent media-bearing requests (#255 DoS control). + release, ok := r.acquireMediaSlot(len(params.Message.FileParts()) > 0) + if !ok { + server.WriteSSEEvent(w, flusher, "error", //nolint:errcheck + a2a.NewErrorResponse(id, a2a.ErrCodeUnavailable, "server is busy processing media requests; retry shortly")) + return + } + defer release() r.logger.Info("tasks/sendSubscribe", map[string]any{"task_id": params.ID}) @@ -2292,17 +2308,38 @@ type restTaskRequest struct { // maxRequestBodyBytes bounds an inbound REST A2A request body — an unbounded // json.Decode on req.Body is a trivial memory-exhaustion vector. Kept in -// parity with the JSON-RPC transport cap in server.handleJSONRPC (both 2 MiB) +// parity with the JSON-RPC transport cap in server.handleJSONRPC (both 32 MiB) // so the inbound-body bound is uniform across transports. // -// Deliberately NOT raised to admit inline media yet: Phase 0 (#255) rejects all -// file parts, so a larger cap would only widen the memory-exhaustion surface -// without admitting anything usable. The media-consuming phase raises this -// (both transports together) alongside the controls that make a large cap safe -// — per-part count/size limits, image decode-dimension bounds, and a -// concurrency semaphore — since a flat MaxBytesReader alone is not media DoS -// protection. -const maxRequestBodyBytes = 2 << 20 // 2 MiB — matches server.handleJSONRPC +// Raised to 32 MiB now that forge consumes inline images (#255) — but ONLY +// alongside the controls that make a large cap safe: per-image size + count +// limits and decode-dimension bounds (checkInboundMedia + CheckImageLimits), +// and the mediaSem concurrency semaphore. A flat MaxBytesReader alone is not +// media DoS protection; these bound total, per-part, and concurrent memory. +const maxRequestBodyBytes = 32 << 20 // 32 MiB — matches server.handleJSONRPC + +// maxConcurrentMediaRequests bounds how many media-bearing requests the runner +// processes at once (#255). Peak media memory ≈ this × maxRequestBodyBytes; a +// request that can't get a slot is shed with a 429/unavailable rather than +// queued, so a burst of large uploads can't exhaust memory. +const maxConcurrentMediaRequests = 4 + +// acquireMediaSlot bounds concurrent media-bearing requests. For a request with +// no media it is a no-op. For a media request it tries (non-blocking) to take a +// semaphore slot: on success it returns a release func and true; when the +// runner is at capacity it returns (nil, false) so the caller sheds the request +// instead of queueing it. +func (r *Runner) acquireMediaSlot(hasMedia bool) (release func(), ok bool) { + if !hasMedia || r.mediaSem == nil { + return func() {}, true + } + select { + case r.mediaSem <- struct{}{}: + return func() { <-r.mediaSem }, true + default: + return nil, false + } +} // checkInboundMedia rejects a message carrying file parts the runtime cannot // forward to the model (#255). Today forge projects only text/data parts into @@ -2325,18 +2362,42 @@ func (r *Runner) checkInboundMedia(ctx context.Context, msg a2a.Message, auditLo var dropped []string reason := "unsupported_media_type" - for _, f := range files { - mt := f.MimeType - if mt == "" { - mt = "application/octet-stream" - } - if coreruntime.IsImageMIME(mt) { - if visionCapable { - continue // accepted → projected to the model as vision input + imageCount := 0 + for _, p := range msg.Parts { + if p.Kind != a2a.PartKindFile { + continue + } + mt := "application/octet-stream" + var data []byte + if p.File != nil { + if p.File.MimeType != "" { + mt = p.File.MimeType } + data = p.File.Bytes + } + if !coreruntime.IsImageMIME(mt) { + reason = "unsupported_media_type" + dropped = append(dropped, "file:"+mt) + continue + } + if !visionCapable { reason = "model_not_vision_capable" + dropped = append(dropped, "file:"+mt) + continue + } + // Image on a vision model: enforce the DoS bounds (#255). + imageCount++ + if imageCount > coreruntime.MaxImagePartsPerMessage { + reason = "too_many_image_parts" + dropped = append(dropped, fmt.Sprintf("file:%s (over the %d-image limit)", mt, coreruntime.MaxImagePartsPerMessage)) + continue } - dropped = append(dropped, "file:"+mt) + if lim := coreruntime.CheckImageLimits(mt, data); lim != "" { + reason = "image_limit_exceeded" + dropped = append(dropped, fmt.Sprintf("file:%s (%s)", mt, lim)) + continue + } + // accepted → projected to the model as vision input } if len(dropped) == 0 { return "" // all file parts are images the model can consume @@ -2352,13 +2413,18 @@ func (r *Runner) checkInboundMedia(ctx context.Context, msg a2a.Message, auditLo }) } var detail string - if reason == "model_not_vision_capable" { + switch reason { + case "model_not_vision_capable": model := "" if r.modelConfig != nil { model = r.modelConfig.Client.Model } detail = fmt.Sprintf("the configured model %q does not support image input", model) - } else { + case "too_many_image_parts": + detail = fmt.Sprintf("a message may carry at most %d image parts", coreruntime.MaxImagePartsPerMessage) + case "image_limit_exceeded": + detail = fmt.Sprintf("each image must be a decodable png/jpeg/gif/webp under %d bytes and %d pixels", coreruntime.MaxImagePartBytes, coreruntime.MaxImagePixels) + default: detail = "only image input (png/jpeg/gif/webp) is accepted; documents and video are not yet supported" } return fmt.Sprintf( @@ -2420,6 +2486,13 @@ func (r *Runner) registerRESTHandlers(srv *server.Server, executor coreruntime.A writeJSON(w, http.StatusBadRequest, map[string]string{"error": reason}) return } + // Bound concurrent media-bearing requests (#255 DoS control). + release, ok := r.acquireMediaSlot(len(params.Message.FileParts()) > 0) + if !ok { + writeJSON(w, http.StatusTooManyRequests, map[string]string{"error": "server is busy processing media requests; retry shortly"}) + return + } + defer release() task, snap, err := r.executeTask(ctx, params, store, executor, guardrails, egressClient, auditLogger) if err != nil { // R4b: a step-up-required error takes priority — we @@ -2475,6 +2548,14 @@ func (r *Runner) registerRESTHandlers(srv *server.Server, executor coreruntime.A writeJSON(w, http.StatusBadRequest, map[string]string{"error": reason}) return } + // Bound concurrent media-bearing requests, while a plain non-stream + // status is still possible — before the SSE Content-Type commits (#255). + release, ok := r.acquireMediaSlot(len(body.Task.Message.FileParts()) > 0) + if !ok { + writeJSON(w, http.StatusTooManyRequests, map[string]string{"error": "server is busy processing media requests; retry shortly"}) + return + } + defer release() w.Header().Set("Content-Type", "text/event-stream") w.Header().Set("Cache-Control", "no-cache") diff --git a/forge-cli/server/a2a_server.go b/forge-cli/server/a2a_server.go index ad93ad8c..f7c08cf6 100644 --- a/forge-cli/server/a2a_server.go +++ b/forge-cli/server/a2a_server.go @@ -264,7 +264,7 @@ func (s *Server) handleHealthz(w http.ResponseWriter, r *http.Request) { } func (s *Server) handleJSONRPC(w http.ResponseWriter, r *http.Request) { - r.Body = http.MaxBytesReader(w, r.Body, 2<<20) // 2 MiB — parity with runtime.maxRequestBodyBytes; raised for inline media in the #255 media phase + r.Body = http.MaxBytesReader(w, r.Body, 32<<20) // 32 MiB — parity with runtime.maxRequestBodyBytes; admits inline media (#255, with per-image + concurrency limits enforced downstream) var req a2a.JSONRPCRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { diff --git a/forge-cli/server/a2a_server_test.go b/forge-cli/server/a2a_server_test.go index a9f7ad33..234ecaa7 100644 --- a/forge-cli/server/a2a_server_test.go +++ b/forge-cli/server/a2a_server_test.go @@ -186,7 +186,7 @@ func TestHandleJSONRPC_OversizedBody(t *testing.T) { // Create a 3 MiB payload that's valid JSON start but oversized (cap 2 MiB). // MaxBytesReader will cut it off, causing a parse error that includes // the "request body too large" message. - huge := `{"data":"` + strings.Repeat("x", 3<<20) + `"}` + huge := `{"data":"` + strings.Repeat("x", 33<<20) + `"}` req := httptest.NewRequest(http.MethodPost, "/", strings.NewReader(huge)) req.Header.Set("Content-Type", "application/json") rec := httptest.NewRecorder() diff --git a/forge-core/runtime/media_limits.go b/forge-core/runtime/media_limits.go new file mode 100644 index 00000000..0740646a --- /dev/null +++ b/forge-core/runtime/media_limits.go @@ -0,0 +1,83 @@ +package runtime + +import ( + "bytes" + "fmt" + "image" + _ "image/gif" // register gif DecodeConfig + _ "image/jpeg" // register jpeg DecodeConfig + _ "image/png" // register png DecodeConfig +) + +// Media ingest limits (#255 DoS controls). These bound what a single inbound +// message can carry so that raising the request-body cap to admit real images +// cannot be turned into a memory/CPU-exhaustion vector. They are enforced at +// the ingest gate (before dispatch), alongside a concurrency semaphore that +// bounds how many media-bearing requests are processed at once. +const ( + // MaxImagePartBytes caps a single image part's raw (pre-base64) bytes. + // 5 MiB matches Anthropic's per-image limit and fits several images plus + // text under the request-body cap. + MaxImagePartBytes = 5 << 20 + + // MaxImagePartsPerMessage caps how many image parts one message may carry, + // so a body full of tiny images can't fan out into a huge provider request. + MaxImagePartsPerMessage = 20 + + // MaxImagePixels bounds decoded image dimensions (width*height) to defuse + // decompression bombs — a tiny file that decodes to a gigapixel canvas. + // ~50 MP covers legitimate high-resolution photos. Enforced only for the + // stdlib-decodable formats (png/jpeg/gif); webp is bounded by the byte cap. + MaxImagePixels = 50_000_000 + + // maxImageDim bounds a single dimension BEFORE the width*height product, so + // that product can't overflow int64. PNG's IHDR carries 32-bit dimensions + // (JPEG/GIF are 16-bit), so a forged PNG can report dims near 2^31; two such + // dims multiplied would wrap int64 negative and slip past the pixel check. + // 100_000 keeps the product ≤ 1e10 (safe) and no legitimate image is that + // large on any single side. + maxImageDim = 100_000 +) + +// CheckImageLimits validates a single image part's bytes against the per-image +// size and decode-dimension bounds. It returns a non-empty reason string when +// the part must be rejected, or "" when it passes. +// +// The dimension check uses image.DecodeConfig, which reads only the header and +// does NOT allocate the pixel buffer, so it is cheap. Formats without a +// registered decoder (webp) can't be dimension-checked here and are bounded by +// the byte cap alone; a malformed image in a format we DO support is rejected. +func CheckImageLimits(mime string, data []byte) string { + if len(data) == 0 { + return "image part has no bytes" + } + if len(data) > MaxImagePartBytes { + return fmt.Sprintf("image exceeds the %d-byte per-image limit (%d bytes)", MaxImagePartBytes, len(data)) + } + cfg, _, err := image.DecodeConfig(bytes.NewReader(data)) + if err != nil { + // No registered decoder for this format (e.g. webp) → rely on the byte + // cap. A decode error for a format we DO support means it's malformed. + if imageFormatIsStdlibDecodable(mime) { + return "image could not be decoded (malformed or truncated)" + } + return "" + } + // Per-dimension bound first, so the product below can't overflow int64 + // (a forged PNG can report dims near 2^31; the raw product would wrap). + if cfg.Width > maxImageDim || cfg.Height > maxImageDim { + return fmt.Sprintf("image dimension %dx%d exceeds the %d-px per-side limit", cfg.Width, cfg.Height, maxImageDim) + } + if int64(cfg.Width)*int64(cfg.Height) > MaxImagePixels { + return fmt.Sprintf("image dimensions %dx%d exceed the %d-pixel limit", cfg.Width, cfg.Height, MaxImagePixels) + } + return "" +} + +func imageFormatIsStdlibDecodable(mime string) bool { + switch NormalizeImageMIME(mime) { + case "image/png", "image/jpeg", "image/gif": + return true + } + return false +} diff --git a/forge-core/runtime/media_limits_test.go b/forge-core/runtime/media_limits_test.go new file mode 100644 index 00000000..e5a1ffb6 --- /dev/null +++ b/forge-core/runtime/media_limits_test.go @@ -0,0 +1,113 @@ +package runtime + +import ( + "bytes" + "encoding/binary" + "hash/crc32" + "strings" + "testing" +) + +// pngWithDims builds a minimal but valid PNG (signature + IHDR chunk) reporting +// the given dimensions. image.DecodeConfig reads only the header, so this is +// enough to exercise the dimension bound WITHOUT allocating a real pixel +// buffer — exactly the decompression-bomb shape (tiny bytes, huge reported +// dimensions). +func pngWithDims(t *testing.T, w, h uint32) []byte { + t.Helper() + var buf bytes.Buffer + buf.Write([]byte{0x89, 'P', 'N', 'G', 0x0d, 0x0a, 0x1a, 0x0a}) // signature + ihdr := make([]byte, 0, 13) + ihdr = binary.BigEndian.AppendUint32(ihdr, w) + ihdr = binary.BigEndian.AppendUint32(ihdr, h) + ihdr = append(ihdr, 8, 2, 0, 0, 0) // bit depth 8, color type 2 (RGB), no compression/filter/interlace + _ = binary.Write(&buf, binary.BigEndian, uint32(len(ihdr))) + buf.WriteString("IHDR") + buf.Write(ihdr) + crc := crc32.ChecksumIEEE(append([]byte("IHDR"), ihdr...)) + _ = binary.Write(&buf, binary.BigEndian, crc) + return buf.Bytes() +} + +func TestCheckImageLimits(t *testing.T) { + small := pngWithDims(t, 100, 100) + + t.Run("valid small png passes", func(t *testing.T) { + if got := CheckImageLimits("image/png", small); got != "" { + t.Errorf("valid png rejected: %q", got) + } + }) + + t.Run("empty bytes rejected", func(t *testing.T) { + if CheckImageLimits("image/png", nil) == "" { + t.Error("empty image bytes must be rejected") + } + }) + + t.Run("oversized bytes rejected before decode", func(t *testing.T) { + big := make([]byte, MaxImagePartBytes+1) + if CheckImageLimits("image/png", big) == "" { + t.Error("image over the byte cap must be rejected") + } + }) + + t.Run("decompression-bomb dimensions rejected", func(t *testing.T) { + bomb := pngWithDims(t, 100_000, 100_000) // 1e10 px, but tiny bytes + if len(bomb) > 1024 { + t.Fatalf("forged png unexpectedly large (%d bytes)", len(bomb)) + } + got := CheckImageLimits("image/png", bomb) + if got == "" { + t.Fatal("a tiny file reporting gigapixel dimensions must be rejected") + } + }) + + t.Run("single oversized dimension rejected (product alone would pass)", func(t *testing.T) { + // 100_001 x 1 = 100_001 px — UNDER MaxImagePixels, so the raw + // width*height check would ACCEPT it. The per-dimension bound is what + // rejects it, and it's also the guard that keeps the product from + // overflowing int64 for a forged 32-bit PNG dimension (#532 review). + strip := pngWithDims(t, 100_001, 1) + got := CheckImageLimits("image/png", strip) + if got == "" { + t.Fatal("an image with a single dimension over the per-side limit must be rejected") + } + if !strings.Contains(got, "per-side") { + t.Errorf("reject reason should cite the per-side limit; got %q", got) + } + }) + + t.Run("large square dimensions rejected before the product multiply", func(t *testing.T) { + // DecodeConfig returns these (verified), so the per-dimension bound + // fires before int64(w)*int64(h) is ever computed. + if CheckImageLimits("image/png", pngWithDims(t, 200_000, 200_000)) == "" { + t.Error("200000x200000 must be rejected") + } + }) + + t.Run("forged overflow dimensions rejected (decoder guards, we don't rely on it)", func(t *testing.T) { + // width=height=2^32-1: Go's png DecodeConfig itself errors on this + // (non-positive after int32 wrap), so it's rejected as malformed — the + // security property holds regardless of the decoder's internal cap. + if CheckImageLimits("image/png", pngWithDims(t, 0xFFFFFFFF, 0xFFFFFFFF)) == "" { + t.Error("a forged gigantic-dimension png must be rejected") + } + }) + + t.Run("malformed png in a supported format rejected", func(t *testing.T) { + if CheckImageLimits("image/png", []byte("not really a png")) == "" { + t.Error("malformed png must be rejected") + } + }) + + t.Run("webp bounded by byte cap only (no stdlib decoder)", func(t *testing.T) { + // Small webp-ish bytes: no registered decoder → passes (byte cap holds). + if got := CheckImageLimits("image/webp", []byte("RIFF....WEBP....")); got != "" { + t.Errorf("small webp should pass (byte-cap-only): %q", got) + } + // Oversized webp still rejected on size. + if CheckImageLimits("image/webp", make([]byte, MaxImagePartBytes+1)) == "" { + t.Error("oversized webp must be rejected on the byte cap") + } + }) +}