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: 4 additions & 2 deletions services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ These are the code-level rules of `services/core` that no contract states. Contr

`internal/persistence/postgres/auditpg` is the one adapter that other adapters call directly. Audit rows are written inside the business transaction, so each adapter calls `RecordWriteAudit`, `RecordAdminMutation` or `RecordDeploymentMutation` with its own transaction's queries; the audit provenance travels in the context. `writeaudit` and `adminaudit` own the sources, their validation, `ErrInvalidSource` and the read models, and `auditpg.Store` serves the audit reads.

The adapter that stores a secret seals and opens it with the credential key `cmd/server` builds it with: domain storage interfaces carry plaintext, domain service constructors never take the key, a missing key is `credentialcrypto.ErrUnavailable`, and a ciphertext that fails to open or authenticate is an internal error, never a missing key, value or row.

Two errors are shared across domains, each with one `api` helper: `textvalue.ErrUnstorable` (400, `writeTextValueError`) for text PostgreSQL cannot store, and `credentialcrypto.ErrUnavailable` (503, `writeCredentialUnavailableError`) for a missing credential key. The audit `ErrInvalidSource` errors pass through adapters unchanged, and `writeAuditSourceError` maps both to 400.

Shared vocabulary has one owner each, and domains use it rather than copy it. `internal/environmentconfig` owns Environment setup, Skills, Plugins and initial files with their validation and public metadata; `Setup.Validate` checks requested configuration, where a Skill may be an unresolved reference, and `Setup.ValidateInstalled` checks frozen, installable configuration. `internal/skills` owns `ParseVersion`, the canonical positive decimal Skill version. `internal/metadata` owns the metadata rules: `Validate` for the pair, key and value limits and U+0000, `ValidateStorable` for U+0000 alone, and `Encode` with its 64 KiB bound. `internal/jsonobject` owns `Normalize`, the stable encoding of stored JSON objects that snapshots and retry identities compare. These packages import no persistence.
Expand Down Expand Up @@ -108,7 +110,7 @@ Reusable Agents are tenant-scoped rows independent of Session snapshots and engi

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.

`modelconfiguration` owns each Harness's deployment default: its service validates a replacement through the Harness declaration and seals the complete bundle, and `modelconfigurationpg` stores it under a private revision UUID generated on every PUT, identical replacements included, with its audit row in the same transaction. Session creation resolves the ciphertext and revision together and freezes them; retries and older Sessions never gain or replace revision metadata. After a successful root terminal commit, the Dispatcher's required `modelconfiguration.Observer` runs one independent pool operation with at most one second to update the matching current revision. Only completed Turns and native provider failures with `engine_failed` count; cancelled work, Core or Runtime errors and input-policy classifications never do. `ShouldObserveProvider` skips the round trip for outcomes that cannot count, and the SQL stays authoritative: it verifies the tenant, root Turn and committed outcome. Both check the same shared cases. The metadata-only transaction sets statement and lock timeouts within the remaining budget and issues one UPDATE that locks only the default and samples database time after the lock. Errors throttle for 30 seconds, ordinary successes throttle for 30 seconds with one immediate recovery write after each accepted error, and an unchanged revision has at most three effective writes in any 30-second window of nondecreasing database time. Observations never change `updated_at`, readiness or execution truth, and can be lost or stale; there is no queue, retry, probe or backfill.
`modelconfiguration` owns each Harness's deployment default: its service validates a replacement through the Harness declaration, and `modelconfigurationpg` seals the complete bundle to the Harness and stores it under a private revision UUID generated on every PUT, identical replacements included, with its audit row in the same transaction. Session creation resolves the opened bundle and its revision together and freezes them; retries and older Sessions never gain or replace revision metadata. After a successful root terminal commit, the Dispatcher's required `modelconfiguration.Observer` runs one independent pool operation with at most one second to update the matching current revision. Only completed Turns and native provider failures with `engine_failed` count; cancelled work, Core or Runtime errors and input-policy classifications never do. `ShouldObserveProvider` skips the round trip for outcomes that cannot count, and the SQL stays authoritative: it verifies the tenant, root Turn and committed outcome. Both check the same shared cases. The metadata-only transaction sets statement and lock timeouts within the remaining budget and issues one UPDATE that locks only the default and samples database time after the lock. Errors throttle for 30 seconds, ordinary successes throttle for 30 seconds with one immediate recovery write after each accepted error, and an unchanged revision has at most three effective writes in any 30-second window of nondecreasing database time. Observations never change `updated_at`, readiness or execution truth, and can be lost or stale; there is no queue, retry, probe or backfill.

