From 826717d9f1b45a479a691f77ae682abe74c1a95b Mon Sep 17 00:00:00 2001 From: MK Date: Mon, 28 Sep 2026 04:40:12 -0400 Subject: [PATCH 1/2] feat(llm): model-generated image output via OpenAI Responses (#255 Phase 5) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Forge can now return model-GENERATED images as file parts in the A2A response, completing the multipart-output direction. Opt-in and Responses-only (Anthropic/ Gemini/Bedrock don't generate images through their chat/messages APIs). - llm: StreamDelta gains Parts (non-text output); ClientConfig gains EnableImageGeneration. - Responses provider: when EnableImageGeneration is set, buildRequest sends the `image_generation` built-in tool; readStream reads the image_generation_call base64 `result` from the authoritative response.completed frame (not fragile partial-image events) and emits it as an image ContentPart; Chat aggregates delta.Parts into resp.Message.Parts. responsesTool.Name is now omitempty so the built-in tool serializes without a name. - Executor: llmMessageToA2A projects response image/document Parts into A2A `file` parts, and the assistant message's generated media is persisted like inbound media (URI in history, no base64 bloat; Bytes stay inline so the image is surfaced this turn). The empty-assistant-turn placeholder now respects a media-only response. - Config: ModelRef.ImageGeneration (yaml image_generation) → ClientConfig, off by default (changes provider behavior + cost). Tests: Responses parses a generated image from a synthesized response.completed SSE frame → resp.Message.Parts; image_generation tool sent only when opted in (and carries no name); llmMessageToA2A emits the generated image as a file part. Docs: runtime-engine output section, forge-yaml-schema image_generation, skill. Note: the OpenAI Responses image wire-shape follows the documented format and is fixture-tested; warrants live-API verification before production reliance. --- .claude/skills/forge.md | 2 +- docs/core-concepts/runtime-engine.md | 2 + docs/reference/forge-yaml-schema.md | 1 + forge-cli/internal/surface/knowledge/forge.md | 2 +- forge-core/llm/client.go | 7 ++ forge-core/llm/providers/responses.go | 30 ++++++- .../llm/providers/responses_image_test.go | 78 +++++++++++++++++++ forge-core/llm/types.go | 11 +-- forge-core/runtime/config.go | 5 ++ forge-core/runtime/loop.go | 21 ++++- forge-core/runtime/loop_projection_test.go | 30 +++++++ forge-core/types/config.go | 6 ++ 12 files changed, 186 insertions(+), 9 deletions(-) create mode 100644 forge-core/llm/providers/responses_image_test.go diff --git a/.claude/skills/forge.md b/.claude/skills/forge.md index df451ff2..2ef29b1c 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` \| `model_not_document_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `too_many_document_parts` \| `image_limit_exceeded` \| `document_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) and **PDFs** on a **doc model** (`ModelSupportsPDF` = Anthropic Sonnet 3.5+/Opus 4+/Haiku 4.5+/Fable 5; Claude 3.0 & 3.5-Haiku excluded, fail-closed), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image`/`document` source blocks, OpenAI `image_url` data URLs. Non-PDF docs/video still rejected; OpenAI-Responses PDF + extraction fallback are follow-ups. **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP/≤100k-per-side (`CheckImageLimits`); per-PDF ≤32 MiB + `%PDF-` sniff (`CheckDocumentLimits`); ≤20 images & ≤5 docs/msg; ≤4 concurrent media requests (excess shed with 429/unavailable). **Persistence/replay** (`persistInboundMedia` + `RehydrateMedia`): inbound media written to `.forge/files/inbound/.`, path stored as `MediaRef.URI`; history keeps URI-only (`Bytes` json:"-"), rehydrated per turn so multi-turn convos replay media without base64 bloat | +| `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` \| `model_not_document_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `too_many_document_parts` \| `image_limit_exceeded` \| `document_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) and **PDFs** on a **doc model** (`ModelSupportsPDF` = Anthropic Sonnet 3.5+/Opus 4+/Haiku 4.5+/Fable 5; Claude 3.0 & 3.5-Haiku excluded, fail-closed), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image`/`document` source blocks, OpenAI `image_url` data URLs. Non-PDF docs/video still rejected; OpenAI-Responses PDF + extraction fallback are follow-ups. **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP/≤100k-per-side (`CheckImageLimits`); per-PDF ≤32 MiB + `%PDF-` sniff (`CheckDocumentLimits`); ≤20 images & ≤5 docs/msg; ≤4 concurrent media requests (excess shed with 429/unavailable). **Persistence/replay** (`persistInboundMedia` + `RehydrateMedia`): inbound media written to `.forge/files/inbound/.`, path stored as `MediaRef.URI`; history keeps URI-only (`Bytes` json:"-"), rehydrated per turn so multi-turn convos replay media without base64 bloat. **Model-generated image OUTPUT** (#255 Phase 5): opt-in `image_generation: true` on an `openai-responses` model → sends the `image_generation` built-in tool; the `image_generation_call` base64 `result` (from `response.completed`) → `resp.Message.Parts` image → `llmMessageToA2A` emits a response `file` part (persisted like inbound). Responses-only (Anthropic/Gemini/Bedrock don't generate images via chat) | | `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 4443b016..d36cb3ef 100644 --- a/docs/core-concepts/runtime-engine.md +++ b/docs/core-concepts/runtime-engine.md @@ -42,6 +42,8 @@ A media `file` part is forwarded to the model as native input when the resolved `a2aMessageToLLM` projects supported parts into `llm.ChatMessage.Parts` (the flattened text stays in `Content` as the text-of-record for the scanners). A text-only message keeps `Parts` empty and marshals byte-identically to before. +**Model-generated image output.** A model can also *produce* an image, which forge returns as a `file` part in the A2A response. For OpenAI **Responses** (`provider: openai-responses`), set `image_generation: true` on the model (`ModelRef.ImageGeneration` → `ClientConfig.EnableImageGeneration`); the client then sends the `image_generation` built-in tool, and the base64 `result` from the `image_generation_call` output (read from the authoritative `response.completed` frame) is decoded into an image `ContentPart` on `resp.Message.Parts`. `llmMessageToA2A` projects those into response `file` parts, and the generated image is persisted like inbound media (URI in history, no base64 bloat). Off by default (it changes provider behavior and cost). Anthropic/Gemini/Bedrock don't generate images through their chat/messages APIs, so this is Responses-only today. + **Persistence & cross-turn replay.** Inbound media is written to `.forge/files/inbound/.` (`persistInboundMedia`) and the on-disk path is recorded as the part's `MediaRef.URI` — content-addressed, so identical uploads dedup and re-writes are idempotent, and the files are available on disk for tools. The bytes stay inline for the current turn's request; because `MediaRef.Bytes` is `json:"-"`, session history persists only the URI (never base64). On a later turn, `RehydrateMedia` reloads the bytes from the URI before the request is built, so a multi-turn conversation keeps seeing earlier images/PDFs without bloating the session file. A file that can't be reloaded (deleted, or an absolute path from another host — remote/distributed session replay is a follow-up) is left byteless and skipped by the provider serializers, degrading to text rather than failing the turn. Persistence and rehydration are best-effort: a failure never fails the turn (media is still fed inline the turn it arrives). Media the model can't consume is **rejected loudly, never silently dropped** (the `checkInboundMedia` ingest gate): an image on a text-only model, a PDF on a non-document model, or an unsupported type (other documents/video) returns a 4xx and emits the `input_media_rejected` audit event. Note media **bytes** are not text-scannable, so guardrail/intent scanning still applies only to the text/data projection; this is an accepted limitation. diff --git a/docs/reference/forge-yaml-schema.md b/docs/reference/forge-yaml-schema.md index f952a15b..d59a3293 100644 --- a/docs/reference/forge-yaml-schema.md +++ b/docs/reference/forge-yaml-schema.md @@ -24,6 +24,7 @@ model: aws_region: "" # Required when auth_scheme: aws_sigv4 — issue #202 auth_header_name: "" # apikey_header[_only] custom header name; default "apikey" — issue #302 disable_store: false # openai-responses only: send store=false so OpenAI doesn't retain responses — issue #383 + image_generation: false # openai-responses only: enable the image_generation tool so the model can return images as file parts — issue #255 fallbacks: # Fallback providers (optional) - provider: "anthropic" name: "claude-sonnet-4-20250514" diff --git a/forge-cli/internal/surface/knowledge/forge.md b/forge-cli/internal/surface/knowledge/forge.md index df451ff2..2ef29b1c 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` \| `model_not_document_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `too_many_document_parts` \| `image_limit_exceeded` \| `document_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) and **PDFs** on a **doc model** (`ModelSupportsPDF` = Anthropic Sonnet 3.5+/Opus 4+/Haiku 4.5+/Fable 5; Claude 3.0 & 3.5-Haiku excluded, fail-closed), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image`/`document` source blocks, OpenAI `image_url` data URLs. Non-PDF docs/video still rejected; OpenAI-Responses PDF + extraction fallback are follow-ups. **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP/≤100k-per-side (`CheckImageLimits`); per-PDF ≤32 MiB + `%PDF-` sniff (`CheckDocumentLimits`); ≤20 images & ≤5 docs/msg; ≤4 concurrent media requests (excess shed with 429/unavailable). **Persistence/replay** (`persistInboundMedia` + `RehydrateMedia`): inbound media written to `.forge/files/inbound/.`, path stored as `MediaRef.URI`; history keeps URI-only (`Bytes` json:"-"), rehydrated per turn so multi-turn convos replay media without base64 bloat | +| `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` \| `model_not_document_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `too_many_document_parts` \| `image_limit_exceeded` \| `document_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) and **PDFs** on a **doc model** (`ModelSupportsPDF` = Anthropic Sonnet 3.5+/Opus 4+/Haiku 4.5+/Fable 5; Claude 3.0 & 3.5-Haiku excluded, fail-closed), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image`/`document` source blocks, OpenAI `image_url` data URLs. Non-PDF docs/video still rejected; OpenAI-Responses PDF + extraction fallback are follow-ups. **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP/≤100k-per-side (`CheckImageLimits`); per-PDF ≤32 MiB + `%PDF-` sniff (`CheckDocumentLimits`); ≤20 images & ≤5 docs/msg; ≤4 concurrent media requests (excess shed with 429/unavailable). **Persistence/replay** (`persistInboundMedia` + `RehydrateMedia`): inbound media written to `.forge/files/inbound/.`, path stored as `MediaRef.URI`; history keeps URI-only (`Bytes` json:"-"), rehydrated per turn so multi-turn convos replay media without base64 bloat. **Model-generated image OUTPUT** (#255 Phase 5): opt-in `image_generation: true` on an `openai-responses` model → sends the `image_generation` built-in tool; the `image_generation_call` base64 `result` (from `response.completed`) → `resp.Message.Parts` image → `llmMessageToA2A` emits a response `file` part (persisted like inbound). Responses-only (Anthropic/Gemini/Bedrock don't generate images via chat) | | `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-core/llm/client.go b/forge-core/llm/client.go index 77fa766e..fda2e5f1 100644 --- a/forge-core/llm/client.go +++ b/forge-core/llm/client.go @@ -136,4 +136,11 @@ type ClientConfig struct { // unset so the API applies its own default (#383). The ChatGPT OAuth // path forces this true regardless (Codex backend requires it). DisableStore bool + + // EnableImageGeneration opts an OpenAI Responses request into the + // `image_generation` built-in tool, letting the model emit images that + // forge surfaces as `file` parts in the A2A response (#255). Off by + // default; only the openai-responses client honors it (other clients + // ignore the field). Opt-in because it changes provider behavior and cost. + EnableImageGeneration bool } diff --git a/forge-core/llm/providers/responses.go b/forge-core/llm/providers/responses.go index b37eb304..a2100efa 100644 --- a/forge-core/llm/providers/responses.go +++ b/forge-core/llm/providers/responses.go @@ -4,6 +4,7 @@ import ( "bufio" "bytes" "context" + "encoding/base64" "encoding/json" "fmt" "io" @@ -26,6 +27,7 @@ type ResponsesClient struct { authHeaderName string client *http.Client disableStore bool // set store=false in requests (required for ChatGPT Codex backend) + enableImageGen bool // send the image_generation built-in tool (#255) } // NewResponsesClient creates a new Responses API client. @@ -57,6 +59,7 @@ func NewResponsesClient(cfg llm.ClientConfig) *ResponsesClient { authScheme: cfg.AuthScheme, authHeaderName: cfg.AuthHeaderName, disableStore: cfg.DisableStore, + enableImageGen: cfg.EnableImageGeneration, client: httpClient, } } @@ -84,6 +87,9 @@ func (c *ResponsesClient) Chat(ctx context.Context, req *llm.ChatRequest) (*llm. if delta.Content != "" { result.Message.Content += delta.Content } + if len(delta.Parts) > 0 { + result.Message.Parts = append(result.Message.Parts, delta.Parts...) + } for _, tc := range delta.ToolCalls { existing, ok := toolCallMap[tc.ID] if !ok { @@ -207,7 +213,7 @@ type responsesInput struct { // responsesTool is the Responses API tool format (flat, not nested under "function"). type responsesTool struct { Type string `json:"type"` - Name string `json:"name"` + Name string `json:"name,omitempty"` // omitted for built-in tools (e.g. image_generation) Description string `json:"description,omitempty"` Parameters json.RawMessage `json:"parameters,omitempty"` } @@ -274,6 +280,11 @@ func (c *ResponsesClient) buildRequest(req *llm.ChatRequest, stream bool) respon Parameters: t.Function.Parameters, }) } + // Opt-in: let the model emit images via the image_generation built-in tool + // (#255). Surfaced back as file parts in the A2A response. + if c.enableImageGen { + tools = append(tools, responsesTool{Type: "image_generation"}) + } // The Responses API requires the instructions field. If no system // message was provided (e.g. summarization calls), use a minimal default @@ -319,6 +330,9 @@ type responsesOutput struct { CallID string `json:"call_id,omitempty"` Name string `json:"name,omitempty"` Arguments string `json:"arguments,omitempty"` + + // For image_generation_call outputs: base64-encoded image bytes (#255). + Result string `json:"result,omitempty"` } type responsesContentPart struct { @@ -473,6 +487,20 @@ func (c *ResponsesClient) readStream(r io.Reader, ch chan<- llm.StreamDelta) { TotalTokens: ev.Response.Usage.TotalTokens, } } + // Surface any model-generated images from the final output array as + // content parts (#255). The completed event carries the authoritative + // output[], so we read the image_generation_call `result` (base64) + // here rather than reassembling partial-image stream events. + for _, out := range ev.Response.Output { + if out.Type == "image_generation_call" && out.Result != "" { + if raw, derr := base64.StdEncoding.DecodeString(out.Result); derr == nil && len(raw) > 0 { + delta.Parts = append(delta.Parts, llm.NewMediaContentPart(llm.ContentPartImage, llm.MediaRef{ + MimeType: "image/png", // Responses image_generation defaults to PNG + Bytes: raw, + })) + } + } + } // Determine finish reason from output for _, out := range ev.Response.Output { if out.Type == "function_call" { diff --git a/forge-core/llm/providers/responses_image_test.go b/forge-core/llm/providers/responses_image_test.go new file mode 100644 index 00000000..db760ec3 --- /dev/null +++ b/forge-core/llm/providers/responses_image_test.go @@ -0,0 +1,78 @@ +package providers + +import ( + "context" + "encoding/base64" + "fmt" + "net/http" + "net/http/httptest" + "testing" + + "github.com/initializ/forge/forge-core/llm" +) + +// TestResponses_ParsesGeneratedImage verifies the Responses client surfaces a +// model-generated image (an image_generation_call output in response.completed) +// as an image content part on the aggregated response (#255 Phase 5). +func TestResponses_ParsesGeneratedImage(t *testing.T) { + raw := []byte("\x89PNG\r\n-fake-image-bytes") + b64 := base64.StdEncoding.EncodeToString(raw) + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "text/event-stream") + _, _ = w.Write([]byte("data: {\"output_index\":0,\"content_index\":0,\"delta\":\"Here you go:\",\"type\":\"response.output_text.delta\"}\n\n")) + _, _ = fmt.Fprintf(w, "data: {\"response\":{\"id\":\"resp-img\",\"status\":\"completed\",\"output\":[{\"type\":\"image_generation_call\",\"id\":\"ig_1\",\"status\":\"completed\",\"result\":%q}],\"usage\":{\"input_tokens\":5,\"output_tokens\":9,\"total_tokens\":14}},\"type\":\"response.completed\"}\n\n", b64) + })) + defer srv.Close() + + client := NewResponsesClient(llm.ClientConfig{APIKey: "sk-test", Model: "gpt-5", BaseURL: srv.URL}) + resp, err := client.Chat(context.Background(), &llm.ChatRequest{ + Messages: []llm.ChatMessage{{Role: llm.RoleUser, Content: "draw a cat"}}, + }) + if err != nil { + t.Fatalf("Chat: %v", err) + } + if resp.Message.Content != "Here you go:" { + t.Errorf("content = %q", resp.Message.Content) + } + if len(resp.Message.Parts) != 1 { + t.Fatalf("Parts = %+v, want one image part", resp.Message.Parts) + } + part := resp.Message.Parts[0] + if part.Type != llm.ContentPartImage || part.Media == nil { + t.Fatalf("part is not an image: %+v", part) + } + if part.Media.MimeType != "image/png" { + t.Errorf("mime = %q, want image/png", part.Media.MimeType) + } + if string(part.Media.Bytes) != string(raw) { + t.Errorf("image bytes not decoded from base64 result") + } +} + +// TestResponses_ImageGenerationToolOptIn: the image_generation built-in tool is +// sent only when EnableImageGeneration is set, and carries no function name. +func TestResponses_ImageGenerationToolOptIn(t *testing.T) { + req := &llm.ChatRequest{Messages: []llm.ChatMessage{{Role: llm.RoleUser, Content: "hi"}}} + + off := NewResponsesClient(llm.ClientConfig{Model: "gpt-5"}).buildRequest(req, false) + for _, tl := range off.Tools { + if tl.Type == "image_generation" { + t.Fatal("image_generation tool must NOT be sent when the opt-in is off") + } + } + + on := NewResponsesClient(llm.ClientConfig{Model: "gpt-5", EnableImageGeneration: true}).buildRequest(req, false) + found := false + for _, tl := range on.Tools { + if tl.Type == "image_generation" { + found = true + if tl.Name != "" { + t.Errorf("built-in image_generation tool must not carry a name; got %q", tl.Name) + } + } + } + if !found { + t.Error("image_generation tool must be sent when EnableImageGeneration is set") + } +} diff --git a/forge-core/llm/types.go b/forge-core/llm/types.go index 6f27ef15..1d38932c 100644 --- a/forge-core/llm/types.go +++ b/forge-core/llm/types.go @@ -138,11 +138,12 @@ type ChatResponse struct { // StreamDelta represents a single chunk in a streaming response. type StreamDelta struct { - Content string `json:"content,omitempty"` - ToolCalls []ToolCall `json:"tool_calls,omitempty"` - FinishReason string `json:"finish_reason,omitempty"` - Done bool `json:"done,omitempty"` - Usage *UsageInfo `json:"usage,omitempty"` + Content string `json:"content,omitempty"` + ToolCalls []ToolCall `json:"tool_calls,omitempty"` + Parts []ContentPart `json:"parts,omitempty"` // non-text output, e.g. a model-generated image (#255) + FinishReason string `json:"finish_reason,omitempty"` + Done bool `json:"done,omitempty"` + Usage *UsageInfo `json:"usage,omitempty"` } // UsageInfo contains token usage information. diff --git a/forge-core/runtime/config.go b/forge-core/runtime/config.go index 98c599cf..0ead7d41 100644 --- a/forge-core/runtime/config.go +++ b/forge-core/runtime/config.go @@ -122,6 +122,11 @@ func ResolveModelConfig(cfg *types.ForgeConfig, envVars map[string]string, provi if cfg.Model.DisableStore { mc.Client.DisableStore = true } + // #255 — opt into the OpenAI Responses image_generation built-in tool. + // Carried unconditionally; only the openai-responses client honors it. + if cfg.Model.ImageGeneration { + mc.Client.EnableImageGeneration = true + } // AWS_REGION env safety-net for the SigV4 path. Mirrors the // OPENAI_BASE_URL / ANTHROPIC_BASE_URL env pattern above — lets // an operator override the region per-deploy without touching diff --git a/forge-core/runtime/loop.go b/forge-core/runtime/loop.go index 4e843f2e..79176d0a 100644 --- a/forge-core/runtime/loop.go +++ b/forge-core/runtime/loop.go @@ -628,9 +628,14 @@ func (e *LLMExecutor) Execute(ctx context.Context, task *a2a.Task, msg *a2a.Mess // names the failure mode so an operator scanning a recovered // session understands what happened. assistantMsg := resp.Message - if assistantMsg.Content == "" && len(assistantMsg.ToolCalls) == 0 { + if assistantMsg.Content == "" && len(assistantMsg.ToolCalls) == 0 && !assistantMsg.HasMedia() { assistantMsg.Content = emptyAssistantPlaceholder } + // Persist any model-generated media (e.g. an image_generation result) + // to disk and record its URI, so history stores the reference not the + // base64 (#255 Phase 5). Bytes stay inline so finalizeResponse still + // surfaces the image as a file part in the A2A response this turn. + _ = persistInboundMedia(ctx, &assistantMsg) mem.Append(assistantMsg) // Check if we're done: the definitive signal is the absence of tool @@ -1341,6 +1346,20 @@ func llmMessageToA2A(msg llm.ChatMessage, extraParts ...a2a.Part) *a2a.Message { } parts := []a2a.Part{a2a.NewTextPart(msg.Content)} + // Surface model-generated media (e.g. an image_generation result) as file + // parts in the response (#255 Phase 5). Text parts are already covered by + // msg.Content above. + for _, p := range msg.Parts { + if p.Media == nil || len(p.Media.Bytes) == 0 { + continue + } + if p.Type == llm.ContentPartImage || p.Type == llm.ContentPartDocument { + parts = append(parts, a2a.NewFilePart(a2a.FileContent{ + MimeType: p.Media.MimeType, + Bytes: p.Media.Bytes, + })) + } + } parts = append(parts, extraParts...) return &a2a.Message{ diff --git a/forge-core/runtime/loop_projection_test.go b/forge-core/runtime/loop_projection_test.go index 314d7561..8e284f28 100644 --- a/forge-core/runtime/loop_projection_test.go +++ b/forge-core/runtime/loop_projection_test.go @@ -81,6 +81,36 @@ func TestA2AMessageToLLM_ProjectsPDFDocument(t *testing.T) { } } +// TestLLMMessageToA2A_EmitsGeneratedImage verifies a model-generated image on +// the response message becomes an A2A file part alongside the text (#255 Phase 5). +func TestLLMMessageToA2A_EmitsGeneratedImage(t *testing.T) { + img := []byte("generated-png-bytes") + msg := llm.ChatMessage{ + Role: llm.RoleAssistant, + Content: "here is your image", + Parts: []llm.ContentPart{ + llm.NewMediaContentPart(llm.ContentPartImage, llm.MediaRef{MimeType: "image/png", Bytes: img}), + }, + } + out := llmMessageToA2A(msg) + if out.Role != a2a.MessageRoleAgent { + t.Errorf("role = %q", out.Role) + } + if len(out.Parts) != 2 { + t.Fatalf("parts = %+v, want [text, file]", out.Parts) + } + if out.Parts[0].Kind != a2a.PartKindText || out.Parts[0].Text != "here is your image" { + t.Errorf("first part should be the text; got %+v", out.Parts[0]) + } + f := out.Parts[1] + if f.Kind != a2a.PartKindFile || f.File == nil || f.File.MimeType != "image/png" { + t.Fatalf("second part should be an image file part; got %+v", f) + } + if string(f.File.Bytes) != string(img) { + t.Errorf("generated image bytes not carried into the file part") + } +} + // TestA2AMessageToLLM_UnsupportedFileNotProjected: an unsupported file part (the // gate would reject it) must not be projected into Parts if it reaches here. func TestA2AMessageToLLM_UnsupportedFileNotProjected(t *testing.T) { diff --git a/forge-core/types/config.go b/forge-core/types/config.go index f528b8ae..c7b1f68a 100644 --- a/forge-core/types/config.go +++ b/forge-core/types/config.go @@ -1136,6 +1136,12 @@ type ModelRef struct { // it; every other provider ignores it. Issue #383. DisableStore bool `yaml:"disable_store,omitempty"` + // ImageGeneration opts an OpenAI Responses model (provider: + // openai-responses) into the `image_generation` built-in tool, letting the + // model emit images that forge returns as `file` parts in the response + // (#255). Off by default; only the openai-responses provider honors it. + ImageGeneration bool `yaml:"image_generation,omitempty"` + Version string `yaml:"version,omitempty"` OrganizationID string `yaml:"organization_id,omitempty"` Fallbacks []ModelFallback `yaml:"fallbacks,omitempty"` From 9cacc9ccadb0d43214de86fa5a861fa925dee032 Mon Sep 17 00:00:00 2001 From: MK Date: Wed, 30 Sep 2026 00:05:36 -0400 Subject: [PATCH 2/2] refactor(runtime): persistMedia with inbound/generated subdirs (#536 review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Phase-5 change routed model-generated (outbound) media through persistInboundMedia into .forge/files/inbound — a naming/semantics mismatch. Rename to persistMedia(ctx, msg, subdir): received uploads persist under inbound/, model-generated output under generated/. Distinct subdirs let the tracked retention/GC follow-up (#255 item-3b) tell received from generated media apart. Purely organizational — same content-addressed store, same best-effort behavior. Adds a generated-subdir test; docs + skill updated. --- .claude/skills/forge.md | 2 +- docs/core-concepts/runtime-engine.md | 6 +-- forge-cli/internal/surface/knowledge/forge.md | 2 +- forge-core/runtime/loop.go | 6 +-- forge-core/runtime/media_persist.go | 41 +++++++++++-------- forge-core/runtime/media_persist_test.go | 26 +++++++++--- 6 files changed, 54 insertions(+), 29 deletions(-) diff --git a/.claude/skills/forge.md b/.claude/skills/forge.md index 2ef29b1c..7e0652a7 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` \| `model_not_document_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `too_many_document_parts` \| `image_limit_exceeded` \| `document_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) and **PDFs** on a **doc model** (`ModelSupportsPDF` = Anthropic Sonnet 3.5+/Opus 4+/Haiku 4.5+/Fable 5; Claude 3.0 & 3.5-Haiku excluded, fail-closed), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image`/`document` source blocks, OpenAI `image_url` data URLs. Non-PDF docs/video still rejected; OpenAI-Responses PDF + extraction fallback are follow-ups. **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP/≤100k-per-side (`CheckImageLimits`); per-PDF ≤32 MiB + `%PDF-` sniff (`CheckDocumentLimits`); ≤20 images & ≤5 docs/msg; ≤4 concurrent media requests (excess shed with 429/unavailable). **Persistence/replay** (`persistInboundMedia` + `RehydrateMedia`): inbound media written to `.forge/files/inbound/.`, path stored as `MediaRef.URI`; history keeps URI-only (`Bytes` json:"-"), rehydrated per turn so multi-turn convos replay media without base64 bloat. **Model-generated image OUTPUT** (#255 Phase 5): opt-in `image_generation: true` on an `openai-responses` model → sends the `image_generation` built-in tool; the `image_generation_call` base64 `result` (from `response.completed`) → `resp.Message.Parts` image → `llmMessageToA2A` emits a response `file` part (persisted like inbound). Responses-only (Anthropic/Gemini/Bedrock don't generate images via chat) | +| `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` \| `model_not_document_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `too_many_document_parts` \| `image_limit_exceeded` \| `document_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) and **PDFs** on a **doc model** (`ModelSupportsPDF` = Anthropic Sonnet 3.5+/Opus 4+/Haiku 4.5+/Fable 5; Claude 3.0 & 3.5-Haiku excluded, fail-closed), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image`/`document` source blocks, OpenAI `image_url` data URLs. Non-PDF docs/video still rejected; OpenAI-Responses PDF + extraction fallback are follow-ups. **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP/≤100k-per-side (`CheckImageLimits`); per-PDF ≤32 MiB + `%PDF-` sniff (`CheckDocumentLimits`); ≤20 images & ≤5 docs/msg; ≤4 concurrent media requests (excess shed with 429/unavailable). **Persistence/replay** (`persistMedia` + `RehydrateMedia`): media written to `.forge/files/{inbound,generated}/.`, path stored as `MediaRef.URI`; history keeps URI-only (`Bytes` json:"-"), rehydrated per turn so multi-turn convos replay media without base64 bloat. **Model-generated image OUTPUT** (#255 Phase 5): opt-in `image_generation: true` on an `openai-responses` model → sends the `image_generation` built-in tool; the `image_generation_call` base64 `result` (from `response.completed`) → `resp.Message.Parts` image → `llmMessageToA2A` emits a response `file` part (persisted like inbound). Responses-only (Anthropic/Gemini/Bedrock don't generate images via chat) | | `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 d36cb3ef..38f39ccc 100644 --- a/docs/core-concepts/runtime-engine.md +++ b/docs/core-concepts/runtime-engine.md @@ -42,9 +42,9 @@ A media `file` part is forwarded to the model as native input when the resolved `a2aMessageToLLM` projects supported parts into `llm.ChatMessage.Parts` (the flattened text stays in `Content` as the text-of-record for the scanners). A text-only message keeps `Parts` empty and marshals byte-identically to before. -**Model-generated image output.** A model can also *produce* an image, which forge returns as a `file` part in the A2A response. For OpenAI **Responses** (`provider: openai-responses`), set `image_generation: true` on the model (`ModelRef.ImageGeneration` → `ClientConfig.EnableImageGeneration`); the client then sends the `image_generation` built-in tool, and the base64 `result` from the `image_generation_call` output (read from the authoritative `response.completed` frame) is decoded into an image `ContentPart` on `resp.Message.Parts`. `llmMessageToA2A` projects those into response `file` parts, and the generated image is persisted like inbound media (URI in history, no base64 bloat). Off by default (it changes provider behavior and cost). Anthropic/Gemini/Bedrock don't generate images through their chat/messages APIs, so this is Responses-only today. +**Model-generated image output.** A model can also *produce* an image, which forge returns as a `file` part in the A2A response. For OpenAI **Responses** (`provider: openai-responses`), set `image_generation: true` on the model (`ModelRef.ImageGeneration` → `ClientConfig.EnableImageGeneration`); the client then sends the `image_generation` built-in tool, and the base64 `result` from the `image_generation_call` output (read from the authoritative `response.completed` frame) is decoded into an image `ContentPart` on `resp.Message.Parts`. `llmMessageToA2A` projects those into response `file` parts, and the generated image is persisted to `.forge/files/generated/` (URI in history, no base64 bloat). Off by default (it changes provider behavior and cost). Anthropic/Gemini/Bedrock don't generate images through their chat/messages APIs, so this is Responses-only today. -**Persistence & cross-turn replay.** Inbound media is written to `.forge/files/inbound/.` (`persistInboundMedia`) and the on-disk path is recorded as the part's `MediaRef.URI` — content-addressed, so identical uploads dedup and re-writes are idempotent, and the files are available on disk for tools. The bytes stay inline for the current turn's request; because `MediaRef.Bytes` is `json:"-"`, session history persists only the URI (never base64). On a later turn, `RehydrateMedia` reloads the bytes from the URI before the request is built, so a multi-turn conversation keeps seeing earlier images/PDFs without bloating the session file. A file that can't be reloaded (deleted, or an absolute path from another host — remote/distributed session replay is a follow-up) is left byteless and skipped by the provider serializers, degrading to text rather than failing the turn. Persistence and rehydration are best-effort: a failure never fails the turn (media is still fed inline the turn it arrives). +**Persistence & cross-turn replay.** Media is written to `.forge/files//.` (`persistMedia`) — received uploads under `inbound/`, model-generated output under `generated/`. The on-disk path is recorded as the part's `MediaRef.URI` — content-addressed, so identical uploads dedup and re-writes are idempotent, and the files are available on disk for tools. The bytes stay inline for the current turn's request; because `MediaRef.Bytes` is `json:"-"`, session history persists only the URI (never base64). On a later turn, `RehydrateMedia` reloads the bytes from the URI before the request is built, so a multi-turn conversation keeps seeing earlier images/PDFs without bloating the session file. A file that can't be reloaded (deleted, or an absolute path from another host — remote/distributed session replay is a follow-up) is left byteless and skipped by the provider serializers, degrading to text rather than failing the turn. Persistence and rehydration are best-effort: a failure never fails the turn (media is still fed inline the turn it arrives). Media the model can't consume is **rejected loudly, never silently dropped** (the `checkInboundMedia` ingest gate): an image on a text-only model, a PDF on a non-document model, or an unsupported type (other documents/video) returns a 4xx and emits the `input_media_rejected` audit event. Note media **bytes** are not text-scannable, so guardrail/intent scanning still applies only to the text/data projection; this is an accepted limitation. @@ -348,7 +348,7 @@ docker run -e KUBECONFIG="$(cat ~/.kube/config)" my-agent ## File Output Directory -The runtime configures a `FilesDir` for tool-generated files (e.g., from `file_create`). This directory defaults to `/.forge/files/` and is injected into the execution context so tools can write files that other tools can reference by path. Inbound media (uploaded images/PDFs) is persisted under `/inbound/` — see [Image and document input](#image-and-document-input-multimodal). +The runtime configures a `FilesDir` for tool-generated files (e.g., from `file_create`). This directory defaults to `/.forge/files/` and is injected into the execution context so tools can write files that other tools can reference by path. Inbound uploads are persisted under `/inbound/` and model-generated media under `/generated/` — see [Image and document input](#image-and-document-input-multimodal). ``` / diff --git a/forge-cli/internal/surface/knowledge/forge.md b/forge-cli/internal/surface/knowledge/forge.md index 2ef29b1c..7e0652a7 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` \| `model_not_document_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `too_many_document_parts` \| `image_limit_exceeded` \| `document_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) and **PDFs** on a **doc model** (`ModelSupportsPDF` = Anthropic Sonnet 3.5+/Opus 4+/Haiku 4.5+/Fable 5; Claude 3.0 & 3.5-Haiku excluded, fail-closed), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image`/`document` source blocks, OpenAI `image_url` data URLs. Non-PDF docs/video still rejected; OpenAI-Responses PDF + extraction fallback are follow-ups. **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP/≤100k-per-side (`CheckImageLimits`); per-PDF ≤32 MiB + `%PDF-` sniff (`CheckDocumentLimits`); ≤20 images & ≤5 docs/msg; ≤4 concurrent media requests (excess shed with 429/unavailable). **Persistence/replay** (`persistInboundMedia` + `RehydrateMedia`): inbound media written to `.forge/files/inbound/.`, path stored as `MediaRef.URI`; history keeps URI-only (`Bytes` json:"-"), rehydrated per turn so multi-turn convos replay media without base64 bloat. **Model-generated image OUTPUT** (#255 Phase 5): opt-in `image_generation: true` on an `openai-responses` model → sends the `image_generation` built-in tool; the `image_generation_call` base64 `result` (from `response.completed`) → `resp.Message.Parts` image → `llmMessageToA2A` emits a response `file` part (persisted like inbound). Responses-only (Anthropic/Gemini/Bedrock don't generate images via chat) | +| `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` \| `model_not_document_capable` \| `unsupported_media_type` \| `too_many_image_parts` \| `too_many_document_parts` \| `image_limit_exceeded` \| `document_limit_exceeded`). Gate: `Runner.checkInboundMedia`. **Images** (png/jpeg/gif/webp) on a **vision model** (`coreruntime.ModelSupportsVision`) and **PDFs** on a **doc model** (`ModelSupportsPDF` = Anthropic Sonnet 3.5+/Opus 4+/Haiku 4.5+/Fable 5; Claude 3.0 & 3.5-Haiku excluded, fail-closed), within the DoS bounds, are NOT rejected — `a2aMessageToLLM` projects them into `llm.ChatMessage.Parts` → Anthropic `image`/`document` source blocks, OpenAI `image_url` data URLs. Non-PDF docs/video still rejected; OpenAI-Responses PDF + extraction fallback are follow-ups. **DoS controls** (`media_limits.go` + `mediaSem`): body cap 32 MiB both transports; per-image ≤5 MiB & ≤50 MP/≤100k-per-side (`CheckImageLimits`); per-PDF ≤32 MiB + `%PDF-` sniff (`CheckDocumentLimits`); ≤20 images & ≤5 docs/msg; ≤4 concurrent media requests (excess shed with 429/unavailable). **Persistence/replay** (`persistMedia` + `RehydrateMedia`): media written to `.forge/files/{inbound,generated}/.`, path stored as `MediaRef.URI`; history keeps URI-only (`Bytes` json:"-"), rehydrated per turn so multi-turn convos replay media without base64 bloat. **Model-generated image OUTPUT** (#255 Phase 5): opt-in `image_generation: true` on an `openai-responses` model → sends the `image_generation` built-in tool; the `image_generation_call` base64 `result` (from `response.completed`) → `resp.Message.Parts` image → `llmMessageToA2A` emits a response `file` part (persisted like inbound). Responses-only (Anthropic/Gemini/Bedrock don't generate images via chat) | | `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-core/runtime/loop.go b/forge-core/runtime/loop.go index 79176d0a..ff28b172 100644 --- a/forge-core/runtime/loop.go +++ b/forge-core/runtime/loop.go @@ -344,7 +344,7 @@ func (e *LLMExecutor) Execute(ctx context.Context, task *a2a.Task, msg *a2a.Mess hm := a2aMessageToLLM(histMsg) // Best-effort persist so any media in replayed history keeps a URI // reference (#255 Phase 4); on failure it stays inline for this turn. - _ = persistInboundMedia(ctx, &hm) + _ = persistMedia(ctx, &hm, mediaSubdirInbound) mem.Append(hm) } } @@ -366,7 +366,7 @@ func (e *LLMExecutor) Execute(ctx context.Context, task *a2a.Task, msg *a2a.Mess // part, so it survives into session history (Bytes are json:"-") and can be // rehydrated on the next turn (#255 Phase 4). Bytes stay inline for THIS // turn's request. Best-effort: on failure media is still fed this turn. - _ = persistInboundMedia(ctx, &newMsg) + _ = persistMedia(ctx, &newMsg, mediaSubdirInbound) if recovered { msgs := mem.Messages() n := len(msgs) @@ -635,7 +635,7 @@ func (e *LLMExecutor) Execute(ctx context.Context, task *a2a.Task, msg *a2a.Mess // to disk and record its URI, so history stores the reference not the // base64 (#255 Phase 5). Bytes stay inline so finalizeResponse still // surfaces the image as a file part in the A2A response this turn. - _ = persistInboundMedia(ctx, &assistantMsg) + _ = persistMedia(ctx, &assistantMsg, mediaSubdirGenerated) mem.Append(assistantMsg) // Check if we're done: the definitive signal is the absence of tool diff --git a/forge-core/runtime/media_persist.go b/forge-core/runtime/media_persist.go index 6624adfd..23cafdfb 100644 --- a/forge-core/runtime/media_persist.go +++ b/forge-core/runtime/media_persist.go @@ -10,25 +10,34 @@ import ( "github.com/initializ/forge/forge-core/llm" ) -// persistInboundMedia writes each inline media part in msg to the context's -// files dir (WithFilesDir → .forge/files/inbound) and records the on-disk path -// as the part's MediaRef.URI, while LEAVING Bytes in place for the current -// turn's request (#255 Phase 4). This is what makes media survive across turns: -// session history persists llm.ChatMessage with MediaRef.Bytes tagged json:"-", -// so only the URI is stored, and RehydrateMedia reloads the bytes on replay. -// The persisted files also give agent tools a real path to open an upload. +// Media store subdirectories under the files dir (#255). Received uploads and +// model-generated output share the same content-addressed store but land in +// distinct subdirs so a future retention/GC pass can tell them apart. +const ( + mediaSubdirInbound = "inbound" // client-uploaded media + mediaSubdirGenerated = "generated" // model-generated media (e.g. image_generation) +) + +// persistMedia writes each inline media part in msg to the context's files dir +// (WithFilesDir → .forge/files/) and records the on-disk path as the +// part's MediaRef.URI, while LEAVING Bytes in place for the current turn's +// request (#255). This is what makes media survive across turns: session +// history persists llm.ChatMessage with MediaRef.Bytes tagged json:"-", so only +// the URI is stored, and RehydrateMedia reloads the bytes on replay. The +// persisted files also give agent tools a real path to open a file. // -// Files are content-addressed (sha256 + a MIME-derived extension), so identical -// uploads dedup across turns/messages and re-writing is idempotent. Best-effort: -// callers ignore the error — on failure the media simply stays inline for this -// turn and won't replay on the next one. Parts with no bytes, or that already -// carry a URI, are skipped. -func persistInboundMedia(ctx context.Context, msg *llm.ChatMessage) error { +// subdir separates received uploads (mediaSubdirInbound) from model-generated +// output (mediaSubdirGenerated). Files are content-addressed (sha256 + a +// MIME-derived extension), so identical media dedups and re-writing is +// idempotent. Best-effort: callers ignore the error — on failure the media +// simply stays inline for this turn and won't replay on the next one. Parts +// with no bytes, or that already carry a URI, are skipped. +func persistMedia(ctx context.Context, msg *llm.ChatMessage, subdir string) error { dir := FilesDirFromContext(ctx) if dir == "" { return nil // no files dir configured → media stays inline (this turn only) } - inboundDir := filepath.Join(dir, "inbound") + mediaDir := filepath.Join(dir, subdir) made := false for i := range msg.Parts { m := msg.Parts[i].Media @@ -36,13 +45,13 @@ func persistInboundMedia(ctx context.Context, msg *llm.ChatMessage) error { continue } if !made { - if err := os.MkdirAll(inboundDir, 0o700); err != nil { + if err := os.MkdirAll(mediaDir, 0o700); err != nil { return err } made = true } sum := sha256.Sum256(m.Bytes) - path := filepath.Join(inboundDir, hex.EncodeToString(sum[:])+extForMIME(m.MimeType)) + path := filepath.Join(mediaDir, hex.EncodeToString(sum[:])+extForMIME(m.MimeType)) if err := os.WriteFile(path, m.Bytes, 0o600); err != nil { return err } diff --git a/forge-core/runtime/media_persist_test.go b/forge-core/runtime/media_persist_test.go index 2772bbb3..05bc2047 100644 --- a/forge-core/runtime/media_persist_test.go +++ b/forge-core/runtime/media_persist_test.go @@ -34,7 +34,7 @@ func TestPersistAndRehydrate_CrossTurnRoundTrip(t *testing.T) { orig := []byte("the-real-image-bytes-xyz") msg := imageMsg("image/png", orig) - if err := persistInboundMedia(ctx, &msg); err != nil { + if err := persistMedia(ctx, &msg, mediaSubdirInbound); err != nil { t.Fatalf("persist: %v", err) } @@ -95,10 +95,10 @@ func TestPersistInboundMedia_ContentAddressedAndIdempotent(t *testing.T) { a := imageMsg("image/png", data) b := imageMsg("image/png", data) - if err := persistInboundMedia(ctx, &a); err != nil { + if err := persistMedia(ctx, &a, mediaSubdirInbound); err != nil { t.Fatal(err) } - if err := persistInboundMedia(ctx, &b); err != nil { + if err := persistMedia(ctx, &b, mediaSubdirInbound); err != nil { t.Fatal(err) } if a.Parts[1].Media.URI != b.Parts[1].Media.URI { @@ -108,7 +108,7 @@ func TestPersistInboundMedia_ContentAddressedAndIdempotent(t *testing.T) { // A part that already has a URI is left untouched (idempotent). preset := imageMsg("image/png", data) preset.Parts[1].Media.URI = "/already/set.png" - if err := persistInboundMedia(ctx, &preset); err != nil { + if err := persistMedia(ctx, &preset, mediaSubdirInbound); err != nil { t.Fatal(err) } if preset.Parts[1].Media.URI != "/already/set.png" { @@ -116,11 +116,27 @@ func TestPersistInboundMedia_ContentAddressedAndIdempotent(t *testing.T) { } } +// TestPersistMedia_GeneratedSubdir: model-generated media lands under the +// generated/ subdir (distinct from inbound/) so a retention pass can tell +// received uploads from generated output (#536 review). +func TestPersistMedia_GeneratedSubdir(t *testing.T) { + dir := t.TempDir() + ctx := WithFilesDir(context.Background(), dir) + msg := imageMsg("image/png", []byte("generated")) + if err := persistMedia(ctx, &msg, mediaSubdirGenerated); err != nil { + t.Fatal(err) + } + uri := msg.Parts[1].Media.URI + if !strings.HasPrefix(uri, filepath.Join(dir, "generated")) { + t.Errorf("generated media should live under generated/, got %q", uri) + } +} + // TestPersistInboundMedia_NoFilesDirIsNoop: without a files dir, persistence is // skipped (media stays inline for the turn) and no error is raised. func TestPersistInboundMedia_NoFilesDirIsNoop(t *testing.T) { msg := imageMsg("image/png", []byte("x")) - if err := persistInboundMedia(context.Background(), &msg); err != nil { + if err := persistMedia(context.Background(), &msg, mediaSubdirInbound); err != nil { t.Fatalf("no files dir should be a no-op, got %v", err) } if msg.Parts[1].Media.URI != "" {