From a934380c9c8fd46dc32a9a125d1aabedc8234335 Mon Sep 17 00:00:00 2001 From: mac Date: Thu, 23 Jul 2026 18:36:24 +0800 Subject: [PATCH] feat: switch base delivery to awx task logs --- frontend/src/api/delivery.ts | 124 +++++ frontend/src/api/deployment.ts | 68 --- frontend/src/api/taskLog.ts | 81 +++ frontend/src/views/service/Catalog.vue | 235 +++++---- frontend/src/views/task/TaskCenter.vue | 336 ++++++++++--- server/.env.example | 5 +- server/internal/config/config.go | 226 +++++---- server/internal/database/database.go | 3 - server/internal/handler/delivery.go | 51 +- server/internal/handler/deployment.go | 249 --------- server/internal/handler/task_log.go | 347 +++++++++++++ server/internal/model/delivery.go | 12 - server/internal/model/models.go | 28 -- server/internal/response/response.go | 4 +- server/internal/router/router.go | 12 +- server/internal/service/awx.go | 270 ++++++++-- server/internal/service/awx_test.go | 36 ++ server/internal/service/container_service.go | 160 ++++++ server/internal/service/delivery.go | 135 ++++- server/internal/service/deployment.go | 499 ------------------- 20 files changed, 1638 insertions(+), 1243 deletions(-) create mode 100644 frontend/src/api/delivery.ts delete mode 100644 frontend/src/api/deployment.ts create mode 100644 frontend/src/api/taskLog.ts delete mode 100644 server/internal/handler/deployment.go create mode 100644 server/internal/handler/task_log.go delete mode 100644 server/internal/service/deployment.go diff --git a/frontend/src/api/delivery.ts b/frontend/src/api/delivery.ts new file mode 100644 index 0000000..b799841 --- /dev/null +++ b/frontend/src/api/delivery.ts @@ -0,0 +1,124 @@ +import { getToken } from '@/utils/auth' + +export interface DeliveryTarget { + id: number + name: string + target_type: string + awx_inventory_id: number + awx_template_id: number + enabled: boolean + metadata?: string +} + +export interface CreateMySQLDeliveryPayload { + business_line_id: number + target_id: number + namespace: string + instance_name: string + mysql_version: string + cpu_milli: number + memory_mi: number + storage_gi: number +} + +export interface DeliveryTask { + id: string + business_line_id: number + requested_by: number + target_type: string + target_id: number + namespace: string + instance_name: string + target_host?: string + target_host_ip?: string + mysql_port?: number + status: string + error_message?: string + created_at: string + updated_at: string + started_at?: string + finished_at?: string +} + +export interface TaskEvent { + id: number + task_id: string + from_state: string + to_state: string + message: string + created_at: string +} + +export const deliveryApi = { + async listTargets(): Promise { + const data = await authRequest('/auth/api/v1/delivery/targets?component=mysql') + return Array.isArray(data.items) ? data.items : [] + }, + + async createMySQL(payload: CreateMySQLDeliveryPayload): Promise { + const data = await authRequest('/auth/api/v1/delivery/mysql', { + method: 'POST', + headers: { + 'Idempotency-Key': createIdempotencyKey(payload), + }, + body: JSON.stringify(payload), + }) + return data.task + }, + + async getTask(taskId: string): Promise<{ task: DeliveryTask; events: TaskEvent[] }> { + const data = await authRequest(`/auth/api/v1/delivery/tasks/${encodeURIComponent(taskId)}`) + return { + task: data.task, + events: Array.isArray(data.events) ? data.events : [], + } + }, + + async cancel(taskId: string): Promise { + await authRequest(`/auth/api/v1/delivery/tasks/${encodeURIComponent(taskId)}/cancel`, { + method: 'POST', + }) + }, +} + +function createIdempotencyKey(payload: CreateMySQLDeliveryPayload) { + return [ + 'mysql', + payload.business_line_id, + payload.target_id, + payload.namespace, + payload.instance_name, + Date.now(), + ].join(':') +} + +async function authRequest(path: string, init: RequestInit = {}) { + const token = getToken() + const response = await fetch(path, { + ...init, + headers: { + Accept: 'application/json', + 'Content-Type': 'application/json', + ...(token ? { Authorization: `Bearer ${token}` } : {}), + ...init.headers, + }, + }) + const text = await response.text() + const data = parseResponseBody(text) + if (!response.ok) { + const message = data?.error || data?.message || text || `HTTP ${response.status}` + throw new Error(message) + } + return data || {} +} + +function parseResponseBody(text: string) { + if (!text.trim()) { + return {} + } + try { + return JSON.parse(text) + } catch { + return { error: text } + } +} diff --git a/frontend/src/api/deployment.ts b/frontend/src/api/deployment.ts deleted file mode 100644 index 99d8118..0000000 --- a/frontend/src/api/deployment.ts +++ /dev/null @@ -1,68 +0,0 @@ -import { getToken } from '@/utils/auth' - -export interface DeploymentCreatePayload { - component: string - business_line_id: number - params: Record -} - -export interface DeploymentCreateResult { - deployment_id: string - status: string -} - -export const deploymentApi = { - async create(payload: DeploymentCreatePayload): Promise { - const data = await authRequest('/auth/api/v1/deployments', { - method: 'POST', - body: JSON.stringify(payload), - }) - return { - deployment_id: String(data.deployment_id || ''), - status: String(data.status || ''), - } - }, - - async cancel(deploymentId: string): Promise { - await authRequest(`/auth/api/v1/deployments/${encodeURIComponent(deploymentId)}/cancel`, { - method: 'POST', - }) - }, - - eventsURL(deploymentId: string): string { - const token = getToken() - const params = token ? `?access_token=${encodeURIComponent(token)}` : '' - return `/auth/api/v1/deployments/${encodeURIComponent(deploymentId)}/events${params}` - }, -} - -async function authRequest(path: string, init: RequestInit = {}) { - const token = getToken() - const response = await fetch(path, { - ...init, - headers: { - Accept: 'application/json', - 'Content-Type': 'application/json', - ...(token ? { Authorization: `Bearer ${token}` } : {}), - ...init.headers, - }, - }) - const text = await response.text() - const data = parseResponseBody(text) - if (!response.ok) { - const message = data?.error || data?.message || text || `HTTP ${response.status}` - throw new Error(message) - } - return data || {} -} - -function parseResponseBody(text: string) { - if (!text.trim()) { - return {} - } - try { - return JSON.parse(text) - } catch { - return { error: text } - } -} diff --git a/frontend/src/api/taskLog.ts b/frontend/src/api/taskLog.ts new file mode 100644 index 0000000..34b7050 --- /dev/null +++ b/frontend/src/api/taskLog.ts @@ -0,0 +1,81 @@ +import { getToken } from '@/utils/auth' + +export interface TaskLogSummary { + id: string + source: 'awx' | 'wayne' + service: string + name: string + runner: string + status: string + status_text: string + status_class: string + business_line_id: number + reference_id: string + created_at: string + updated_at: string +} + +export interface TaskLogLine { + time: string + message: string + class: string +} + +export interface TaskLogListParams { + source?: string + businessLineId?: number +} + +export const taskLogApi = { + async list(params: TaskLogListParams = {}): Promise { + const query = new URLSearchParams() + if (params.source && params.source !== 'all') { + query.set('source', params.source) + } + if (params.businessLineId) { + query.set('business_line_id', String(params.businessLineId)) + } + const suffix = query.toString() ? `?${query.toString()}` : '' + const data = await authRequest(`/auth/api/v1/task-logs${suffix}`) + return Array.isArray(data.items) ? data.items : [] + }, + + async get(id: string): Promise<{ task: TaskLogSummary; lines: TaskLogLine[] }> { + const data = await authRequest(`/auth/api/v1/task-logs/${encodeURIComponent(id)}`) + return { + task: data.task, + lines: Array.isArray(data.lines) ? data.lines : [], + } + }, +} + +async function authRequest(path: string, init: RequestInit = {}) { + const token = getToken() + const response = await fetch(path, { + ...init, + headers: { + Accept: 'application/json', + 'Content-Type': 'application/json', + ...(token ? { Authorization: `Bearer ${token}` } : {}), + ...init.headers, + }, + }) + const text = await response.text() + const data = parseResponseBody(text) + if (!response.ok) { + const message = data?.error || data?.message || text || `HTTP ${response.status}` + throw new Error(message) + } + return data || {} +} + +function parseResponseBody(text: string) { + if (!text.trim()) { + return {} + } + try { + return JSON.parse(text) + } catch { + return { error: text } + } +} diff --git a/frontend/src/views/service/Catalog.vue b/frontend/src/views/service/Catalog.vue index d34be53..f0574b4 100644 --- a/frontend/src/views/service/Catalog.vue +++ b/frontend/src/views/service/Catalog.vue @@ -94,6 +94,12 @@ +
@@ -219,8 +225,8 @@
-

Ansible 自动交付

-

{{ deploymentId || '创建任务后生成 deployment_id' }} · {{ activeService.runner }} · {{ activeService.template }}

+

AWX 自动交付

+

{{ deploymentId || '创建任务后生成 task_id' }} · {{ activeService.runner }} · {{ activeService.template }}

