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
3 changes: 2 additions & 1 deletion services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ Shared vocabulary has one owner each, and domains use it rather than copy it. `i
Domain owners, each with its PostgreSQL adapter under `internal/persistence/postgres`:

- `agents` (`agentpg`): saved Agents, their configuration merge and bounds, and the encrypted model-provider bundle bound to each Agent.
- `files` (`filepg`): source Files.

## Request handling

Expand All @@ -42,7 +43,7 @@ On the Beta group, the OpenAI-Beta check (exactly one `agents=v1` value) runs be

## Source Files and Artifacts

Source Files are Project resources with a lifecycle independent of copied workspace files. Store immutable source metadata and PostgreSQL large objects in Core's database with the pinned pgx driver. Upload validation, metadata insertion and the bytes commit atomically; deletion removes the metadata and unlinks the object in one transaction. Keep OIDs private and authorize every metadata, content and delete lookup by tenant before opening a body. Stream bounded chunks; never hold an entire upload in memory or use a filename as a filesystem path. A direct download of the `user_data` purpose is rejected after the tenant-scoped metadata lookup, while initialization and workspace copies keep their authorized store read. A read-only repeatable-read transaction preserves an admitted source across concurrent deletion; resolve that snapshot before entering the Environment write path, and a later deletion never undoes a completed workspace copy. Bound request and transaction lifetimes, roll back incomplete bodies and never retry an ambiguous commit automatically. Backups must include PostgreSQL large objects, and a schema rollback must not orphan them.
Source Files are Project resources with a lifecycle independent of copied workspace files. `files` owns their vocabulary, the upload envelope and 512 MiB content bound, the list rules and the create and delete use cases; `filepg` stores immutable source metadata and PostgreSQL large objects in Core's database with the pinned pgx driver. Upload validation, metadata insertion, the bytes and the write audit commit atomically; deletion removes the metadata, unlinks the object and records the audit in one transaction. Keep OIDs private and authorize every metadata, content and delete lookup by tenant before opening a body. Write content through `pgunit`'s large-object writer, which streams bounded chunks and reports the size and SHA-256 digest; never hold an entire upload in memory or use a filename as a filesystem path. A direct download of the `user_data` purpose is rejected after the tenant-scoped metadata lookup, while workspace copies keep their authorized `filepg` read and Session initialization copies read the object inside Session creation's store transaction. A read-only repeatable-read transaction preserves an admitted source across concurrent deletion; resolve that snapshot before entering the Environment write path, and a later deletion never undoes a completed workspace copy. Bound request and transaction lifetimes, roll back incomplete bodies and never retry an ambiguous commit automatically. Backups must include PostgreSQL large objects, and a schema rollback must not orphan them.

Session Artifacts are immutable published copies, separate from live workspace files and source Files. The private output exporter reuses the authorized workspace path boundary and streams bounded bytes; publication requires complete capture and confirmed helper and transport success, not merely valid archive syntax. The daemon owns and drains the exporter's stdout pipe separately from child reaping, so pull-transport backpressure cannot consume the process-exit I/O deadline; after helper exit, each pipe read has one second, reset after consumer delays, which rejects inherited pipes that never close. Cancellation closes the owned reader and the dispatch consumer, then waits for the child. Never extract an output archive into Core's filesystem or hold the execution lease through a large transfer.

Expand Down
2 changes: 1 addition & 1 deletion services/core/cmd/server/http_routes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,7 @@ func daemonComposition(t testing.TB) http.Handler {
apiHandler, err := api.NewHandler(api.Dependencies{
Engine: "codex", CoreKeys: admin, InstallationBindings: struct{ api.InstallationBindings }{},
Projects: trapProjects{keys: keys}, Vaults: struct{ api.Vaults }{}, ModelProviders: struct{ api.ModelProviders }{},
Files: struct{ api.Files }{}, Skills: struct{ api.Skills }{}, EnvironmentTemplates: struct{ api.EnvironmentTemplates }{},
Files: struct{ api.Files }{}, FilesReader: struct{ api.FilesReader }{}, Skills: struct{ api.Skills }{}, EnvironmentTemplates: struct{ api.EnvironmentTemplates }{},
Agents: struct{ api.Agents }{}, AgentsReader: struct{ api.AgentsReader }{}, Sessions: struct{ api.Sessions }{}, SessionEvents: struct{ api.SessionEvents }{},
SessionHistory: struct{ api.SessionHistory }{}, Subagents: struct{ api.Subagents }{}, Artifacts: struct{ api.Artifacts }{},
SessionAdmin: struct{ api.SessionAdmin }{}, Environments: struct{ api.Environments }{}, ExecutorConnections: struct{ api.ExecutorConnections }{},
Expand Down
13 changes: 11 additions & 2 deletions services/core/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,11 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/coremetrics"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/databaseurl"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/files"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/nativeinstaller"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/agentpg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/auditpg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/filepg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeenrollment"
Expand Down Expand Up @@ -299,11 +301,18 @@ func run() error {
metricsDone := make(chan struct{})
go func() { defer close(metricsDone); metrics.Run(metricsCtx) }()
defer func() { cancelMetrics(); <-metricsDone }()
fileStore := filepg.New(units)
fileService, err := files.NewService(fileStore)
if err != nil {
return err
}
deps := api.Dependencies{
Engine: engine, Harnesses: kinds, CoreKeys: keyAdmin,
Installation: installation, InstallationBindings: executionStore,
Projects: executionStore, Vaults: executionStore, ModelProviders: executionStore, Files: executionStore,
Skills: executionStore, EnvironmentTemplates: executionStore, Agents: agentService, AgentsReader: agentStore,
Projects: executionStore, Vaults: executionStore, ModelProviders: executionStore,
Skills: executionStore, EnvironmentTemplates: executionStore,
Files: fileService, FilesReader: fileStore,
Agents: agentService, AgentsReader: agentStore,
Sessions: executionStore, SessionEvents: executionStore, SessionHistory: executionStore,
Subagents: executionStore, Artifacts: executionStore, SessionAdmin: executionStore,
Environments: executionStore, ExecutorConnections: executorConnections{store: executionStore, registry: registry},
Expand Down
8 changes: 5 additions & 3 deletions services/core/internal/api/dependencies.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ type Dependencies struct {
Vaults Vaults
ModelProviders ModelProviders
Files Files
FilesReader FilesReader
Skills Skills
EnvironmentTemplates EnvironmentTemplates
Agents Agents
Expand Down Expand Up @@ -111,9 +112,10 @@ func (d Dependencies) validate() error {
}
if err := required(
field{"InstallationBindings", d.InstallationBindings}, field{"Projects", d.Projects}, field{"Vaults", d.Vaults},
field{"ModelProviders", d.ModelProviders}, field{"Files", d.Files}, field{"Skills", d.Skills},
field{"EnvironmentTemplates", d.EnvironmentTemplates}, field{"Agents", d.Agents}, field{"AgentsReader", d.AgentsReader},
field{"Sessions", d.Sessions},
field{"ModelProviders", d.ModelProviders}, field{"Skills", d.Skills},
field{"Files", d.Files}, field{"FilesReader", d.FilesReader},
field{"Agents", d.Agents}, field{"AgentsReader", d.AgentsReader},
field{"EnvironmentTemplates", d.EnvironmentTemplates}, field{"Sessions", d.Sessions},
field{"SessionEvents", d.SessionEvents}, field{"SessionHistory", d.SessionHistory}, field{"Subagents", d.Subagents},
field{"Artifacts", d.Artifacts}, field{"SessionAdmin", d.SessionAdmin}, field{"Environments", d.Environments},
field{"ExecutorConnections", d.ExecutorConnections}, field{"Admin", d.Admin}, field{"AdminAudit", d.AdminAudit}, field{"WriteAudit", d.WriteAudit},
Expand Down
13 changes: 9 additions & 4 deletions services/core/internal/api/dependencies_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ type testFakes struct {
vaults *fakeVaults
modelProviders *fakeModelProviders
files *fakeFiles
filesReader *fakeFilesReader
skills *fakeSkills
environmentTemplates *fakeEnvironmentTemplates
agents *fakeAgents
Expand Down Expand Up @@ -55,8 +56,10 @@ func testDependencies(t testing.TB) (Dependencies, *testFakes) {
t.Helper()
f := &testFakes{
projects: &fakeProjects{t: t}, vaults: &fakeVaults{t: t}, modelProviders: &fakeModelProviders{t: t},
files: &fakeFiles{t: t}, skills: &fakeSkills{t: t}, environmentTemplates: &fakeEnvironmentTemplates{t: t},
agents: &fakeAgents{t: t}, agentsReader: &fakeAgentsReader{t: t}, sessions: &fakeSessions{t: t}, sessionEvents: &fakeSessionEvents{t: t},
skills: &fakeSkills{t: t}, environmentTemplates: &fakeEnvironmentTemplates{t: t},
files: &fakeFiles{t: t}, filesReader: &fakeFilesReader{t: t},
agents: &fakeAgents{t: t}, agentsReader: &fakeAgentsReader{t: t},
sessions: &fakeSessions{t: t}, sessionEvents: &fakeSessionEvents{t: t},
sessionHistory: &fakeSessionHistory{t: t}, subagents: &fakeSubagents{t: t}, artifacts: &fakeArtifacts{t: t},
sessionAdmin: &fakeSessionAdmin{t: t}, environments: &fakeEnvironments{t: t}, executorConnections: &fakeExecutorConnections{t: t},
admin: &fakeAdmin{t: t}, adminAudit: &fakeAdminAudit{t: t}, writeAudit: &fakeWriteAudit{t: t}, metrics: &fakeMetrics{t: t},
Expand All @@ -66,8 +69,10 @@ func testDependencies(t testing.TB) (Dependencies, *testFakes) {
}
return Dependencies{
Engine: "codex", CoreKeys: coreKeys(t, "admin"), InstallationBindings: f.installationBindings,
Projects: f.projects, Vaults: f.vaults, ModelProviders: f.modelProviders, Files: f.files, Skills: f.skills,
EnvironmentTemplates: f.environmentTemplates, Agents: f.agents, AgentsReader: f.agentsReader, Sessions: f.sessions, SessionEvents: f.sessionEvents,
Projects: f.projects, Vaults: f.vaults, ModelProviders: f.modelProviders, Skills: f.skills,
Files: f.files, FilesReader: f.filesReader,
Agents: f.agents, AgentsReader: f.agentsReader,
EnvironmentTemplates: f.environmentTemplates, Sessions: f.sessions, SessionEvents: f.sessionEvents,
SessionHistory: f.sessionHistory, Subagents: f.subagents, Artifacts: f.artifacts, SessionAdmin: f.sessionAdmin,
Environments: f.environments, ExecutorConnections: f.executorConnections, Admin: f.admin, AdminAudit: f.adminAudit, WriteAudit: f.writeAudit,
Metrics: f.metrics, RuntimeObservations: f.runtimeObservations, RuntimeHistory: f.runtimeHistory,
Expand Down
7 changes: 4 additions & 3 deletions services/core/internal/api/environment_files_create.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/echotext"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/files"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/go-chi/chi/v5"
)
Expand Down Expand Up @@ -101,9 +102,9 @@ func (h *Handler) createEnvironmentFile(w http.ResponseWriter, r *http.Request)
if request.Type == "file_id" {
ctx, cancel := context.WithTimeout(r.Context(), 30*time.Second)
defer cancel()
err = h.Files.ReadSourceFile(ctx, tenantID(r), *request.FileID, func(file store.SourceFile, body io.Reader) error {
err = h.FilesReader.Read(ctx, tenantID(r), *request.FileID, func(file files.File, body io.Reader) error {
if file.SizeBytes > proto.WorkspaceWriteMaxBytes {
return store.ErrSourceFileTooLarge
return files.ErrTooLarge
}
data, err = io.ReadAll(io.LimitReader(body, proto.WorkspaceWriteMaxBytes+1))
if err == nil && int64(len(data)) != file.SizeBytes {
Expand All @@ -112,7 +113,7 @@ func (h *Handler) createEnvironmentFile(w http.ResponseWriter, r *http.Request)
return err
})
if err != nil {
writeStoreError(w, r, err)
writeFilesError(w, r, err)
return
}
}
Expand Down
3 changes: 2 additions & 1 deletion services/core/internal/api/environment_files_write_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"time"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/files"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/google/uuid"
)
Expand Down Expand Up @@ -47,7 +48,7 @@ func TestEnvironmentFileCreateSourceCopyKeepsDestinationBound(t *testing.T) {
h, f := environmentFileCreateHandler(t, sources.wire)
data := bytes.Repeat([]byte{9}, 6<<20)
sources.tenant, sources.data = f.environment.TenantID, data
sources.file = store.SourceFile{ID: "file-" + uuid.NewString(), SizeBytes: int64(len(data)), CreatedAt: time.Unix(1, 0)}
sources.file = files.File{ID: "file-" + uuid.NewString(), SizeBytes: int64(len(data)), CreatedAt: time.Unix(1, 0)}
body := `{"type":"file_id","file_id":"` + sources.file.ID + `","path":"/workspace/copy.bin"}`
if w := requestCreateEnvironmentFile(h, f.environment.ID, body, "files-key"); w.Code != 201 || f.writes != 1 || !bytes.Equal(f.data, data) {
t.Fatal("6 MiB source copy rejected", w.Code, w.Body)
Expand Down
12 changes: 8 additions & 4 deletions services/core/internal/api/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,12 @@ func writeInputError(w http.ResponseWriter, r *http.Request, err error) {
}
}

// writeContentTooLarge reports uploaded or copied content beyond the
// operation's limit.
func writeContentTooLarge(w http.ResponseWriter) {
writeError(w, http.StatusRequestEntityTooLarge, "request_too_large", "File exceeds this operation's content limit.")
}

const unstorableTextMessage = "Request text contains characters this service cannot store or compare, such as U+0000 or invalid UTF-8."

// fieldError is a request validation failure reported with the official
Expand Down Expand Up @@ -185,8 +191,6 @@ func writeStoreError(w http.ResponseWriter, r *http.Request, err error, notFound

case errors.Is(err, store.ErrDefaultSkillVersion):
writeError(w, http.StatusBadRequest, "invalid_value", "Cannot delete the default skill version.", "version")
case errors.Is(err, store.ErrSourceFileTooLarge):
writeError(w, http.StatusRequestEntityTooLarge, "request_too_large", "File exceeds this operation's content limit.")
case errors.Is(err, store.ErrModelProviderRequired):
writeError(w, http.StatusBadRequest, "model_provider_required", "This Session was created without a model provider and cannot run. Create a new Session with x_agents_core.model_provider or an Agent that has one saved.")
case errors.Is(err, store.ErrHostedEnvironmentFailed):
Expand Down Expand Up @@ -220,8 +224,8 @@ func writeStoreError(w http.ResponseWriter, r *http.Request, err error, notFound
}
case errors.Is(err, store.ErrNotFound):
code := "not_found_error"
// Files and Skills retain their non-beta error envelope.
if strings.HasPrefix(r.URL.Path, "/v1/files/") || strings.HasPrefix(r.URL.Path, "/v1/skills/") || r.URL.Path == "/v1/files" || r.URL.Path == "/v1/skills" {
// Skills retain their non-beta error envelope.
if strings.HasPrefix(r.URL.Path, "/v1/skills/") || r.URL.Path == "/v1/skills" {
code = ""
}
writeError(w, http.StatusNotFound, code, "Resource not found.", notFoundParam...)
Expand Down
35 changes: 35 additions & 0 deletions services/core/internal/api/errors_files.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
package api

import (
"errors"
"net/http"
"strings"

"github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/files"
)

// writeFilesError maps a files error to its public response. notFoundParam
// names the parameter a missing File came from.
func writeFilesError(w http.ResponseWriter, r *http.Request, err error, notFoundParam ...string) {
if writeAuditSourceError(w, r, err) || writeTextValueError(w, r, err) || writeCredentialUnavailableError(w, r, err) {
return
}
switch {
case errors.Is(err, files.ErrNotFound):
code := "not_found_error"
// The Files routes keep their non-beta error envelope.
if r.URL.Path == "/v1/files" || strings.HasPrefix(r.URL.Path, "/v1/files/") {
code = ""
}
writeError(w, http.StatusNotFound, code, "Resource not found.", notFoundParam...)
case errors.Is(err, files.ErrTooLarge):
writeContentTooLarge(w)
case errors.Is(err, files.ErrInvalidInput):
writeError(w, http.StatusBadRequest, "invalid_request", "Invalid resource identifier or request limits.")
default:
// Storage failures can include submitted values; do not log the error.
log.Ctx(r.Context()).Error("oac-core persistence operation failed")
writeError(w, http.StatusInternalServerError, "internal_error", "The operation could not be completed.")
}
}
7 changes: 6 additions & 1 deletion services/core/internal/api/errors_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/credentialcrypto"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/files"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/textvalue"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/writeaudit"
Expand All @@ -27,7 +28,11 @@ func TestResourceNotFoundErrorSurfaces(t *testing.T) {
t.Run(path, func(t *testing.T) {
response := httptest.NewRecorder()
request := httptest.NewRequest(http.MethodGet, path, nil)
writeStoreError(response, request, fmt.Errorf("lookup: %w", store.ErrNotFound))
if strings.HasPrefix(path, "/v1/files") {
writeFilesError(response, request, fmt.Errorf("lookup: %w", files.ErrNotFound))
} else {
writeStoreError(response, request, fmt.Errorf("lookup: %w", store.ErrNotFound))
}
var body v1.ErrorResponse
if response.Code != http.StatusNotFound || json.Unmarshal(response.Body.Bytes(), &body) != nil {
t.Fatalf("response = %d %s", response.Code, response.Body)
Expand Down
Loading