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
38 changes: 38 additions & 0 deletions services/core/internal/api/errors_sessions.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,50 @@ import (
"errors"
"net/http"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/agents"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmentconfig"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmenttemplates"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/projects"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/vaults"
)

// configurationError marks an error from resolving the model, Harness and
// execution selection of a Session request. Its message is the 400 answer
// unless the error is a selection, stored-data or field error, which map first.
type configurationError struct{ err error }

func (e *configurationError) Error() string { return e.err.Error() }
func (e *configurationError) Unwrap() error { return e.err }

// writeSessionCreationError maps a Session creation failure. Template, saved
// Agent (including its provider bundle) and Vault reads map by their domain.
func writeSessionCreationError(w http.ResponseWriter, r *http.Request, err error) {
var required *modelProviderRequiredError
var configuration *configurationError
var provider *v1.ModelProviderError
var credentials *vaults.MCPCredentialSelectionError
switch {
case errors.As(err, &required):
writeError(w, http.StatusBadRequest, "model_provider_required", required.message, "x_agents_core.model_provider")
case writeSelectionError(w, err):
case writeStoredDataError(w, r, err):
case writeFieldError(w, err):
case errors.As(err, &configuration):
writeError(w, http.StatusBadRequest, "unsupported_or_invalid_configuration", err.Error())
case errors.Is(err, environmenttemplates.ErrNotFound), errors.Is(err, environmenttemplates.ErrInvalidInput):
writeEnvironmentTemplatesError(w, r, err)
case errors.Is(err, agents.ErrNotFound), errors.Is(err, agents.ErrInvalidInput), errors.As(err, &provider):
writeAgentsError(w, r, err)
case errors.As(err, &credentials), errors.Is(err, vaults.ErrNotFound), errors.Is(err, vaults.ErrInvalidInput):
writeVaultsError(w, r, err)
default:
writeOperationError(w, r, err)
}
}

