From 9053c53e8e95f09f9b0af15ec13512febc7491c1 Mon Sep 17 00:00:00 2001 From: actiontech-zihan Date: Wed, 29 Jul 2026 13:01:35 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20DMS=E2=86=94ODC=20=E7=94=A8=E6=88=B7?= =?UTF-8?q?=E5=90=8C=E6=AD=A5=E8=87=AA=E6=84=88=EF=BC=88=E7=AD=96=E7=95=A5?= =?UTF-8?q?=20A+B=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 进台缓存未命中时先查后建并冲突兜底;删除用户时清理工作台缓存并对 ODC 用户禁用。 关联 https://github.com/actiontech/dms-ee/issues/948 Co-authored-by: Cursor --- internal/apiserver/service/router.go | 3 + internal/dms/biz/user.go | 20 ++ internal/dms/storage/sql_workbench.go | 18 ++ .../client/sql_workbench_client.go | 296 +++++++++++++++++- .../service/sql_workbench_service.go | 214 ++++++++++++- 5 files changed, 533 insertions(+), 18 deletions(-) diff --git a/internal/apiserver/service/router.go b/internal/apiserver/service/router.go index 9b1617b5d..e1b1573ff 100644 --- a/internal/apiserver/service/router.go +++ b/internal/apiserver/service/router.go @@ -532,6 +532,9 @@ func (s *APIServer) installController() error { } s.SqlWorkbenchController.SqlWorkbenchService.SetSqlResultMasker(masker) + // 策略 B:DMS 删用户时清理 SqlWorkbench 缓存/会话并禁用 ODC 对应用户 + s.DMSController.DMS.UserUsecase.SetSqlWorkbenchLifecycle(s.SqlWorkbenchController.SqlWorkbenchService) + // s.AuthController.RegisterPlugin(s.DMSController.GetRegisterPluginFn()) return nil } diff --git a/internal/dms/biz/user.go b/internal/dms/biz/user.go index e05be0d53..035750787 100644 --- a/internal/dms/biz/user.go +++ b/internal/dms/biz/user.go @@ -217,6 +217,7 @@ type SqlWorkbenchUser struct { type SqlWorkbenchUserRepo interface { GetSqlWorkbenchUserByDMSUserID(ctx context.Context, dmsUserID string) (*SqlWorkbenchUser, bool, error) SaveSqlWorkbenchUserCache(ctx context.Context, user *SqlWorkbenchUser) error + DeleteSqlWorkbenchUserCache(ctx context.Context, dmsUserID string) error } // SqlWorkbenchDatasource SqlWorkbench数据源缓存 @@ -233,9 +234,15 @@ type SqlWorkbenchDatasourceRepo interface { GetSqlWorkbenchDatasourceByDMSDBServiceID(ctx context.Context, dmsDBServiceID, dmsUserID, purpose string) (*SqlWorkbenchDatasource, bool, error) SaveSqlWorkbenchDatasourceCache(ctx context.Context, datasource *SqlWorkbenchDatasource) error DeleteSqlWorkbenchDatasourceCache(ctx context.Context, dmsDBServiceID, dmsUserID, purpose string) error + DeleteSqlWorkbenchDatasourceCachesByUserID(ctx context.Context, dmsUserID string) error GetSqlWorkbenchDatasourcesByUserID(ctx context.Context, dmsUserID string) ([]*SqlWorkbenchDatasource, error) } +// SqlWorkbenchLifecycle 删除 DMS 用户时的 SqlWorkbench/ODC 生命周期清理(策略 B) +type SqlWorkbenchLifecycle interface { + CleanupOnDMSUserDelete(ctx context.Context, dmsUserID, dmsUserName string) error +} + type UserUsecase struct { tx TransactionGenerator repo UserRepo @@ -247,9 +254,15 @@ type UserUsecase struct { ldapConfigurationUsecase *LDAPConfigurationUsecase cloudBeaverRepo CloudbeaverRepo gatewayUsecase *GatewayUsecase + sqlWorkbenchLifecycle SqlWorkbenchLifecycle log *utilLog.Helper } +// SetSqlWorkbenchLifecycle 注入策略 B 清理能力(由 apiserver 在 SqlWorkbenchService 就绪后接线) +func (d *UserUsecase) SetSqlWorkbenchLifecycle(lifecycle SqlWorkbenchLifecycle) { + d.sqlWorkbenchLifecycle = lifecycle +} + func NewUserUsecase(log utilLog.Logger, tx TransactionGenerator, repo UserRepo, userGroupRepo UserGroupRepo, pluginUsecase *PluginUsecase, opPermissionUsecase *OpPermissionUsecase, OpPermissionVerifyUsecase *OpPermissionVerifyUsecase, loginConfigurationUsecase *LoginConfigurationUsecase, ldapConfigurationUsecase *LDAPConfigurationUsecase, cloudBeaverRepo CloudbeaverRepo, gatewayUsecase *GatewayUsecase, ) *UserUsecase { @@ -719,6 +732,13 @@ func (d *UserUsecase) DelUser(ctx context.Context, currentUserUid, UserUid strin return fmt.Errorf("delete cloudbeaver cache failed: %v", err) } + // 策略 B:清 SqlWorkbench 缓存 + 会话,并对 ODC 对应用户禁用(ODC 失败不阻断) + if d.sqlWorkbenchLifecycle != nil { + if err := d.sqlWorkbenchLifecycle.CleanupOnDMSUserDelete(tx, UserUid, ds.Name); err != nil { + return fmt.Errorf("delete sql workbench cache failed: %v", err) + } + } + if err := d.repo.DelUser(tx, UserUid); nil != err { return fmt.Errorf("delete user error: %v", err) } diff --git a/internal/dms/storage/sql_workbench.go b/internal/dms/storage/sql_workbench.go index 613163ac4..49ffd6a3a 100644 --- a/internal/dms/storage/sql_workbench.go +++ b/internal/dms/storage/sql_workbench.go @@ -45,6 +45,15 @@ func (sr *SqlWorkbenchRepo) SaveSqlWorkbenchUserCache(ctx context.Context, user }) } +func (sr *SqlWorkbenchRepo) DeleteSqlWorkbenchUserCache(ctx context.Context, dmsUserID string) error { + return transaction(sr.log, ctx, sr.db, func(tx *gorm.DB) error { + if err := tx.WithContext(ctx).Where("dms_user_id = ?", dmsUserID).Delete(&model.SqlWorkbenchUserCache{}).Error; err != nil { + return fmt.Errorf("failed to delete sql workbench user cache: %v", err) + } + return nil + }) +} + func convertModelSqlWorkbenchUser(user *model.SqlWorkbenchUserCache) *biz.SqlWorkbenchUser { return &biz.SqlWorkbenchUser{ DMSUserID: user.DMSUserID, @@ -110,6 +119,15 @@ func (sr *SqlWorkbenchDatasourceRepo) DeleteSqlWorkbenchDatasourceCache(ctx cont }) } +func (sr *SqlWorkbenchDatasourceRepo) DeleteSqlWorkbenchDatasourceCachesByUserID(ctx context.Context, dmsUserID string) error { + return transaction(sr.log, ctx, sr.db, func(tx *gorm.DB) error { + if err := tx.WithContext(ctx).Where("dms_user_id = ?", dmsUserID).Delete(&model.SqlWorkbenchDatasourceCache{}).Error; err != nil { + return fmt.Errorf("failed to delete sql workbench datasource caches by user id: %v", err) + } + return nil + }) +} + func (sr *SqlWorkbenchDatasourceRepo) GetSqlWorkbenchDatasourcesByUserID(ctx context.Context, dmsUserID string) ([]*biz.SqlWorkbenchDatasource, error) { var datasources []model.SqlWorkbenchDatasourceCache err := transaction(sr.log, ctx, sr.db, func(tx *gorm.DB) error { diff --git a/internal/sql_workbench/client/sql_workbench_client.go b/internal/sql_workbench/client/sql_workbench_client.go index 399d2e923..356a49a9e 100644 --- a/internal/sql_workbench/client/sql_workbench_client.go +++ b/internal/sql_workbench/client/sql_workbench_client.go @@ -7,10 +7,12 @@ import ( "crypto/x509" "encoding/base64" "encoding/json" + "errors" "fmt" "io" "mime/multipart" "net/http" + "net/url" "regexp" "strings" "time" @@ -612,14 +614,19 @@ func (c *SqlWorkbenchClient) CreateUsers(users []CreateUserRequest, publicKey st return nil, fmt.Errorf("failed to read create users response: %v", err) } + var createUsersResp CreateUsersResponse + _ = json.Unmarshal(body, &createUsersResp) + // 检查HTTP状态码 if resp.StatusCode != http.StatusOK { c.log.Errorf("Create users failed with status code: %d, response: %s", resp.StatusCode, string(body)) + if isCreateUsersDuplicatedResponse(resp.StatusCode, body, &createUsersResp) { + return nil, newCreateUsersConflictError(resp.StatusCode, &createUsersResp, body) + } return nil, fmt.Errorf("create users failed with status code: %d", resp.StatusCode) } - // 解析响应 - var createUsersResp CreateUsersResponse + // 解析响应(HTTP 200 时要求结构合法) if err := json.Unmarshal(body, &createUsersResp); err != nil { c.log.Errorf("Failed to parse create users response: %v", err) return nil, fmt.Errorf("failed to parse create users response: %v", err) @@ -632,6 +639,9 @@ func (c *SqlWorkbenchClient) CreateUsers(users []CreateUserRequest, publicKey st errorMsg = *createUsersResp.Message } c.log.Errorf("Create users failed: %s", errorMsg) + if isCreateUsersDuplicatedResponse(resp.StatusCode, body, &createUsersResp) { + return nil, newCreateUsersConflictError(resp.StatusCode, &createUsersResp, body) + } return nil, fmt.Errorf("create users failed: %s", errorMsg) } @@ -639,6 +649,231 @@ func (c *SqlWorkbenchClient) CreateUsers(users []CreateUserRequest, publicKey st return &createUsersResp, nil } +// CreateUsersConflictError 表示 ODC 侧账号已存在冲突,调用方可走再查绑定兜底 +type CreateUsersConflictError struct { + StatusCode int + Code string + Message string +} + +func (e *CreateUsersConflictError) Error() string { + if e.Message != "" { + return fmt.Sprintf("create users conflict: %s", e.Message) + } + if e.Code != "" { + return fmt.Sprintf("create users conflict: code=%s status=%d", e.Code, e.StatusCode) + } + return fmt.Sprintf("create users conflict: status=%d", e.StatusCode) +} + +// IsCreateUsersConflict 判断错误是否为账号已存在冲突 +func IsCreateUsersConflict(err error) bool { + var conflict *CreateUsersConflictError + return errors.As(err, &conflict) +} + +func newCreateUsersConflictError(statusCode int, resp *CreateUsersResponse, body []byte) *CreateUsersConflictError { + err := &CreateUsersConflictError{StatusCode: statusCode} + if resp != nil { + if resp.Code != nil { + err.Code = *resp.Code + } + if resp.Message != nil { + err.Message = *resp.Message + } + } + if err.Message == "" { + err.Message = string(body) + } + return err +} + +func isCreateUsersDuplicatedResponse(statusCode int, body []byte, resp *CreateUsersResponse) bool { + bodyStr := string(body) + if resp != nil && resp.Code != nil && strings.Contains(*resp.Code, "DuplicatedExists") { + return true + } + if strings.Contains(bodyStr, "DuplicatedExists") { + return true + } + if resp != nil && resp.Message != nil { + msg := *resp.Message + if strings.Contains(msg, "已存在") || strings.Contains(strings.ToLower(msg), "already exist") { + return true + } + } + if strings.Contains(bodyStr, "已存在") { + return true + } + _ = statusCode + return false +} + +// ListUsers 按 accountName 查询 ODC 用户列表 +func (c *SqlWorkbenchClient) ListUsers(accountName string, cookie string) (*ListUsersResponse, error) { + c.log.Infof("Attempting to list users by accountName: %s", accountName) + + listUsersURL := fmt.Sprintf("%s/api/v2/iam/users?accountName=%s¤tOrganizationId=1&size=100", + c.baseURL, url.QueryEscape(accountName)) + + req, err := http.NewRequest("GET", listUsersURL, nil) + if err != nil { + c.log.Errorf("Failed to create list users request: %v", err) + return nil, fmt.Errorf("failed to create list users request: %v", err) + } + + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Cookie", cookie) + req.Header.Set("X-Xsrf-Token", c.ExtractCookieValue(cookie, "XSRF-TOKEN")) + + resp, err := c.httpClient.Do(req) + if err != nil { + c.log.Errorf("Failed to send list users request: %v", err) + return nil, fmt.Errorf("failed to send list users request: %v", err) + } + defer resp.Body.Close() + + body, err := io.ReadAll(resp.Body) + if err != nil { + c.log.Errorf("Failed to read list users response: %v", err) + return nil, fmt.Errorf("failed to read list users response: %v", err) + } + + if resp.StatusCode != http.StatusOK { + c.log.Errorf("List users failed with status code: %d, response: %s", resp.StatusCode, string(body)) + return nil, fmt.Errorf("list users failed with status code: %d", resp.StatusCode) + } + + var listUsersResp ListUsersResponse + if err := json.Unmarshal(body, &listUsersResp); err != nil { + c.log.Errorf("Failed to parse list users response: %v", err) + return nil, fmt.Errorf("failed to parse list users response: %v", err) + } + + if !listUsersResp.Successful { + errorMsg := "list users failed" + if listUsersResp.Message != nil { + errorMsg = *listUsersResp.Message + } + c.log.Errorf("List users failed: %s", errorMsg) + return nil, fmt.Errorf("list users failed: %s", errorMsg) + } + + c.log.Infof("Successfully listed %d users for accountName %s", len(listUsersResp.Data.Contents), accountName) + return &listUsersResp, nil +} + +// SetUserEnabled 启用或禁用 ODC 用户 +func (c *SqlWorkbenchClient) SetUserEnabled(userID int64, enabled bool, cookie string) (*SetUserEnabledResponse, error) { + c.log.Infof("Attempting to set user enabled=%v for id=%d", enabled, userID) + + setEnabledURL := fmt.Sprintf("%s/api/v2/iam/users/%d/setEnabled?currentOrganizationId=1", c.baseURL, userID) + payload := SetEnabledRequest{Enabled: enabled} + jsonData, err := json.Marshal(payload) + if err != nil { + return nil, fmt.Errorf("failed to marshal setEnabled request: %v", err) + } + + req, err := http.NewRequest("POST", setEnabledURL, bytes.NewBuffer(jsonData)) + if err != nil { + return nil, fmt.Errorf("failed to create setEnabled request: %v", err) + } + + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Cookie", cookie) + req.Header.Set("X-Xsrf-Token", c.ExtractCookieValue(cookie, "XSRF-TOKEN")) + + resp, err := c.httpClient.Do(req) + if err != nil { + c.log.Errorf("Failed to send setEnabled request: %v", err) + return nil, fmt.Errorf("failed to send setEnabled request: %v", err) + } + defer resp.Body.Close() + + body, err := io.ReadAll(resp.Body) + if err != nil { + return nil, fmt.Errorf("failed to read setEnabled response: %v", err) + } + + if resp.StatusCode != http.StatusOK { + c.log.Errorf("setEnabled failed with status code: %d, response: %s", resp.StatusCode, string(body)) + return nil, fmt.Errorf("setEnabled failed with status code: %d", resp.StatusCode) + } + + var setEnabledResp SetUserEnabledResponse + if err := json.Unmarshal(body, &setEnabledResp); err != nil { + return nil, fmt.Errorf("failed to parse setEnabled response: %v", err) + } + if !setEnabledResp.Successful { + errorMsg := "setEnabled failed" + if setEnabledResp.Message != nil { + errorMsg = *setEnabledResp.Message + } + return nil, fmt.Errorf("%s", errorMsg) + } + + c.log.Infof("Successfully set user enabled=%v for id=%d", enabled, userID) + return &setEnabledResp, nil +} + +// ResetUserPassword 管理员重置用户密码(newPassword 明文,内部 RSA 加密) +func (c *SqlWorkbenchClient) ResetUserPassword(userID int64, newPassword, publicKey, cookie string) (*ResetUserPasswordResponse, error) { + c.log.Infof("Attempting to reset password for user id=%d", userID) + + encryptedPassword, err := c.EncryptPasswordWithRSA(newPassword, publicKey) + if err != nil { + return nil, fmt.Errorf("failed to encrypt new password: %v", err) + } + + resetURL := fmt.Sprintf("%s/api/v2/iam/users/resetPassword?id=%d¤tOrganizationId=1", c.baseURL, userID) + payload := ResetPasswordRequest{NewPassword: encryptedPassword} + jsonData, err := json.Marshal(payload) + if err != nil { + return nil, fmt.Errorf("failed to marshal resetPassword request: %v", err) + } + + req, err := http.NewRequest("POST", resetURL, bytes.NewBuffer(jsonData)) + if err != nil { + return nil, fmt.Errorf("failed to create resetPassword request: %v", err) + } + + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Cookie", cookie) + req.Header.Set("X-Xsrf-Token", c.ExtractCookieValue(cookie, "XSRF-TOKEN")) + + resp, err := c.httpClient.Do(req) + if err != nil { + c.log.Errorf("Failed to send resetPassword request: %v", err) + return nil, fmt.Errorf("failed to send resetPassword request: %v", err) + } + defer resp.Body.Close() + + body, err := io.ReadAll(resp.Body) + if err != nil { + return nil, fmt.Errorf("failed to read resetPassword response: %v", err) + } + + if resp.StatusCode != http.StatusOK { + c.log.Errorf("resetPassword failed with status code: %d, response: %s", resp.StatusCode, string(body)) + return nil, fmt.Errorf("resetPassword failed with status code: %d", resp.StatusCode) + } + + var resetResp ResetUserPasswordResponse + if err := json.Unmarshal(body, &resetResp); err != nil { + return nil, fmt.Errorf("failed to parse resetPassword response: %v", err) + } + if !resetResp.Successful { + errorMsg := "resetPassword failed" + if resetResp.Message != nil { + errorMsg = *resetResp.Message + } + return nil, fmt.Errorf("%s", errorMsg) + } + + c.log.Infof("Successfully reset password for user id=%d", userID) + return &resetResp, nil +} + // ActivateUser 激活用户 func (c *SqlWorkbenchClient) ActivateUser(username, currentPassword, newPassword, publicKey, cookie string) (*ActivateUserResponse, error) { c.log.Infof("Attempting to activate user: %s", username) @@ -921,6 +1156,63 @@ type CreateUsersResponse struct { Error interface{} `json:"error,omitempty"` } +// ListUsersResponse 查询用户列表响应 +type ListUsersResponse struct { + Data struct { + Contents []User `json:"contents"` + } `json:"data"` + DurationMillis int64 `json:"durationMillis"` + HTTPStatus string `json:"httpStatus"` + RequestID string `json:"requestId"` + Server string `json:"server"` + Successful bool `json:"successful"` + Timestamp float64 `json:"timestamp"` + TraceID string `json:"traceId"` + Code *string `json:"code,omitempty"` + Message *string `json:"message,omitempty"` + Error interface{} `json:"error,omitempty"` +} + +// SetEnabledRequest 启用/禁用用户请求 +type SetEnabledRequest struct { + Enabled bool `json:"enabled"` +} + +// SetUserEnabledResponse 启用/禁用用户响应 +type SetUserEnabledResponse struct { + Data User `json:"data"` + DurationMillis int64 `json:"durationMillis"` + HTTPStatus string `json:"httpStatus"` + RequestID string `json:"requestId"` + Server string `json:"server"` + Successful bool `json:"successful"` + Timestamp float64 `json:"timestamp"` + TraceID string `json:"traceId"` + Code *string `json:"code,omitempty"` + Message *string `json:"message,omitempty"` + Error interface{} `json:"error,omitempty"` +} + +// ResetPasswordRequest 管理员重置密码请求 +type ResetPasswordRequest struct { + NewPassword string `json:"newPassword"` +} + +// ResetUserPasswordResponse 管理员重置密码响应 +type ResetUserPasswordResponse struct { + Data User `json:"data"` + DurationMillis int64 `json:"durationMillis"` + HTTPStatus string `json:"httpStatus"` + RequestID string `json:"requestId"` + Server string `json:"server"` + Successful bool `json:"successful"` + Timestamp float64 `json:"timestamp"` + TraceID string `json:"traceId"` + Code *string `json:"code,omitempty"` + Message *string `json:"message,omitempty"` + Error interface{} `json:"error,omitempty"` +} + // Organization 组织结构 type Organization struct { Builtin bool `json:"builtin"` diff --git a/internal/sql_workbench/service/sql_workbench_service.go b/internal/sql_workbench/service/sql_workbench_service.go index a4176ff1b..ee177152f 100644 --- a/internal/sql_workbench/service/sql_workbench_service.go +++ b/internal/sql_workbench/service/sql_workbench_service.go @@ -62,8 +62,16 @@ type odcSession struct { var ( dmsUserIdODCSessionMap = make(map[string]odcSession) odcSessionMutex = &sync.Mutex{} + ensureUserLocks sync.Map // accountName -> *sync.Mutex,串行化同账号 ensure,避免并发重置密码 ) +func lockEnsureSqlWorkbenchUser(accountName string) func() { + v, _ := ensureUserLocks.LoadOrStore(accountName, &sync.Mutex{}) + mu := v.(*sync.Mutex) + mu.Lock() + return mu.Unlock +} + // generateSqlWorkbenchUsername 生成 SQL Workbench 用户名 func (s *SqlWorkbenchService) generateSqlWorkbenchUsername(dmsUserName string) string { return SQL_WORKBENCH_PREFIX + dmsUserName @@ -283,14 +291,14 @@ func (sqlWorkbenchService *SqlWorkbenchService) Login() echo.MiddlewareFunc { return err } - // 3. 如果用户不存在,调用sqlworkbench创建用户接口进行创建 + // 3. 缓存未命中:策略 A — 先查 ODC 再创建,冲突再查绑定 if !exists { - err = sqlWorkbenchService.createSqlWorkbenchUser(c.Request().Context(), user) + err = sqlWorkbenchService.ensureSqlWorkbenchUser(c.Request().Context(), user) if err != nil { - sqlWorkbenchService.log.Errorf("Failed to create sql workbench user: %v", err) + sqlWorkbenchService.log.Errorf("Failed to ensure sql workbench user: %v", err) return err } - // 重新获取创建后的用户信息 + // 重新获取绑定/创建后的用户信息 sqlWorkbenchUser, _, err = sqlWorkbenchService.sqlWorkbenchUserRepo.GetSqlWorkbenchUserByDMSUserID(c.Request().Context(), dmsUserId) if err != nil { sqlWorkbenchService.log.Errorf("Failed to get created sql workbench user: %v", err) @@ -320,18 +328,29 @@ func (sqlWorkbenchService *SqlWorkbenchService) Login() echo.MiddlewareFunc { } } -// createSqlWorkbenchUser 创建SqlWorkbench用户 -func (sqlWorkbenchService *SqlWorkbenchService) createSqlWorkbenchUser(ctx context.Context, dmsUser *biz.User) error { +// ensureSqlWorkbenchUser 策略 A:映射缺失时先查 ODC,命中则复用绑定;无则创建;创建冲突再查绑定 +func (sqlWorkbenchService *SqlWorkbenchService) ensureSqlWorkbenchUser(ctx context.Context, dmsUser *biz.User) error { + accountName := sqlWorkbenchService.generateSqlWorkbenchUsername(dmsUser.Name) + unlock := lockEnsureSqlWorkbenchUser(accountName) + defer unlock() + cookie, _, publicKey, err := sqlWorkbenchService.getUserCookie(sqlWorkbenchService.cfg.AdminUser, sqlWorkbenchService.cfg.AdminPassword) if err != nil { return err } - // 创建用户请求 - sqlWorkbenchUsername := sqlWorkbenchService.generateSqlWorkbenchUsername(dmsUser.Name) + odcUser, err := sqlWorkbenchService.findExactSqlWorkbenchUser(accountName, cookie) + if err != nil { + return err + } + if odcUser != nil { + sqlWorkbenchService.log.Infof("Found existing sql workbench user %s (ID: %d), binding for DMS user %s", accountName, odcUser.ID, dmsUser.Name) + return sqlWorkbenchService.bindExistingSqlWorkbenchUser(ctx, dmsUser, odcUser, publicKey, cookie) + } + createUserReq := []client.CreateUserRequest{ { - AccountName: sqlWorkbenchUsername, + AccountName: accountName, Name: dmsUser.Name, Password: SQL_WORKBENCH_DEFAULT_PASSWORD, Enabled: true, @@ -339,9 +358,22 @@ func (sqlWorkbenchService *SqlWorkbenchService) createSqlWorkbenchUser(ctx conte }, } - // 调用创建用户接口 createUserResp, err := sqlWorkbenchService.client.CreateUsers(createUserReq, publicKey, cookie) if err != nil { + // 冲突(DuplicatedExists)或并发创建写冲突(如 DataAccessError):再查绑定 + sqlWorkbenchService.log.Infof("Create sql workbench user failed for %s: %v; re-list and try bind", accountName, err) + odcUser, listErr := sqlWorkbenchService.findExactSqlWorkbenchUser(accountName, cookie) + if listErr != nil { + return fmt.Errorf("failed to create user in sql workbench: %v (re-list: %v)", err, listErr) + } + if odcUser != nil { + if client.IsCreateUsersConflict(err) { + sqlWorkbenchService.log.Infof("Create conflict for %s, binding existing user", accountName) + } else { + sqlWorkbenchService.log.Infof("Create failed but user %s now exists, binding", accountName) + } + return sqlWorkbenchService.bindExistingSqlWorkbenchUser(ctx, dmsUser, odcUser, publicKey, cookie) + } return fmt.Errorf("failed to create user in sql workbench: %v", err) } @@ -349,9 +381,8 @@ func (sqlWorkbenchService *SqlWorkbenchService) createSqlWorkbenchUser(ctx conte return fmt.Errorf("no user created in sql workbench") } - // 激活用户 activateUserResp, err := sqlWorkbenchService.client.ActivateUser( - sqlWorkbenchService.generateSqlWorkbenchUsername(dmsUser.Name), + accountName, SQL_WORKBENCH_DEFAULT_PASSWORD, SQL_WORKBENCH_REAL_PASSWORD, publicKey, @@ -361,22 +392,99 @@ func (sqlWorkbenchService *SqlWorkbenchService) createSqlWorkbenchUser(ctx conte return fmt.Errorf("failed to activate user in sql workbench: %v", err) } - // 保存用户缓存 sqlWorkbenchUser := &biz.SqlWorkbenchUser{ - SqlWorkbenchUsername: sqlWorkbenchService.generateSqlWorkbenchUsername(dmsUser.Name), + SqlWorkbenchUsername: accountName, DMSUserID: dmsUser.UID, SqlWorkbenchUserId: activateUserResp.Data.ID, } + if err = sqlWorkbenchService.sqlWorkbenchUserRepo.SaveSqlWorkbenchUserCache(ctx, sqlWorkbenchUser); err != nil { + return fmt.Errorf("failed to save sql workbench user cache: %v", err) + } + + sqlWorkbenchService.log.Infof("Successfully created and activated sql workbench user for DMS user %s (ID: %d)", dmsUser.Name, activateUserResp.Data.ID) + return nil +} - err = sqlWorkbenchService.sqlWorkbenchUserRepo.SaveSqlWorkbenchUserCache(ctx, sqlWorkbenchUser) +// findExactSqlWorkbenchUser 按 accountName 精确匹配;0 返回 nil;>1 硬失败 +func (sqlWorkbenchService *SqlWorkbenchService) findExactSqlWorkbenchUser(accountName, cookie string) (*client.User, error) { + listResp, err := sqlWorkbenchService.client.ListUsers(accountName, cookie) if err != nil { + return nil, fmt.Errorf("failed to list sql workbench users: %v", err) + } + + matches := make([]client.User, 0) + for _, u := range listResp.Data.Contents { + if u.AccountName == accountName { + matches = append(matches, u) + } + } + switch len(matches) { + case 0: + return nil, nil + case 1: + user := matches[0] + return &user, nil + default: + return nil, fmt.Errorf("ambiguous sql workbench users for accountName %s: found %d", accountName, len(matches)) + } +} + +// bindExistingSqlWorkbenchUser 复用已有 ODC 账户:ensureUsable 后写本地映射 +func (sqlWorkbenchService *SqlWorkbenchService) bindExistingSqlWorkbenchUser(ctx context.Context, dmsUser *biz.User, odcUser *client.User, publicKey, cookie string) error { + if err := sqlWorkbenchService.ensureSqlWorkbenchUserUsable(odcUser, publicKey, cookie); err != nil { + return err + } + + sqlWorkbenchUser := &biz.SqlWorkbenchUser{ + SqlWorkbenchUsername: sqlWorkbenchService.generateSqlWorkbenchUsername(dmsUser.Name), + DMSUserID: dmsUser.UID, + SqlWorkbenchUserId: odcUser.ID, + } + if err := sqlWorkbenchService.sqlWorkbenchUserRepo.SaveSqlWorkbenchUserCache(ctx, sqlWorkbenchUser); err != nil { return fmt.Errorf("failed to save sql workbench user cache: %v", err) } - sqlWorkbenchService.log.Infof("Successfully created and activated sql workbench user for DMS user %s (ID: %d)", dmsUser.Name, activateUserResp.Data.ID) + sqlWorkbenchService.log.Infof("Successfully bound existing sql workbench user %s (ID: %d) for DMS user %s", sqlWorkbenchUser.SqlWorkbenchUsername, odcUser.ID, dmsUser.Name) return nil } +// ensureSqlWorkbenchUserUsable 启用(若禁用)并强制对齐代管登录密码。 +// ODC admin resetPassword 会将 active 置为 false,故先重置为 DEFAULT 再 Activate(DEFAULT→REAL), +// 最终密码为 SQL_WORKBENCH_REAL_PASSWORD 且账户可用(与新建路径一致)。 +// 若代管密码已可登录则跳过重置,避免并发绑定互相踩密码。 +func (sqlWorkbenchService *SqlWorkbenchService) ensureSqlWorkbenchUserUsable(odcUser *client.User, publicKey, cookie string) error { + if !odcUser.Enabled { + if _, err := sqlWorkbenchService.client.SetUserEnabled(odcUser.ID, true, cookie); err != nil { + return fmt.Errorf("failed to enable sql workbench user: %v", err) + } + } + + if _, err := sqlWorkbenchService.client.Login(odcUser.AccountName, SQL_WORKBENCH_REAL_PASSWORD, publicKey); err == nil { + sqlWorkbenchService.log.Infof("sql workbench user %s already usable with managed password, skip reset", odcUser.AccountName) + return nil + } + + if _, err := sqlWorkbenchService.client.ResetUserPassword(odcUser.ID, SQL_WORKBENCH_DEFAULT_PASSWORD, publicKey, cookie); err != nil { + return fmt.Errorf("failed to reset sql workbench user password: %v", err) + } + + if _, err := sqlWorkbenchService.client.ActivateUser( + odcUser.AccountName, + SQL_WORKBENCH_DEFAULT_PASSWORD, + SQL_WORKBENCH_REAL_PASSWORD, + publicKey, + cookie, + ); err != nil { + return fmt.Errorf("failed to activate sql workbench user after password reset: %v", err) + } + return nil +} + +// createSqlWorkbenchUser 保留为 ensure 新建路径的语义别名(测试/兼容调用) +func (sqlWorkbenchService *SqlWorkbenchService) createSqlWorkbenchUser(ctx context.Context, dmsUser *biz.User) error { + return sqlWorkbenchService.ensureSqlWorkbenchUser(ctx, dmsUser) +} + // loginSqlWorkbenchUser 使用SqlWorkbench用户登录并设置Cookie // 返回 jsessionID 和 xsrfToken,并缓存会话 func (sqlWorkbenchService *SqlWorkbenchService) loginSqlWorkbenchUser(c echo.Context, dmsUser *biz.User, dmsUserId, dmsToken string) (string, string, error) { @@ -462,6 +570,80 @@ func (sqlWorkbenchService *SqlWorkbenchService) clearODCSession(dmsUserId string delete(dmsUserIdODCSessionMap, dmsUserId) } +// CleanupOnDMSUserDelete 策略 B:删除 DMS 用户时清理 SqlWorkbench 缓存/会话并对 ODC 用户禁用(非删除) +// B1–B3 缓存失败返回 error(阻断 DMS 删除);B4 忽略不存在;B5 ODC 禁用失败仅记日志不返回 error。 +func (sqlWorkbenchService *SqlWorkbenchService) CleanupOnDMSUserDelete(ctx context.Context, dmsUserID, dmsUserName string) error { + // B1: 读用户缓存(记下 ODC user id,供 B5) + cachedUser, _, err := sqlWorkbenchService.sqlWorkbenchUserRepo.GetSqlWorkbenchUserByDMSUserID(ctx, dmsUserID) + if err != nil { + return fmt.Errorf("get sql workbench user cache failed: %v", err) + } + + var odcUserID int64 + var odcUsername string + if cachedUser != nil { + odcUserID = cachedUser.SqlWorkbenchUserId + odcUsername = cachedUser.SqlWorkbenchUsername + } + + // B2: 删数据源缓存 + if err := sqlWorkbenchService.sqlWorkbenchDatasourceRepo.DeleteSqlWorkbenchDatasourceCachesByUserID(ctx, dmsUserID); err != nil { + return fmt.Errorf("delete sql workbench datasource caches failed: %v", err) + } + + // B3: 删用户缓存 + if err := sqlWorkbenchService.sqlWorkbenchUserRepo.DeleteSqlWorkbenchUserCache(ctx, dmsUserID); err != nil { + return fmt.Errorf("delete sql workbench user cache failed: %v", err) + } + + // B4: 清内存会话 + sqlWorkbenchService.clearODCSession(dmsUserID) + + // B5: ODC 禁用(失败不阻断) + sqlWorkbenchService.disableODCUserOnDMSDelete(odcUserID, odcUsername, dmsUserName) + + sqlWorkbenchService.log.Infof("SqlWorkbench cleanup on DMS user delete done: dmsUserId=%s name=%s", dmsUserID, dmsUserName) + return nil +} + +// disableODCUserOnDMSDelete 对 ODC 对应用户执行 setEnabled(false);任何失败只记日志 +func (sqlWorkbenchService *SqlWorkbenchService) disableODCUserOnDMSDelete(odcUserID int64, cachedUsername, dmsUserName string) { + if !sqlWorkbenchService.IsConfigured() || sqlWorkbenchService.client == nil { + sqlWorkbenchService.log.Infof("SqlWorkbench not configured, skip ODC disable for DMS user %s", dmsUserName) + return + } + + cookie, _, _, err := sqlWorkbenchService.getUserCookie(sqlWorkbenchService.cfg.AdminUser, sqlWorkbenchService.cfg.AdminPassword) + if err != nil { + sqlWorkbenchService.log.Errorf("ODC disable skipped: admin login failed for DMS user %s: %v", dmsUserName, err) + return + } + + targetID := odcUserID + if targetID == 0 { + accountName := cachedUsername + if accountName == "" { + accountName = sqlWorkbenchService.generateSqlWorkbenchUsername(dmsUserName) + } + odcUser, listErr := sqlWorkbenchService.findExactSqlWorkbenchUser(accountName, cookie) + if listErr != nil { + sqlWorkbenchService.log.Errorf("ODC disable skipped: list user %s failed: %v", accountName, listErr) + return + } + if odcUser == nil { + sqlWorkbenchService.log.Infof("ODC disable no-op: no unique user for %s", accountName) + return + } + targetID = odcUser.ID + } + + if _, err := sqlWorkbenchService.client.SetUserEnabled(targetID, false, cookie); err != nil { + sqlWorkbenchService.log.Errorf("ODC disable failed for user id=%d (DMS user %s): %v", targetID, dmsUserName, err) + return + } + sqlWorkbenchService.log.Infof("ODC user id=%d disabled after DMS user %s delete", targetID, dmsUserName) +} + // validateODCSession 验证 ODC 会话是否有效 // 通过调用 GetOrganizations API 来验证会话 func (sqlWorkbenchService *SqlWorkbenchService) validateODCSession(jsessionID, xsrfToken string) bool {