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
2 changes: 1 addition & 1 deletion .claude/skills/forge.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:<mime>"]`), `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:<mime>"]`), `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/<sha256>.<ext>`, 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 |
Expand Down
4 changes: 3 additions & 1 deletion docs/core-concepts/runtime-engine.md
Original file line number Diff line number Diff line change
Expand Up @@ -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/<sha256>.<ext>` (`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:
Expand Down Expand Up @@ -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 `<WorkDir>/.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 `<WorkDir>/.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 `<FilesDir>/inbound/` — see [Image and document input](#image-and-document-input-multimodal).

```
<WorkDir>/
Expand Down
2 changes: 1 addition & 1 deletion forge-cli/internal/surface/knowledge/forge.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:<mime>"]`), `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:<mime>"]`), `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/<sha256>.<ext>`, 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 |
Expand Down
18 changes: 17 additions & 1 deletion forge-core/runtime/loop.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}

Expand All @@ -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)
Expand Down Expand Up @@ -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{
Expand Down
71 changes: 71 additions & 0 deletions forge-core/runtime/media_persist.go
Original file line number Diff line number Diff line change
@@ -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 {

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.

Follow-up — LOW/MEDIUM (not blocking): no cleanup/GC for inbound/. Content-addressed writes accumulate indefinitely — there's no TTL, size cap, or refcount, and content-addressing makes session-scoped cleanup non-trivial (files are shared across sessions/turns, so deleting on one session's end can break another). Per-request DoS is bounded (body cap + per-part limits, #532), but a long-running agent receiving media grows inbound/ without bound → eventual disk exhaustion. This is the realization of #255 checklist item 3 (storage). The PR's scope-boundaries flag remote replay + tool-path surfacing but not retention — recommend tracking a retention policy (size cap / TTL / GC) as a follow-up.

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"
}
}
129 changes: 129 additions & 0 deletions forge-core/runtime/media_persist_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
Loading