diff --git a/services/core/internal/api/errors_sessions.go b/services/core/internal/api/errors_sessions.go index cefdf683f..91206a437 100644 --- a/services/core/internal/api/errors_sessions.go +++ b/services/core/internal/api/errors_sessions.go @@ -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. diff --git a/services/core/internal/api/handler.go b/services/core/internal/api/handler.go index a3fcb0508..8fa9b128a 100644 --- a/services/core/internal/api/handler.go +++ b/services/core/internal/api/handler.go @@ -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" @@ -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) { diff --git a/services/core/internal/api/session_creation.go b/services/core/internal/api/session_creation.go new file mode 100644 index 000000000..a5db974f1 --- /dev/null +++ b/services/core/internal/api/session_creation.go @@ -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}) +} diff --git a/services/core/internal/api/session_creation_identity.go b/services/core/internal/api/session_creation_identity.go deleted file mode 100644 index b174220db..000000000 --- a/services/core/internal/api/session_creation_identity.go +++ /dev/null @@ -1,63 +0,0 @@ -package api - -import ( - "encoding/json" - "errors" - "net/http" - - v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" - - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" -) - -// 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}) -} - -func (h *Handler) recoverSessionCreation(w http.ResponseWriter, r *http.Request, key string, request json.RawMessage, stream bool) bool { - result, err := h.SessionCreation.FindSessionCreation(r.Context(), tenantID(r), key, request, sessionCreator(r)) - if errors.Is(err, sessions.ErrNotFound) { - return false - } - if err != nil { - writeSessionsError(w, r, err) - return true - } - if stream { - if !h.auditSessionOperation(w, r, result.Session.ID, "create") { - return true - } - // Recorded-intent lookup finds an existing creation, which sends no events. - h.respondSessionCreationStream(w, r, result) - } else { - session, err := h.SessionsReader.GetSession(r.Context(), tenantID(r), result.Session.ID) - if err != nil { - writeSessionsError(w, r, err) - } else if h.auditSessionOperation(w, r, session.ID, "create") { - h.respondSessionStatus(w, r, session, http.StatusCreated) - } - } - return true -} diff --git a/services/core/internal/api/session_model_configuration.go b/services/core/internal/api/session_model_configuration.go index 1b5019066..de8fa75a3 100644 --- a/services/core/internal/api/session_model_configuration.go +++ b/services/core/internal/api/session_model_configuration.go @@ -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} } @@ -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) -} diff --git a/services/core/internal/api/session_write_audit_test.go b/services/core/internal/api/session_write_audit_test.go index a341ce362..1a36694b4 100644 --- a/services/core/internal/api/session_write_audit_test.go +++ b/services/core/internal/api/session_write_audit_test.go @@ -9,6 +9,7 @@ import ( "strings" "testing" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmenttemplates" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/identity" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/writeaudit" @@ -59,13 +60,18 @@ func TestSessionAuditOnlyRoutesFailClosed(t *testing.T) { routeContext := chi.NewRouteContext() routeContext.URLParams.Add("session_id", session) ctx = context.WithValue(ctx, chi.RouteCtxKey, routeContext) - r := httptest.NewRequest(http.MethodPost, "/v1/agents/sessions", strings.NewReader(`{"events":[]}`)).WithContext(ctx) + body := map[string]string{ + "empty-events": `{"events":[]}`, + "creation-replay": `{"agent":{"model":"test"},"environment":{"type":"none"},"input":"hello"}`, + "stream-replay": `{"agent":{"model":"test"},"environment":{"type":"none"},"input":"hello","stream":true}`, + }[route] + r := httptest.NewRequest(http.MethodPost, "/v1/agents/sessions", strings.NewReader(body)).WithContext(ctx) r.Header.Set("Content-Type", "application/json") w := httptest.NewRecorder() if route == "empty-events" { h.createEvents(w, r) - } else if !h.recoverSessionCreation(w, r, "key", json.RawMessage(`{}`), route == "stream-replay") { - t.Fatal("replay not handled") + } else { + h.createSession(w, r) } wantStatus, wantAction := 201, "create" if route == "empty-events" { @@ -84,3 +90,37 @@ func TestSessionAuditOnlyRoutesFailClosed(t *testing.T) { } } } + +// A retry whose Template is gone answers with the Session that a concurrent +// request with the same key and intent committed after the entry lookup. +func TestSessionCreationReplaysAfterResolutionFails(t *testing.T) { + session, lookups, audits := uuid.NewString(), 0, 0 + h, _, _ := testHandler(t, func(_ *Dependencies, f *testFakes) { + f.environmentTemplatesReader.resolve = func(context.Context, string, string) (environmenttemplates.Resolved, error) { + return environmenttemplates.Resolved{}, environmenttemplates.ErrNotFound + } + f.sessionCreation.findSessionCreation = func(context.Context, string, string, json.RawMessage, identity.Subject) (sessions.Creation, error) { + if lookups++; lookups == 1 { + return sessions.Creation{}, sessions.ErrNotFound + } + return sessions.Creation{Session: sessions.Session{ID: session}}, nil + } + f.sessions.auditSessionOperation = func(context.Context, sessions.AuditSessionOperationCommand) error { + audits++ + return nil + } + f.sessionsReader.getSession = func(_ context.Context, tenant, id string) (sessions.Session, error) { + return sessions.Session{ID: id, TenantID: tenant, Configuration: json.RawMessage(`{"environment":{"type":"none"},"agent":{"id":"agent","model":"test"}}`)}, nil + } + }) + r := httptest.NewRequest(http.MethodPost, "/v1/agents/sessions", strings.NewReader(`{"agent":{"model":"test"},"environment":{"type":"openai_hosted","environment_template_id":"saved"},"input":"Start."}`)) + r.Header.Set("Authorization", "Bearer test-api-key") + r.Header.Set("OpenAI-Beta", "agents=v1") + r.Header.Set("Content-Type", "application/json") + r.Header.Set("Idempotency-Key", "retry") + w := httptest.NewRecorder() + h.ServeHTTP(w, r) + if w.Code != http.StatusCreated || !strings.Contains(w.Body.String(), session) || lookups != 2 || audits != 1 { + t.Fatal(w.Code, w.Body.String(), lookups, audits) + } +}