取消任务 @@ -265,7 +271,7 @@
健康状态 - {{ deliveryFailed ? '健康检查失败 · 已回退' : activeService.healthText }} + {{ deliveryFailed ? resultSubtitle : activeService.healthText }}
@@ -310,11 +316,11 @@ diff --git a/server/.env.example b/server/.env.example index be1c64f..7ce4835 100644 --- a/server/.env.example +++ b/server/.env.example @@ -11,9 +11,6 @@ WAYNE_API_BASE_URL=http://wayne-backend:8080 WAYNE_ADMIN_USERNAME=admin WAYNE_ADMIN_PASSWORD=change-this-wayne-admin-password WAYNE_TOKEN_TTL_MINUTES=1440 -ANSIBLE_SERVICE_BASE_URL=http://ansible-runner:8084 -ANSIBLE_INTERNAL_TOKEN=change-this-ansible-internal-token -DEPLOYMENT_CALLBACK_BASE_URL=http://authserver-backend:8080 MYSQL_DSN=auth:auth@tcp(127.0.0.1:3000)/authserver?charset=utf8mb4&parseTime=True&loc=Local AUTO_MIGRATE=true @@ -27,6 +24,8 @@ DELIVERY_TARGET_LIMIT=2 DELIVERY_BUSINESS_LIMIT=1 AWX_BASE_URL= AWX_TOKEN= +AWX_USERNAME= +AWX_PASSWORD= DELIVERY_SERVICE_TOKEN= CLOUDDM_REGISTER_URL= CLOUDDM_API_TOKEN= diff --git a/server/internal/config/config.go b/server/internal/config/config.go index 7fdc7b6..98a10e7 100644 --- a/server/internal/config/config.go +++ b/server/internal/config/config.go @@ -14,63 +14,62 @@ type OAuthClient struct { } type Config struct { - AppEnv string - HTTPAddr string - PublicBaseURL string - MySQLDSN string - AutoMigrate bool - SSOEnabled bool - JWTSecret string - JWTIssuer string - JWTTTLMinutes int - SAMLEntityID string - SAMLACSURL string - SAMLSPCert string - SAMLSPKey string - SAMLIDPMetaURL string - SAMLLogoutURL string - WayenLoginURL string - WayenTargetURL string - WayenUsernameKey string - WayenPasswordKey string - WayenLoginFormat string - WayenLoginValue string - WayenOAuthRef string - WayenOAuthLoginURL string - WayneAPIBaseURL string - WayneAdminUsername string - WayneAdminPassword string - WayneTokenTTLMinutes int - WayneInternalAPIBaseURL string - WayneServiceName string - WayneServiceAPISecretKey string - OAuthClientID string - OAuthClientSecret string - OAuthRedirectURI string - OAuthCodeTTLSeconds int - OIDCIssuer string - OIDCAuthorizeURL string - OIDCTokenURL string - OIDCUserInfoURL string - OIDCJWKSURL string - CloudDMClientID string - CloudDMClientSecret string - CloudDMRedirectURI string - CloudDMTargetURL string - CloudDMRegisterURL string - CloudDMAPIToken string - AnsibleServiceBaseURL string - AnsibleInternalToken string - DeploymentCallbackBaseURL string - AWXBaseURL string - AWXToken string - DeliveryServiceToken string - DeliverySchedulerEnabled bool - DeliveryPollSeconds int - ReservationTTLMinutes int - DeliveryGlobalLimit int - DeliveryTargetLimit int - DeliveryBusinessLimit int + AppEnv string + HTTPAddr string + PublicBaseURL string + MySQLDSN string + AutoMigrate bool + SSOEnabled bool + JWTSecret string + JWTIssuer string + JWTTTLMinutes int + SAMLEntityID string + SAMLACSURL string + SAMLSPCert string + SAMLSPKey string + SAMLIDPMetaURL string + SAMLLogoutURL string + WayenLoginURL string + WayenTargetURL string + WayenUsernameKey string + WayenPasswordKey string + WayenLoginFormat string + WayenLoginValue string + WayenOAuthRef string + WayenOAuthLoginURL string + WayneAPIBaseURL string + WayneAdminUsername string + WayneAdminPassword string + WayneTokenTTLMinutes int + WayneInternalAPIBaseURL string + WayneServiceName string + WayneServiceAPISecretKey string + OAuthClientID string + OAuthClientSecret string + OAuthRedirectURI string + OAuthCodeTTLSeconds int + OIDCIssuer string + OIDCAuthorizeURL string + OIDCTokenURL string + OIDCUserInfoURL string + OIDCJWKSURL string + CloudDMClientID string + CloudDMClientSecret string + CloudDMRedirectURI string + CloudDMTargetURL string + CloudDMRegisterURL string + CloudDMAPIToken string + AWXBaseURL string + AWXToken string + AWXUsername string + AWXPassword string + DeliveryServiceToken string + DeliverySchedulerEnabled bool + DeliveryPollSeconds int + ReservationTTLMinutes int + DeliveryGlobalLimit int + DeliveryTargetLimit int + DeliveryBusinessLimit int } func Load() Config { @@ -83,63 +82,62 @@ func Load() Config { oidcIssuer = strings.TrimRight(oidcIssuer, "/") return Config{ - AppEnv: env("APP_ENV", "dev"), - HTTPAddr: httpAddr, - PublicBaseURL: publicBaseURL, - MySQLDSN: env("MYSQL_DSN", "auth:auth@tcp(127.0.0.1:3306)/authserver?charset=utf8mb4&parseTime=True&loc=Local"), - AutoMigrate: envBool("AUTO_MIGRATE", true), - SSOEnabled: envBool("SSO_ENABLED", true), - JWTSecret: env("JWT_SECRET", "change-this-secret"), - JWTIssuer: env("JWT_ISSUER", "authserver"), - JWTTTLMinutes: envInt("JWT_TTL_MINUTES", 120), - SAMLEntityID: samlEntityID, - SAMLACSURL: samlACSURL, - SAMLSPCert: env("SAML_SP_CERT_FILE", "certs/sp.crt"), - SAMLSPKey: env("SAML_SP_KEY_FILE", "certs/sp.key"), - SAMLIDPMetaURL: env("SAML_IDP_METADATA_URL", "http://sso-internal.dev.qiniu.io/saml2/meta"), - SAMLLogoutURL: trimURL(env("SAML_LOGOUT_URL", "")), - WayenLoginURL: env("WAYEN_LOGIN_URL", ""), - WayenTargetURL: env("WAYEN_TARGET_URL", ""), - WayenUsernameKey: env("WAYEN_USERNAME_KEY", "email"), - WayenPasswordKey: env("WAYEN_PASSWORD_KEY", "password"), - WayenLoginFormat: env("WAYEN_LOGIN_FORMAT", "form"), - WayenLoginValue: env("WAYEN_LOGIN_VALUE", "email"), - WayenOAuthRef: env("WAYEN_OAUTH_REF", "/portal/namespace/1/app"), - WayenOAuthLoginURL: trimURL(env("WAYEN_OAUTH_LOGIN_URL", "")), - WayneAPIBaseURL: trimURL(env("WAYNE_API_BASE_URL", env("WAYNE_INTERNAL_API_BASE_URL", ""))), - WayneAdminUsername: env("WAYNE_ADMIN_USERNAME", ""), - WayneAdminPassword: env("WAYNE_ADMIN_PASSWORD", ""), - WayneTokenTTLMinutes: envInt("WAYNE_TOKEN_TTL_MINUTES", 1440), - WayneInternalAPIBaseURL: trimURL(env("WAYNE_INTERNAL_API_BASE_URL", "")), - WayneServiceName: env("WAYNE_SERVICE_NAME", "xinfra"), - WayneServiceAPISecretKey: env("WAYNE_SERVICE_API_SECRET_KEY", ""), - OAuthClientID: env("OAUTH_WAYNE_CLIENT_ID", "wayne"), - OAuthClientSecret: env("OAUTH_WAYNE_CLIENT_SECRET", "wayne-secret"), - OAuthRedirectURI: env("OAUTH_WAYNE_REDIRECT_URI", ""), - OAuthCodeTTLSeconds: envInt("OAUTH_CODE_TTL_SECONDS", 120), - OIDCIssuer: oidcIssuer, - OIDCAuthorizeURL: trimURL(env("OIDC_AUTHORIZATION_ENDPOINT", oidcIssuer+"/oauth/authorize")), - OIDCTokenURL: trimURL(env("OIDC_TOKEN_ENDPOINT", oidcIssuer+"/oauth/token")), - OIDCUserInfoURL: trimURL(env("OIDC_USERINFO_ENDPOINT", oidcIssuer+"/oauth/userinfo")), - OIDCJWKSURL: trimURL(env("OIDC_JWKS_URI", oidcIssuer+"/oauth/jwks")), - CloudDMClientID: env("OIDC_CLOUDDM_CLIENT_ID", "clouddm"), - CloudDMClientSecret: env("OIDC_CLOUDDM_CLIENT_SECRET", ""), - CloudDMRedirectURI: env("OIDC_CLOUDDM_REDIRECT_URI", ""), - CloudDMTargetURL: env("CLOUDDM_TARGET_URL", ""), - CloudDMRegisterURL: trimURL(env("CLOUDDM_REGISTER_URL", "")), - CloudDMAPIToken: env("CLOUDDM_API_TOKEN", ""), - AnsibleServiceBaseURL: trimURL(env("ANSIBLE_SERVICE_BASE_URL", "")), - AnsibleInternalToken: env("ANSIBLE_INTERNAL_TOKEN", ""), - DeploymentCallbackBaseURL: trimURL(env("DEPLOYMENT_CALLBACK_BASE_URL", publicBaseURL)), - AWXBaseURL: trimURL(env("AWX_BASE_URL", "")), - AWXToken: env("AWX_TOKEN", ""), - DeliveryServiceToken: env("DELIVERY_SERVICE_TOKEN", ""), - DeliverySchedulerEnabled: envBool("DELIVERY_SCHEDULER_ENABLED", false), - DeliveryPollSeconds: envInt("DELIVERY_POLL_SECONDS", 5), - ReservationTTLMinutes: envInt("DELIVERY_RESERVATION_TTL_MINUTES", 120), - DeliveryGlobalLimit: envInt("DELIVERY_GLOBAL_LIMIT", 2), - DeliveryTargetLimit: envInt("DELIVERY_TARGET_LIMIT", 2), - DeliveryBusinessLimit: envInt("DELIVERY_BUSINESS_LIMIT", 1), + AppEnv: env("APP_ENV", "dev"), + HTTPAddr: httpAddr, + PublicBaseURL: publicBaseURL, + MySQLDSN: env("MYSQL_DSN", "auth:auth@tcp(127.0.0.1:3306)/authserver?charset=utf8mb4&parseTime=True&loc=Local"), + AutoMigrate: envBool("AUTO_MIGRATE", true), + SSOEnabled: envBool("SSO_ENABLED", true), + JWTSecret: env("JWT_SECRET", "change-this-secret"), + JWTIssuer: env("JWT_ISSUER", "authserver"), + JWTTTLMinutes: envInt("JWT_TTL_MINUTES", 120), + SAMLEntityID: samlEntityID, + SAMLACSURL: samlACSURL, + SAMLSPCert: env("SAML_SP_CERT_FILE", "certs/sp.crt"), + SAMLSPKey: env("SAML_SP_KEY_FILE", "certs/sp.key"), + SAMLIDPMetaURL: env("SAML_IDP_METADATA_URL", "http://sso-internal.dev.qiniu.io/saml2/meta"), + SAMLLogoutURL: trimURL(env("SAML_LOGOUT_URL", "")), + WayenLoginURL: env("WAYEN_LOGIN_URL", ""), + WayenTargetURL: env("WAYEN_TARGET_URL", ""), + WayenUsernameKey: env("WAYEN_USERNAME_KEY", "email"), + WayenPasswordKey: env("WAYEN_PASSWORD_KEY", "password"), + WayenLoginFormat: env("WAYEN_LOGIN_FORMAT", "form"), + WayenLoginValue: env("WAYEN_LOGIN_VALUE", "email"), + WayenOAuthRef: env("WAYEN_OAUTH_REF", "/portal/namespace/1/app"), + WayenOAuthLoginURL: trimURL(env("WAYEN_OAUTH_LOGIN_URL", "")), + WayneAPIBaseURL: trimURL(env("WAYNE_API_BASE_URL", env("WAYNE_INTERNAL_API_BASE_URL", ""))), + WayneAdminUsername: env("WAYNE_ADMIN_USERNAME", ""), + WayneAdminPassword: env("WAYNE_ADMIN_PASSWORD", ""), + WayneTokenTTLMinutes: envInt("WAYNE_TOKEN_TTL_MINUTES", 1440), + WayneInternalAPIBaseURL: trimURL(env("WAYNE_INTERNAL_API_BASE_URL", "")), + WayneServiceName: env("WAYNE_SERVICE_NAME", "xinfra"), + WayneServiceAPISecretKey: env("WAYNE_SERVICE_API_SECRET_KEY", ""), + OAuthClientID: env("OAUTH_WAYNE_CLIENT_ID", "wayne"), + OAuthClientSecret: env("OAUTH_WAYNE_CLIENT_SECRET", "wayne-secret"), + OAuthRedirectURI: env("OAUTH_WAYNE_REDIRECT_URI", ""), + OAuthCodeTTLSeconds: envInt("OAUTH_CODE_TTL_SECONDS", 120), + OIDCIssuer: oidcIssuer, + OIDCAuthorizeURL: trimURL(env("OIDC_AUTHORIZATION_ENDPOINT", oidcIssuer+"/oauth/authorize")), + OIDCTokenURL: trimURL(env("OIDC_TOKEN_ENDPOINT", oidcIssuer+"/oauth/token")), + OIDCUserInfoURL: trimURL(env("OIDC_USERINFO_ENDPOINT", oidcIssuer+"/oauth/userinfo")), + OIDCJWKSURL: trimURL(env("OIDC_JWKS_URI", oidcIssuer+"/oauth/jwks")), + CloudDMClientID: env("OIDC_CLOUDDM_CLIENT_ID", "clouddm"), + CloudDMClientSecret: env("OIDC_CLOUDDM_CLIENT_SECRET", ""), + CloudDMRedirectURI: env("OIDC_CLOUDDM_REDIRECT_URI", ""), + CloudDMTargetURL: env("CLOUDDM_TARGET_URL", ""), + CloudDMRegisterURL: trimURL(env("CLOUDDM_REGISTER_URL", "")), + CloudDMAPIToken: env("CLOUDDM_API_TOKEN", ""), + AWXBaseURL: trimURL(env("AWX_BASE_URL", "")), + AWXToken: env("AWX_TOKEN", ""), + AWXUsername: env("AWX_USERNAME", ""), + AWXPassword: env("AWX_PASSWORD", ""), + DeliveryServiceToken: env("DELIVERY_SERVICE_TOKEN", ""), + DeliverySchedulerEnabled: envBool("DELIVERY_SCHEDULER_ENABLED", false), + DeliveryPollSeconds: envInt("DELIVERY_POLL_SECONDS", 5), + ReservationTTLMinutes: envInt("DELIVERY_RESERVATION_TTL_MINUTES", 120), + DeliveryGlobalLimit: envInt("DELIVERY_GLOBAL_LIMIT", 2), + DeliveryTargetLimit: envInt("DELIVERY_TARGET_LIMIT", 2), + DeliveryBusinessLimit: envInt("DELIVERY_BUSINESS_LIMIT", 1), } } diff --git a/server/internal/database/database.go b/server/internal/database/database.go index 0a296c5..f8daaa8 100644 --- a/server/internal/database/database.go +++ b/server/internal/database/database.go @@ -20,10 +20,7 @@ func AutoMigrate(db *gorm.DB) error { &model.BusinessLineWayneNamespace{}, &model.AccessToken{}, &model.WayneToken{}, - &model.Deployment{}, - &model.DeploymentEvent{}, &model.AuditLog{}, - &model.DeploymentTarget{}, &model.ResourceQuota{}, &model.DeliveryTask{}, &model.ResourceReservation{}, diff --git a/server/internal/handler/delivery.go b/server/internal/handler/delivery.go index f4f2a7d..f36b948 100644 --- a/server/internal/handler/delivery.go +++ b/server/internal/handler/delivery.go @@ -2,7 +2,6 @@ package handler import ( "crypto/subtle" - "encoding/json" "errors" "net/http" "strconv" @@ -178,65 +177,23 @@ func (h *DeliveryHandler) Cancel(c *gin.Context) { // Targets 获取可用部署目标 // @Summary 获取可用部署目标 -// @Description 返回所有已启用的部署目标(如 k8s 集群、主机池) +// @Description 从 AWX 动态返回可用 Job Template 及其 Inventory hosts // @Tags delivery // @Produce json +// @Param component query string false "组件过滤,例如 mysql" // @Success 200 {object} map[string]any "items: 部署目标数组" // @Failure 500 {object} map[string]any "内部错误" // @Router /auth/api/v1/delivery/targets [get] // @Security BearerAuth func (h *DeliveryHandler) Targets(c *gin.Context) { - var items []model.DeploymentTarget - if err := h.service.DB().WithContext(c.Request.Context()).Where("enabled = ?", true).Order("id ASC").Find(&items).Error; err != nil { + items, err := h.service.ListTargets(c.Request.Context(), c.Query("component")) + if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } c.JSON(http.StatusOK, gin.H{"items": items}) } -type targetPayload struct { - Name string `json:"name" binding:"required"` - TargetType string `json:"target_type" binding:"required"` - AWXInventoryID uint64 `json:"awx_inventory_id" binding:"required"` - AWXTemplateID uint64 `json:"awx_template_id" binding:"required"` - Metadata map[string]any `json:"metadata"` -} - -// CreateTarget 创建部署目标 -// @Summary 创建部署目标 -// @Description 管理员创建新的部署目标(目前仅支持 k8s 类型) -// @Tags delivery -// @Accept json -// @Produce json -// @Param body body targetPayload true "目标配置" -// @Success 201 {object} model.DeploymentTarget "目标已创建" -// @Failure 400 {object} map[string]any "参数错误" -// @Failure 401 {object} map[string]any "未授权" -// @Failure 409 {object} map[string]any "名称冲突" -// @Router /auth/api/v1/delivery/targets [post] -// @Security BearerAuth -func (h *DeliveryHandler) CreateTarget(c *gin.Context) { - if !requirePlatformAdmin(c) { - return - } - var req targetPayload - if err := c.ShouldBindJSON(&req); err != nil { - c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) - return - } - if req.TargetType != "k8s" { - c.JSON(http.StatusBadRequest, gin.H{"error": "only k8s targets are supported"}) - return - } - raw, _ := json.Marshal(req.Metadata) - item := model.DeploymentTarget{Name: req.Name, TargetType: req.TargetType, AWXInventoryID: req.AWXInventoryID, AWXTemplateID: req.AWXTemplateID, Enabled: true, Metadata: string(raw)} - if err := h.service.DB().Create(&item).Error; err != nil { - c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) - return - } - c.JSON(http.StatusCreated, item) -} - type quotaPayload struct { BusinessLineID uint64 `json:"business_line_id" binding:"required"` TargetID uint64 `json:"target_id" binding:"required"` diff --git a/server/internal/handler/deployment.go b/server/internal/handler/deployment.go deleted file mode 100644 index b2a535b..0000000 --- a/server/internal/handler/deployment.go +++ /dev/null @@ -1,249 +0,0 @@ -package handler - -import ( - "encoding/json" - "errors" - "fmt" - "net/http" - "strings" - - "github.com/1024XEngineer/xinfra/server/internal/config" - "github.com/1024XEngineer/xinfra/server/internal/model" - "github.com/1024XEngineer/xinfra/server/internal/service" - - "github.com/gin-gonic/gin" - "gorm.io/gorm" -) - -type DeploymentHandler struct { - cfg config.Config - db *gorm.DB - deployments *service.DeploymentService -} - -func NewDeploymentHandler(cfg config.Config, db *gorm.DB, deployments *service.DeploymentService) *DeploymentHandler { - return &DeploymentHandler{cfg: cfg, db: db, deployments: deployments} -} - -func (h *DeploymentHandler) Create(c *gin.Context) { - claims, ok := CurrentClaims(c) - if !ok { - c.JSON(http.StatusUnauthorized, gin.H{"error": "missing auth claims"}) - return - } - var req service.DeploymentCreateRequest - if err := c.ShouldBindJSON(&req); err != nil { - c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) - return - } - if !h.ensureBusinessLineMember(c, req.BusinessLineID, claims.UserID, claims.IsAdmin) { - return - } - deployment, err := h.deployments.Create(c.Request.Context(), req, claims.UserID, claims.Username) - if err != nil { - status := http.StatusBadGateway - if errors.Is(err, service.ErrDeploymentNotConfigured) { - status = http.StatusServiceUnavailable - } - c.JSON(status, gin.H{ - "error": err.Error(), - "deployment_id": deployment.DeploymentID, - "status": deployment.Status, - }) - return - } - c.JSON(http.StatusAccepted, gin.H{"deployment_id": deployment.DeploymentID, "status": deployment.Status}) -} - -func (h *DeploymentHandler) Get(c *gin.Context) { - deployment, ok := h.getAuthorizedDeployment(c) - if !ok { - return - } - events, err := h.deployments.Events(c.Request.Context(), deployment.DeploymentID) - if err != nil { - c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) - return - } - c.JSON(http.StatusOK, gin.H{"deployment": deployment, "events": events}) -} - -func (h *DeploymentHandler) Events(c *gin.Context) { - deployment, ok := h.getAuthorizedDeployment(c) - if !ok { - return - } - events, err := h.deployments.Events(c.Request.Context(), deployment.DeploymentID) - if err != nil { - c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) - return - } - w := c.Writer - w.Header().Set("Content-Type", "text/event-stream") - w.Header().Set("Cache-Control", "no-cache") - w.Header().Set("Connection", "keep-alive") - w.Header().Set("X-Accel-Buffering", "no") - - for _, event := range events { - if err := writeSSE(w, event.Type, service.DeploymentEventData(event)); err != nil { - return - } - } - if isTerminalDeploymentStatus(deployment.Status) { - _ = writeSSE(w, service.DeploymentEventDone, gin.H{"status": deployment.Status}) - return - } - - ch, unsubscribe := h.deployments.Subscribe(deployment.DeploymentID) - defer unsubscribe() - flusher, _ := w.(http.Flusher) - if flusher != nil { - flusher.Flush() - } - for { - select { - case <-c.Request.Context().Done(): - return - case event := <-ch: - if err := writeSSE(w, event.Event.Type, event.Data); err != nil { - return - } - if event.Event.Type == service.DeploymentEventDone { - return - } - } - } -} - -func (h *DeploymentHandler) Cancel(c *gin.Context) { - deployment, ok := h.getAuthorizedDeployment(c) - if !ok { - return - } - if err := h.deployments.Cancel(c.Request.Context(), deployment.DeploymentID); err != nil { - writeDeploymentError(c, err) - return - } - c.JSON(http.StatusAccepted, gin.H{"deployment_id": deployment.DeploymentID, "status": service.DeploymentStatusCanceling}) -} - -func (h *DeploymentHandler) InternalEvent(c *gin.Context) { - if !h.authorizeInternal(c) { - return - } - deploymentID := strings.TrimSpace(c.Param("id")) - var req service.DeploymentEventRequest - if err := c.ShouldBindJSON(&req); err != nil { - c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) - return - } - event, err := h.deployments.AppendEvent(c.Request.Context(), deploymentID, req) - if err != nil { - writeDeploymentError(c, err) - return - } - c.JSON(http.StatusAccepted, gin.H{"event_id": event.ID, "seq": event.Seq}) -} - -func (h *DeploymentHandler) InternalFinish(c *gin.Context) { - if !h.authorizeInternal(c) { - return - } - deploymentID := strings.TrimSpace(c.Param("id")) - var req service.DeploymentFinishRequest - if err := c.ShouldBindJSON(&req); err != nil { - c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) - return - } - deployment, err := h.deployments.Finish(c.Request.Context(), deploymentID, req) - if err != nil { - writeDeploymentError(c, err) - return - } - c.JSON(http.StatusAccepted, gin.H{"deployment_id": deployment.DeploymentID, "status": deployment.Status}) -} - -func (h *DeploymentHandler) getAuthorizedDeployment(c *gin.Context) (model.Deployment, bool) { - claims, ok := CurrentClaims(c) - if !ok { - c.JSON(http.StatusUnauthorized, gin.H{"error": "missing auth claims"}) - return model.Deployment{}, false - } - deploymentID := strings.TrimSpace(c.Param("id")) - deployment, err := h.deployments.Get(c.Request.Context(), deploymentID) - if err != nil { - writeDeploymentError(c, err) - return model.Deployment{}, false - } - if !h.ensureBusinessLineMember(c, deployment.BusinessLineID, claims.UserID, claims.IsAdmin) { - return model.Deployment{}, false - } - return deployment, true -} - -func (h *DeploymentHandler) ensureBusinessLineMember(c *gin.Context, businessLineID uint64, userID uint64, isAdmin bool) bool { - if businessLineID == 0 { - c.JSON(http.StatusBadRequest, gin.H{"error": "business_line_id is required"}) - return false - } - if isAdmin { - return true - } - var binding model.BusinessLineUser - err := h.db.WithContext(c.Request.Context()).Where("business_line_id = ? AND user_id = ?", businessLineID, userID).First(&binding).Error - if errors.Is(err, gorm.ErrRecordNotFound) { - c.JSON(http.StatusForbidden, gin.H{"error": "current user is not assigned to this business line"}) - return false - } - if err != nil { - c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) - return false - } - return true -} - -func (h *DeploymentHandler) authorizeInternal(c *gin.Context) bool { - if h.cfg.AnsibleInternalToken == "" { - c.JSON(http.StatusServiceUnavailable, gin.H{"error": "ansible internal token is not configured"}) - return false - } - value := c.GetHeader("Authorization") - if strings.TrimPrefix(value, "Bearer ") != h.cfg.AnsibleInternalToken { - c.JSON(http.StatusUnauthorized, gin.H{"error": "invalid internal token"}) - return false - } - return true -} - -func writeSSE(w gin.ResponseWriter, event string, data any) error { - body, err := json.Marshal(data) - if err != nil { - return err - } - if _, err := fmt.Fprintf(w, "event: %s\ndata: %s\n\n", event, body); err != nil { - return err - } - if flusher, ok := w.(http.Flusher); ok { - flusher.Flush() - } - return nil -} - -func writeDeploymentError(c *gin.Context, err error) { - switch { - case errors.Is(err, service.ErrDeploymentNotFound): - c.JSON(http.StatusNotFound, gin.H{"error": err.Error()}) - case errors.Is(err, service.ErrDeploymentForbidden): - c.JSON(http.StatusForbidden, gin.H{"error": err.Error()}) - case errors.Is(err, service.ErrDeploymentInvalidState): - c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) - case errors.Is(err, service.ErrDeploymentNotConfigured): - c.JSON(http.StatusServiceUnavailable, gin.H{"error": err.Error()}) - default: - c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) - } -} - -func isTerminalDeploymentStatus(status string) bool { - return status == service.DeploymentStatusSuccess || status == service.DeploymentStatusFailed || status == service.DeploymentStatusCanceled -} diff --git a/server/internal/handler/task_log.go b/server/internal/handler/task_log.go new file mode 100644 index 0000000..ee0ad27 --- /dev/null +++ b/server/internal/handler/task_log.go @@ -0,0 +1,347 @@ +package handler + +import ( + "errors" + "net/http" + "strconv" + "strings" + "time" + + "github.com/1024XEngineer/xinfra/server/internal/model" + "github.com/1024XEngineer/xinfra/server/internal/service" + "github.com/gin-gonic/gin" + "gorm.io/gorm" +) + +type TaskLogHandler struct { + db *gorm.DB + delivery *service.DeliveryService + wayne *service.WayneRoleBindingService +} + +func NewTaskLogHandler(db *gorm.DB, delivery *service.DeliveryService, wayne *service.WayneRoleBindingService) *TaskLogHandler { + return &TaskLogHandler{db: db, delivery: delivery, wayne: wayne} +} + +type taskLogSummary struct { + ID string `json:"id"` + Source string `json:"source"` + Service string `json:"service"` + Name string `json:"name"` + Runner string `json:"runner"` + Status string `json:"status"` + StatusText string `json:"status_text"` + StatusClass string `json:"status_class"` + BusinessLineID uint64 `json:"business_line_id"` + ReferenceID string `json:"reference_id"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +type taskLogLine struct { + Time string `json:"time"` + Message string `json:"message"` + Class string `json:"class"` +} + +// List 聚合 AWX 交付任务和 Wayne 部署服务任务。 +func (h *TaskLogHandler) List(c *gin.Context) { + claims, ok := CurrentClaims(c) + if !ok { + c.JSON(http.StatusUnauthorized, gin.H{"error": "missing current user"}) + return + } + source := strings.ToLower(strings.TrimSpace(c.DefaultQuery("source", "all"))) + businessLineID, _ := strconv.ParseUint(c.Query("business_line_id"), 10, 64) + + items := make([]taskLogSummary, 0) + if source == "" || source == "all" || source == "awx" { + awxItems, err := h.listAWXTasks(c, claims.UserID, claims.IsAdmin, businessLineID) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + items = append(items, awxItems...) + } + if source == "" || source == "all" || source == "wayne" { + wayneItems, err := h.listWayneTasks(c, claims.UserID, claims.IsAdmin, businessLineID) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + items = append(items, wayneItems...) + } + sortTaskLogSummaries(items) + if len(items) > 100 { + items = items[:100] + } + c.JSON(http.StatusOK, gin.H{"items": items}) +} + +func (h *TaskLogHandler) Get(c *gin.Context) { + claims, ok := CurrentClaims(c) + if !ok { + c.JSON(http.StatusUnauthorized, gin.H{"error": "missing current user"}) + return + } + id := c.Param("id") + switch { + case strings.HasPrefix(id, "awx:"): + h.getAWXTask(c, strings.TrimPrefix(id, "awx:"), claims.UserID, claims.IsAdmin) + case strings.HasPrefix(id, "wayne:publish:"): + h.getWayneTask(c, id, claims.UserID, claims.IsAdmin) + default: + c.JSON(http.StatusBadRequest, gin.H{"error": "unknown task log source"}) + } +} + +func (h *TaskLogHandler) listAWXTasks(c *gin.Context, userID uint64, isAdmin bool, businessLineID uint64) ([]taskLogSummary, error) { + tasks, err := h.delivery.ListTasks(c.Request.Context(), userID, isAdmin, businessLineID) + if err != nil { + return nil, err + } + items := make([]taskLogSummary, 0, len(tasks)) + for _, task := range tasks { + items = append(items, awxTaskSummary(task)) + } + return items, nil +} + +func (h *TaskLogHandler) listWayneTasks(c *gin.Context, userID uint64, isAdmin bool, businessLineID uint64) ([]taskLogSummary, error) { + namespaces, err := h.visibleWayneNamespaces(c, userID, isAdmin, businessLineID) + if err != nil { + return nil, err + } + histories, err := h.wayne.ListDeploymentHistories(c.Request.Context(), namespaces, 100) + if err != nil { + return nil, err + } + items := make([]taskLogSummary, 0, len(histories)) + for _, history := range histories { + items = append(items, wayneTaskSummary(history)) + } + return items, nil +} + +func (h *TaskLogHandler) getAWXTask(c *gin.Context, taskID string, userID uint64, isAdmin bool) { + task, events, err := h.delivery.GetTask(c.Request.Context(), taskID, userID, isAdmin) + if errors.Is(err, gorm.ErrRecordNotFound) { + c.JSON(http.StatusNotFound, gin.H{"error": "task log not found"}) + return + } + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + lines := make([]taskLogLine, 0, len(events)+16) + for _, event := range events { + lines = append(lines, taskLogLine{Time: formatTaskLogTime(event.CreatedAt), Message: "[" + event.ToState + "] " + event.Message, Class: classForTaskStatus(event.ToState)}) + } + var execution model.ExecutionJob + if err := h.db.WithContext(c.Request.Context()).Where("task_id = ?", task.ID).First(&execution).Error; err == nil && execution.ExecutorJobID != "" && execution.ExecutorJobID != "pending" { + stdout, stdoutErr := h.delivery.AWXJobStdout(c.Request.Context(), execution.ExecutorJobID) + if stdoutErr != nil { + lines = append(lines, taskLogLine{Time: formatTaskLogTime(time.Now()), Message: "[awx] stdout fetch failed: " + stdoutErr.Error(), Class: "err"}) + } else { + lines = append(lines, splitStdoutLines(stdout)...) + } + } + c.JSON(http.StatusOK, gin.H{"task": awxTaskSummary(*task), "lines": lines}) +} + +func (h *TaskLogHandler) getWayneTask(c *gin.Context, id string, userID uint64, isAdmin bool) { + resourceID, historyID, ok := parseWaynePublishTaskID(id) + if !ok { + c.JSON(http.StatusBadRequest, gin.H{"error": "invalid Wayne task log id"}) + return + } + namespaces, err := h.visibleWayneNamespaces(c, userID, isAdmin, 0) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + history, err := h.wayne.GetDeploymentHistory(c.Request.Context(), namespaces, resourceID, historyID) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": err.Error()}) + return + } + lines := []taskLogLine{ + {Time: formatTaskLogTime(history.CreatedAt), Message: "[wayne] publish history #" + strconv.FormatInt(history.ID, 10), Class: classForWaynePublishStatus(history.Status)}, + {Time: formatTaskLogTime(history.CreatedAt), Message: "[deployment] " + history.ResourceName + " resource_id=" + strconv.FormatInt(history.ResourceID, 10), Class: ""}, + {Time: formatTaskLogTime(history.CreatedAt), Message: "[cluster] " + history.Cluster + " template_id=" + strconv.FormatInt(history.TemplateID, 10), Class: ""}, + {Time: formatTaskLogTime(history.CreatedAt), Message: "[user] " + history.User, Class: ""}, + } + if strings.TrimSpace(history.Message) != "" { + lines = append(lines, taskLogLine{Time: formatTaskLogTime(history.CreatedAt), Message: "[message] " + history.Message, Class: classForWaynePublishStatus(history.Status)}) + } + c.JSON(http.StatusOK, gin.H{"task": wayneTaskSummary(*history), "lines": lines}) +} + +func awxTaskSummary(task model.DeliveryTask) taskLogSummary { + return taskLogSummary{ + ID: "awx:" + task.ID, + Source: "awx", + Service: "mysql", + Name: "MySQL 标准化交付 · " + task.InstanceName, + Runner: "AWX Job Template #" + strconv.FormatUint(task.TargetID, 10), + Status: task.Status, + StatusText: textForTaskStatus(task.Status), + StatusClass: classForTaskStatus(task.Status), + BusinessLineID: task.BusinessLineID, + ReferenceID: task.ID, + CreatedAt: task.CreatedAt, + UpdatedAt: task.UpdatedAt, + } +} + +func wayneTaskSummary(history service.WayneDeploymentHistory) taskLogSummary { + return taskLogSummary{ + ID: waynePublishTaskID(history), + Source: "wayne", + Service: "wayne-deployment", + Name: wayneDeploymentTaskName(history), + Runner: "Wayne Native API", + Status: strconv.Itoa(history.Status), + StatusText: textForWaynePublishStatus(history.Status), + StatusClass: classForWaynePublishStatus(history.Status), + BusinessLineID: history.BusinessLineID, + ReferenceID: strconv.FormatInt(history.ID, 10), + CreatedAt: history.CreatedAt, + UpdatedAt: history.CreatedAt, + } +} + +func (h *TaskLogHandler) visibleWayneNamespaces(c *gin.Context, userID uint64, isAdmin bool, businessLineID uint64) ([]model.BusinessLineWayneNamespace, error) { + query := h.db.WithContext(c.Request.Context()).Order("business_line_id ASC, wayne_namespace_id ASC") + if businessLineID != 0 { + query = query.Where("business_line_id = ?", businessLineID) + } + if !isAdmin { + query = query.Where("business_line_id IN (?)", h.db.Model(&model.BusinessLineUser{}).Select("business_line_id").Where("user_id = ?", userID)) + } + var namespaces []model.BusinessLineWayneNamespace + if err := query.Find(&namespaces).Error; err != nil { + return nil, err + } + return namespaces, nil +} + +func waynePublishTaskID(history service.WayneDeploymentHistory) string { + return "wayne:publish:" + strconv.FormatInt(history.ResourceID, 10) + ":" + strconv.FormatInt(history.ID, 10) +} + +func parseWaynePublishTaskID(id string) (int64, int64, bool) { + parts := strings.Split(id, ":") + if len(parts) != 4 || parts[0] != "wayne" || parts[1] != "publish" { + return 0, 0, false + } + resourceID, resourceErr := strconv.ParseInt(parts[2], 10, 64) + historyID, historyErr := strconv.ParseInt(parts[3], 10, 64) + if resourceErr != nil || historyErr != nil { + return 0, 0, false + } + return resourceID, historyID, true +} + +func wayneDeploymentTaskName(history service.WayneDeploymentHistory) string { + name := strings.TrimSpace(history.ResourceName) + if name == "" { + name = strconv.FormatInt(history.ResourceID, 10) + } + return "Wayne 服务部署 · " + name +} + +func textForTaskStatus(status string) string { + switch status { + case model.TaskPending, model.TaskValidating, model.TaskDispatching: + return "等待" + case model.TaskRunning, model.TaskRegistering, model.TaskCanceling: + return "执行中" + case model.TaskFinished: + return "成功" + case model.TaskCanceled: + return "已取消" + default: + return "失败" + } +} + +func classForTaskStatus(status string) string { + switch status { + case model.TaskFinished: + return "ok" + case model.TaskExecutionFailed, model.TaskValidationFailed, model.TaskRegisterFailed, model.TaskCanceled: + return "err" + case model.TaskRunning, model.TaskDispatching, model.TaskRegistering, model.TaskCanceling: + return "warn" + default: + return "" + } +} + +func textForWaynePublishStatus(status int) string { + switch status { + case 1: + return "成功" + case 0: + return "失败" + default: + return "未知" + } +} + +func classForWaynePublishStatus(status int) string { + switch status { + case 1: + return "ok" + case 0: + return "err" + default: + return "" + } +} + +func splitStdoutLines(stdout string) []taskLogLine { + lines := make([]taskLogLine, 0) + for _, line := range strings.Split(stdout, "\n") { + line = strings.TrimRight(line, "\r") + if strings.TrimSpace(line) == "" { + continue + } + lines = append(lines, taskLogLine{Time: "", Message: line, Class: classForOutputLine(line)}) + } + return lines +} + +func classForOutputLine(line string) string { + lower := strings.ToLower(line) + switch { + case strings.Contains(lower, "failed") || strings.Contains(lower, "fatal") || strings.Contains(lower, "error"): + return "err" + case strings.Contains(lower, "ok:") || strings.Contains(lower, "successful") || strings.Contains(lower, "success"): + return "ok" + case strings.Contains(lower, "changed:"): + return "tag-ok" + default: + return "" + } +} + +func formatTaskLogTime(t time.Time) string { + if t.IsZero() { + return "" + } + return t.Format("15:04:05") +} + +func sortTaskLogSummaries(items []taskLogSummary) { + for i := 1; i < len(items); i++ { + item := items[i] + j := i - 1 + for j >= 0 && items[j].UpdatedAt.Before(item.UpdatedAt) { + items[j+1] = items[j] + j-- + } + items[j+1] = item + } +} diff --git a/server/internal/model/delivery.go b/server/internal/model/delivery.go index 6c7262d..e30b228 100644 --- a/server/internal/model/delivery.go +++ b/server/internal/model/delivery.go @@ -16,18 +16,6 @@ const ( TaskCanceled = "canceled" ) -type DeploymentTarget struct { - ID uint64 `gorm:"primaryKey" json:"id"` - Name string `gorm:"size:128;not null;uniqueIndex" json:"name"` - TargetType string `gorm:"size:32;not null;index" json:"target_type"` - AWXInventoryID uint64 `gorm:"not null" json:"awx_inventory_id"` - AWXTemplateID uint64 `gorm:"not null" json:"awx_template_id"` - Enabled bool `gorm:"not null;default:true;index" json:"enabled"` - Metadata string `gorm:"type:json" json:"metadata"` - CreatedAt time.Time `json:"created_at"` - UpdatedAt time.Time `json:"updated_at"` -} - type ResourceQuota struct { ID uint64 `gorm:"primaryKey" json:"id"` BusinessLineID uint64 `gorm:"not null;uniqueIndex:idx_resource_quota_scope,priority:1" json:"business_line_id"` diff --git a/server/internal/model/models.go b/server/internal/model/models.go index 2113125..1fbcbc6 100644 --- a/server/internal/model/models.go +++ b/server/internal/model/models.go @@ -80,34 +80,6 @@ type WayneToken struct { UpdatedAt time.Time `json:"updated_at"` } -type Deployment struct { - ID uint64 `gorm:"primaryKey" json:"id"` - DeploymentID string `gorm:"size:64;not null;uniqueIndex" json:"deployment_id"` - Component string `gorm:"size:64;not null;index" json:"component"` - BusinessLineID uint64 `gorm:"not null;index" json:"business_line_id"` - BusinessLine string `gorm:"size:128;not null;default:''" json:"business_line"` - Status string `gorm:"size:32;not null;index" json:"status"` - RequestPayload string `gorm:"type:longtext" json:"request_payload"` - ResultPayload string `gorm:"type:longtext" json:"result_payload"` - CreatedBy uint64 `gorm:"not null;index" json:"created_by"` - CreatedByName string `gorm:"size:128;not null;default:''" json:"created_by_name"` - StartedAt *time.Time `json:"started_at"` - FinishedAt *time.Time `json:"finished_at"` - CreatedAt time.Time `gorm:"index" json:"created_at"` - UpdatedAt time.Time `json:"updated_at"` -} - -type DeploymentEvent struct { - ID uint64 `gorm:"primaryKey" json:"id"` - DeploymentID string `gorm:"size:64;not null;index:idx_deployment_events_deployment_seq,priority:1" json:"deployment_id"` - Seq uint64 `gorm:"not null;index:idx_deployment_events_deployment_seq,priority:2" json:"seq"` - Type string `gorm:"size:32;not null;index" json:"type"` - Level string `gorm:"size:32;not null;default:''" json:"level"` - Message string `gorm:"type:longtext" json:"message"` - Payload string `gorm:"type:longtext" json:"payload"` - CreatedAt time.Time `gorm:"index" json:"created_at"` -} - type AuditLog struct { ID uint64 `gorm:"primaryKey" json:"id"` RequestID string `gorm:"size:128;not null;default:''" json:"request_id"` diff --git a/server/internal/response/response.go b/server/internal/response/response.go index b67281c..02abd4a 100644 --- a/server/internal/response/response.go +++ b/server/internal/response/response.go @@ -38,8 +38,8 @@ var codeMessages = map[int]string{ CodeTokenExpired: "Token 已过期", CodeSSONotConfigured: "子系统未配置 SSO", CodeSSOGenerateFailed: "SSO 跳转生成失败", - CodeTaskCreateFailed: "Ansible 任务创建失败", - CodeTaskExecTimeout: "Ansible 任务执行超时", + CodeTaskCreateFailed: "交付任务创建失败", + CodeTaskExecTimeout: "交付任务执行超时", } func Message(code int) string { diff --git a/server/internal/router/router.go b/server/internal/router/router.go index 0408fbe..a1bb599 100644 --- a/server/internal/router/router.go +++ b/server/internal/router/router.go @@ -71,7 +71,6 @@ func registerAuthServerRoutes(r *gin.Engine, deps Dependencies) { authService := service.NewAuthService(deps.Config, deps.DB, auditService) wayenService := service.NewWayenService(deps.Config, deps.DB) wayneRoleBindingService := service.NewWayneRoleBindingService(deps.Config, deps.DB) - deploymentService := service.NewDeploymentService(deps.Config, deps.DB) deliveryService := service.NewDeliveryService(deps.Config, deps.DB, auditService) if deps.Config.DeliverySchedulerEnabled { go deliveryService.Run(context.Background()) @@ -87,10 +86,10 @@ func registerAuthServerRoutes(r *gin.Engine, deps Dependencies) { clouddmHandler := handler.NewCloudDMHandler(deps.Config, auditService) samlHandler := handler.NewSAMLHandler(deps.Config, authService) oauthHandler := handler.NewOAuthHandler(deps.Config, deps.DB, auditService) - deploymentHandler := handler.NewDeploymentHandler(deps.Config, deps.DB, deploymentService) deliveryHandler := handler.NewDeliveryHandler(deliveryService) executionHandler := handler.NewExecutionHandler(deliveryService, deps.Config.DeliveryServiceToken) containerServiceHandler := handler.NewContainerServiceHandler(deps.DB, wayneRoleBindingService) + taskLogHandler := handler.NewTaskLogHandler(deps.DB, deliveryService, wayneRoleBindingService) r.GET("/healthz", healthHandler.Healthz) r.GET("/readyz", healthHandler.Readyz) @@ -100,8 +99,6 @@ func registerAuthServerRoutes(r *gin.Engine, deps Dependencies) { r.POST("/auth/oauth/token", oauthHandler.Token) r.GET("/auth/oauth/jwks", oauthHandler.JWKS) r.GET("/auth/oauth/userinfo", oauthHandler.UserInfo) - r.POST("/auth/internal/deployments/:id/events", deploymentHandler.InternalEvent) - r.POST("/auth/internal/deployments/:id/finish", deploymentHandler.InternalFinish) v1 := r.Group("/auth/api/v1") { @@ -145,17 +142,14 @@ func registerAuthServerRoutes(r *gin.Engine, deps Dependencies) { protected.POST("/subsystem-auth/wayne/business-lines/:id/users/:userid/init", subsystemAuthHandler.InitWayneBusinessLineUser) protected.GET("/container-services/business-lines/:id/summary", containerServiceHandler.Summary) protected.GET("/container-services/business-lines/:id/workloads", containerServiceHandler.Workloads) - protected.POST("/deployments", deploymentHandler.Create) - protected.GET("/deployments/:id", deploymentHandler.Get) - protected.GET("/deployments/:id/events", deploymentHandler.Events) - protected.POST("/deployments/:id/cancel", deploymentHandler.Cancel) protected.GET("/clouddm/login", clouddmHandler.Login) protected.GET("/delivery/targets", deliveryHandler.Targets) - protected.POST("/delivery/targets", deliveryHandler.CreateTarget) protected.PUT("/delivery/quotas", deliveryHandler.UpsertQuota) protected.POST("/delivery/mysql", deliveryHandler.CreateMySQL) protected.GET("/delivery/tasks", deliveryHandler.List) protected.GET("/delivery/tasks/:id", deliveryHandler.Get) protected.POST("/delivery/tasks/:id/cancel", deliveryHandler.Cancel) + protected.GET("/task-logs", taskLogHandler.List) + protected.GET("/task-logs/:id", taskLogHandler.Get) } } diff --git a/server/internal/service/awx.go b/server/internal/service/awx.go index 8b919ab..e719798 100644 --- a/server/internal/service/awx.go +++ b/server/internal/service/awx.go @@ -7,15 +7,22 @@ import ( "fmt" "io" "net/http" + "regexp" "strconv" "strings" + "sync" "time" ) type AWXClient struct { - baseURL string - token string - client *http.Client + baseURL string + staticToken string + username string + password string + cachedToken string + tokenExpires time.Time + tokenMu sync.Mutex + client *http.Client } type AWXLaunchRequest struct { @@ -30,16 +37,50 @@ type AWXJob struct { Failed bool `json:"failed"` } -func NewAWXClient(baseURL, token string) *AWXClient { +type AWXJobTemplate struct { + ID uint64 `json:"id"` + Name string `json:"name"` + Description string `json:"description"` + Inventory uint64 `json:"inventory"` +} + +type AWXInventoryHost struct { + ID uint64 `json:"id"` + Name string `json:"name"` + Enabled bool `json:"enabled"` + Variables string `json:"variables"` +} + +type awxListResponse[T any] struct { + Next string `json:"next"` + Results []T `json:"results"` +} + +type awxTokenResponse struct { + Token string `json:"token"` + Expires string `json:"expires"` +} + +func NewAWXClient(baseURL, token string, credentials ...string) *AWXClient { + username := "" + password := "" + if len(credentials) > 0 { + username = credentials[0] + } + if len(credentials) > 1 { + password = credentials[1] + } return &AWXClient{ - baseURL: strings.TrimRight(baseURL, "/"), - token: strings.TrimSpace(token), - client: &http.Client{Timeout: 30 * time.Second}, + baseURL: strings.TrimRight(baseURL, "/"), + staticToken: strings.TrimSpace(token), + username: strings.TrimSpace(username), + password: strings.TrimSpace(password), + client: &http.Client{Timeout: 30 * time.Second}, } } func (c *AWXClient) Configured() bool { - return c.baseURL != "" && c.token != "" + return c.baseURL != "" && (c.staticToken != "" || (c.username != "" && c.password != "")) } func (c *AWXClient) Launch(ctx context.Context, templateID uint64, input AWXLaunchRequest) (*AWXJob, error) { @@ -68,6 +109,17 @@ func (c *AWXClient) GetJob(ctx context.Context, jobID string) (*AWXJob, error) { return &job, nil } +func (c *AWXClient) JobStdout(ctx context.Context, jobID string) (string, error) { + if _, err := strconv.ParseUint(jobID, 10, 64); err != nil { + return "", fmt.Errorf("invalid AWX job id %q", jobID) + } + raw, err := c.requestRaw(ctx, http.MethodGet, "/api/v2/jobs/"+jobID+"/stdout/?format=txt", nil, "text/plain") + if err != nil { + return "", err + } + return string(raw), nil +} + func (c *AWXClient) Cancel(ctx context.Context, jobID string) error { if _, err := strconv.ParseUint(jobID, 10, 64); err != nil { return fmt.Errorf("invalid AWX job id %q", jobID) @@ -75,39 +127,56 @@ func (c *AWXClient) Cancel(ctx context.Context, jobID string) error { return c.request(ctx, http.MethodPost, "/api/v2/jobs/"+jobID+"/cancel/", map[string]any{}, nil) } -func (c *AWXClient) request(ctx context.Context, method, path string, payload any, output any) error { - if !c.Configured() { - return fmt.Errorf("AWX_BASE_URL and AWX_TOKEN must be configured") - } - var body io.Reader - if payload != nil { - raw, err := json.Marshal(payload) - if err != nil { - return err +func (c *AWXClient) ListJobTemplates(ctx context.Context) ([]AWXJobTemplate, error) { + var out []AWXJobTemplate + path := "/api/v2/job_templates/?page_size=200" + for path != "" { + var page awxListResponse[AWXJobTemplate] + if err := c.request(ctx, http.MethodGet, path, nil, &page); err != nil { + return nil, err } - body = bytes.NewReader(raw) + out = append(out, page.Results...) + path = awxNextPath(page.Next) } - req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, body) + return out, nil +} + +func (c *AWXClient) GetJobTemplate(ctx context.Context, templateID uint64) (*AWXJobTemplate, error) { + if templateID == 0 { + return nil, fmt.Errorf("invalid AWX job template id") + } + var item AWXJobTemplate + if err := c.request(ctx, http.MethodGet, fmt.Sprintf("/api/v2/job_templates/%d/", templateID), nil, &item); err != nil { + return nil, err + } + if item.ID == 0 { + return nil, fmt.Errorf("AWX job template %d was not found", templateID) + } + return &item, nil +} + +func (c *AWXClient) ListInventoryHosts(ctx context.Context, inventoryID uint64) ([]AWXInventoryHost, error) { + if inventoryID == 0 { + return nil, fmt.Errorf("AWX job template does not bind an inventory") + } + var out []AWXInventoryHost + path := fmt.Sprintf("/api/v2/inventories/%d/hosts/?page_size=200", inventoryID) + for path != "" { + var page awxListResponse[AWXInventoryHost] + if err := c.request(ctx, http.MethodGet, path, nil, &page); err != nil { + return nil, err + } + out = append(out, page.Results...) + path = awxNextPath(page.Next) + } + return out, nil +} + +func (c *AWXClient) request(ctx context.Context, method, path string, payload any, output any) error { + raw, err := c.requestRaw(ctx, method, path, payload, "application/json") if err != nil { return err } - req.Header.Set("Authorization", "Bearer "+c.token) - req.Header.Set("Accept", "application/json") - if payload != nil { - req.Header.Set("Content-Type", "application/json") - } - resp, err := c.client.Do(req) - if err != nil { - return fmt.Errorf("AWX request: %w", err) - } - defer resp.Body.Close() - raw, err := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) - if err != nil { - return err - } - if resp.StatusCode < 200 || resp.StatusCode >= 300 { - return fmt.Errorf("AWX returned %s: %s", resp.Status, strings.TrimSpace(string(raw))) - } if output != nil && len(raw) > 0 { if err := json.Unmarshal(raw, output); err != nil { return fmt.Errorf("decode AWX response: %w", err) @@ -115,3 +184,132 @@ func (c *AWXClient) request(ctx context.Context, method, path string, payload an } return nil } + +func (c *AWXClient) requestRaw(ctx context.Context, method, path string, payload any, accept string) ([]byte, error) { + if !c.Configured() { + return nil, fmt.Errorf("AWX_BASE_URL and AWX_TOKEN or AWX_USERNAME/AWX_PASSWORD must be configured") + } + token, err := c.bearerToken(ctx) + if err != nil { + return nil, err + } + var body io.Reader + if payload != nil { + raw, err := json.Marshal(payload) + if err != nil { + return nil, err + } + body = bytes.NewReader(raw) + } + req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, body) + if err != nil { + return nil, err + } + req.Header.Set("Authorization", "Bearer "+token) + if accept == "" { + accept = "application/json" + } + req.Header.Set("Accept", accept) + if payload != nil { + req.Header.Set("Content-Type", "application/json") + } + resp, err := c.client.Do(req) + if err != nil { + return nil, fmt.Errorf("AWX request: %w", err) + } + defer resp.Body.Close() + raw, err := io.ReadAll(io.LimitReader(resp.Body, 4<<20)) + if err != nil { + return nil, err + } + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return nil, fmt.Errorf("AWX returned %s: %s", resp.Status, strings.TrimSpace(string(raw))) + } + return raw, nil +} + +func (c *AWXClient) bearerToken(ctx context.Context) (string, error) { + if c.staticToken != "" { + return c.staticToken, nil + } + c.tokenMu.Lock() + defer c.tokenMu.Unlock() + if c.cachedToken != "" && time.Until(c.tokenExpires) > 5*time.Minute { + return c.cachedToken, nil + } + token, expires, err := c.createToken(ctx) + if err != nil { + return "", err + } + c.cachedToken = token + c.tokenExpires = expires + return token, nil +} + +func (c *AWXClient) createToken(ctx context.Context) (string, time.Time, error) { + raw, err := json.Marshal(map[string]any{"description": "xinfra delivery"}) + if err != nil { + return "", time.Time{}, err + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/api/v2/tokens/", bytes.NewReader(raw)) + if err != nil { + return "", time.Time{}, err + } + req.SetBasicAuth(c.username, c.password) + req.Header.Set("Accept", "application/json") + req.Header.Set("Content-Type", "application/json") + resp, err := c.client.Do(req) + if err != nil { + return "", time.Time{}, fmt.Errorf("AWX token request: %w", err) + } + defer resp.Body.Close() + body, err := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) + if err != nil { + return "", time.Time{}, err + } + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return "", time.Time{}, fmt.Errorf("AWX token returned %s: %s", resp.Status, strings.TrimSpace(string(body))) + } + var data awxTokenResponse + if err := json.Unmarshal(body, &data); err != nil { + return "", time.Time{}, fmt.Errorf("decode AWX token response: %w", err) + } + if strings.TrimSpace(data.Token) == "" { + return "", time.Time{}, fmt.Errorf("AWX token response did not include a token") + } + expires := time.Now().Add(24 * time.Hour) + if data.Expires != "" { + if parsed, err := time.Parse(time.RFC3339, data.Expires); err == nil { + expires = parsed + } + } + return strings.TrimSpace(data.Token), expires, nil +} + +func awxNextPath(next string) string { + if next == "" { + return "" + } + if strings.HasPrefix(next, "http://") || strings.HasPrefix(next, "https://") { + if idx := strings.Index(next, "/api/"); idx >= 0 { + return next[idx:] + } + return "" + } + return next +} + +var ansibleHostPattern = regexp.MustCompile(`(?m)^\s*ansible_host\s*:\s*"?([^"\s]+)"?\s*$`) + +func AWXHostIP(host AWXInventoryHost) string { + var parsed map[string]any + if err := json.Unmarshal([]byte(host.Variables), &parsed); err == nil { + if value, ok := parsed["ansible_host"].(string); ok && strings.TrimSpace(value) != "" { + return strings.TrimSpace(value) + } + } + if match := ansibleHostPattern.FindStringSubmatch(host.Variables); len(match) == 2 { + return strings.TrimSpace(match[1]) + } + return host.Name +} diff --git a/server/internal/service/awx_test.go b/server/internal/service/awx_test.go index e074231..09c8706 100644 --- a/server/internal/service/awx_test.go +++ b/server/internal/service/awx_test.go @@ -7,6 +7,7 @@ import ( "io" "net/http" "testing" + "time" ) type roundTripFunc func(*http.Request) (*http.Response, error) @@ -37,3 +38,38 @@ func TestAWXClientLaunchAndGetJob(t *testing.T) { t.Fatalf("get job = %#v, err=%v", job, err) } } + +func TestAWXClientCreatesTokenFromCredentials(t *testing.T) { + client := NewAWXClient("https://awx.example", "", "admin", "password") + tokenRequests := 0 + client.client.Transport = roundTripFunc(func(r *http.Request) (*http.Response, error) { + switch r.URL.Path { + case "/api/v2/tokens/": + tokenRequests++ + username, password, ok := r.BasicAuth() + if !ok || username != "admin" || password != "password" { + t.Fatalf("missing basic auth for token request") + } + raw, _ := json.Marshal(awxTokenResponse{Token: "generated-token", Expires: time.Now().Add(time.Hour).Format(time.RFC3339)}) + return &http.Response{StatusCode: http.StatusCreated, Status: "201 Created", Body: io.NopCloser(bytes.NewReader(raw)), Header: make(http.Header)}, nil + case "/api/v2/jobs/42/": + if r.Header.Get("Authorization") != "Bearer generated-token" { + t.Fatalf("missing generated bearer token") + } + raw, _ := json.Marshal(AWXJob{ID: 42, Status: "successful"}) + return &http.Response{StatusCode: http.StatusOK, Status: "200 OK", Body: io.NopCloser(bytes.NewReader(raw)), Header: make(http.Header)}, nil + default: + t.Fatalf("unexpected path %s", r.URL.Path) + return nil, nil + } + }) + if _, err := client.GetJob(context.Background(), "42"); err != nil { + t.Fatalf("get job: %v", err) + } + if _, err := client.GetJob(context.Background(), "42"); err != nil { + t.Fatalf("get job with cached token: %v", err) + } + if tokenRequests != 1 { + t.Fatalf("token requests = %d, want 1", tokenRequests) + } +} diff --git a/server/internal/service/container_service.go b/server/internal/service/container_service.go index 9c783f2..7218931 100644 --- a/server/internal/service/container_service.go +++ b/server/internal/service/container_service.go @@ -8,6 +8,7 @@ import ( "net/url" "strconv" "strings" + "time" "github.com/1024XEngineer/xinfra/server/internal/model" ) @@ -57,6 +58,22 @@ type ContainerBusinessLineQuery struct { Namespaces []model.BusinessLineWayneNamespace } +type WayneDeploymentHistory struct { + ID int64 `json:"id"` + Type int `json:"type"` + ResourceID int64 `json:"resourceId"` + ResourceName string `json:"resourceName"` + TemplateID int64 `json:"templateId"` + Cluster string `json:"cluster"` + Status int `json:"status"` + Message string `json:"message"` + User string `json:"user"` + CreateTime string `json:"createTime"` + CreatedAt time.Time `json:"-"` + BusinessLineID uint64 `json:"businessLineId"` + NamespaceID uint64 `json:"namespaceId"` +} + func (s *WayneRoleBindingService) ContainerServiceData(ctx context.Context, query ContainerBusinessLineQuery) (*ContainerServiceData, error) { summary := ContainerServiceSummary{ BusinessLineID: query.BusinessLineID, @@ -143,6 +160,65 @@ func (s *WayneRoleBindingService) ContainerServiceData(ctx context.Context, quer return &ContainerServiceData{Summary: summary, Workloads: workloads}, nil } +func (s *WayneRoleBindingService) ListDeploymentHistories(ctx context.Context, namespaces []model.BusinessLineWayneNamespace, limit int) ([]WayneDeploymentHistory, error) { + if limit <= 0 { + limit = 100 + } + histories := make([]WayneDeploymentHistory, 0) + seenDeployments := map[int64]struct{}{} + for _, binding := range namespaces { + apps, err := s.listWayneNamespaceApps(ctx, binding.WayneNamespaceID) + if err != nil { + return nil, err + } + for _, app := range apps { + deployments, err := s.listWayneAppDeployments(ctx, app.ID) + if err != nil { + return nil, err + } + for _, deployment := range deployments { + if deployment.ID == 0 { + continue + } + if _, ok := seenDeployments[deployment.ID]; ok { + continue + } + seenDeployments[deployment.ID] = struct{}{} + items, err := s.listWaynePublishHistories(ctx, deployment.ID, 20) + if err != nil { + return nil, err + } + for _, item := range items { + item.BusinessLineID = binding.BusinessLineID + item.NamespaceID = binding.WayneNamespaceID + if item.ResourceName == "" { + item.ResourceName = deployment.Name + } + histories = append(histories, item) + } + } + } + } + sortWayneDeploymentHistories(histories) + if len(histories) > limit { + histories = histories[:limit] + } + return histories, nil +} + +func (s *WayneRoleBindingService) GetDeploymentHistory(ctx context.Context, namespaces []model.BusinessLineWayneNamespace, resourceID, historyID int64) (*WayneDeploymentHistory, error) { + histories, err := s.ListDeploymentHistories(ctx, namespaces, 500) + if err != nil { + return nil, err + } + for _, history := range histories { + if history.ResourceID == resourceID && history.ID == historyID { + return &history, nil + } + } + return nil, fmt.Errorf("wayne deployment history %d was not found", historyID) +} + func (s *WayneRoleBindingService) getWayneNamespace(ctx context.Context, id uint64) (wayneNamespaceDetail, error) { result, err := s.callRaw(ctx, http.MethodGet, fmt.Sprintf("/api/v1/namespaces/%d", id), nil) if err != nil { @@ -162,6 +238,32 @@ func (s *WayneRoleBindingService) listWayneNamespaceApps(ctx context.Context, na return parseWayneApps(result.Body) } +func (s *WayneRoleBindingService) listWayneAppDeployments(ctx context.Context, appID uint64) ([]wayneDeployment, error) { + values := url.Values{} + values.Set("pageNo", "1") + values.Set("pageSize", "500") + values.Set("deleted", "false") + result, err := s.callRaw(ctx, http.MethodGet, fmt.Sprintf("/api/v1/apps/%d/deployments?%s", appID, values.Encode()), nil) + if err != nil { + return nil, err + } + return parseWayneDeployments(result.Body) +} + +func (s *WayneRoleBindingService) listWaynePublishHistories(ctx context.Context, resourceID int64, pageSize int) ([]WayneDeploymentHistory, error) { + values := url.Values{} + values.Set("pageNo", "1") + values.Set("pageSize", strconv.Itoa(pageSize)) + values.Set("type", "0") + values.Set("resourceId", strconv.FormatInt(resourceID, 10)) + values.Set("sortby", "-createTime") + result, err := s.callRaw(ctx, http.MethodGet, "/api/v1/publish/histories?"+values.Encode(), nil) + if err != nil { + return nil, err + } + return parseWaynePublishHistories(result.Body) +} + func (s *WayneRoleBindingService) getWayneNodes(ctx context.Context, cluster string) (wayneNodeSummary, error) { result, err := s.callRaw(ctx, http.MethodGet, fmt.Sprintf("/api/v1/kubernetes/nodes/clusters/%s", url.PathEscape(cluster)), nil) if err != nil { @@ -216,6 +318,11 @@ type wayneApp struct { Name string `json:"name"` } +type wayneDeployment struct { + ID int64 `json:"id"` + Name string `json:"name"` +} + type wayneNodeSummary struct { Total int Ready int @@ -248,6 +355,33 @@ func parseWayneApps(body []byte) ([]wayneApp, error) { return wrapped.Data.List, nil } +func parseWayneDeployments(body []byte) ([]wayneDeployment, error) { + var wrapped struct { + Data struct { + List []wayneDeployment `json:"list"` + } `json:"data"` + } + if err := json.Unmarshal(body, &wrapped); err != nil { + return nil, err + } + return wrapped.Data.List, nil +} + +func parseWaynePublishHistories(body []byte) ([]WayneDeploymentHistory, error) { + var wrapped struct { + Data struct { + List []WayneDeploymentHistory `json:"list"` + } `json:"data"` + } + if err := json.Unmarshal(body, &wrapped); err != nil { + return nil, err + } + for i := range wrapped.Data.List { + wrapped.Data.List[i].CreatedAt = parseWayneTime(wrapped.Data.List[i].CreateTime) + } + return wrapped.Data.List, nil +} + func parseWayneMapList(body []byte) ([]map[string]any, error) { var wrapped struct { Data struct { @@ -260,6 +394,32 @@ func parseWayneMapList(body []byte) ([]map[string]any, error) { return wrapped.Data.List, nil } +func parseWayneTime(raw string) time.Time { + raw = strings.TrimSpace(raw) + if raw == "" { + return time.Time{} + } + layouts := []string{time.RFC3339, "2006-01-02T15:04:05Z07:00", "2006-01-02 15:04:05", "2006-01-02T15:04:05"} + for _, layout := range layouts { + if parsed, err := time.Parse(layout, raw); err == nil { + return parsed + } + } + return time.Time{} +} + +func sortWayneDeploymentHistories(items []WayneDeploymentHistory) { + for i := 1; i < len(items); i++ { + item := items[i] + j := i - 1 + for j >= 0 && items[j].CreatedAt.Before(item.CreatedAt) { + items[j+1] = items[j] + j-- + } + items[j+1] = item + } +} + func parseWayneNodeSummary(body []byte) (wayneNodeSummary, error) { var wrapped struct { Data struct { diff --git a/server/internal/service/delivery.go b/server/internal/service/delivery.go index 3135a42..1186212 100644 --- a/server/internal/service/delivery.go +++ b/server/internal/service/delivery.go @@ -40,7 +40,17 @@ type deliveryPayload struct { TargetType string `json:"target_type"` } -// targetMetadata describes the native VM候选节点池以及部署形态,存储在 DeploymentTarget.Metadata (JSON)。 +type DeliveryTarget struct { + ID uint64 `json:"id"` + Name string `json:"name"` + TargetType string `json:"target_type"` + AWXInventoryID uint64 `json:"awx_inventory_id"` + AWXTemplateID uint64 `json:"awx_template_id"` + Enabled bool `json:"enabled"` + Metadata string `json:"metadata"` +} + +// targetMetadata describes the native VM候选节点池以及部署形态,由 AWX inventory hosts 动态组装。 type targetMetadata struct { Topology string `json:"topology"` MySQLPort int `json:"mysql_port"` @@ -91,7 +101,68 @@ type DeliveryService struct { func (s *DeliveryService) DB() *gorm.DB { return s.db } func NewDeliveryService(cfg config.Config, db *gorm.DB, audit *AuditService) *DeliveryService { - return &DeliveryService{db: db, cfg: cfg, awx: NewAWXClient(cfg.AWXBaseURL, cfg.AWXToken), audit: audit} + return &DeliveryService{db: db, cfg: cfg, awx: NewAWXClient(cfg.AWXBaseURL, cfg.AWXToken, cfg.AWXUsername, cfg.AWXPassword), audit: audit} +} + +func (s *DeliveryService) ListTargets(ctx context.Context, component string) ([]DeliveryTarget, error) { + templates, err := s.awx.ListJobTemplates(ctx) + if err != nil { + return nil, err + } + component = strings.ToLower(strings.TrimSpace(component)) + var targets []DeliveryTarget + for _, template := range templates { + if template.Inventory == 0 { + continue + } + if component != "" && component != "all" { + text := strings.ToLower(template.Name + " " + template.Description) + if !strings.Contains(text, component) { + continue + } + } + target, err := s.awxDeliveryTarget(ctx, template) + if err != nil { + continue + } + targets = append(targets, target) + } + return targets, nil +} + +func (s *DeliveryService) getTarget(ctx context.Context, templateID uint64) (DeliveryTarget, error) { + template, err := s.awx.GetJobTemplate(ctx, templateID) + if err != nil { + return DeliveryTarget{}, fmt.Errorf("deployment target is unavailable: %w", err) + } + return s.awxDeliveryTarget(ctx, *template) +} + +func (s *DeliveryService) awxDeliveryTarget(ctx context.Context, template AWXJobTemplate) (DeliveryTarget, error) { + hosts, err := s.awx.ListInventoryHosts(ctx, template.Inventory) + if err != nil { + return DeliveryTarget{}, err + } + meta := targetMetadata{Topology: "standalone", MySQLPort: 3307} + for _, host := range hosts { + if !host.Enabled { + continue + } + meta.Hosts = append(meta.Hosts, targetHost{Name: host.Name, IP: AWXHostIP(host)}) + } + raw, err := json.Marshal(meta) + if err != nil { + return DeliveryTarget{}, err + } + return DeliveryTarget{ + ID: template.ID, + Name: template.Name, + TargetType: "k8s", + AWXInventoryID: template.Inventory, + AWXTemplateID: template.ID, + Enabled: true, + Metadata: string(raw), + }, nil } func (s *DeliveryService) CreateTask(ctx context.Context, userID uint64, isAdmin bool, idempotencyKey string, input MySQLDeliveryInput) (*model.DeliveryTask, bool, error) { @@ -113,9 +184,9 @@ func (s *DeliveryService) CreateTask(ctx context.Context, userID uint64, isAdmin return nil, false, err } - var target model.DeploymentTarget - if err := s.db.WithContext(ctx).First(&target, "id = ? AND enabled = ?", input.TargetID, true).Error; err != nil { - return nil, false, fmt.Errorf("deployment target is unavailable: %w", err) + target, err := s.getTarget(ctx, input.TargetID) + if err != nil { + return nil, false, err } if target.TargetType != "k8s" { return nil, false, fmt.Errorf("target type %q is not supported in the MVP", target.TargetType) @@ -208,6 +279,10 @@ func (s *DeliveryService) GetTask(ctx context.Context, taskID string, userID uin return &task, events, nil } +func (s *DeliveryService) AWXJobStdout(ctx context.Context, jobID string) (string, error) { + return s.awx.JobStdout(ctx, jobID) +} + func (s *DeliveryService) Cancel(ctx context.Context, taskID string, userID uint64, isAdmin bool) error { task, _, err := s.GetTask(ctx, taskID, userID, isAdmin) if err != nil { @@ -232,14 +307,16 @@ func (s *DeliveryService) Cancel(ctx context.Context, taskID string, userID uint func (s *DeliveryService) claimAndReserve(ctx context.Context) (*model.DeliveryTask, error) { var task model.DeliveryTask var payload deliveryPayload - var target model.DeploymentTarget + var target DeliveryTarget dispatchable := false err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { if err := tx.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}).Where("status = ?", model.TaskPending).Order("created_at ASC").First(&task).Error; err != nil { return err } - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&target, "id = ? AND enabled = ?", task.TargetID, true).Error; err != nil { - return s.failInTransaction(tx, &task, model.TaskValidationFailed, "deployment target is unavailable") + var targetErr error + target, targetErr = s.getTarget(ctx, task.TargetID) + if targetErr != nil { + return s.failInTransaction(tx, &task, model.TaskValidationFailed, targetErr.Error()) } if err := json.Unmarshal([]byte(task.ImmutablePayload), &payload); err != nil { return s.failInTransaction(tx, &task, model.TaskValidationFailed, "stored deployment payload is invalid") @@ -267,19 +344,11 @@ func (s *DeliveryService) claimAndReserve(ctx context.Context) (*model.DeliveryT } } } - var quota model.ResourceQuota - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("business_line_id = ? AND target_id = ?", task.BusinessLineID, task.TargetID).First("a).Error; err != nil { - return s.failInTransaction(tx, &task, model.TaskValidationFailed, "resource quota is not configured") - } - type totals struct{ CPU, Memory, Storage, Instances int64 } - var used, reserved totals - if err := tx.Model(&model.ResourceUsage{}).Select("COALESCE(SUM(cpu_milli),0) cpu, COALESCE(SUM(memory_mi),0) memory, COALESCE(SUM(storage_gi),0) storage, COALESCE(SUM(instance_count),0) instances").Where("business_line_id = ? AND target_id = ? AND status = ?", task.BusinessLineID, task.TargetID, "active").Scan(&used).Error; err != nil { + quotaOK, err := checkResourceQuota(tx, task.BusinessLineID, task.TargetID, payload) + if err != nil { return err } - if err := tx.Model(&model.ResourceReservation{}).Select("COALESCE(SUM(cpu_milli),0) cpu, COALESCE(SUM(memory_mi),0) memory, COALESCE(SUM(storage_gi),0) storage, COALESCE(SUM(instance_count),0) instances").Where("business_line_id = ? AND target_id = ? AND status = ? AND expires_at > ?", task.BusinessLineID, task.TargetID, "reserved", time.Now()).Scan(&reserved).Error; err != nil { - return err - } - if used.CPU+reserved.CPU+payload.CPUMilli > quota.CPUMilli || used.Memory+reserved.Memory+payload.MemoryMi > quota.MemoryMi || used.Storage+reserved.Storage+payload.StorageGi > quota.StorageGi || used.Instances+reserved.Instances+1 > quota.InstanceLimit { + if !quotaOK { return s.failInTransaction(tx, &task, model.TaskValidationFailed, "resource quota is insufficient") } meta := parseTargetMetadata(target.Metadata) @@ -317,6 +386,28 @@ func (s *DeliveryService) claimAndReserve(ctx context.Context) (*model.DeliveryT return &task, err } +func checkResourceQuota(tx *gorm.DB, businessLineID, targetID uint64, payload deliveryPayload) (bool, error) { + var quota model.ResourceQuota + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("business_line_id = ? AND target_id = ?", businessLineID, targetID).First("a).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return true, nil + } + return false, err + } + type totals struct{ CPU, Memory, Storage, Instances int64 } + var used, reserved totals + if err := tx.Model(&model.ResourceUsage{}).Select("COALESCE(SUM(cpu_milli),0) cpu, COALESCE(SUM(memory_mi),0) memory, COALESCE(SUM(storage_gi),0) storage, COALESCE(SUM(instance_count),0) instances").Where("business_line_id = ? AND target_id = ? AND status = ?", businessLineID, targetID, "active").Scan(&used).Error; err != nil { + return false, err + } + if err := tx.Model(&model.ResourceReservation{}).Select("COALESCE(SUM(cpu_milli),0) cpu, COALESCE(SUM(memory_mi),0) memory, COALESCE(SUM(storage_gi),0) storage, COALESCE(SUM(instance_count),0) instances").Where("business_line_id = ? AND target_id = ? AND status = ? AND expires_at > ?", businessLineID, targetID, "reserved", time.Now()).Scan(&reserved).Error; err != nil { + return false, err + } + return used.CPU+reserved.CPU+payload.CPUMilli <= quota.CPUMilli && + used.Memory+reserved.Memory+payload.MemoryMi <= quota.MemoryMi && + used.Storage+reserved.Storage+payload.StorageGi <= quota.StorageGi && + used.Instances+reserved.Instances+1 <= quota.InstanceLimit, nil +} + func (s *DeliveryService) failInTransaction(tx *gorm.DB, task *model.DeliveryTask, status, message string) error { if err := s.transitionTx(tx, task, status, message, message); err != nil { return err @@ -410,8 +501,8 @@ func (s *DeliveryService) CreateExecution(ctx context.Context, taskID, payloadHa if task.Status != model.TaskDispatching { return nil, false, fmt.Errorf("task in state %q is not ready for execution", task.Status) } - var target model.DeploymentTarget - if err := s.db.WithContext(ctx).First(&target, "id = ? AND enabled = ?", task.TargetID, true).Error; err != nil { + target, err := s.getTarget(ctx, task.TargetID) + if err != nil { return nil, false, err } var payload deliveryPayload @@ -455,6 +546,8 @@ func (s *DeliveryService) PollOnce(ctx context.Context) error { for _, execution := range jobs { job, err := s.awx.GetJob(ctx, execution.ExecutorJobID) if err != nil { + s.finishExecution(ctx, &execution, "failed") + _ = s.failTask(ctx, &model.DeliveryTask{ID: execution.TaskID, Status: model.TaskRunning}, model.TaskExecutionFailed, "poll AWX job "+execution.ExecutorJobID+": "+err.Error()) continue } switch strings.ToLower(job.Status) { diff --git a/server/internal/service/deployment.go b/server/internal/service/deployment.go deleted file mode 100644 index 8abee34..0000000 --- a/server/internal/service/deployment.go +++ /dev/null @@ -1,499 +0,0 @@ -package service - -import ( - "bytes" - "context" - "crypto/rand" - "encoding/hex" - "encoding/json" - "errors" - "fmt" - "io" - "net/http" - "sort" - "strings" - "sync" - "time" - - "github.com/1024XEngineer/xinfra/server/internal/config" - "github.com/1024XEngineer/xinfra/server/internal/model" - - "gorm.io/gorm" - "gorm.io/gorm/clause" -) - -const ( - DeploymentStatusPending = "pending" - DeploymentStatusRunning = "running" - DeploymentStatusSuccess = "success" - DeploymentStatusFailed = "failed" - DeploymentStatusCanceling = "canceling" - DeploymentStatusCanceled = "canceled" - - DeploymentEventLog = "log" - DeploymentEventStatus = "status" - DeploymentEventResult = "result" - DeploymentEventDone = "done" - DeploymentEventError = "error" -) - -var ( - ErrDeploymentNotConfigured = errors.New("ansible deployment service is not configured") - ErrDeploymentNotFound = errors.New("deployment not found") - ErrDeploymentForbidden = errors.New("deployment permission denied") - ErrDeploymentInvalidState = errors.New("deployment state does not allow this operation") -) - -type DeploymentCreateRequest struct { - Component string `json:"component"` - BusinessLineID uint64 `json:"business_line_id"` - Params map[string]any `json:"params"` -} - -type DeploymentEventRequest struct { - Type string `json:"type"` - Level string `json:"level"` - Status string `json:"status"` - Seq uint64 `json:"seq"` - Message string `json:"message"` - Payload map[string]any `json:"payload"` -} - -type DeploymentFinishRequest struct { - Status string `json:"status"` - ExitCode int `json:"exit_code"` - Error string `json:"error"` - Summary map[string]any `json:"summary"` -} - -type DeploymentEventEnvelope struct { - Event model.DeploymentEvent - Data map[string]any -} - -type DeploymentService struct { - cfg config.Config - db *gorm.DB - client *http.Client - - mu sync.Mutex - subscribers map[string]map[chan DeploymentEventEnvelope]struct{} -} - -func NewDeploymentService(cfg config.Config, db *gorm.DB) *DeploymentService { - return &DeploymentService{ - cfg: cfg, - db: db, - client: &http.Client{Timeout: 8 * time.Second}, - subscribers: make(map[string]map[chan DeploymentEventEnvelope]struct{}), - } -} - -func (s *DeploymentService) Create(ctx context.Context, req DeploymentCreateRequest, actorID uint64, actorName string) (model.Deployment, error) { - component := strings.TrimSpace(strings.ToLower(req.Component)) - if !isSupportedDeploymentComponent(component) { - return model.Deployment{}, fmt.Errorf("unsupported deployment component: %s", req.Component) - } - if req.BusinessLineID == 0 { - return model.Deployment{}, errors.New("business_line_id is required") - } - payload, err := json.Marshal(req) - if err != nil { - return model.Deployment{}, err - } - - var businessLine model.BusinessLine - if err := s.db.WithContext(ctx).First(&businessLine, req.BusinessLineID).Error; err != nil { - return model.Deployment{}, err - } - - deployment := model.Deployment{ - DeploymentID: newDeploymentID(), - Component: component, - BusinessLineID: req.BusinessLineID, - BusinessLine: businessLine.Name, - Status: DeploymentStatusPending, - RequestPayload: string(payload), - CreatedBy: actorID, - CreatedByName: actorName, - } - - if err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - if err := tx.Create(&deployment).Error; err != nil { - return err - } - _, err := s.appendEventTx(ctx, tx, deployment.DeploymentID, DeploymentEventStatus, "info", "deployment created", map[string]any{"status": DeploymentStatusPending}) - return err - }); err != nil { - return model.Deployment{}, err - } - - if err := s.startPythonDeployment(ctx, deployment, req.Params); err != nil { - _ = s.FailStart(ctx, deployment.DeploymentID, err) - return deployment, err - } - return deployment, nil -} - -func (s *DeploymentService) Get(ctx context.Context, deploymentID string) (model.Deployment, error) { - var deployment model.Deployment - if err := s.db.WithContext(ctx).Where("deployment_id = ?", deploymentID).First(&deployment).Error; err != nil { - if errors.Is(err, gorm.ErrRecordNotFound) { - return deployment, ErrDeploymentNotFound - } - return deployment, err - } - return deployment, nil -} - -func (s *DeploymentService) Events(ctx context.Context, deploymentID string) ([]model.DeploymentEvent, error) { - var events []model.DeploymentEvent - err := s.db.WithContext(ctx).Where("deployment_id = ?", deploymentID).Order("seq ASC").Find(&events).Error - return events, err -} - -func (s *DeploymentService) AppendEvent(ctx context.Context, deploymentID string, req DeploymentEventRequest) (model.DeploymentEvent, error) { - eventType := normalizeDeploymentEventType(req.Type) - level := strings.TrimSpace(req.Level) - if level == "" { - level = "info" - } - payload := req.Payload - if payload == nil { - payload = map[string]any{} - } - if req.Status != "" { - payload["status"] = normalizeDeploymentStatus(req.Status) - } - event, err := s.appendEvent(ctx, deploymentID, eventType, level, req.Message, payload) - if err != nil { - return event, err - } - if status, _ := payload["status"].(string); status != "" { - _ = s.updateStatus(ctx, deploymentID, status, "") - } - return event, nil -} - -func (s *DeploymentService) Finish(ctx context.Context, deploymentID string, req DeploymentFinishRequest) (model.Deployment, error) { - status := normalizeDeploymentStatus(req.Status) - if status == "" { - status = DeploymentStatusFailed - } - payload := map[string]any{ - "status": status, - "exit_code": req.ExitCode, - "summary": req.Summary, - } - if req.Error != "" { - payload["error"] = req.Error - } - result, _ := json.Marshal(payload) - now := time.Now() - var deployment model.Deployment - err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("deployment_id = ?", deploymentID).First(&deployment).Error; err != nil { - if errors.Is(err, gorm.ErrRecordNotFound) { - return ErrDeploymentNotFound - } - return err - } - updates := map[string]any{"status": status, "result_payload": string(result), "finished_at": &now} - if err := tx.Model(&model.Deployment{}).Where("deployment_id = ?", deploymentID).Updates(updates).Error; err != nil { - return err - } - eventType := DeploymentEventResult - if status == DeploymentStatusFailed || status == DeploymentStatusCanceled { - eventType = DeploymentEventError - } - if _, err := s.appendEventTx(ctx, tx, deploymentID, eventType, eventLevelForStatus(status), finishMessage(status, req.Error), payload); err != nil { - return err - } - _, err := s.appendEventTx(ctx, tx, deploymentID, DeploymentEventDone, eventLevelForStatus(status), status, map[string]any{"status": status}) - return err - }) - if err != nil { - return deployment, err - } - deployment.Status = status - deployment.ResultPayload = string(result) - deployment.FinishedAt = &now - return deployment, nil -} - -func (s *DeploymentService) Cancel(ctx context.Context, deploymentID string) error { - deployment, err := s.Get(ctx, deploymentID) - if err != nil { - return err - } - if !deploymentCancelable(deployment.Status) { - return ErrDeploymentInvalidState - } - if err := s.updateStatus(ctx, deploymentID, DeploymentStatusCanceling, "cancel requested"); err != nil { - return err - } - if s.cfg.AnsibleServiceBaseURL == "" { - return ErrDeploymentNotConfigured - } - body, _ := json.Marshal(map[string]any{"deployment_id": deploymentID}) - path := s.cfg.AnsibleServiceBaseURL + "/internal/ansible/deployments/" + deploymentID + "/cancel" - httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, path, bytes.NewReader(body)) - if err != nil { - return err - } - httpReq.Header.Set("Content-Type", "application/json") - if s.cfg.AnsibleInternalToken != "" { - httpReq.Header.Set("Authorization", "Bearer "+s.cfg.AnsibleInternalToken) - } - resp, err := s.client.Do(httpReq) - if err != nil { - return err - } - defer resp.Body.Close() - if resp.StatusCode < 200 || resp.StatusCode >= 300 { - data, _ := io.ReadAll(io.LimitReader(resp.Body, 2048)) - return fmt.Errorf("ansible cancel failed: status %d: %s", resp.StatusCode, strings.TrimSpace(string(data))) - } - return nil -} - -func (s *DeploymentService) Subscribe(deploymentID string) (chan DeploymentEventEnvelope, func()) { - ch := make(chan DeploymentEventEnvelope, 64) - s.mu.Lock() - if s.subscribers[deploymentID] == nil { - s.subscribers[deploymentID] = make(map[chan DeploymentEventEnvelope]struct{}) - } - s.subscribers[deploymentID][ch] = struct{}{} - s.mu.Unlock() - return ch, func() { - s.mu.Lock() - if subscribers := s.subscribers[deploymentID]; subscribers != nil { - delete(subscribers, ch) - if len(subscribers) == 0 { - delete(s.subscribers, deploymentID) - } - } - s.mu.Unlock() - close(ch) - } -} - -func (s *DeploymentService) FailStart(ctx context.Context, deploymentID string, cause error) error { - _, err := s.Finish(ctx, deploymentID, DeploymentFinishRequest{ - Status: DeploymentStatusFailed, - Error: cause.Error(), - }) - return err -} - -func (s *DeploymentService) startPythonDeployment(ctx context.Context, deployment model.Deployment, params map[string]any) error { - if s.cfg.AnsibleServiceBaseURL == "" { - return ErrDeploymentNotConfigured - } - callbackBaseURL := strings.TrimRight(s.cfg.DeploymentCallbackBaseURL, "/") - body, err := json.Marshal(map[string]any{ - "deployment_id": deployment.DeploymentID, - "component": deployment.Component, - "callback_url": callbackBaseURL + "/auth/internal/deployments/" + deployment.DeploymentID + "/events", - "finish_url": callbackBaseURL + "/auth/internal/deployments/" + deployment.DeploymentID + "/finish", - "params": params, - }) - if err != nil { - return err - } - httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, s.cfg.AnsibleServiceBaseURL+"/internal/ansible/deploy", bytes.NewReader(body)) - if err != nil { - return err - } - httpReq.Header.Set("Content-Type", "application/json") - if s.cfg.AnsibleInternalToken != "" { - httpReq.Header.Set("Authorization", "Bearer "+s.cfg.AnsibleInternalToken) - } - resp, err := s.client.Do(httpReq) - if err != nil { - return err - } - defer resp.Body.Close() - if resp.StatusCode < 200 || resp.StatusCode >= 300 { - data, _ := io.ReadAll(io.LimitReader(resp.Body, 2048)) - return fmt.Errorf("ansible deploy failed: status %d: %s", resp.StatusCode, strings.TrimSpace(string(data))) - } - return nil -} - -func (s *DeploymentService) appendEvent(ctx context.Context, deploymentID, eventType, level, message string, payload map[string]any) (model.DeploymentEvent, error) { - var event model.DeploymentEvent - err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - var err error - event, err = s.appendEventTx(ctx, tx, deploymentID, eventType, level, message, payload) - return err - }) - return event, err -} - -func (s *DeploymentService) appendEventTx(ctx context.Context, tx *gorm.DB, deploymentID, eventType, level, message string, payload map[string]any) (model.DeploymentEvent, error) { - var deployment model.Deployment - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("deployment_id = ?", deploymentID).First(&deployment).Error; err != nil { - if errors.Is(err, gorm.ErrRecordNotFound) { - return model.DeploymentEvent{}, ErrDeploymentNotFound - } - return model.DeploymentEvent{}, err - } - var maxSeq uint64 - if err := tx.Model(&model.DeploymentEvent{}).Where("deployment_id = ?", deploymentID).Select("COALESCE(MAX(seq), 0)").Scan(&maxSeq).Error; err != nil { - return model.DeploymentEvent{}, err - } - payloadBody, _ := json.Marshal(payload) - event := model.DeploymentEvent{ - DeploymentID: deploymentID, - Seq: maxSeq + 1, - Type: eventType, - Level: level, - Message: message, - Payload: string(payloadBody), - } - if err := tx.Create(&event).Error; err != nil { - return model.DeploymentEvent{}, err - } - go s.broadcast(event, payload) - return event, nil -} - -func (s *DeploymentService) updateStatus(ctx context.Context, deploymentID string, status string, message string) error { - status = normalizeDeploymentStatus(status) - if status == "" { - return nil - } - now := time.Now() - updates := map[string]any{"status": status} - if status == DeploymentStatusRunning { - updates["started_at"] = &now - } - if isDeploymentTerminal(status) { - updates["finished_at"] = &now - } - if err := s.db.WithContext(ctx).Model(&model.Deployment{}).Where("deployment_id = ?", deploymentID).Updates(updates).Error; err != nil { - return err - } - if message != "" { - _, err := s.appendEvent(ctx, deploymentID, DeploymentEventStatus, eventLevelForStatus(status), message, map[string]any{"status": status}) - return err - } - return nil -} - -func (s *DeploymentService) broadcast(event model.DeploymentEvent, data map[string]any) { - envelope := DeploymentEventEnvelope{Event: event, Data: eventData(event, data)} - s.mu.Lock() - subscribers := make([]chan DeploymentEventEnvelope, 0, len(s.subscribers[event.DeploymentID])) - for ch := range s.subscribers[event.DeploymentID] { - subscribers = append(subscribers, ch) - } - s.mu.Unlock() - for _, ch := range subscribers { - select { - case ch <- envelope: - default: - } - } -} - -func eventData(event model.DeploymentEvent, data map[string]any) map[string]any { - out := map[string]any{ - "deployment_id": event.DeploymentID, - "seq": event.Seq, - "type": event.Type, - "level": event.Level, - "message": event.Message, - "created_at": event.CreatedAt, - } - for key, value := range data { - out[key] = value - } - return out -} - -func DeploymentEventData(event model.DeploymentEvent) map[string]any { - payload := map[string]any{} - if strings.TrimSpace(event.Payload) != "" { - _ = json.Unmarshal([]byte(event.Payload), &payload) - } - return eventData(event, payload) -} - -func isSupportedDeploymentComponent(component string) bool { - switch component { - case "mysql", "openresty": - return true - default: - return false - } -} - -func normalizeDeploymentEventType(value string) string { - switch strings.TrimSpace(strings.ToLower(value)) { - case DeploymentEventStatus: - return DeploymentEventStatus - case DeploymentEventResult: - return DeploymentEventResult - case DeploymentEventDone: - return DeploymentEventDone - case DeploymentEventError: - return DeploymentEventError - default: - return DeploymentEventLog - } -} - -func normalizeDeploymentStatus(value string) string { - switch strings.TrimSpace(strings.ToLower(value)) { - case DeploymentStatusPending, DeploymentStatusRunning, DeploymentStatusSuccess, DeploymentStatusFailed, DeploymentStatusCanceling, DeploymentStatusCanceled: - return strings.TrimSpace(strings.ToLower(value)) - default: - return "" - } -} - -func isDeploymentTerminal(status string) bool { - return status == DeploymentStatusSuccess || status == DeploymentStatusFailed || status == DeploymentStatusCanceled -} - -func deploymentCancelable(status string) bool { - return status == DeploymentStatusPending || status == DeploymentStatusRunning || status == DeploymentStatusCanceling -} - -func eventLevelForStatus(status string) string { - if status == DeploymentStatusFailed || status == DeploymentStatusCanceled { - return "error" - } - return "info" -} - -func finishMessage(status, fallback string) string { - if fallback != "" { - return fallback - } - switch status { - case DeploymentStatusSuccess: - return "deployment completed" - case DeploymentStatusCanceled: - return "deployment canceled" - default: - return "deployment failed" - } -} - -func newDeploymentID() string { - now := time.Now() - buf := make([]byte, 3) - if _, err := rand.Read(buf); err != nil { - return fmt.Sprintf("CMP-%s-%d", now.Format("20060102"), now.UnixNano()%1000000) - } - return fmt.Sprintf("CMP-%s-%s", now.Format("20060102"), strings.ToUpper(hex.EncodeToString(buf))) -} - -func SortedDeploymentEvents(events []model.DeploymentEvent) { - sort.Slice(events, func(i, j int) bool { - return events[i].Seq < events[j].Seq - }) -}