Session execution-configuration reads use a separate immutable safe projection written with its provenance in the Session's creation transaction. It reads no ciphertext, never recomputes sources from current Agents or defaults, never touches activity or wakes a sandbox, and does not affect retry identity.

Expand All @@ -121,7 +123,7 @@ Provider input validation uses the adapter rules in `internal/harnessconfig`: on
[Vaults and credentials](../../contracts/agents-api/vaults.md) describes the resources, selection rules, refresh and deletion. `vaults` implements them, with `vaultpg` as its storage, under these rules:

- Credentials are children of tenant-owned Vaults. Creation admits the owner in the same SQL statement as the insert; retrieval joins the owning Vault; listing enforces Project and Vault ownership on the parent, cursor and row query. Metadata queries never select ciphertext and need no encryption key.
- Secret values are encrypted before they reach SQL, with Core's separately configured random 32-byte key and the standard library's random-nonce AES-GCM. The versioned authenticated binding covers tenant, Vault, Credential, authentication purpose and exact destination. Never reuse daemon transport encryption for this storage. A missing key disables credential writes with `credentialcrypto.ErrUnavailable`; a malformed configured key fails startup.
- `vaultpg` seals secret values before they reach SQL, with Core's separately configured random 32-byte key and the standard library's random-nonce AES-GCM. The versioned authenticated binding covers tenant, Vault, Credential, authentication purpose and exact destination. Never reuse daemon transport encryption for this storage. A missing key disables credential writes with `credentialcrypto.ErrUnavailable`; a malformed configured key fails startup.
- A static replacement is one SQL mutation scoped by tenant, Vault, Credential, auth type and destination, reusing the safe metadata for the immutable binding; it never decrypts the previous token, and a failed write keeps the old row.
- OAuth refresh and replacement serialize on the Credential row lock, authenticate the stored grant metadata against its encrypted copy before using an endpoint, and persist the refreshed grant before returning an access token, so a stale refresh cannot undo a deletion or Vault cascade.
- Credential deletion is one mutation checked by tenant, Vault and ID; Vault deletion removes the parent and its Credentials through the foreign-key cascade in one SQL statement, without decrypting, needing the key or calling providers.
Expand Down
8 changes: 4 additions & 4 deletions services/core/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -128,8 +128,8 @@ func run() error {
if err != nil {
return err
}
vaultStore := vaultpg.New(units)
vaultService, err := vaults.NewService(vaultStore, credentialKey, oauthClient)
vaultStore := vaultpg.New(units, credentialKey)
vaultService, err := vaults.NewService(vaultStore, oauthClient)
if err != nil {
return err
}
Expand All @@ -138,8 +138,8 @@ func run() error {
if err != nil {
return err
}
modelConfigurationStore := modelconfigurationpg.New(units)
modelConfigurationService, err := modelconfiguration.NewService(modelConfigurationStore, credentialKey)
modelConfigurationStore := modelconfigurationpg.New(units, credentialKey)
modelConfigurationService, err := modelconfiguration.NewService(modelConfigurationStore)
if err != nil {
return err
}
Expand Down
5 changes: 3 additions & 2 deletions services/core/internal/agents/storage.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,8 @@ type Reader interface {
GetAgent(ctx context.Context, tenantID, agentID string) (Agent, error)
ListAgents(context.Context, ListQuery) (Page, error)
// GetAgentWithModelProvider reads the Agent and its opened model provider
// bundle from one snapshot. The bundle is nil when the Agent has none; a
// bundle that cannot be opened is credentialcrypto.ErrUnavailable.
// bundle from one snapshot. The bundle is nil when the Agent has none.
// Opening one without a credential key is credentialcrypto.ErrUnavailable;
// a bundle that fails to open is an internal error.
GetAgentWithModelProvider(ctx context.Context, tenantID, agentID string) (Agent, *v1.ModelProviderInput, error)
}
14 changes: 4 additions & 10 deletions services/core/internal/api/core_model_provider_validation_test.go
Original file line number Diff line number Diff line change
@@ -1,15 +1,13 @@
package api

import (
"bytes"
"context"
"encoding/json"
"net/http"
"reflect"
"strings"
"testing"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/credentialcrypto"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/modelconfiguration"
)

Expand All @@ -30,17 +28,13 @@ func (s *coreProviderValidationStore) Delete(context.Context, string) error {
return nil
}

func (s *coreProviderValidationStore) LoadSealed(context.Context, string) (modelconfiguration.Sealed, error) {
unexpectedCall(s.t, "LoadSealed")
return modelconfiguration.Sealed{}, nil
func (s *coreProviderValidationStore) LoadBundle(context.Context, string) (modelconfiguration.Bundle, error) {
unexpectedCall(s.t, "LoadBundle")
return modelconfiguration.Bundle{}, nil
}

func (s *coreProviderValidationStore) configure(d *Dependencies, _ *testFakes) {
cipher, err := credentialcrypto.New(bytes.Repeat([]byte{3}, 32))
if err != nil {
s.t.Fatal(err)
}
service, err := modelconfiguration.NewService(s, cipher)
service, err := modelconfiguration.NewService(s)
if err != nil {
s.t.Fatal(err)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,8 +51,8 @@ func newFinishObservationFixture(t *testing.T, maxConnections int32) finishObser
if err != nil {
t.Fatal(err)
}
defaults := modelconfigurationpg.New(pgunit.NewPool(pool))
service, err := modelconfiguration.NewService(defaults, cipher)
defaults := modelconfigurationpg.New(pgunit.NewPool(pool), cipher)
service, err := modelconfiguration.NewService(defaults)
if err != nil {
t.Fatal(err)
}
Expand Down Expand Up @@ -194,7 +194,7 @@ func TestFinishRunObservationLockTimeoutAndFailureKeepLease(t *testing.T) {
if err != nil {
t.Fatal(err)
}
d := Dispatcher{Observer: modelconfigurationpg.New(pgunit.NewPool(pool))}
d := Dispatcher{Observer: modelconfigurationpg.New(pgunit.NewPool(pool), nil)}
started := time.Now()
d.observeDeploymentProvider(f.tenant, f.session.ID, turn)
if elapsed := time.Since(started); elapsed < 900*time.Millisecond || elapsed > 2*time.Second {
Expand Down
14 changes: 7 additions & 7 deletions services/core/internal/modelconfiguration/configuration.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,18 +40,18 @@ type Snapshot struct {
}

// Record is a validated deployment default ready to store: its safe columns
// and the sealed complete bundle. Storage assigns a new revision to each
// Record it stores.
// and the complete bundle, provider key included, which storage seals to the
// Harness. Storage assigns a new revision to each Record it stores.
type Record struct {
Harness string
Provider v1.ModelProviderView
Model string
HarnessConfig json.RawMessage
Sealed []byte
Configuration v1.ModelConfigurationInput
}

// Sealed is a stored bundle and the revision read with it.
type Sealed struct {
Bundle []byte
Revision uuid.UUID
// Bundle is an opened stored bundle and the revision read with it.
type Bundle struct {
Configuration v1.ModelConfigurationInput
Revision uuid.UUID
}
6 changes: 3 additions & 3 deletions services/core/internal/modelconfiguration/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,9 @@ import "errors"

// Replace and Resolve report a configuration the Harness declaration rejects
// with the contract's *v1.ModelProviderError, which names the field, and a
// bundle they cannot seal or open with credentialcrypto.ErrUnavailable. Storage
// passes textvalue.ErrUnstorable and adminaudit.ErrInvalidSource through
// unchanged.
// missing credential key with credentialcrypto.ErrUnavailable. A bundle that
// fails to open is an internal error. Storage passes textvalue.ErrUnstorable
// and adminaudit.ErrInvalidSource through unchanged.
var (
// ErrNotFound reports a Harness without a deployment default.
ErrNotFound = errors.New("the harness has no deployment default model configuration")
Expand Down
45 changes: 13 additions & 32 deletions services/core/internal/modelconfiguration/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,47 +2,35 @@ package modelconfiguration

import (
"context"
"encoding/json"
"errors"

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

// Service replaces, removes and resolves deployment defaults. The complete
// bundle, key included, is sealed to its Harness before it reaches storage.
// Service replaces, removes and resolves deployment defaults. Storage seals
// the complete bundle, key included, to its Harness.
type Service struct {
storage Storage
cipher *credentialcrypto.Cipher
}

// NewService requires storage. A nil cipher means credential encryption is not
// configured: Replace and Resolve then return credentialcrypto.ErrUnavailable.
func NewService(storage Storage, cipher *credentialcrypto.Cipher) (*Service, error) {
// NewService requires storage.
func NewService(storage Storage) (*Service, error) {
if storage == nil {
return nil, errors.New("model configuration requires storage")
}
return &Service{storage: storage, cipher: cipher}, nil
return &Service{storage: storage}, nil
}

// Replace validates the complete configuration through the Harness
// declaration, seals it and stores it under a new revision. Sessions that
// already froze a default keep theirs.
// declaration and stores it under a new revision. Sessions that already froze
// a default keep theirs.
func (s *Service) Replace(ctx context.Context, replacement Replacement) (Configuration, error) {
configuration := replacement.Configuration
if err := configuration.ValidateHarness(replacement.Harness); err != nil {
return Configuration{}, err
}
raw, err := json.Marshal(configuration)
if err != nil {
return Configuration{}, err
}
sealed, err := s.cipher.SealDeploymentModelProvider(raw, replacement.Harness)
if err != nil {
return Configuration{}, credentialcrypto.ErrUnavailable
}
view := configuration.SafeView()
return s.storage.Replace(ctx, Record{Harness: replacement.Harness, Provider: *view.ModelProvider, Model: view.Model, HarnessConfig: view.HarnessConfig, Sealed: sealed})
return s.storage.Replace(ctx, Record{Harness: replacement.Harness, Provider: *view.ModelProvider, Model: view.Model, HarnessConfig: view.HarnessConfig, Configuration: configuration})
}

// Delete removes the Harness's default. It is idempotent. Sessions that
Expand All @@ -52,28 +40,21 @@ func (s *Service) Delete(ctx context.Context, harness string) error {
}

// Resolve opens the Harness's default for Session creation. It returns nil
// when the Harness has none, and credentialcrypto.ErrUnavailable when the
// bundle does not open.
// when the Harness has none, and credentialcrypto.ErrUnavailable without a
// credential key.
func (s *Service) Resolve(ctx context.Context, harness string) (*Snapshot, error) {
sealed, err := s.storage.LoadSealed(ctx, harness)
bundle, err := s.storage.LoadBundle(ctx, harness)
if errors.Is(err, ErrNotFound) {
return nil, nil
}
if err != nil {
return nil, err
}
raw, err := s.cipher.OpenDeploymentModelProvider(sealed.Bundle, harness)
if err != nil {
return nil, credentialcrypto.ErrUnavailable
}
var configuration v1.ModelConfigurationInput
if json.Unmarshal(raw, &configuration) != nil {
return nil, credentialcrypto.ErrUnavailable
}
configuration := bundle.Configuration
// A bundle that opens but no longer validates is not a credential failure.
// Its stored row stays intact so the operator can inspect and replace it.
if err := configuration.ValidateHarness(harness); err != nil {
return nil, err
}
return &Snapshot{Provider: &configuration.ModelProvider, Model: configuration.Model, HarnessConfig: v1.ResolvedHarnessConfig(configuration.HarnessConfig), Revision: sealed.Revision}, nil
return &Snapshot{Provider: &configuration.ModelProvider, Model: configuration.Model, HarnessConfig: v1.ResolvedHarnessConfig(configuration.HarnessConfig), Revision: bundle.Revision}, nil
}
Loading