// writeInputError reports Session input admission failures. Input that the
// Session cannot accept in its current state is the official conflict_error;
// Idempotency-Key reuse keeps Core's local idempotency_conflict code.
Expand Down
114 changes: 15 additions & 99 deletions services/core/internal/api/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,14 +4,11 @@ import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"reflect"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
"github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/identity"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/go-chi/chi/v5"
Expand Down Expand Up @@ -169,106 +166,25 @@ func (h *Handler) createSession(w http.ResponseWriter, r *http.Request) {
writeError(w, http.StatusBadRequest, "invalid_request_error", "streaming session creation requires initial input")
return
}
creationRequest, err := sessionCreationRequest(input, initialInputs)
if err != nil {
writeSessionsError(w, r, err)
return
}
if h.recoverSessionCreation(w, r, key, creationRequest, input.Stream) {
return
}
if input.templateID != "" {
template, err := h.EnvironmentTemplatesReader.Resolve(r.Context(), tenantID(r), input.templateID)
if err != nil {
if !h.recoverSessionCreation(w, r, key, creationRequest, input.Stream) {
writeEnvironmentTemplatesError(w, r, err)
}
return
}
if err := applyTemplateEnvironment(&input, template); err != nil {
if !h.recoverSessionCreation(w, r, key, creationRequest, input.Stream) && !writeFieldError(w, err) {
writeSessionsError(w, r, err)
}
return
creation, replayed, err := h.createSessionFrom(r.Context(), tenantID(r), sessionCreator(r), key, input, initialInputs)
switch {
case err != nil:
writeSessionCreationError(w, r, err)
case input.Stream:
// A replayed creation sends no events.
if !replayed || h.auditSessionOperation(w, r, creation.Session.ID, "create") {
h.respondSessionCreationStream(w, r, creation)
}
}
saved, inheritedProvider, err := h.sessionAgentDefaults(r.Context(), tenantID(r), input)
if err != nil {
if !h.recoverSessionCreation(w, r, key, creationRequest, input.Stream) {
writeAgentsError(w, r, err)
}
return
}
err = h.prepareSessionModelConfiguration(r.Context(), &input, saved, inheritedProvider)
var configuration json.RawMessage
if err == nil {
configuration, err = resolve(input, tenantID(r), key, saved)
}
if err == nil {
configuration, err = freezeSessionHarnessConfig(configuration, input.resolvedHarnessConfig)
}
if err == nil {
configuration, err = h.bindSessionCredentials(r.Context(), tenantID(r), configuration)
case replayed:
session, err := h.SessionsReader.GetSession(r.Context(), tenantID(r), creation.Session.ID)
if err != nil {
if !h.recoverSessionCreation(w, r, key, creationRequest, input.Stream) {
writeVaultsError(w, r, err)
}
return
}
}
selectedEngine := h.Engine
var provider *v1.ModelProviderInput
var providerSource string
var deploymentRevision uuid.UUID
if err == nil {
selectedEngine, provider, providerSource, deploymentRevision, err = h.resolveSessionExecution(r.Context(), input, inheritedProvider, configuration)
}
if err == nil {
if invalid := execution.ValidateSessionConfiguration(selectedEngine, configuration); invalid != nil {
err = fmt.Errorf("Harness %s does not support the requested Agent/environment configuration: %w", selectedEngine, invalid)
}
}
if err == nil {
err = validateSessionModelConfiguration(selectedEngine, provider, configuration)
}
if err != nil {
if h.recoverSessionCreation(w, r, key, creationRequest, input.Stream) {
return
}
var required *modelProviderRequiredError
switch {
case errors.As(err, &required):
writeError(w, http.StatusBadRequest, "model_provider_required", required.message, "x_agents_core.model_provider")
case writeSelectionError(w, err):
case writeStoredDataError(w, r, err):
case !writeFieldError(w, err):
writeError(w, http.StatusBadRequest, "unsupported_or_invalid_configuration", err.Error())
writeSessionsError(w, r, err)
} else if h.auditSessionOperation(w, r, session.ID, "create") {
h.respondSessionStatus(w, r, session, http.StatusCreated)
}
return
}
executionConfiguration := sessionExecutionProjection(input, saved, inheritedProvider, provider, selectedEngine, configuration)
createInput := sessions.CreateSession{
ExecutionConfiguration: &executionConfiguration,
ModelProvider: provider,
ModelProviderSource: providerSource,
DeploymentProviderRevision: deploymentRevision,
Creator: sessionCreator(r), InitialFiles: input.initialFiles, Initialization: input.initialization,
Engine: selectedEngine, IdempotencyKey: key, Metadata: input.Metadata, Configuration: configuration, InitialInputs: initialInputs, CreationRequest: creationRequest,
}
create := h.SessionCreation.CreateSession
if len(initialInputs) > 0 || input.Environment.Type == "openai_hosted" {
create = h.Execution.SessionAdmission.CreateSession
}
result, err := create(r.Context(), tenantID(r), createInput)
if err != nil {
writeOperationError(w, r, err)
return
}
if input.Stream {
h.respondSessionCreationStream(w, r, result)
return
default:
h.respondSessionStatus(w, r, creation.Session, http.StatusCreated)
}
h.respondSessionStatus(w, r, result.Session, http.StatusCreated)
}

func (h *Handler) getSession(w http.ResponseWriter, r *http.Request) {
Expand Down
117 changes: 117 additions & 0 deletions services/core/internal/api/session_creation.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
package api

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

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

// createSessionFrom creates the Session a decoded request describes or replays
// the one an earlier request with the same key and intent created. A replay
// holds only the Session ID and admits nothing.
func (h *Handler) createSessionFrom(ctx context.Context, tenant string, creator identity.Subject, key string, input sessionRequest, initial []sessions.Input) (creation sessions.Creation, replayed bool, err error) {
intent, err := sessionCreationRequest(input, initial)
if err != nil {
return sessions.Creation{}, false, err
}
if creation, err := h.SessionCreation.FindSessionCreation(ctx, tenant, key, intent, creator); !errors.Is(err, sessions.ErrNotFound) {
return creation, true, err
}
command, err := h.resolveSessionCreation(ctx, tenant, key, input, initial)
if err != nil {
// A concurrent request with the same intent may have committed since the
// first lookup; its Session answers this retry instead of the failure.
if creation, findErr := h.SessionCreation.FindSessionCreation(ctx, tenant, key, intent, creator); !errors.Is(findErr, sessions.ErrNotFound) {
return creation, true, findErr
}
return sessions.Creation{}, false, err
}
command.Creator, command.IdempotencyKey, command.CreationRequest = creator, key, intent
create := h.SessionCreation.CreateSession
if len(initial) > 0 || input.Environment.Type == "openai_hosted" {
create = h.Execution.SessionAdmission.CreateSession
}
creation, err = create(ctx, tenant, command)
return creation, false, err
}

// resolveSessionCreation resolves and validates the Session's configuration.
// Template, saved Agent and credential errors keep their own type; errors of
// the model, Harness and execution selection are a configurationError.
func (h *Handler) resolveSessionCreation(ctx context.Context, tenant, key string, input sessionRequest, initial []sessions.Input) (sessions.CreateSession, error) {
if input.templateID != "" {
template, err := h.EnvironmentTemplatesReader.Resolve(ctx, tenant, input.templateID)
if err != nil {
return sessions.CreateSession{}, err
}
if err := applyTemplateEnvironment(&input, template); err != nil {
return sessions.CreateSession{}, err
}
}
saved, inheritedProvider, err := h.sessionAgentDefaults(ctx, tenant, input)
if err != nil {
return sessions.CreateSession{}, err
}
err = h.prepareSessionModelConfiguration(ctx, &input, saved, inheritedProvider)
var configuration json.RawMessage
if err == nil {
configuration, err = resolve(input, tenant, key, saved)
}
if err == nil {
configuration, err = freezeSessionHarnessConfig(configuration, input.resolvedHarnessConfig)
}
if err != nil {
return sessions.CreateSession{}, &configurationError{err}
}
if configuration, err = h.bindSessionCredentials(ctx, tenant, configuration); err != nil {
return sessions.CreateSession{}, err
}
engine, provider, providerSource, deploymentRevision, err := h.resolveSessionExecution(ctx, input, inheritedProvider, configuration)
if err != nil {
return sessions.CreateSession{}, &configurationError{err}
}
if err := execution.ValidateSessionConfiguration(engine, configuration); err != nil {
return sessions.CreateSession{}, &configurationError{fmt.Errorf("Harness %s does not support the requested Agent/environment configuration: %w", engine, err)}
}
executionConfiguration := sessionExecutionProjection(input, saved, inheritedProvider, provider, engine, configuration)
return sessions.CreateSession{
ExecutionConfiguration: &executionConfiguration,
ModelProvider: provider,
ModelProviderSource: providerSource,
DeploymentProviderRevision: deploymentRevision,
InitialFiles: input.initialFiles, Initialization: input.initialization,
Engine: engine, Metadata: input.Metadata, Configuration: configuration, InitialInputs: initial,
}, nil
}

// sessionCreationRequest records caller intent before mutable sources resolve:
// saved Agents, templates, credentials and deployment defaults. A retry then
// returns the committed Session even after those sources change or go away.
func sessionCreationRequest(input sessionRequest, initial []sessions.Input) (json.RawMessage, error) {
agentID := ""
if input.AgentID != nil {
agentID = *input.AgentID
}
var environment any = input.Environment
if input.templateID != "" || len(input.initialFiles) > 0 || !input.initialization.Empty() || (input.XAgentsCore != nil && len(input.XAgentsCore.Environment) > 0) {
environment = input.originalEnvironment
if len(input.originalEnvironment) == 0 {
environment = input.templateEnvironment
}
}
return json.Marshal(struct {
Execution *v1.SessionExecutionInput `json:"x_agents_core,omitempty"`
AgentID string `json:"agent_id"`
Agent map[string]json.RawMessage `json:"agent,omitempty"`
Environment any `json:"environment"`
Metadata map[string]string `json:"metadata,omitempty"`
VaultIDs []string `json:"vault_ids,omitempty"`
InitialInputs []sessions.Input `json:"initial_inputs,omitempty"`
}{input.XAgentsCore, agentID, input.agentFields, environment, input.Metadata, input.VaultIDs, initial})
}
63 changes: 0 additions & 63 deletions services/core/internal/api/session_creation_identity.go

This file was deleted.

18 changes: 3 additions & 15 deletions services/core/internal/api/session_model_configuration.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,9 +54,9 @@ func (h *Handler) prepareSessionModelConfiguration(ctx context.Context, input *s
if (inherited == nil && !explicitProvider) || needsModel {
input.deploymentDefaults, err = h.ModelProviders.Resolve(ctx, engine)
if err != nil {
var configurationError *v1.ModelProviderError
if errors.As(err, &configurationError) {
return configurationError
var invalid *v1.ModelProviderError
if errors.As(err, &invalid) {
return invalid
}
return &storedDataError{err}
}
Expand Down Expand Up @@ -122,15 +122,3 @@ func freezeSessionHarnessConfig(raw json.RawMessage, native json.RawMessage) (js
cfg.Agent.XAgentsCore.HarnessConfig = v1.ResolvedHarnessConfig(native)
return json.Marshal(cfg)
}

func validateSessionModelConfiguration(engine string, provider *v1.ModelProviderInput, raw json.RawMessage) error {
var cfg configuration
if err := json.Unmarshal(raw, &cfg); err != nil {
return err
}
var native json.RawMessage
if cfg.Agent.XAgentsCore != nil {
native = cfg.Agent.XAgentsCore.HarnessConfig
}
return provider.ValidateConfiguration(engine, native)
}
Loading