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
1 change: 1 addition & 0 deletions .claude/skills/forge.md
Original file line number Diff line number Diff line change
Expand Up @@ -1188,6 +1188,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`) |
| `EventAuthFail` | `auth_fail` | Inbound request rejected (`reason`, `token_kind`) |
| `AuditInputMediaRejected` | `input_media_rejected` | Inbound message carried `file` parts the runtime can't forward to the model → rejected 4xx instead of silently dropped (#255). Fields: `dropped` (`["file:<mime>"]`), `count`, `reason`. Gate: `Runner.checkInboundMedia` at the four send handlers; `a2a.Message.FileParts()` |
| `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 |
Expand Down
1 change: 1 addition & 0 deletions docs/security/audit-logging.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +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 one or more `file` parts (image/document/…) 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:<mimeType>", …]`), `fields.count`, and `fields.reason` (`file_input_unsupported`). Multimodal input arrives in later slices; until then this makes the accepted-but-dropped case loud. |
| `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). |
Expand Down
1 change: 1 addition & 0 deletions forge-cli/internal/surface/knowledge/forge.md
Original file line number Diff line number Diff line change
Expand Up @@ -1188,6 +1188,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`) |
| `EventAuthFail` | `auth_fail` | Inbound request rejected (`reason`, `token_kind`) |
| `AuditInputMediaRejected` | `input_media_rejected` | Inbound message carried `file` parts the runtime can't forward to the model → rejected 4xx instead of silently dropped (#255). Fields: `dropped` (`["file:<mime>"]`), `count`, `reason`. Gate: `Runner.checkInboundMedia` at the four send handlers; `a2a.Message.FileParts()` |
| `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 |
Expand Down
80 changes: 80 additions & 0 deletions forge-cli/runtime/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -1679,6 +1679,12 @@ func (r *Runner) registerHandlers(srv *server.Server, executor coreruntime.Agent
})
return a2a.NewErrorResponse(id, a2a.ErrCodeInvalidParams, "invalid message: "+err.Error())
}
// Reject file/image parts the runtime can't forward to the model
// instead of silently dropping them (#255).
if reason := r.checkInboundMedia(ctx, params.Message, auditLogger); reason != "" {
r.logger.Warn("tasks/send rejected: unsupported media", map[string]any{"task_id": params.ID, "reason": reason})
return a2a.NewErrorResponse(id, a2a.ErrCodeInvalidParams, reason)
}
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 /
Expand Down Expand Up @@ -1729,6 +1735,14 @@ func (r *Runner) registerHandlers(srv *server.Server, executor coreruntime.Agent
a2a.NewErrorResponse(id, a2a.ErrCodeInvalidParams, "invalid message: "+err.Error()))
return
}
// Reject unsupported file/image parts before the stream carries any
// content (#255).
if reason := r.checkInboundMedia(ctx, params.Message, auditLogger); reason != "" {
r.logger.Warn("tasks/sendSubscribe rejected: unsupported media", map[string]any{"task_id": params.ID, "reason": reason})
server.WriteSSEEvent(w, flusher, "error", //nolint:errcheck
a2a.NewErrorResponse(id, a2a.ErrCodeInvalidParams, reason))
return
}

r.logger.Info("tasks/sendSubscribe", map[string]any{"task_id": params.ID})

Expand Down Expand Up @@ -2276,6 +2290,57 @@ type restTaskRequest struct {
} `json:"task"`
}

// 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)
// 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

// checkInboundMedia rejects a message carrying file parts the runtime cannot
// forward to the model (#255). Today forge projects only text/data parts into
// the prompt (a2a.Message.PromptText); file parts (images/documents) are not
// yet consumed by any provider, so accepting them would silently drop the
// attachment and return a plausible answer that ignored it — the footgun this
// closes. It returns a human-readable reason for the 4xx / JSON-RPC error and
// emits an input_media_rejected audit event, or "" when the message is
// acceptable. Centralized so all four send handlers share one contract.
func (r *Runner) checkInboundMedia(ctx context.Context, msg a2a.Message, auditLogger *coreruntime.AuditLogger) string {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

LOW (future-proofing): per-handler gating, no single choke point. This gate is invoked from 4 handlers; current coverage is complete (verified: 2 executeTask callers + 2 ExecuteStream handlers, all gated). But there's no one shared point below them in the runtime — the non-streaming paths funnel through executeTask, the streaming ones through ExecuteStream — so a future 5th ingest path (a new endpoint, an in-process channel→executor call, a batch/replay) would silently reintroduce the footgun unless it remembers this gate. Consider a structural backstop at the executor boundary (forge-core Execute/ExecuteStream) in a later #255 phase so media rejection is guaranteed regardless of caller. Not needed for Phase 0 — just flagging the maintenance risk of a 4-call-site invariant.

files := msg.FileParts()
if len(files) == 0 {
return ""
}
dropped := make([]string, 0, len(files))
for _, f := range files {
mt := f.MimeType
if mt == "" {
mt = "application/octet-stream"
}
dropped = append(dropped, "file:"+mt)
}
if auditLogger != nil {
auditLogger.EmitFromContext(ctx, coreruntime.AuditEvent{
Event: coreruntime.AuditInputMediaRejected,
Fields: map[string]any{
"dropped": dropped,
"count": len(files),
"reason": "file_input_unsupported",
},
})
}
return fmt.Sprintf(
"this agent does not accept file/image input: %d file part(s) rejected (%s). Send text or data parts instead.",
len(files), strings.Join(dropped, ", "),
)
}

// registerRESTHandlers registers REST-style HTTP endpoints on the server.
func (r *Runner) registerRESTHandlers(srv *server.Server, executor coreruntime.AgentExecutor, guardrails coreruntime.GuardrailChecker, egressClient *http.Client, auditLogger *coreruntime.AuditLogger) {
store := srv.TaskStore()
Expand All @@ -2286,6 +2351,7 @@ func (r *Runner) registerRESTHandlers(srv *server.Server, executor coreruntime.A
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "agent disabled by kill switch: not accepting new tasks"})
return
}
req.Body = http.MaxBytesReader(w, req.Body, maxRequestBodyBytes)
var body restTaskRequest
if err := json.NewDecoder(req.Body).Decode(&body); err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request body: " + err.Error()})
Expand Down Expand Up @@ -2322,6 +2388,12 @@ func (r *Runner) registerRESTHandlers(srv *server.Server, executor coreruntime.A
// Same for tenancy override headers (#157).
ctx = coreruntime.WithTenancyContext(ctx,
coreruntime.TenancyContextFromHTTPHeaders(req.Header))
// Reject file/image parts the runtime can't forward to the model (#255).
if reason := r.checkInboundMedia(ctx, params.Message, auditLogger); reason != "" {
r.logger.Warn("REST /tasks/send rejected: unsupported media", map[string]any{"task_id": params.ID, "reason": reason, "remote_addr": req.RemoteAddr})
writeJSON(w, http.StatusBadRequest, map[string]string{"error": reason})
return
}
task, snap, err := r.executeTask(ctx, params, store, executor, guardrails, egressClient, auditLogger)
if err != nil {
// R4b: a step-up-required error takes priority — we
Expand Down Expand Up @@ -2349,6 +2421,7 @@ func (r *Runner) registerRESTHandlers(srv *server.Server, executor coreruntime.A
return
}

req.Body = http.MaxBytesReader(w, req.Body, maxRequestBodyBytes)
var body restTaskRequest
if err := json.NewDecoder(req.Body).Decode(&body); err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request body: " + err.Error()})
Expand All @@ -2369,6 +2442,13 @@ func (r *Runner) registerRESTHandlers(srv *server.Server, executor coreruntime.A
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid message: " + err.Error()})
return
}
// Reject unsupported file/image parts while a plain 400 is still
// possible — before the SSE Content-Type is committed (#255).
if reason := r.checkInboundMedia(req.Context(), body.Task.Message, auditLogger); reason != "" {
r.logger.Warn("REST /tasks/sendSubscribe rejected: unsupported media", map[string]any{"task_id": body.Task.ID, "reason": reason, "remote_addr": req.RemoteAddr})
writeJSON(w, http.StatusBadRequest, map[string]string{"error": reason})
return
}

w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
Expand Down
117 changes: 117 additions & 0 deletions forge-cli/runtime/runner_media_reject_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
package runtime

import (
"context"
"encoding/json"
"fmt"
"net/http"
"strings"
"testing"
"time"

"github.com/initializ/forge/forge-core/a2a"
"github.com/initializ/forge/forge-core/auth"
"github.com/initializ/forge/forge-core/types"
)

// TestRunner_RejectsInboundFileParts is the #255 footgun regression: a client
// attaching a file (image/document) part used to get a 200 and a plausible
// answer that silently ignored the attachment, because the runtime projects
// only text/data parts into the prompt. Now every send path rejects a file
// part loudly (JSON-RPC error / HTTP 400) naming the mime type, while a
// text-only message still succeeds.
func TestRunner_RejectsInboundFileParts(t *testing.T) {
dir := t.TempDir()
cfg := &types.ForgeConfig{
AgentID: "media-reject-test",
Version: "0.1.0",
Framework: "forge",
Entrypoint: "python main.py",
Tools: []types.ToolRef{{Name: "search"}},
}
port, err := findFreePort()
if err != nil {
t.Fatal(err)
}
runner, err := NewRunner(RunnerConfig{Config: cfg, WorkDir: dir, Port: port, MockTools: true})
if err != nil {
t.Fatalf("NewRunner: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go func() { _ = runner.Run(ctx) }()

baseURL := fmt.Sprintf("http://localhost:%d", port)
waitForServer(t, baseURL, 5*time.Second)
token, err := auth.LoadToken(dir)
if err != nil {
t.Fatalf("loading auth token: %v", err)
}

fileMsg := a2a.Message{
Role: a2a.MessageRoleUser,
Parts: []a2a.Part{
a2a.NewTextPart("describe this"),
a2a.NewFilePart(a2a.FileContent{Name: "photo.png", MimeType: "image/png", Bytes: []byte{0x89, 0x50, 0x4e, 0x47}}),
},
}

t.Run("json-rpc tasks/send rejects file part", func(t *testing.T) {
rpcReq := a2a.JSONRPCRequest{
JSONRPC: "2.0", ID: "1", Method: "tasks/send",
Params: mustMarshal(a2a.SendTaskParams{ID: "t-media-1", Message: fileMsg}),
}
body, _ := json.Marshal(rpcReq)
resp, err := authPost(baseURL+"/", token, body)
if err != nil {
t.Fatalf("send: %v", err)
}
defer func() { _ = resp.Body.Close() }()

var rpcResp a2a.JSONRPCResponse
if err := json.NewDecoder(resp.Body).Decode(&rpcResp); err != nil {
t.Fatalf("decode: %v", err)
}
if rpcResp.Error == nil {
t.Fatalf("expected a JSON-RPC error rejecting the file part, got result: %+v", rpcResp.Result)
}
if !strings.Contains(rpcResp.Error.Message, "image/png") {
t.Errorf("error message should name the rejected mime type; got %q", rpcResp.Error.Message)
}
})

t.Run("REST /tasks/send rejects file part with 400", func(t *testing.T) {
body, _ := json.Marshal(restBody("t-media-2", fileMsg))
resp, err := authPost(baseURL+"/tasks/send", token, body)
if err != nil {
t.Fatalf("send: %v", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode != http.StatusBadRequest {
t.Fatalf("status = %d, want 400", resp.StatusCode)
}
var errBody map[string]string
_ = json.NewDecoder(resp.Body).Decode(&errBody)
if !strings.Contains(errBody["error"], "image/png") {
t.Errorf("400 body should name the rejected mime type; got %q", errBody["error"])
}
})

t.Run("text-only message still succeeds", func(t *testing.T) {
textMsg := a2a.Message{Role: a2a.MessageRoleUser, Parts: []a2a.Part{a2a.NewTextPart("hello")}}
body, _ := json.Marshal(restBody("t-text-ok", textMsg))
resp, err := authPost(baseURL+"/tasks/send", token, body)
if err != nil {
t.Fatalf("send: %v", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode != http.StatusOK {
t.Fatalf("text-only status = %d, want 200 (media gate must not affect text)", resp.StatusCode)
}
})
}

// restBody builds the REST {"task":{"id","message"}} envelope.
func restBody(id string, msg a2a.Message) map[string]any {
return map[string]any{"task": map[string]any{"id": id, "message": msg}}
}
2 changes: 1 addition & 1 deletion forge-cli/server/a2a_server.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
r.Body = http.MaxBytesReader(w, r.Body, 2<<20) // 2 MiB — parity with runtime.maxRequestBodyBytes; raised for inline media in the #255 media phase

var req a2a.JSONRPCRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
Expand Down
2 changes: 1 addition & 1 deletion forge-cli/server/a2a_server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -183,7 +183,7 @@ func TestSecurityHeadersOnErrorResponses(t *testing.T) {

func TestHandleJSONRPC_OversizedBody(t *testing.T) {
s := NewServer(ServerConfig{Port: 0})
// Create a 3 MiB payload that's valid JSON start but oversized.
// 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) + `"}`
Expand Down
63 changes: 63 additions & 0 deletions forge-core/a2a/file_parts_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
package a2a

