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
5 changes: 3 additions & 2 deletions config_example.yml
Original file line number Diff line number Diff line change
Expand Up @@ -51,12 +51,13 @@ BOOTSTRAP_TOKEN: <PleaseChangeMe>
# 业务请求沿用用户 Cookie(或已有 Authorization),写请求同时携带 CSRF token。
# 用户凭据只按 Run 保存在进程内存中,不写入 Journal 或模型输入。
# 启用前必须确认 Core 已应用 Runtime Store migration,且组件签名可访问
# /api/swagger.json,并提供 x-jms-required-permissions 与 x-jms-permission-dynamic。
# /api/swagger.json,并提供 x-jms-required-permissions、
# x-jms-permission-dynamic 与 x-jms-rbac-protected。
PLATFORM_GATEWAY_ENABLED: true
# PLATFORM_CA_CERT: ""
# PLATFORM_CLIENT_CERT: ""
# PLATFORM_CLIENT_KEY: ""
# PLATFORM_ALLOWED_METHODS: [GET, POST, PUT, PATCH]
# PLATFORM_ALLOWED_METHODS: [GET, POST, PUT, PATCH, DELETE]
# PLATFORM_REGISTRY_TTL: 1h
# PLATFORM_TIMEOUT: 15s
# PLATFORM_MAX_RESPONSE_BYTES: 1048576
1 change: 1 addition & 0 deletions internal/component/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -290,6 +290,7 @@ func (c *Client) OpenAPISchema(ctx context.Context) (map[string]any, error) {
request.Header.Set("Accept", "application/json")
request.Header.Set("Content-Type", "application/json")
request.Header.Set("X-JMS-ORG", "ROOT")
request.Header.Set("X-JMS-AI-Schema", "1")
if err = (&httplib.SigAuth{KeyID: c.accessKeyID, SecretID: c.accessKeySecret}).Sign(request); err != nil {
return nil, fmt.Errorf("load Core OpenAPI schema: sign request: %w", err)
}
Expand Down
2 changes: 1 addition & 1 deletion internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -137,7 +137,7 @@ func setDefaults(v *viper.Viper) {
v.SetDefault("TERMINAL_AI_MAX_BYTES", int64(1<<30))
v.SetDefault("TERMINAL_AI_MIN_FREE_BYTES", int64(1<<30))
v.SetDefault("PLATFORM_GATEWAY_ENABLED", true)
v.SetDefault("PLATFORM_ALLOWED_METHODS", []string{"GET", "POST", "PUT", "PATCH"})
v.SetDefault("PLATFORM_ALLOWED_METHODS", []string{"GET", "POST", "PUT", "PATCH", "DELETE"})
v.SetDefault("PLATFORM_REGISTRY_TTL", "1h")
v.SetDefault("PLATFORM_TIMEOUT", "15s")
v.SetDefault("PLATFORM_MAX_RESPONSE_BYTES", 1024*1024)
Expand Down
314 changes: 286 additions & 28 deletions internal/platformgateway/gateway.go

Large diffs are not rendered by default.

4 changes: 2 additions & 2 deletions internal/policy/profiles.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,10 @@ type Profile struct {

var profiles = map[string]Profile{
"general": {
ID: "general", Version: "7", Name: "JumpServer assistant", Kind: "general", MaxRisk: "dangerous", CoreAPIEnabled: true,
ID: "general", Version: "10", Name: "JumpServer assistant", Kind: "general", MaxRisk: "dangerous", CoreAPIEnabled: true,
AllowedNamespaces: []string{"general"},
Description: "Unified JumpServer assistant for product questions and authorized live Core operations.",
Instructions: "Act as the unified JumpServer assistant. Answer stable product concepts and usage questions directly, but never claim that a current product option or API capability is absent without checking the authorized Core operations and their schemas. For live environment work, immediately search only server-authorized Core operations and call only an operation returned by that search. Follow operation guidance and server-declared, gateway-enforced create-reuse policies returned by capability search. Do not offer to update, delete or perform another mutation unless its authorized operation appeared in the current search results. Resolve related identifiers with authorized read operations before asking the user. Ask one concise question only for irreducible missing values and combine them in that question. Once required values are available, call the write operation immediately: the trusted approval UI provides confirmation, so never ask the user to type confirmation in chat. Never claim execution without a successful tool result.",
Instructions: "Act as the unified JumpServer assistant. Answer stable product concepts and usage questions directly, but never claim that a current product option or API capability is absent without checking the authorized Core operations and their schemas. For live environment work, immediately search RBAC-protected Core operations and call only an operation returned by that search. Static permissions are checked at discovery; Core checks dynamic permissions against the actual request, so a dynamic search result is not proof that execution is allowed. Follow operation guidance and server-declared, gateway-enforced create-reuse policies returned by capability search. Do not offer to update, delete or perform another mutation unless its authorized operation appeared in the current search results. Resolve related identifiers with authorized read operations before asking the user. Do not repeat passwords or secrets supplied in chat, and never place them in tool arguments. For supported write-only credential fields, omit the field from tool arguments; the trusted approval form collects it directly from the user. If a request combines host creation with an account, search for both authorized operations before any write, then create the host without accounts and add its account separately. For account creation, provide account metadata only; the trusted approval form collects the password. If account creation is unavailable to this assistant, describe the capability limit rather than claiming Core lacks the API, and direct the user to JumpServer account management. Ask one concise question only for irreducible missing values and combine them in that question. Once required values are available, call the write operation immediately: the trusted approval UI provides confirmation, so never ask the user to type confirmation in chat. Never claim execution without a successful tool result.",
StarterPrompts: []string{"What can JumpServer help me manage?", "Explain how to use JumpServer safely."},
},
"platform.management": {
Expand Down
6 changes: 6 additions & 0 deletions internal/ports/capability.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,12 @@ type CapabilityRequest struct {
Profile string
Registration domain.Registration
Arguments json.RawMessage
// AccountPassword is supplied by the approval UI only. It must never be
// included in Arguments, previews, tool results, or the runtime store.
AccountPassword string
// SecretInputs are supplied by the approval UI for write-only Core fields.
// They are never part of model arguments, previews, or persisted tool calls.
SecretInputs map[string]string
}

type CapabilityPolicy struct {
Expand Down
54 changes: 54 additions & 0 deletions internal/service/credentials.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,66 @@ package service
import (
"context"
"errors"
"time"

"github.com/jumpserver/kael/internal/domain"
"github.com/jumpserver/kael/internal/identity"
"github.com/jumpserver/kael/internal/ports"
)

// Approval passwords live only in this process until the waiting call takes
// them, and are never written to the journal or audit store.
func (s *Service) keepAccountPassword(approvalID, password string) {
s.accountPasswordsMu.Lock()
if s.accountPasswords == nil {
s.accountPasswords = make(map[string]string)
}
s.accountPasswords[approvalID] = password
s.accountPasswordsMu.Unlock()
time.AfterFunc(10*time.Minute, func() {
s.accountPasswordsMu.Lock()
delete(s.accountPasswords, approvalID)
s.accountPasswordsMu.Unlock()
})
}

func (s *Service) takeAccountPassword(approvalID string) string {
s.accountPasswordsMu.Lock()
defer s.accountPasswordsMu.Unlock()
password := s.accountPasswords[approvalID]
delete(s.accountPasswords, approvalID)
return password
}

func (s *Service) keepSecretInputs(approvalID string, inputs map[string]string) {
if len(inputs) == 0 {
return
}
copyOfInputs := make(map[string]string, len(inputs))
for name, value := range inputs {
copyOfInputs[name] = value
}
s.accountPasswordsMu.Lock()
if s.approvalSecrets == nil {
s.approvalSecrets = make(map[string]map[string]string)
}
s.approvalSecrets[approvalID] = copyOfInputs
s.accountPasswordsMu.Unlock()
time.AfterFunc(10*time.Minute, func() {
s.accountPasswordsMu.Lock()
delete(s.approvalSecrets, approvalID)
s.accountPasswordsMu.Unlock()
})
}

func (s *Service) takeSecretInputs(approvalID string) map[string]string {
s.accountPasswordsMu.Lock()
defer s.accountPasswordsMu.Unlock()
inputs := s.approvalSecrets[approvalID]
delete(s.approvalSecrets, approvalID)
return inputs
}

// The caller holds credentialsMu across queueing and binding so a worker
// cannot execute a committed run before its request credentials are available.
func (s *Service) bindRunCredentials(ctx context.Context, run *domain.Run) {
Expand Down
68 changes: 63 additions & 5 deletions internal/service/execution.go
Original file line number Diff line number Diff line change
Expand Up @@ -915,10 +915,12 @@ func (s *Service) SubmitToolResult(ctx context.Context, principal domain.Princip
}

type ApprovalDecisionRequest struct {
Decision string `json:"decision"`
RunID string `json:"run_id"`
ArgumentsDigest string `json:"arguments_digest"`
Remember bool `json:"remember,omitempty"`
Decision string `json:"decision"`
RunID string `json:"run_id"`
ArgumentsDigest string `json:"arguments_digest"`
Remember bool `json:"remember,omitempty"`
AccountPassword string `json:"account_password,omitempty"`
SecretInputs map[string]string `json:"secret_inputs,omitempty"`
}

func (s *Service) DecideApproval(ctx context.Context, principal domain.Principal, id string, request ApprovalDecisionRequest) (*domain.Approval, bool, error) {
Expand All @@ -928,10 +930,18 @@ func (s *Service) DecideApproval(ctx context.Context, principal domain.Principal
if request.Remember && request.Decision != "approve" {
return nil, false, serviceError(Invalid, "invalid_remembered_decision", "only an approved command can be remembered", nil)
}
digest, _ := domain.HashValue(request)
// The decision digest must not contain even a hash of the password.
digest, _ := domain.HashValue(struct {
Decision string
RunID string
ArgumentsDigest string
Remember bool
}{request.Decision, request.RunID, request.ArgumentsDigest, request.Remember})
var approval *domain.Approval
duplicate := false
expired := false
passwordStored := false
secretsStored := false
var notify []string
err := s.store.Transaction(ctx, func(tx ports.Tx) error {
now := time.Now().UTC()
Expand Down Expand Up @@ -964,6 +974,40 @@ func (s *Service) DecideApproval(ctx context.Context, principal domain.Principal
if request.RunID != "" && request.RunID != approval.RunID || request.ArgumentsDigest != "" && request.ArgumentsDigest != approval.ArgumentsDigest {
return serviceError(Forbidden, "approval_binding_mismatch", "approval binding is invalid", nil)
}
var preview struct {
AccountPasswordRequired bool `json:"account_password_required"`
SecretInputFields []struct {
Name string `json:"name"`
Required bool `json:"required"`
} `json:"secret_input_fields"`
}
if len(approval.Preview) > 0 {
if err := json.Unmarshal(approval.Preview, &preview); err != nil {
return serviceError(Invalid, "invalid_approval_preview", "approval preview is invalid", nil)
}
}
if preview.AccountPasswordRequired {
if request.Decision == "approve" && (request.AccountPassword == "" || len(request.AccountPassword) > 40960) {
return serviceError(Invalid, "account_password_required", "enter an account password in the approval form", nil)
}
} else if request.AccountPassword != "" {
return serviceError(Invalid, "account_password_unexpected", "this operation does not accept an account password", nil)
}
allowedSecrets := make(map[string]bool, len(preview.SecretInputFields))
for _, field := range preview.SecretInputFields {
allowedSecrets[field.Name] = true
if request.Decision == "approve" && field.Required && request.SecretInputs[field.Name] == "" {
return serviceError(Invalid, "secret_input_required", "enter the required secret in the approval form", nil)
}
}
if len(request.SecretInputs) > 20 {
return serviceError(Invalid, "secret_input_invalid", "too many approval secret inputs", nil)
}
for name, value := range request.SecretInputs {
if request.Decision != "approve" || !allowedSecrets[name] || len(value) > 40960 {
return serviceError(Invalid, "secret_input_invalid", "approval secret input is invalid", nil)
}
}
panel, err := tx.Panel(approval.PanelSessionID, principal, true)
if err != nil {
return err
Expand All @@ -980,6 +1024,14 @@ func (s *Service) DecideApproval(ctx context.Context, principal domain.Principal
return serviceError(Forbidden, "remembered_approval_forbidden", "this approval cannot be remembered", nil)
}
}
if preview.AccountPasswordRequired && request.Decision == "approve" {
s.keepAccountPassword(id, request.AccountPassword)
passwordStored = true
}
if len(request.SecretInputs) > 0 && request.Decision == "approve" {
s.keepSecretInputs(id, request.SecretInputs)
secretsStored = true
}
approval.State = "approved"
if request.Decision == "reject" {
approval.State = "rejected"
Expand Down Expand Up @@ -1019,6 +1071,12 @@ func (s *Service) DecideApproval(ctx context.Context, principal domain.Principal
return s.audit(tx, principal, "approval.decided", approval.ConversationID, approval.PanelSessionID, approval.RunID, map[string]any{"approval_id": approval.ID, "decision": request.Decision, "remembered": approval.Remembered})
})
if err != nil {
if passwordStored {
s.takeAccountPassword(id)
}
if secretsStored {
s.takeSecretInputs(id)
}
return nil, false, translateOrService(err)
}
if expired {
Expand Down
59 changes: 31 additions & 28 deletions internal/service/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -89,34 +89,37 @@ type Options struct {
}

type Service struct {
store ports.Store
engine agentruntime.Engine
bus *event.Bus
logger *zap.Logger
instanceID string
workers int
runTimeout time.Duration
toolResultTimeout time.Duration
panelLease time.Duration
registrationLease time.Duration
eventRetention time.Duration
artifactDir string
maxArtifactBytes int64
capability ports.CapabilityProvider
storageKind string
storageDurable bool
wake chan struct{}
stop chan struct{}
done chan struct{}
startOnce sync.Once
stopOnce sync.Once
lifecycleMu sync.Mutex
started bool
startErr error
activeMu sync.Mutex
active map[string]context.CancelFunc
credentialsMu sync.Mutex
runCredentials map[string]identity.CoreCredentials
store ports.Store
engine agentruntime.Engine
bus *event.Bus
logger *zap.Logger
instanceID string
workers int
runTimeout time.Duration
toolResultTimeout time.Duration
panelLease time.Duration
registrationLease time.Duration
eventRetention time.Duration
artifactDir string
maxArtifactBytes int64
capability ports.CapabilityProvider
storageKind string
storageDurable bool
wake chan struct{}
stop chan struct{}
done chan struct{}
startOnce sync.Once
stopOnce sync.Once
lifecycleMu sync.Mutex
started bool
startErr error
activeMu sync.Mutex
active map[string]context.CancelFunc
credentialsMu sync.Mutex
runCredentials map[string]identity.CoreCredentials
accountPasswordsMu sync.Mutex
accountPasswords map[string]string
approvalSecrets map[string]map[string]string
}

func New(options Options) (*Service, error) {
Expand Down
2 changes: 2 additions & 0 deletions internal/service/service_capability.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,8 @@ func (s *Service) callServiceCapability(ctx context.Context, run *domain.Run, re
}
return agentruntime.ToolObservation{}, err
}
request.AccountPassword = s.takeAccountPassword(approval.ID)
request.SecretInputs = s.takeSecretInputs(approval.ID)
}
if err = s.startServiceCapability(ctx, run, call, approval); err != nil {
return agentruntime.ToolObservation{}, err
Expand Down
Loading