From 58ca598af7e870a21185947ab06c1acb3ce91b30 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 8 Oct 2026 08:11:45 +0000 Subject: [PATCH] Create Sessions through one use case createSession ran the whole creation flow inside the HTTP handler and re-checked replay on six branches. createSessionFrom now owns the flow without HTTP types: one replay lookup at entry, one more at the single resolution failure point, then admission or plain creation. The handler decodes, calls it, maps errors in one switch and writes the response. A configurationError marks rejections of the resolved configuration so the switch keeps the stage-dependent mapping. validateSessionModelConfiguration repeated the provider check resolveSessionExecution already runs on the same harness_config, so it is deleted. /v1 responses, error bodies, check order, stream events and audit rows are unchanged. --- services/core/internal/api/errors_sessions.go | 38 ++++++ services/core/internal/api/handler.go | 114 +++-------------- .../core/internal/api/session_creation.go | 117 ++++++++++++++++++ .../internal/api/session_creation_identity.go | 63 ---------- .../api/session_model_configuration.go | 18 +-- .../internal/api/session_write_audit_test.go | 46 ++++++- 6 files changed, 216 insertions(+), 180 deletions(-) create mode 100644 services/core/internal/api/session_creation.go delete mode 100644 services/core/internal/api/session_creation_identity.go 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) + } +}