diff --git a/.claude/skills/forge.md b/.claude/skills/forge.md index 44dce16b..df451ff2 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) | +| `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 | | `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 c2d48ea6..4443b016 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. +**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. **DoS bounds.** Because inline media raises the inbound-body cap to 32 MiB (both transports), the gate also enforces per-part 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: @@ -344,7 +346,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. +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). ``` / diff --git a/forge-cli/internal/surface/knowledge/forge.md b/forge-cli/internal/surface/knowledge/forge.md index 44dce16b..df451ff2 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) | +| `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 | | `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 59235bbd..4e843f2e 100644 --- a/forge-core/runtime/loop.go +++ b/forge-core/runtime/loop.go @@ -341,7 +341,11 @@ func (e *LLMExecutor) Execute(ctx context.Context, task *a2a.Task, msg *a2a.Mess historyToLoad = historyToLoad[:n-1] } for _, histMsg := range historyToLoad { - mem.Append(a2aMessageToLLM(histMsg)) + 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) + mem.Append(hm) } } @@ -358,6 +362,11 @@ func (e *LLMExecutor) Execute(ctx context.Context, task *a2a.Task, msg *a2a.Mess // replayed the poisoned transcript and repeated the last answer without // re-attempting tools (#378). newMsg := a2aMessageToLLM(*msg) + // Persist inbound media to .forge/files/inbound and record the URI on each + // 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) if recovered { msgs := mem.Messages() n := len(msgs) @@ -447,6 +456,13 @@ func (e *LLMExecutor) Execute(ctx context.Context, task *a2a.Task, msg *a2a.Mess } messages := mem.Messages() + // Reload media bytes from their URI for any part replayed from a + // persisted session (history stores URI-only; Bytes are json:"-"). + // Freshly-ingested parts already hold Bytes and are skipped. Best-effort: + // a media file that can't be reloaded (e.g. deleted) is left byteless and + // the provider serializers skip it, degrading to text rather than failing + // the turn (#255 Phase 4). + _ = RehydrateMedia(ctx, messages) // Fire BeforeLLMCall hook if err := e.hooks.Fire(ctx, BeforeLLMCall, &HookContext{ diff --git a/forge-core/runtime/media_persist.go b/forge-core/runtime/media_persist.go new file mode 100644 index 00000000..6624adfd --- /dev/null +++ b/forge-core/runtime/media_persist.go @@ -0,0 +1,71 @@ +package runtime + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "os" + "path/filepath" + + "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. +// +// 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 { + dir := FilesDirFromContext(ctx) + if dir == "" { + return nil // no files dir configured → media stays inline (this turn only) + } + inboundDir := filepath.Join(dir, "inbound") + made := false + for i := range msg.Parts { + m := msg.Parts[i].Media + if m == nil || len(m.Bytes) == 0 || m.URI != "" { + continue + } + if !made { + if err := os.MkdirAll(inboundDir, 0o700); err != nil { + return err + } + made = true + } + sum := sha256.Sum256(m.Bytes) + path := filepath.Join(inboundDir, hex.EncodeToString(sum[:])+extForMIME(m.MimeType)) + if err := os.WriteFile(path, m.Bytes, 0o600); err != nil { + return err + } + m.URI = path + } + return nil +} + +// extForMIME returns a file extension for a media MIME type (for readable +// on-disk names). Falls back to .bin for anything unexpected. +func extForMIME(mime string) string { + switch NormalizeImageMIME(mime) { + case "image/png": + return ".png" + case "image/jpeg": + return ".jpg" + case "image/gif": + return ".gif" + case "image/webp": + return ".webp" + case "application/pdf": + return ".pdf" + default: + return ".bin" + } +} diff --git a/forge-core/runtime/media_persist_test.go b/forge-core/runtime/media_persist_test.go new file mode 100644 index 00000000..2772bbb3 --- /dev/null +++ b/forge-core/runtime/media_persist_test.go @@ -0,0 +1,129 @@ +package runtime + +import ( + "context" + "encoding/base64" + "encoding/json" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/initializ/forge/forge-core/llm" +) + +func imageMsg(mime string, data []byte) llm.ChatMessage { + return llm.ChatMessage{ + Role: llm.RoleUser, + Content: "hi", + Parts: []llm.ContentPart{ + llm.NewTextContentPart("hi"), + llm.NewMediaContentPart(llm.ContentPartImage, llm.MediaRef{MimeType: mime, Bytes: data}), + }, + } +} + +// TestPersistAndRehydrate_CrossTurnRoundTrip is the core Phase 4 guarantee: +// inbound media is persisted to .forge/files with a URI; the session marshal +// drops the inline Bytes (json:"-") keeping only the URI, and RehydrateMedia +// reloads the Bytes on the next turn — so multi-turn conversations replay the +// media without ever storing base64 in history (#255). +func TestPersistAndRehydrate_CrossTurnRoundTrip(t *testing.T) { + dir := t.TempDir() + ctx := WithFilesDir(context.Background(), dir) + orig := []byte("the-real-image-bytes-xyz") + + msg := imageMsg("image/png", orig) + if err := persistInboundMedia(ctx, &msg); err != nil { + t.Fatalf("persist: %v", err) + } + + // URI points under .forge/files/inbound and the file holds the bytes. + uri := msg.Parts[1].Media.URI + if uri == "" { + t.Fatal("URI not set after persist") + } + if !strings.HasPrefix(uri, filepath.Join(dir, "inbound")) { + t.Errorf("URI %q not under the inbound dir", uri) + } + if got, _ := os.ReadFile(uri); string(got) != string(orig) { + t.Errorf("persisted file content mismatch") + } + if !strings.HasSuffix(uri, ".png") { + t.Errorf("URI should carry a mime-derived extension; got %q", uri) + } + // Bytes stay inline for THIS turn's request. + if string(msg.Parts[1].Media.Bytes) != string(orig) { + t.Error("Bytes must stay inline for the current turn") + } + + // Simulate session persistence: marshal drops Bytes, keeps URI. + data, err := json.Marshal([]llm.ChatMessage{msg}) + if err != nil { + t.Fatal(err) + } + if strings.Contains(string(data), base64.StdEncoding.EncodeToString(orig)) { + t.Errorf("inline bytes leaked into persisted history:\n%s", data) + } + var replayed []llm.ChatMessage + if err := json.Unmarshal(data, &replayed); err != nil { + t.Fatal(err) + } + if len(replayed[0].Parts[1].Media.Bytes) != 0 { + t.Fatal("Bytes must be dropped from persisted history") + } + if replayed[0].Parts[1].Media.URI != uri { + t.Fatalf("URI must survive persistence: got %q", replayed[0].Parts[1].Media.URI) + } + + // Next turn: rehydrate reloads the bytes from the URI. + if err := RehydrateMedia(ctx, replayed); err != nil { + t.Fatalf("rehydrate: %v", err) + } + if string(replayed[0].Parts[1].Media.Bytes) != string(orig) { + t.Errorf("rehydrated bytes = %q, want %q", replayed[0].Parts[1].Media.Bytes, orig) + } +} + +// TestPersistInboundMedia_ContentAddressedAndIdempotent: identical bytes across +// two messages map to the same path (dedup), and re-persisting is a no-op once +// a URI is set. +func TestPersistInboundMedia_ContentAddressedAndIdempotent(t *testing.T) { + dir := t.TempDir() + ctx := WithFilesDir(context.Background(), dir) + data := []byte("same-bytes") + + a := imageMsg("image/png", data) + b := imageMsg("image/png", data) + if err := persistInboundMedia(ctx, &a); err != nil { + t.Fatal(err) + } + if err := persistInboundMedia(ctx, &b); err != nil { + t.Fatal(err) + } + if a.Parts[1].Media.URI != b.Parts[1].Media.URI { + t.Errorf("identical bytes should map to the same content-addressed path: %q vs %q", a.Parts[1].Media.URI, b.Parts[1].Media.URI) + } + + // 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 { + t.Fatal(err) + } + if preset.Parts[1].Media.URI != "/already/set.png" { + t.Error("a part with an existing URI must not be re-persisted") + } +} + +// 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 { + t.Fatalf("no files dir should be a no-op, got %v", err) + } + if msg.Parts[1].Media.URI != "" { + t.Error("no files dir → URI must stay empty") + } +}