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
6 changes: 5 additions & 1 deletion services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,10 @@ Shared vocabulary has one owner each, and domains use it rather than copy it. `i

`cmd/server` owns the execution lease. It acquires one `pgunit.Lease`, builds every lease-bound adapter on it, and passes the lease and those adapters together as one `execution.Owner` to `execution.StartWorker`. If anything fails before that call, `cmd/server` closes the lease. From that call the Worker owns cleanup: a failed start closes the lease before it returns, and a started Worker closes it after `Run` has cancelled and drained its work. Each close runs under its own bounded deadline, independent of the cancelled request or run. Lease-bound adapters and the store writer borrow the lease and never close it, and the Worker uses the lease only through `Owner.Lease`, never through an adapter. Store integration tests start the Worker the same way through `startWorker`.

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.

## Request handling

Every Agents API JSON route reads its body through `readJSONObject` before decoding, validation or lookup. The gate requires a JSON Content-Type, applies the route's body limit and rejects invalid UTF-8, malformed JSON (including unpaired surrogate escapes), repeated keys and non-object roots with the official messages; an empty body or `null` becomes `{}`. DELETE, multipart, Core extension and internal routes keep their own readers. Member names match exactly: decode request objects with `decodeInputObject`, or check `inexactMember` before another decoder, so `encoding/json` never matches a case variant.
Expand Down Expand Up @@ -92,7 +96,7 @@ Session status and last activity use the public projection in [`internal/api/ses

## Agents and model providers

Reusable Agents are tenant-scoped rows independent of Session snapshots and engine bindings. The store persists caller-validated configuration without applying harness restrictions or model defaults, with internal limits of 512 KiB for configuration and the `metadata.Encode` bound of 64 KiB for metadata. An update locks the Agent row while merging the supplied fields and enforcing the configuration bound, then commits configuration, metadata and update time together, so a stale full snapshot never overwrites another update. An empty update preserves the saved fields and advances `updated_at` through the same SQL update. Deletion is one tenant-scoped `DELETE … RETURNING id`. A Session copies the saved configuration into its immutable snapshot and never looks up its source again.
Reusable Agents are tenant-scoped rows independent of Session snapshots and engine bindings. `agents` accepts caller-validated configuration without applying harness restrictions or model defaults, with internal limits of 512 KiB for configuration and the `metadata.Encode` bound of 64 KiB for metadata. An update runs inside `agentpg`'s Agent row lock: `agents` merges the supplied fields over the locked Agent and enforces the configuration bound, then `agentpg` commits configuration, metadata and update time together, so a stale full snapshot never overwrites another update. An empty update preserves the saved fields and advances `updated_at` through the same SQL update. Deletion is one tenant-scoped `DELETE … RETURNING id`. A Session copies the saved configuration into its immutable snapshot and never looks up its source again.

