-
Notifications
You must be signed in to change notification settings - Fork 13
feat(runtime): persist inbound media for cross-turn replay (#255 Phase 4) #535
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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 { | ||
| 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" | ||
| } | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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") | ||
| } | ||
| } |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
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 growsinbound/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.