import "testing"

// TestMessageFileParts verifies FileParts surfaces exactly the file parts (with
// index + mime + name) and ignores text/data parts — the basis for the #255
// ingest gate that rejects media the runtime can't forward to the model.
func TestMessageFileParts(t *testing.T) {
tests := []struct {
name string
parts []Part
want []MediaPartInfo
}{
{
name: "text only",
parts: []Part{NewTextPart("hi")},
want: nil,
},
{
name: "data only",
parts: []Part{NewDataPart(map[string]any{"k": "v"})},
want: nil,
},
{
name: "single file part",
parts: []Part{
NewTextPart("look at this"),
NewFilePart(FileContent{Name: "photo.png", MimeType: "image/png", Bytes: []byte{1, 2, 3}}),
},
want: []MediaPartInfo{{Index: 1, MimeType: "image/png", Name: "photo.png"}},
},
{
name: "multiple file parts keep their indices",
parts: []Part{
NewFilePart(FileContent{Name: "a.pdf", MimeType: "application/pdf"}),
NewTextPart("between"),
NewFilePart(FileContent{Name: "b.mp4", MimeType: "video/mp4"}),
},
want: []MediaPartInfo{
{Index: 0, MimeType: "application/pdf", Name: "a.pdf"},
{Index: 2, MimeType: "video/mp4", Name: "b.mp4"},
},
},
{
name: "file part with nil FileContent still reported",
parts: []Part{{Kind: PartKindFile}},
want: []MediaPartInfo{{Index: 0}},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got := Message{Role: MessageRoleUser, Parts: tt.parts}.FileParts()
if len(got) != len(tt.want) {
t.Fatalf("FileParts() = %+v, want %+v", got, tt.want)
}
for i := range got {
if got[i] != tt.want[i] {
t.Errorf("FileParts()[%d] = %+v, want %+v", i, got[i], tt.want[i])
}
}
})
}
}
31 changes: 31 additions & 0 deletions forge-core/a2a/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,37 @@ func (m Message) PromptText() string {
return strings.Join(segments, "\n")
}

// MediaPartInfo describes a non-text "file" part for capability gating and
// diagnostics. PromptText does not project file parts, so a caller that cannot
// forward them to the model must reject the request rather than drop the
// attachment silently (#255).
type MediaPartInfo struct {
Index int
MimeType string
Name string
}

// FileParts returns info for every file part in the message. File parts carry
// bytes/URIs (images, documents, …) that only reach the model through a
// multimodal path; because PromptText omits them, accepting a message with
// file parts on a path that can't consume them silently discards the
// attachment — the caller uses this to reject loudly instead.
func (m Message) FileParts() []MediaPartInfo {
var out []MediaPartInfo
for i, p := range m.Parts {
if p.Kind != PartKindFile {
continue
}
info := MediaPartInfo{Index: i}
if p.File != nil {
info.MimeType = p.File.MimeType
info.Name = p.File.Name
}
out = append(out, info)
}
return out
}

// PartKind discriminates the content type of a Part.
type PartKind string

Expand Down
Loading
Loading