Merge pull request #106 from 1024XEngineer/feat/base-service-delivery

Fix(base service delivery)任务日志
This commit is contained in:
ztkkOip
2026-07-27 09:47:46 +08:00
committed by GitHub
22 changed files with 1655 additions and 1244 deletions
+16
View File
@@ -0,0 +1,16 @@
.git/
.DS_Store
authserver/
frontend/node_modules/
frontend/dist/
frontend/.vite/
server/.env
server/.env.*
server/bin/
server/data/
server/*.tar
*.log
+124
View File
@@ -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<DeliveryTarget[]> {
const data = await authRequest('/auth/api/v1/delivery/targets?component=mysql')
return Array.isArray(data.items) ? data.items : []
},
async createMySQL(payload: CreateMySQLDeliveryPayload): Promise<DeliveryTask> {
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<void> {
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 }
}
}
-68
View File
@@ -1,68 +0,0 @@
import { getToken } from '@/utils/auth'
export interface DeploymentCreatePayload {
component: string
business_line_id: number
params: Record<string, unknown>
}
export interface DeploymentCreateResult {
deployment_id: string
status: string
}
export const deploymentApi = {
async create(payload: DeploymentCreatePayload): Promise<DeploymentCreateResult> {
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<void> {
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 }
}
}
+81
View File
@@ -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<TaskLogSummary[]> {
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 }
}
}
+146 -89
View File
@@ -94,6 +94,12 @@
<el-option v-for="version in activeService.versions" :key="version" :label="version" :value="version" />
</el-select>
</label>
<label class="form-field">
部署目标
<el-select v-model="selectedTargetId" placeholder="选择 AWX 部署目标" :loading="targetsLoading">
<el-option v-for="target in deliveryTargets" :key="target.id" :label="target.name" :value="target.id" />
</el-select>
</label>
</div>
<div class="config-list">
@@ -219,8 +225,8 @@
<section v-else-if="activeView === 'execution'" class="execution-view">
<div class="execution-head">
<div>
<h3>Ansible 自动交付</h3>
<p>{{ deploymentId || '创建任务后生成 deployment_id' }} · {{ activeService.runner }} · {{ activeService.template }}</p>
<h3>AWX 自动交付</h3>
<p>{{ deploymentId || '创建任务后生成 task_id' }} · {{ activeService.runner }} · {{ activeService.template }}</p>
</div>
<div class="head-actions">
<el-button :disabled="!canCancelDeployment" :loading="canceling" @click="cancelDeployment">取消任务</el-button>
@@ -265,7 +271,7 @@
</div>
<div class="result-item">
<span>健康状态</span>
<strong>{{ deliveryFailed ? '健康检查失败 · 已回退' : activeService.healthText }}</strong>
<strong>{{ deliveryFailed ? resultSubtitle : activeService.healthText }}</strong>
</div>
</div>
@@ -310,11 +316,11 @@
</template>
<script setup lang="ts">
import { computed, reactive, ref, watch } from 'vue'
import { computed, onBeforeUnmount, onMounted, reactive, ref, watch } from 'vue'
import { ElMessage } from 'element-plus'
import { Back, CircleCheck, Promotion, Refresh } from '@element-plus/icons-vue'
import { useRoute, useRouter } from 'vue-router'
import { deploymentApi } from '@/api/deployment'
import { deliveryApi, type DeliveryTarget, type TaskEvent } from '@/api/delivery'
import { useBusinessLineStore } from '@/stores/businessLine'
import { useBusinessLineMockProfile } from '@/utils/businessLineMock'
@@ -376,20 +382,18 @@ const basicServices = ref<Service[]>([
icon: 'My',
bgColor: 'var(--logo-w-bg)',
iconColor: 'var(--tag-blue-text)',
description: '一主多从标准化交付,完成后自动注册进 CloudDM / open-cdm 数据源,统一审核 SQL 上线。',
description: '单实例标准化交付,完成后自动注册进 CloudDM / open-cdm 数据源,统一审核 SQL 上线。',
playbook: 'roles/mysql-deploy',
version: 'v8.0.36',
version: 'v8.0',
status: '可交付',
template: 'mysql-delivery@v1.8.3',
runner: 'runner-02',
versions: ['MySQL 8.0.36', 'MySQL 8.0.40', 'MySQL 5.7.44'],
versions: ['MySQL 8.0'],
modes: [
{ value: 'single', label: '单实例' },
{ value: 'replica', label: '一主一从' },
{ value: 'mgr', label: 'MGR 三节点' },
],
specs: ['4C / 16G', '8C / 32G', '16C / 64G'],
disks: ['200 GB', '500 GB', '1 TB', '2 TB'],
specs: ['1C / 1G', '2C / 2G', '4C / 4G'],
disks: ['10 GB', '20 GB', '50 GB', '100 GB'],
charsets: ['utf8mb4', 'utf8'],
defaultPort: 3306,
defaultPaths: { install: '/opt/mysql', data: '/data/mysql', log: '/data/mysql-log' },
@@ -407,7 +411,7 @@ const basicServices = ref<Service[]>([
description: '容器化部署在 K8s 内,作为业务边缘网关,支持灰度路由配置。',
playbook: 'roles/openresty-deploy',
version: 'v1.25',
status: '可交付',
status: '待接入 AWX 模板',
template: 'nginx-delivery@v1.3.0',
runner: 'runner-02',
versions: ['OpenResty 1.25.3', 'OpenResty 1.21.4'],
@@ -424,6 +428,9 @@ const basicServices = ref<Service[]>([
configId: 'CFG-NGX-02841',
assetId: 'CMP-NGX-1042',
healthText: 'HTTP 200 · 23 ms',
statusColor: 'var(--text-dim)',
disabled: true,
opacity: 0.75,
},
{
key: 'pgsql',
@@ -508,7 +515,7 @@ const basicServices = ref<Service[]>([
const flowTabs = [
{ key: 'config' as FlowView, title: '01 参数配置', subtitle: '模板与资源规划' },
{ key: 'execution' as FlowView, title: '02 自动交付', subtitle: 'Ansible + 健康探测' },
{ key: 'execution' as FlowView, title: '02 自动交付', subtitle: 'AWX + 健康探测' },
{ key: 'result' as FlowView, title: '03 交付结果', subtitle: '资产与访问地址' },
]
@@ -520,15 +527,22 @@ const deliveryDone = ref(false)
const deliveryFailed = ref(false)
const canceling = ref(false)
const deploymentId = ref('')
const eventSource = ref<EventSource | null>(null)
const deliveryError = ref('')
const deliveredHost = ref('')
const deliveredPort = ref<number>()
const deliveryTargets = ref<DeliveryTarget[]>([])
const selectedTargetId = ref<number>()
const targetsLoading = ref(false)
const pollTimer = ref<number | undefined>()
const seenEventIds = ref(new Set<number>())
const deliveryLog = ref('[ready] 等待创建交付任务...')
const deliveryForm = reactive({
instanceName: `mysql-${currentName.value}-billing-02`,
version: 'MySQL 8.0.36',
mode: 'replica',
spec: '8C / 32G',
disk: '500 GB',
version: 'MySQL 8.0',
mode: 'single',
spec: '2C / 2G',
disk: '20 GB',
port: 3306,
charset: 'utf8mb4',
pool: 'auto',
@@ -567,7 +581,7 @@ const topologyNodes = computed(() => {
const runnerPreview = computed(() => {
const playbook = activeService.value?.playbook || '-'
return [
'# ansible-runner 调用预览',
'# AWX 调用预览',
`playbook: ${playbook}`,
`tenant: ${currentName.value}`,
`component: ${activeService.value?.key || '-'}`,
@@ -579,13 +593,13 @@ const runnerPreview = computed(() => {
`register_to: ${activeService.value?.registerTo || '-'}`,
].join('\n')
})
const resultAddress = computed(() => `${topologyNodes.value[0]?.ip || '10.24.18.21'}:${deliveryForm.port}`)
const resultAddress = computed(() => `${deliveredHost.value || topologyNodes.value[0]?.ip || '10.24.18.21'}:${deliveredPort.value || deliveryForm.port}`)
const resultTitle = computed(() => {
if (deliveryFailed.value) return '交付失败,已自动回退'
return activeServiceKey.value === 'mysql' ? 'MySQL 实例已交付' : 'OpenResty 集群已交付'
})
const resultSubtitle = computed(() => {
if (deliveryFailed.value) return '健康检查未通过 · 资源已释放 · 变更未交付'
if (deliveryFailed.value) return deliveryError.value || '交付失败 · 资源已释放 · 变更未交付'
return activeServiceKey.value === 'mysql' ? '全部步骤执行成功 · 用时 06:42' : '全部步骤执行成功 · 用时 02:18'
})
const deliveryStateText = computed(() => {
@@ -620,6 +634,8 @@ watch(
)
hydrateServiceDefaults()
onMounted(loadDeliveryTargets)
onBeforeUnmount(stopPolling)
function handleCardClick(service: Service) {
if (service.disabled) {
@@ -662,13 +678,16 @@ function resetWorkbench() {
}
function resetExecutionState() {
closeEventSource()
stopPolling()
precheckPassed.value = false
running.value = false
deliveryDone.value = false
deliveryFailed.value = false
deliveryError.value = ''
canceling.value = false
deploymentId.value = ''
deliveredHost.value = ''
deliveredPort.value = undefined
deliveryLog.value = '[ready] 等待创建交付任务...'
steps.value = defaultSteps().map((step) => ({ ...step, state: 'pending' }))
}
@@ -716,20 +735,25 @@ async function createTask() {
return
}
if (!activeService.value) return
if (activeService.value.key !== 'mysql') {
ElMessage.warning(`${activeService.value.name} 的 AWX 交付模板尚未接入`)
return
}
if (!selectedTargetId.value) {
ElMessage.warning('请先选择部署目标')
return
}
running.value = true
deliveryDone.value = false
deliveryFailed.value = false
activeView.value = 'execution'
steps.value = defaultSteps().map((step) => ({ ...step, state: 'pending' }))
try {
const result = await deploymentApi.create({
component: activeService.value.key,
business_line_id: businessLineId,
params: deploymentParams(),
})
deploymentId.value = result.deployment_id
deliveryLog.value += `\n[task] ${result.deployment_id} created by ${currentName.value}`
connectDeploymentEvents(result.deployment_id)
const task = await deliveryApi.createMySQL(mysqlDeliveryPayload(businessLineId))
deploymentId.value = task.id
deliveryLog.value += `\n[task] ${task.id} created by ${currentName.value}`
applyDeliveryStatus(task.status, '')
startTaskPolling(task.id)
ElMessage.success('交付任务已创建')
} catch (error) {
running.value = false
@@ -743,7 +767,7 @@ async function cancelDeployment() {
if (!deploymentId.value) return
canceling.value = true
try {
await deploymentApi.cancel(deploymentId.value)
await deliveryApi.cancel(deploymentId.value)
appendLog('[cancel] cancel requested')
} catch (error) {
ElMessage.error(error instanceof Error ? error.message : '取消任务失败')
@@ -752,98 +776,131 @@ async function cancelDeployment() {
}
}
function deploymentParams() {
function mysqlDeliveryPayload(businessLineId: number) {
const resources = parseSpec(deliveryForm.spec)
return {
business_line: currentName.value,
instance_name: deliveryForm.instanceName,
version: deliveryForm.version,
mode: deliveryForm.mode,
mode_label: currentModeLabel.value,
spec: deliveryForm.spec,
disk: deliveryForm.disk,
port: deliveryForm.port,
charset: deliveryForm.charset,
pool: deliveryForm.pool,
install_path: deliveryForm.installPath,
data_path: deliveryForm.dataPath,
log_path: deliveryForm.logPath,
business_line_id: businessLineId,
target_id: selectedTargetId.value || 0,
namespace: normalizeDNSLabel(currentName.value),
instance_name: normalizeDNSLabel(deliveryForm.instanceName),
mysql_version: mysqlVersionValue(deliveryForm.version),
cpu_milli: resources.cpuMilli,
memory_mi: resources.memoryMi,
storage_gi: parseStorageGi(deliveryForm.disk),
}
}
function connectDeploymentEvents(id: string) {
closeEventSource()
const source = new EventSource(deploymentApi.eventsURL(id))
eventSource.value = source
source.addEventListener('log', (event) => {
const data = parseSSEData(event)
appendLog(data.message)
markStepRunning()
})
source.addEventListener('status', (event) => {
const data = parseSSEData(event)
applyDeploymentStatus(String(data.status || ''), data.message ? String(data.message) : '')
})
source.addEventListener('error', (event) => {
const data = parseSSEData(event)
if (data.message) appendLog(String(data.message))
applyDeploymentStatus(String(data.status || 'failed'), '')
})
source.addEventListener('result', (event) => {
const data = parseSSEData(event)
applyDeploymentStatus(String(data.status || 'success'), data.message ? String(data.message) : '')
})
source.addEventListener('done', (event) => {
const data = parseSSEData(event)
applyDeploymentStatus(String(data.status || ''), '')
closeEventSource()
activeView.value = 'result'
})
source.onerror = () => {
appendLog('[sse] connection interrupted')
}
function startTaskPolling(id: string) {
stopPolling()
seenEventIds.value = new Set()
pollTask(id)
pollTimer.value = window.setInterval(() => pollTask(id), 3000)
}
function closeEventSource() {
eventSource.value?.close()
eventSource.value = null
function stopPolling() {
if (pollTimer.value) window.clearInterval(pollTimer.value)
pollTimer.value = undefined
}
function parseSSEData(event: Event): Record<string, unknown> {
const message = event as MessageEvent<string>
async function pollTask(id: string) {
try {
return JSON.parse(message.data || '{}')
} catch {
return {}
const data = await deliveryApi.getTask(id)
deliveredHost.value = data.task.target_host_ip || deliveredHost.value
deliveredPort.value = data.task.mysql_port || deliveredPort.value
applyTaskEvents(data.events)
applyDeliveryStatus(data.task.status, data.task.error_message || '')
if (isTerminalDeliveryStatus(data.task.status)) {
stopPolling()
activeView.value = 'result'
}
} catch (error) {
appendLog(`[poll] ${error instanceof Error ? error.message : '获取任务状态失败'}`)
}
}
function applyTaskEvents(events: TaskEvent[]) {
events.forEach((event) => {
if (seenEventIds.value.has(event.id)) return
seenEventIds.value.add(event.id)
appendLog(`[${event.to_state}] ${event.message}`)
})
}
function appendLog(message: unknown) {
if (!message) return
deliveryLog.value += `\n${String(message)}`
}
function applyDeploymentStatus(status: string, message: string) {
function applyDeliveryStatus(status: string, message: string) {
if (message) appendLog(message)
if (status === 'running') {
if (['pending', 'validating', 'dispatching', 'running', 'registering', 'canceling'].includes(status)) {
running.value = true
markStepRunning()
return
}
if (status === 'success') {
if (status === 'finished') {
running.value = false
deliveryDone.value = true
deliveryFailed.value = false
steps.value = steps.value.map((step) => ({ ...step, state: 'done' }))
return
}
if (status === 'failed' || status === 'canceled') {
if (['execution_failed', 'validation_failed', 'register_failed', 'canceled'].includes(status)) {
running.value = false
deliveryDone.value = false
deliveryFailed.value = true
deliveryError.value = message || deliveryError.value
markCurrentStepFailed()
}
}
async function loadDeliveryTargets() {
targetsLoading.value = true
try {
deliveryTargets.value = await deliveryApi.listTargets()
selectedTargetId.value = deliveryTargets.value[0]?.id
} catch (error) {
ElMessage.error(error instanceof Error ? error.message : '获取部署目标失败')
} finally {
targetsLoading.value = false
}
}
function parseSpec(spec: string) {
const cpu = Number(spec.match(/(\d+)\s*C/i)?.[1] || 1)
const memory = Number(spec.match(/\/\s*(\d+)\s*G/i)?.[1] || 1)
return {
cpuMilli: cpu * 1000,
memoryMi: memory * 1024,
}
}
function parseStorageGi(disk: string) {
const value = Number(disk.match(/(\d+)/)?.[1] || 1)
if (/TB/i.test(disk)) return value * 1024
return value
}
function mysqlVersionValue(version: string) {
const matched = version.match(/\d+(?:\.\d+)+/)?.[0] || version || '8.0'
return matched.split('.').slice(0, 2).join('.')
}
function normalizeDNSLabel(value: string) {
const normalized = value
.toLowerCase()
.replace(/[^a-z0-9-]+/g, '-')
.replace(/^-+|-+$/g, '')
.replace(/-{2,}/g, '-')
.slice(0, 63)
.replace(/-+$/g, '')
return normalized || 'default'
}
function isTerminalDeliveryStatus(status: string) {
return ['finished', 'execution_failed', 'validation_failed', 'register_failed', 'canceled'].includes(status)
}
function markStepRunning() {
const index = steps.value.findIndex((step) => step.state === 'pending' || step.state === 'running')
if (index < 0) return
+273 -63
View File
@@ -3,52 +3,67 @@
<div class="page-head">
<div>
<h1>任务中心</h1>
<p>所有 ansible-playbook / Wayne 发布任务的执行记录与实时日志</p>
<p>AWX 交付任务与 Wayne 部署服务的执行记录</p>
</div>
<div class="task-filters">
<el-select v-model="sourceFilter" size="small" class="filter-select" @change="loadTasks">
<el-option label="全部来源" value="all" />
<el-option label="AWX" value="awx" />
<el-option label="Wayne" value="wayne" />
</el-select>
<el-button size="small" :loading="loadingTasks || loadingLogs" @click="refreshCurrent">刷新</el-button>
</div>
</div>
<div class="cols-2-equal">
<div class="panel">
<div class="task-layout">
<div class="panel task-list-panel">
<div class="panel-head">
<h3>任务列表</h3>
<span class="meta">{{ lastLoadedText }}</span>
</div>
<div class="panel-body">
<table>
<tbody>
<tr v-for="task in visibleTasks" :key="task.id" class="tr-hover" :class="{ active: task.id === selectedTask }" @click="selectedTask = task.id">
<td :class="['status-text', task.statusClass]">● {{ task.status }}</td>
<td class="strong">{{ task.name }}</td>
<td class="mono text-xs text-dim">{{ task.playbook }}</td>
</tr>
</tbody>
</table>
<div class="task-list">
<button
v-for="task in tasks"
:key="task.id"
type="button"
class="task-row"
:class="{ active: task.id === selectedTaskId }"
@click="selectTask(task.id)"
>
<span :class="['task-status', task.status_class]">● {{ task.status_text }}</span>
<span class="task-main">
<span class="task-name">{{ task.name }}</span>
<span class="task-runner mono">{{ task.runner }}</span>
</span>
</button>
</div>
<div v-if="!loadingTasks && !tasks.length" class="empty-state">
暂无任务记录
</div>
<div class="pagination">
<span>共 {{ visibleTasks.length }} 条 · 当前业务线:{{ currentName }}</span>
<div class="pg-btns">
<span class="pg-btn disabled">‹</span>
<span class="pg-btn active">1</span>
<span class="pg-btn">2</span>
<span class="pg-btn">3</span>
<span class="pg-sep">…</span>
<span class="pg-btn">15</span>
<span class="pg-btn">›</span>
</div>
<span>共 {{ tasks.length }} 条 · 当前业务线:{{ currentName }}</span>
</div>
</div>
</div>
<div class="panel">
<div class="panel task-log-panel">
<div class="panel-head">
<h3>实时日志 · {{ currentName }} RKE2 节点加入</h3>
<span class="meta">WebSocket 流式输出</span>
<h3>
<span>任务日志</span>
<span class="selected-name">{{ selectedTaskName }}</span>
</h3>
<span class="meta">{{ selectedTaskMeta }}</span>
</div>
<div class="panel-body log-stream">
<div v-for="(log, index) in visibleLogs" :key="index" :class="['task-log-line', log.class]">
<div v-for="(log, index) in logs" :key="index" :class="['task-log-line', log.class]">
<span class="t">{{ log.time }}</span>{{ log.message }}
</div>
<div class="task-log-line">
<span class="t">10:42:33</span>
<span class="blink">▌</span> 等待节点 Ready...
<div v-if="loadingLogs" class="task-log-line">
<span class="t">...</span>加载中
</div>
<div v-if="!loadingLogs && !logs.length" class="empty-state">
请选择任务查看日志
</div>
</div>
</div>
@@ -57,58 +72,253 @@
</template>
<script setup lang="ts">
import { computed, ref } from 'vue'
import { computed, onMounted, ref, watch } from 'vue'
import { taskLogApi, type TaskLogLine, type TaskLogSummary } from '@/api/taskLog'
import { useBusinessLineStore } from '@/stores/businessLine'
import { useBusinessLineMockProfile } from '@/utils/businessLineMock'
const { currentName } = useBusinessLineMockProfile()
const businessLineStore = useBusinessLineStore()
const selectedTask = ref(1)
const sourceFilter = ref('all')
const selectedTaskId = ref('')
const tasks = ref<TaskLogSummary[]>([])
const logs = ref<TaskLogLine[]>([])
const loadingTasks = ref(false)
const loadingLogs = ref(false)
const lastLoadedAt = ref<Date | null>(null)
const tasks = ref([
{ id: 1, name: 'RKE2 节点加入 · bj-node-061', playbook: 'roles/rke2-node-join', status: '执行中', statusClass: 'warn' },
{ id: 2, name: 'MySQL 主从部署 · kodo', playbook: 'roles/mysql-deploy', status: '成功', statusClass: 'ok' },
{ id: 3, name: 'Redis Cluster 部署 · linxi', playbook: 'roles/redis-deploy', status: '成功', statusClass: 'ok' },
{ id: 4, name: 'openresty 网关部署', playbook: 'roles/openresty-deploy', status: '失败', statusClass: 'err' },
{ id: 5, name: '业务线标签同步 · LDAP', playbook: 'internal/label-sync', status: '成功', statusClass: 'ok' },
{ id: 6, name: 'MySQL 主从部署 · xinfra', playbook: 'roles/mysql-deploy', status: '成功', statusClass: 'ok' },
])
const selectedTask = computed(() => tasks.value.find((task) => task.id === selectedTaskId.value))
const selectedTaskName = computed(() => selectedTask.value?.name || '未选择')
const selectedTaskMeta = computed(() => selectedTask.value ? `${selectedTask.value.source.toUpperCase()} · ${selectedTask.value.runner}` : '接口日志')
const lastLoadedText = computed(() => lastLoadedAt.value ? `更新于 ${formatTime(lastLoadedAt.value)}` : '')
const logs = ref([
{ time: '10:42:01', message: 'PLAY [rke2-node-join] **********************', class: '' },
{ time: '10:42:02', message: 'TASK [初始化内核参数] ...', class: '' },
{ time: '10:42:04', message: 'ok: [bj-node-061]', class: 'ok' },
{ time: '10:42:05', message: 'TASK [安装 containerd] ...', class: '' },
{ time: '10:42:18', message: 'ok: [bj-node-061]', class: 'ok' },
{ time: '10:42:19', message: 'TASK [写入 RKE2 node-labels: business-line=kodo] ...', class: '' },
{ time: '10:42:20', message: 'changed: [bj-node-061] => labels applied', class: 'tag-ok' },
{ time: '10:42:21', message: 'TASK [加入集群 rke2-bj-prod-01] ...', class: '' },
])
async function loadTasks() {
loadingTasks.value = true
try {
tasks.value = await taskLogApi.list({
source: sourceFilter.value,
businessLineId: businessLineStore.current?.id,
})
lastLoadedAt.value = new Date()
if (!tasks.value.some((task) => task.id === selectedTaskId.value)) {
selectedTaskId.value = tasks.value[0]?.id || ''
}
if (selectedTaskId.value) {
await loadLogs(selectedTaskId.value)
} else {
logs.value = []
}
} catch (error) {
logs.value = [{ time: formatTime(new Date()), message: error instanceof Error ? error.message : '任务日志加载失败', class: 'err' }]
} finally {
loadingTasks.value = false
}
}
const visibleTasks = computed(() => tasks.value.map((task) => ({
...task,
name: task.name.replace(/ · (kodo|linxi|xinfra|las)|$/, ` · ${currentName.value}`),
})))
async function loadLogs(taskId: string) {
loadingLogs.value = true
try {
const data = await taskLogApi.get(taskId)
logs.value = data.lines
lastLoadedAt.value = new Date()
} catch (error) {
logs.value = [{ time: formatTime(new Date()), message: error instanceof Error ? error.message : '任务日志加载失败', class: 'err' }]
} finally {
loadingLogs.value = false
}
}
const visibleLogs = computed(() => logs.value.map((log) => ({
...log,
message: log.message.replace(/business-line=(kodo|linxi|xinfra|las)/, `business-line=${currentName.value}`),
})))
function selectTask(taskId: string) {
if (selectedTaskId.value === taskId) {
return
}
selectedTaskId.value = taskId
void loadLogs(taskId)
}
function refreshCurrent() {
void loadTasks()
}
function formatTime(date: Date) {
return date.toLocaleTimeString('zh-CN', { hour12: false })
}
watch(
() => businessLineStore.current?.id,
() => {
void loadTasks()
},
)
onMounted(() => {
void loadTasks()
})
</script>
<style scoped>
/* 公共样式已在 global.css 中定义 */
.log-stream {
max-height: 340px;
.task-layout {
display: grid;
grid-template-columns: minmax(320px, 380px) minmax(0, 1fr);
gap: 18px;
align-items: start;
}
.task-list-panel,
.task-log-panel {
min-width: 0;
}
.task-list-panel .panel-body {
padding-top: 0;
}
.task-list {
max-height: calc(100vh - 245px);
min-height: 360px;
overflow-y: auto;
}
.blink {
animation: blink 1s step-start infinite;
.task-row {
width: 100%;
min-height: 62px;
display: grid;
grid-template-columns: 72px minmax(0, 1fr);
gap: 12px;
align-items: center;
padding: 10px 16px;
border: 0;
border-bottom: 1px solid var(--line-soft);
background: transparent;
color: inherit;
text-align: left;
cursor: pointer;
}
.task-row:hover,
.task-row.active {
background: var(--bg-panel-2);
}
.task-status {
font-size: 12px;
white-space: nowrap;
}
.task-status.ok {
color: var(--accent);
}
@keyframes blink {
50% { opacity: 0; }
.task-status.warn {
color: var(--warn);
}
.task-status.err {
color: var(--err);
}
.task-main {
min-width: 0;
display: flex;
flex-direction: column;
gap: 5px;
}
.task-name,
.task-runner,
.selected-name {
overflow: hidden;
text-overflow: ellipsis;
white-space: nowrap;
}
.task-name {
color: var(--text-hi);
font-size: 12.5px;
font-weight: 600;
}
.task-runner {
color: var(--text-dim);
font-size: 11px;
}
.task-log-panel .panel-head {
gap: 14px;
}
.task-log-panel .panel-head h3 {
min-width: 0;
display: flex;
align-items: center;
gap: 10px;
}
.selected-name {
color: var(--text-dim);
font-weight: 500;
}
.log-stream {
max-height: calc(100vh - 210px);
min-height: 520px;
overflow-y: auto;
overflow-x: auto;
}
.task-filters {
display: flex;
gap: 10px;
align-items: center;
flex-wrap: wrap;
}
.filter-select {
width: 160px;
}
.empty-state {
padding: 18px 0;
color: var(--text-dim);
font-size: 13px;
}
@media (max-width: 1180px) {
.task-layout {
grid-template-columns: minmax(280px, 340px) minmax(0, 1fr);
}
}
@media (max-width: 1024px) {
.task-layout {
grid-template-columns: 1fr;
}
.task-list,
.log-stream {
max-height: none;
}
.log-stream {
min-height: 420px;
}
}
@media (max-width: 640px) {
.filter-select {
width: 100%;
}
.task-filters {
width: 100%;
}
.task-row {
grid-template-columns: 64px minmax(0, 1fr);
padding: 10px 12px;
}
}
</style>
+1 -1
View File
@@ -2,8 +2,8 @@
.git/
.DS_Store
.env
.env.*
bin/
data/
web/node_modules/
web/dist/
+2 -3
View File
@@ -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=
+112 -114
View File
@@ -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),
}
}
-3
View File
@@ -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{},
+4 -47
View File
@@ -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"`
-249
View File
@@ -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
}
+347
View File
@@ -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
}
}
-12
View File
@@ -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"`
-28
View File
@@ -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"`
+2 -2
View File
@@ -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 {
+3 -9
View File
@@ -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)
}
}
+234 -36
View File
@@ -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
}
+36
View File
@@ -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)
}
}
@@ -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 {
+114 -21
View File
@@ -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(&quota).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(&quota).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) {
-499
View File
@@ -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
})
}