Saved execution defaults keep a model-provider bundle whole at every replacement boundary: endpoint, key, protocol and limits are never inherited separately. Agent JSON holds only the safe provider fields and an output-only configured flag; the complete bundle is encrypted separately with a tenant and Agent binding and its own purpose, and configuration and secret changes commit together under the Agent row lock. Model-only edits need no key. Merged harness, protocol and limits are validated without reading keys. Session creation reads safe defaults and ciphertext in one snapshot, and a complete Session override does not decrypt the inherited bundle. The resolved bundle is frozen in an encrypted Session-owned row, and dispatch fails closed when that snapshot is missing or cannot be decrypted; later Agent edits, default changes, restarts and suspension never resolve it again.

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 @@ -89,7 +89,7 @@ func daemonComposition(t testing.TB) http.Handler {
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 }{},
Agents: struct{ api.Agents }{}, Sessions: struct{ api.Sessions }{}, SessionEvents: struct{ api.SessionEvents }{},
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 }{},
Admin: struct{ api.Admin }{}, AdminAudit: struct{ api.AdminAudit }{}, WriteAudit: struct{ api.WriteAudit }{}, Metrics: struct{ api.Metrics }{},
Expand Down
9 changes: 8 additions & 1 deletion services/core/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,11 +31,13 @@ import (
"time"

"github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/agents"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/api"
"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/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/pgunit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime"
Expand Down Expand Up @@ -109,6 +111,11 @@ func run() error {
executionStore.SetPublicURL(public)
units := pgunit.NewPool(pool)
auditStore := auditpg.New(units)
agentStore := agentpg.New(units, credentialKey)
agentService, err := agents.NewService(agentStore)
if err != nil {
return err
}
installation, err := installationFacts(public)
if err != nil {
return err
Expand Down Expand Up @@ -296,7 +303,7 @@ func run() error {
Engine: engine, Harnesses: kinds, CoreKeys: keyAdmin,
Installation: installation, InstallationBindings: executionStore,
Projects: executionStore, Vaults: executionStore, ModelProviders: executionStore, Files: executionStore,
Skills: executionStore, EnvironmentTemplates: executionStore, Agents: executionStore,
Skills: executionStore, EnvironmentTemplates: executionStore, 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
56 changes: 56 additions & 0 deletions services/core/internal/agents/agent.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
package agents

import (
"encoding/json"
"fmt"
"time"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
)

// Agent is saved configuration owned by a tenant. It has no Harness binding or
// live execution state; Session snapshots are separate objects.
type Agent struct {
ID string
TenantID string
Metadata map[string]string
Configuration json.RawMessage
CreatedAt time.Time
UpdatedAt time.Time
}

// MaxConfigurationBytes bounds a saved configuration and an update's patch.
const MaxConfigurationBytes = 512 * 1024

// MaxPageSize bounds one page of ListAgents.
const MaxPageSize = 100

// ListQuery selects one page of a tenant's Agents in creation order. After is
// the last Agent of the previous page; an After that names no Agent of the
// tenant is ErrNotFound.
type ListQuery struct {
TenantID string
After string
Limit int
Ascending bool
}

// Validate checks the page size.
func (q ListQuery) Validate() error {
if q.Limit < 1 || q.Limit > MaxPageSize {
return fmt.Errorf("%w: page size must be 1..%d", ErrInvalidInput, MaxPageSize)
}
return nil
}

// Page is one page of Agents. NextCursor is empty on the last page.
type Page struct {
Agents []Agent
NextCursor string
}

// ModelProviderChange replaces an Agent's saved model provider bundle. A nil
// Provider removes it.
type ModelProviderChange struct {
Provider *v1.ModelProviderInput
}
149 changes: 149 additions & 0 deletions services/core/internal/agents/configuration.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
package agents

import (
"encoding/json"
"fmt"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/jsonobject"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/metadata"
)

// encodeMetadata encodes caller metadata for storage; nil is the empty object.
func encodeMetadata(values map[string]string) (json.RawMessage, error) {
encoded, err := metadata.Encode(values)
if err != nil {
return nil, fmt.Errorf("%w: metadata: %v", ErrInvalidInput, err)
}
return encoded, nil
}

// createConfiguration normalizes the configuration of a new Agent.
func createConfiguration(raw json.RawMessage) (json.RawMessage, error) {
if len(raw) == 0 || len(raw) > MaxConfigurationBytes {
return nil, fmt.Errorf("%w: configuration must be an object of at most 512 KiB", ErrInvalidInput)
}
return normalizeConfiguration(raw)
}

// decodePatch normalizes an update's supplied fields; empty supplies none.
func decodePatch(raw json.RawMessage) (map[string]json.RawMessage, error) {
if len(raw) == 0 {
raw = json.RawMessage(`{}`)
}
if len(raw) > MaxConfigurationBytes {
return nil, fmt.Errorf("%w: configuration patch exceeds 512 KiB", ErrInvalidInput)
}
normalized, err := normalizeConfiguration(raw)
if err != nil {
return nil, err
}
var patch map[string]json.RawMessage
if err := json.Unmarshal(normalized, &patch); err != nil {
return nil, fmt.Errorf("%w: configuration patch: %v", ErrInvalidInput, err)
}
return patch, nil
}

func normalizeConfiguration(raw json.RawMessage) (json.RawMessage, error) {
normalized, err := jsonobject.Normalize(raw)
if err != nil {
return nil, fmt.Errorf("%w: configuration: %v", ErrInvalidInput, err)
}
return normalized, nil
}

// validateModelExecution checks the saved Harness configuration and that the
// model provider bundle suits the saved Harness. provider is the bundle being
// saved, if any.
func validateModelExecution(configuration json.RawMessage, provider *v1.ModelProviderInput) error {
if provider != nil {
if err := provider.Validate(); err != nil {
return fmt.Errorf("%w: %s", ErrInvalidInput, err)
}
}
var config struct {
Core *v1.SavedAgentCore `json:"x_agents_core"`
}
if err := json.Unmarshal(configuration, &config); err != nil {
return fmt.Errorf("%w: x_agents_core: %v", ErrInvalidInput, err)
}
if config.Core != nil {
if err := v1.ValidateHarnessConfig(config.Core.Harness, config.Core.HarnessConfig); err != nil {
return fmt.Errorf("%w: %s", ErrInvalidInput, err)
}
}
if config.Core == nil || config.Core.ModelProvider == nil || config.Core.Harness == "" {
return nil
}
if err := config.Core.ModelProvider.ValidateHarness(config.Core.Harness); err != nil {
return fmt.Errorf("%w: %s", ErrInvalidInput, err)
}
return nil
}

// mergeConfiguration applies an update's supplied fields to the saved
// configuration. Supplied top-level fields replace saved ones; x_agents_core
// subfields merge one by one, and a null x_agents_core removes it. Changing the
// model, the model provider or the Harness without supplying harness_config
// resets harness_config to {}, because native options belong to one model and
// Harness.
func mergeConfiguration(saved json.RawMessage, patch map[string]json.RawMessage) (json.RawMessage, error) {
configuration := map[string]json.RawMessage{}
if err := json.Unmarshal(saved, &configuration); err != nil {
return nil, fmt.Errorf("decode saved Agent configuration: %w", err)
}
coreNull := string(patch["x_agents_core"]) == "null"
var corePatch map[string]json.RawMessage
if raw := patch["x_agents_core"]; len(raw) > 0 && !coreNull {
if err := json.Unmarshal(raw, &corePatch); err != nil {
return nil, fmt.Errorf("%w: x_agents_core: %v", ErrInvalidInput, err)
}
}
_, modelChanged := patch["model"]
_, providerChanged := corePatch["model_provider"]
_, harnessChanged := corePatch["harness"]
_, nativeSupplied := corePatch["harness_config"]
if (modelChanged || providerChanged || harnessChanged) && !nativeSupplied && !coreNull {
if corePatch == nil {
corePatch = map[string]json.RawMessage{}
}
corePatch["harness_config"] = json.RawMessage(`{}`)
}
for field, value := range patch {
if field != "x_agents_core" {
configuration[field] = value
}
}
switch {
case coreNull:
configuration["x_agents_core"] = patch["x_agents_core"]
case corePatch != nil:
core := map[string]json.RawMessage{}
if old := configuration["x_agents_core"]; len(old) != 0 && string(old) != "null" {
if err := json.Unmarshal(old, &core); err != nil {
return nil, fmt.Errorf("decode saved Agent configuration: %w", err)
}
}
for key, replacement := range corePatch {
core[key] = replacement
}
merged, err := json.Marshal(core)
if err != nil {
return nil, err
}
configuration["x_agents_core"] = merged
}
merged, err := json.Marshal(configuration)
if err != nil {
return nil, err
}
merged, err = normalizeConfiguration(merged)
if err != nil {
return nil, err
}
if len(merged) > MaxConfigurationBytes {
return nil, fmt.Errorf("%w: configuration exceeds 512 KiB", ErrInvalidInput)
}
return merged, nil
}
Loading
Loading