3666 lines
165 KiB
Python
3666 lines
165 KiB
Python
from __future__ import annotations
|
||
|
||
import json
|
||
import os
|
||
import re
|
||
import socket
|
||
import subprocess
|
||
from datetime import datetime
|
||
from math import ceil
|
||
from pathlib import Path
|
||
from threading import Lock
|
||
from urllib.parse import urlsplit, urlunsplit
|
||
from uuid import uuid4
|
||
|
||
from psycopg2 import errors
|
||
|
||
from app.core.db import get_db
|
||
from app.services.build_info_service import get_runtime_build_info
|
||
from app.services.ops_command_service import build_bash_command
|
||
from app.services.ops_execution_mode_service import execution_mode_label
|
||
|
||
|
||
_RELEASE_SCHEMA_SQL = """
|
||
CREATE TABLE IF NOT EXISTS ops_releases (
|
||
id BIGSERIAL PRIMARY KEY,
|
||
release_version VARCHAR(128) NOT NULL UNIQUE,
|
||
channel VARCHAR(64) NOT NULL DEFAULT 'stable',
|
||
commit_sha VARCHAR(64) NOT NULL DEFAULT '',
|
||
status VARCHAR(32) NOT NULL DEFAULT 'draft',
|
||
artifact_url TEXT NOT NULL DEFAULT '',
|
||
checksum VARCHAR(128) NOT NULL DEFAULT '',
|
||
notes TEXT NOT NULL DEFAULT '',
|
||
metadata_json JSONB,
|
||
created_by VARCHAR(64) NOT NULL DEFAULT 'api',
|
||
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||
activated_at TIMESTAMP,
|
||
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||
);
|
||
|
||
CREATE INDEX IF NOT EXISTS idx_ops_releases_channel_status
|
||
ON ops_releases(channel, status, created_at DESC);
|
||
|
||
CREATE TABLE IF NOT EXISTS ops_release_rollouts (
|
||
id BIGSERIAL PRIMARY KEY,
|
||
release_id BIGINT NOT NULL REFERENCES ops_releases(id) ON DELETE CASCADE,
|
||
rollout_code VARCHAR(64) NOT NULL DEFAULT '',
|
||
target_selector_json JSONB,
|
||
target_nodes_json JSONB,
|
||
policy_json JSONB,
|
||
status VARCHAR(32) NOT NULL DEFAULT 'planned',
|
||
batch_cursor INTEGER NOT NULL DEFAULT 0,
|
||
batches_total INTEGER NOT NULL DEFAULT 0,
|
||
jobs_total INTEGER NOT NULL DEFAULT 0,
|
||
jobs_created INTEGER NOT NULL DEFAULT 0,
|
||
result_summary_json JSONB,
|
||
created_by VARCHAR(64) NOT NULL DEFAULT 'api',
|
||
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||
);
|
||
|
||
CREATE INDEX IF NOT EXISTS idx_ops_release_rollouts_release
|
||
ON ops_release_rollouts(release_id, created_at DESC);
|
||
|
||
ALTER TABLE ops_release_rollouts ADD COLUMN IF NOT EXISTS rollout_code VARCHAR(64) NOT NULL DEFAULT '';
|
||
ALTER TABLE ops_release_rollouts ADD COLUMN IF NOT EXISTS target_nodes_json JSONB;
|
||
ALTER TABLE ops_release_rollouts ADD COLUMN IF NOT EXISTS batch_cursor INTEGER NOT NULL DEFAULT 0;
|
||
ALTER TABLE ops_release_rollouts ADD COLUMN IF NOT EXISTS batches_total INTEGER NOT NULL DEFAULT 0;
|
||
ALTER TABLE ops_release_rollouts ADD COLUMN IF NOT EXISTS jobs_total INTEGER NOT NULL DEFAULT 0;
|
||
ALTER TABLE ops_release_rollouts ADD COLUMN IF NOT EXISTS jobs_created INTEGER NOT NULL DEFAULT 0;
|
||
ALTER TABLE ops_release_rollouts ADD COLUMN IF NOT EXISTS result_summary_json JSONB;
|
||
"""
|
||
|
||
_ACTIVE_JOB_STATUSES = {"queued", "dispatching", "running", "awaiting_approval"}
|
||
_ISSUE_JOB_STATUSES = {"failed", "cancelled", "blocked", "partially_succeeded"}
|
||
_ROLLOUT_INSPECTION_ACTION_KEYS = ("health.snapshot", "logs.collect", "diagnostics.collect")
|
||
_SAFE_RELEASE_PACKAGE_NAME_RE = re.compile(r"^[A-Za-z0-9._-]+$")
|
||
_RELEASE_SCHEMA_LOCK = Lock()
|
||
_RELEASE_SCHEMA_READY = False
|
||
_RELEASE_SCHEMA_ADVISORY_LOCK_KEY = 90421803
|
||
_RELEASE_REQUIRED_TABLES = ("ops_releases", "ops_release_rollouts")
|
||
_RELEASE_REQUIRED_COLUMNS = {
|
||
"ops_release_rollouts": (
|
||
"rollout_code",
|
||
"target_nodes_json",
|
||
"batch_cursor",
|
||
"batches_total",
|
||
"jobs_total",
|
||
"jobs_created",
|
||
"result_summary_json",
|
||
),
|
||
}
|
||
|
||
|
||
def ensure_ops_release_schema() -> None:
|
||
global _RELEASE_SCHEMA_READY
|
||
|
||
if _RELEASE_SCHEMA_READY:
|
||
return
|
||
|
||
with _RELEASE_SCHEMA_LOCK:
|
||
if _RELEASE_SCHEMA_READY:
|
||
return
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
if _ops_release_schema_basics_present(cur):
|
||
_RELEASE_SCHEMA_READY = True
|
||
return
|
||
conn.autocommit = False
|
||
try:
|
||
with conn.cursor() as cur:
|
||
cur.execute("SELECT pg_advisory_xact_lock(%s)", (_RELEASE_SCHEMA_ADVISORY_LOCK_KEY,))
|
||
cur.execute(_RELEASE_SCHEMA_SQL)
|
||
conn.commit()
|
||
except Exception as exc:
|
||
recoverable = isinstance(exc, (errors.DeadlockDetected, errors.LockNotAvailable))
|
||
try:
|
||
conn.rollback()
|
||
except Exception:
|
||
pass
|
||
if not recoverable:
|
||
raise
|
||
with conn.cursor() as cur:
|
||
if not _ops_release_schema_basics_present(cur):
|
||
raise
|
||
_RELEASE_SCHEMA_READY = True
|
||
|
||
|
||
def _ops_release_schema_basics_present(cur) -> bool:
|
||
for table_name in _RELEASE_REQUIRED_TABLES:
|
||
cur.execute("SELECT to_regclass(%s)", (f"public.{table_name}",))
|
||
row = cur.fetchone()
|
||
if not row or not row[0]:
|
||
return False
|
||
|
||
for table_name, required_columns in _RELEASE_REQUIRED_COLUMNS.items():
|
||
cur.execute(
|
||
"""
|
||
SELECT column_name
|
||
FROM information_schema.columns
|
||
WHERE table_schema = 'public' AND table_name = %s
|
||
""",
|
||
(table_name,),
|
||
)
|
||
existing_columns = {str(row[0] or "").strip() for row in list(cur.fetchall() or [])}
|
||
if not set(required_columns).issubset(existing_columns):
|
||
return False
|
||
return True
|
||
|
||
|
||
def _decode_json(value: object) -> dict:
|
||
if isinstance(value, dict):
|
||
return value
|
||
if value in (None, ""):
|
||
return {}
|
||
try:
|
||
return json.loads(value)
|
||
except Exception:
|
||
return {}
|
||
|
||
|
||
def _decode_items(value: object) -> list[dict]:
|
||
decoded = _decode_json(value)
|
||
items = decoded.get("items") if isinstance(decoded, dict) else []
|
||
return [item for item in list(items or []) if isinstance(item, dict)]
|
||
|
||
|
||
def _ts(value: object) -> str:
|
||
if isinstance(value, datetime):
|
||
return value.isoformat(sep=" ", timespec="seconds")
|
||
return ""
|
||
|
||
|
||
def _release_status_label(status: object) -> str:
|
||
normalized_status = str(status or "").strip()
|
||
mapping = {
|
||
"draft": "草稿",
|
||
"ready": "已就绪",
|
||
"active": "已激活",
|
||
"archived": "已归档",
|
||
}
|
||
return mapping.get(normalized_status, normalized_status or "未知")
|
||
|
||
|
||
def _rollout_status_label(status: object) -> str:
|
||
normalized_status = str(status or "").strip()
|
||
mapping = {
|
||
"planned": "待推进",
|
||
"running": "执行中",
|
||
"awaiting_approval": "待审批",
|
||
"ready_for_next_batch": "可推进下一批",
|
||
"halted": "已暂停",
|
||
"completed": "已完成",
|
||
"completed_with_issues": "带问题完成",
|
||
"failed": "失败",
|
||
}
|
||
return mapping.get(normalized_status, normalized_status or "未知")
|
||
|
||
|
||
def _build_release_focus_ref(
|
||
release: dict | None = None,
|
||
*,
|
||
rollout_id: int = 0,
|
||
rollout_code: str = "",
|
||
section: str = "",
|
||
) -> dict:
|
||
normalized_release = dict(release or {})
|
||
return {
|
||
"kind": "release_hub",
|
||
"release_id": int(normalized_release.get("id") or 0),
|
||
"release_version": str(normalized_release.get("release_version") or "").strip(),
|
||
"channel": str(normalized_release.get("channel") or "").strip(),
|
||
"rollout_id": int(rollout_id or 0),
|
||
"rollout_code": str(rollout_code or "").strip(),
|
||
"section": str(section or "").strip(),
|
||
}
|
||
|
||
|
||
def _build_rollout_focus_ref(rollout: dict | None = None, *, section: str = "rollout_jobs") -> dict:
|
||
normalized_rollout = dict(rollout or {})
|
||
return {
|
||
"kind": "release_rollout",
|
||
"rollout_id": int(normalized_rollout.get("id") or 0),
|
||
"rollout_code": str(normalized_rollout.get("rollout_code") or "").strip(),
|
||
"release_id": int(normalized_rollout.get("release_id") or 0),
|
||
"section": str(section or "").strip(),
|
||
}
|
||
|
||
|
||
def _build_release_summary_text(release: dict | None = None) -> str:
|
||
normalized_release = dict(release or {})
|
||
status = str(normalized_release.get("status") or "").strip()
|
||
release_version = str(normalized_release.get("release_version") or "").strip() or "当前版本"
|
||
artifact_url = str(normalized_release.get("artifact_url") or "").strip()
|
||
checksum = str(normalized_release.get("checksum") or "").strip()
|
||
|
||
if status == "active":
|
||
return f"{release_version} 当前为激活版本,可继续进入 Rollout 评估。"
|
||
if status == "ready":
|
||
if artifact_url and checksum:
|
||
return f"{release_version} 已就绪,制品与校验信息齐备,可进入 Rollout。"
|
||
if artifact_url:
|
||
return f"{release_version} 已就绪,但建议补齐 checksum 后再推进正式发布。"
|
||
return f"{release_version} 已就绪,但仍缺少可直接下发的制品地址。"
|
||
if status == "draft":
|
||
return f"{release_version} 仍为草稿,建议先补齐制品、校验信息与说明。"
|
||
if status == "archived":
|
||
return f"{release_version} 已归档,通常不再作为默认发布入口。"
|
||
return f"{release_version} 当前状态为 {_release_status_label(status)}。"
|
||
|
||
|
||
def _build_rollout_summary_text(rollout: dict | None = None) -> str:
|
||
normalized_rollout = dict(rollout or {})
|
||
result_summary = dict(normalized_rollout.get("result_summary") or {})
|
||
status = str(normalized_rollout.get("status") or result_summary.get("status") or "").strip()
|
||
target_nodes = list(normalized_rollout.get("target_nodes") or [])
|
||
total_targets = len(target_nodes) or int(normalized_rollout.get("jobs_total") or 0)
|
||
launched_targets = int(result_summary.get("launched_targets", normalized_rollout.get("batch_cursor", 0)) or 0)
|
||
remaining_targets = int(result_summary.get("remaining_targets", max(total_targets - launched_targets, 0)) or 0)
|
||
active_jobs = int(result_summary.get("active_jobs", 0) or 0)
|
||
issue_jobs = int(result_summary.get("issue_jobs", 0) or 0)
|
||
batches_total = int(normalized_rollout.get("batches_total") or result_summary.get("batches_total") or 0)
|
||
batch_index = int(result_summary.get("current_batch_index", normalized_rollout.get("batch_cursor", 0)) or 0)
|
||
|
||
if status == "running":
|
||
return f"当前已覆盖 {launched_targets}/{total_targets} 台节点,仍有 {active_jobs} 条任务执行中。"
|
||
if status == "awaiting_approval":
|
||
return f"当前批次等待审批,已覆盖 {launched_targets}/{total_targets} 台节点。"
|
||
if status == "ready_for_next_batch":
|
||
return f"当前批次已收口,可继续推进下一批;剩余节点 {remaining_targets} 台。"
|
||
if status == "halted":
|
||
return f"当前 rollout 已暂停,存在 {issue_jobs} 条待处理问题任务。"
|
||
if status == "completed":
|
||
return f"全部 {total_targets} 台目标节点已完成发布收口。"
|
||
if status == "completed_with_issues":
|
||
return f"发布已完成,但仍残留 {issue_jobs} 条问题任务,建议先复盘再继续。"
|
||
if status == "failed":
|
||
return f"当前 rollout 执行失败,建议进入任务列表查看失败节点与阻断原因。"
|
||
if status == "planned":
|
||
batch_text = f"{batch_index}/{batches_total}" if batches_total > 0 else "0/0"
|
||
return f"当前 rollout 已创建,目标节点 {total_targets} 台,批次 {batch_text},等待首批推进。"
|
||
return f"当前 rollout 状态为 {_rollout_status_label(status)}。"
|
||
|
||
|
||
def _extract_release_manifest(payload: dict | None = None) -> dict:
|
||
normalized_payload = dict(payload or {})
|
||
metadata = dict(normalized_payload.get("metadata") or {})
|
||
manifest = dict(metadata.get("manifest") or normalized_payload.get("manifest") or {})
|
||
artifact_kind = (
|
||
str(manifest.get("artifact_kind") or "").strip()
|
||
or str(normalized_payload.get("package_type") or manifest.get("package_type") or "").strip()
|
||
)
|
||
package_name = (
|
||
str(manifest.get("package_name") or "").strip()
|
||
or str(normalized_payload.get("package_name") or "").strip()
|
||
or str(normalized_payload.get("release_version") or "").strip()
|
||
)
|
||
checksum = (
|
||
str(normalized_payload.get("checksum") or "").strip()
|
||
or str(normalized_payload.get("sha256") or "").strip()
|
||
or str(manifest.get("checksum") or "").strip()
|
||
)
|
||
commit_sha = (
|
||
str(normalized_payload.get("commit_sha") or "").strip()
|
||
or str(manifest.get("commit_sha") or "").strip()
|
||
)
|
||
built_at = (
|
||
str(manifest.get("built_at") or "").strip()
|
||
or str(normalized_payload.get("generated_at") or "").strip()
|
||
or str(normalized_payload.get("created_at") or "").strip()
|
||
)
|
||
normalized_manifest = {
|
||
"artifact_kind": artifact_kind,
|
||
"package_name": package_name,
|
||
"package_size_bytes": int(manifest.get("package_size_bytes", 0) or 0),
|
||
"checksum": checksum,
|
||
"build_id": str(manifest.get("build_id") or "").strip(),
|
||
"commit_sha": commit_sha,
|
||
"built_at": built_at,
|
||
"included_services": list(manifest.get("included_services") or []),
|
||
"required_env_keys": list(manifest.get("required_env_keys") or []),
|
||
"health_checks": list(manifest.get("health_checks") or []),
|
||
"rollback_hint": str(manifest.get("rollback_hint") or "").strip(),
|
||
}
|
||
if any(
|
||
[
|
||
normalized_manifest["artifact_kind"],
|
||
normalized_manifest["package_name"],
|
||
normalized_manifest["checksum"],
|
||
normalized_manifest["build_id"],
|
||
normalized_manifest["commit_sha"],
|
||
normalized_manifest["built_at"],
|
||
normalized_manifest["included_services"],
|
||
normalized_manifest["required_env_keys"],
|
||
normalized_manifest["health_checks"],
|
||
normalized_manifest["rollback_hint"],
|
||
normalized_manifest["package_size_bytes"] > 0,
|
||
]
|
||
):
|
||
return normalized_manifest
|
||
return {}
|
||
|
||
|
||
def _extract_node_release_version(*candidates: dict) -> str:
|
||
keys = (
|
||
"current_release_version",
|
||
"release_version",
|
||
"active_release_version",
|
||
"current_version",
|
||
"deployed_release_version",
|
||
)
|
||
for candidate in candidates:
|
||
normalized_candidate = dict(candidate or {})
|
||
for key in keys:
|
||
value = str(normalized_candidate.get(key) or "").strip()
|
||
if value:
|
||
return value
|
||
metadata = dict(normalized_candidate.get("metadata") or {})
|
||
for key in keys:
|
||
value = str(metadata.get(key) or "").strip()
|
||
if value:
|
||
return value
|
||
return ""
|
||
|
||
|
||
def _build_rollout_readiness_focus_ref(release: dict | None = None, *, node_code: str = "") -> dict:
|
||
focus_ref = _build_release_focus_ref(release, section="default_rollout_gate")
|
||
if str(node_code or "").strip():
|
||
focus_ref["node_code"] = str(node_code).strip()
|
||
return focus_ref
|
||
|
||
|
||
def _rollout_readiness_risk_level(*, blocking_reasons: list[str], inspection_status: str) -> str:
|
||
if list(blocking_reasons or []):
|
||
return "blocked"
|
||
normalized_status = str(inspection_status or "").strip()
|
||
if normalized_status == "attention":
|
||
return "high"
|
||
if normalized_status in {"missing", "running"}:
|
||
return "medium"
|
||
return "low"
|
||
|
||
|
||
def _parse_time(value: object) -> datetime | None:
|
||
if isinstance(value, datetime):
|
||
return value
|
||
raw_value = str(value or "").strip()
|
||
if not raw_value:
|
||
return None
|
||
normalized = raw_value.replace("Z", "+00:00")
|
||
try:
|
||
return datetime.fromisoformat(normalized)
|
||
except Exception:
|
||
pass
|
||
for fmt in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%d %H:%M:%S.%f"):
|
||
try:
|
||
return datetime.strptime(raw_value, fmt)
|
||
except ValueError:
|
||
continue
|
||
return None
|
||
|
||
|
||
def _serialize_release_row(row: tuple) -> dict:
|
||
serialized = {
|
||
"id": int(row[0]),
|
||
"release_version": str(row[1] or ""),
|
||
"channel": str(row[2] or ""),
|
||
"commit_sha": str(row[3] or ""),
|
||
"status": str(row[4] or ""),
|
||
"artifact_url": str(row[5] or ""),
|
||
"checksum": str(row[6] or ""),
|
||
"notes": str(row[7] or ""),
|
||
"metadata": _decode_json(row[8]),
|
||
"created_by": str(row[9] or ""),
|
||
"created_at": _ts(row[10]),
|
||
"activated_at": _ts(row[11]),
|
||
"updated_at": _ts(row[12]),
|
||
}
|
||
serialized["manifest"] = _extract_release_manifest(serialized)
|
||
serialized["status_label"] = _release_status_label(serialized.get("status"))
|
||
serialized["summary"] = _build_release_summary_text(serialized)
|
||
serialized["summary_text"] = str(serialized.get("summary") or "")
|
||
serialized["focus_ref"] = _build_release_focus_ref(serialized, section="release_detail")
|
||
return serialized
|
||
|
||
|
||
def _serialize_rollout_row(row: tuple) -> dict:
|
||
target_nodes = _decode_items(row[4])
|
||
policy = _decode_json(row[5])
|
||
serialized = {
|
||
"id": int(row[0]),
|
||
"release_id": int(row[1]),
|
||
"rollout_code": str(row[2] or ""),
|
||
"target_selector": _decode_json(row[3]),
|
||
"target_nodes": target_nodes,
|
||
"target_node_codes": [
|
||
str(item.get("node_code") or "").strip()
|
||
for item in target_nodes
|
||
if str(item.get("node_code") or "").strip()
|
||
],
|
||
"policy": policy,
|
||
"status": str(row[6] or ""),
|
||
"batch_cursor": int(row[7] or 0),
|
||
"batches_total": int(row[8] or 0),
|
||
"jobs_total": int(row[9] or 0),
|
||
"jobs_created": int(row[10] or 0),
|
||
"result_summary": _decode_json(row[11]),
|
||
"created_by": str(row[12] or ""),
|
||
"created_at": _ts(row[13]),
|
||
"updated_at": _ts(row[14]),
|
||
}
|
||
serialized["status_label"] = _rollout_status_label(serialized.get("status"))
|
||
serialized["summary"] = _build_rollout_summary_text(serialized)
|
||
serialized["summary_text"] = str(serialized.get("summary") or "")
|
||
serialized["focus_ref"] = _build_rollout_focus_ref(serialized)
|
||
return serialized
|
||
|
||
|
||
def _build_rollout_code() -> str:
|
||
return f"rollout-{datetime.now().strftime('%Y%m%d%H%M%S')}-{uuid4().hex[:4]}"
|
||
|
||
|
||
def _normalize_rollout_policy(policy: dict, *, total_targets: int) -> dict:
|
||
input_policy = dict(policy or {})
|
||
safe_total = max(int(total_targets or 0), 0)
|
||
requested_batch_size = int(input_policy.get("batch_size") or safe_total or 1)
|
||
batch_size = max(1, min(requested_batch_size, safe_total or requested_batch_size))
|
||
requested_first_batch_size = int(
|
||
input_policy.get("first_batch_size")
|
||
or input_policy.get("canary_size")
|
||
or batch_size
|
||
)
|
||
first_batch_size = max(1, min(requested_first_batch_size, safe_total or requested_first_batch_size))
|
||
if safe_total == 0:
|
||
first_batch_size = 0
|
||
batch_size = 0
|
||
elif batch_size < first_batch_size:
|
||
batch_size = first_batch_size
|
||
|
||
deploy_payload = dict(input_policy.get("deploy_payload") or {})
|
||
deploy_payload.setdefault("restart_services", [])
|
||
deploy_payload.setdefault("health_check_urls", [])
|
||
deploy_payload.setdefault("health_check_timeout_seconds", 10)
|
||
deploy_payload.setdefault("health_check_retries", 2)
|
||
deploy_payload.setdefault("health_check_interval_seconds", 2)
|
||
deploy_payload.setdefault("rollback_on_failure", True)
|
||
deploy_payload.setdefault("switch_current", True)
|
||
|
||
normalized = {
|
||
"auto_start": bool(input_policy.get("auto_start", True)),
|
||
"auto_dispatch": bool(input_policy.get("auto_dispatch", False)),
|
||
"auto_approve": bool(input_policy.get("auto_approve", False)),
|
||
"auto_advance": bool(input_policy.get("auto_advance", True)),
|
||
"continue_on_failure": bool(input_policy.get("continue_on_failure", False)),
|
||
"continue_on_partial": bool(input_policy.get("continue_on_partial", False)),
|
||
"execution_mode": str(input_policy.get("execution_mode") or "remote-agent").strip() or "remote-agent",
|
||
"batch_size": batch_size,
|
||
"first_batch_size": first_batch_size,
|
||
"pause_seconds": max(0, int(input_policy.get("pause_seconds") or 0)),
|
||
"deploy_payload": deploy_payload,
|
||
}
|
||
return normalized
|
||
|
||
|
||
def _calculate_total_batches(total_targets: int, policy: dict) -> int:
|
||
safe_total = max(int(total_targets or 0), 0)
|
||
if safe_total <= 0:
|
||
return 0
|
||
first_batch_size = max(1, int(policy.get("first_batch_size") or 1))
|
||
batch_size = max(1, int(policy.get("batch_size") or first_batch_size))
|
||
if safe_total <= first_batch_size:
|
||
return 1
|
||
remaining = safe_total - first_batch_size
|
||
return 1 + int(ceil(remaining / batch_size))
|
||
|
||
|
||
def _decode_payload_json(value: object) -> dict:
|
||
return _decode_json(value)
|
||
|
||
|
||
def _release_package_root() -> Path:
|
||
configured = str(os.environ.get("DOMAINCHECK_RELEASE_ROOT") or "").strip()
|
||
if configured:
|
||
return Path(configured).expanduser().resolve()
|
||
return Path(__file__).resolve().parents[3] / "release"
|
||
|
||
|
||
def _release_repo_root() -> Path:
|
||
return Path(__file__).resolve().parents[3]
|
||
|
||
|
||
def _trim_release_command_output(value: object, limit: int = 4000) -> str:
|
||
text = str(value or "").strip()
|
||
if len(text) <= int(limit):
|
||
return text
|
||
return text[-int(limit):]
|
||
|
||
|
||
def _safe_release_package_name(package_name: str) -> str:
|
||
normalized = str(package_name or "").strip()
|
||
if not normalized or not _SAFE_RELEASE_PACKAGE_NAME_RE.fullmatch(normalized):
|
||
raise ValueError("invalid package_name")
|
||
return normalized
|
||
|
||
|
||
def _normalize_public_base_url(raw_value: str) -> str:
|
||
normalized = str(raw_value or "").strip().rstrip("/")
|
||
if normalized.endswith("/api/v1"):
|
||
normalized = normalized[: -len("/api/v1")]
|
||
return _rewrite_loopback_control_plane_url(normalized)
|
||
|
||
|
||
def _is_loopback_hostname(hostname: str) -> bool:
|
||
normalized = str(hostname or "").strip().lower().strip("[]")
|
||
return normalized in {"127.0.0.1", "localhost", "0.0.0.0", "::1"}
|
||
|
||
|
||
def _build_url_with_host(raw_url: str, *, host: str, scheme: str = "", port: int | None = None) -> str:
|
||
normalized_url = str(raw_url or "").strip()
|
||
if not normalized_url:
|
||
return ""
|
||
parsed = urlsplit(normalized_url)
|
||
if not parsed.scheme or not parsed.netloc:
|
||
return normalized_url
|
||
normalized_host = str(host or "").strip().strip("[]")
|
||
if not normalized_host:
|
||
return normalized_url
|
||
final_scheme = str(scheme or parsed.scheme or "http").strip() or "http"
|
||
final_port = parsed.port if port is None else int(port)
|
||
netloc = f"{normalized_host}:{final_port}" if final_port else normalized_host
|
||
return urlunsplit((final_scheme, netloc, parsed.path, parsed.query, parsed.fragment))
|
||
|
||
|
||
def _resolve_public_control_plane_origin(loopback_url: str) -> str:
|
||
normalized_loopback_url = str(loopback_url or "").strip()
|
||
if not normalized_loopback_url:
|
||
return ""
|
||
parsed_loopback = urlsplit(normalized_loopback_url)
|
||
default_scheme = str(parsed_loopback.scheme or "http").strip() or "http"
|
||
default_port = parsed_loopback.port
|
||
|
||
env_candidates = [
|
||
os.getenv("OPS_CONTROL_PLANE_PUBLIC_BASE_URL", ""),
|
||
os.getenv("CONTROL_PLANE_PUBLIC_BASE_URL", ""),
|
||
os.getenv("OPS_CONTROL_PLANE_BASE_URL", ""),
|
||
]
|
||
for candidate in env_candidates:
|
||
normalized_candidate = str(candidate or "").strip().rstrip("/")
|
||
if not normalized_candidate:
|
||
continue
|
||
parsed_candidate = urlsplit(
|
||
normalized_candidate if "://" in normalized_candidate else f"{default_scheme}://{normalized_candidate}"
|
||
)
|
||
candidate_host = str(parsed_candidate.hostname or "").strip()
|
||
if candidate_host and not _is_loopback_hostname(candidate_host):
|
||
return _build_url_with_host(
|
||
normalized_loopback_url,
|
||
host=candidate_host,
|
||
scheme=str(parsed_candidate.scheme or default_scheme),
|
||
port=parsed_candidate.port if parsed_candidate.port is not None else default_port,
|
||
)
|
||
|
||
try:
|
||
from app.services.cluster_runtime_service import get_cluster_snapshot
|
||
|
||
snapshot = get_cluster_snapshot()
|
||
local_hostnames = {
|
||
str(socket.gethostname() or "").strip().lower(),
|
||
str(socket.getfqdn() or "").strip().lower(),
|
||
}
|
||
fallback_control_hosts: list[str] = []
|
||
for item in list(snapshot.get("nodes") or []):
|
||
if str(item.get("role") or "").strip() != "control":
|
||
continue
|
||
control_host = str(item.get("hostname") or "").strip().lower()
|
||
control_ip = str(item.get("ip") or "").strip()
|
||
if not control_ip or _is_loopback_hostname(control_ip):
|
||
continue
|
||
if control_host and control_host in local_hostnames:
|
||
return _build_url_with_host(
|
||
normalized_loopback_url,
|
||
host=control_ip,
|
||
scheme=default_scheme,
|
||
port=default_port,
|
||
)
|
||
fallback_control_hosts.append(control_ip)
|
||
for control_ip in fallback_control_hosts:
|
||
if control_ip and not _is_loopback_hostname(control_ip):
|
||
return _build_url_with_host(
|
||
normalized_loopback_url,
|
||
host=control_ip,
|
||
scheme=default_scheme,
|
||
port=default_port,
|
||
)
|
||
except Exception:
|
||
pass
|
||
|
||
try:
|
||
resolved_host = str(socket.gethostbyname(socket.gethostname()) or "").strip()
|
||
if resolved_host and not _is_loopback_hostname(resolved_host):
|
||
return _build_url_with_host(
|
||
normalized_loopback_url,
|
||
host=resolved_host,
|
||
scheme=default_scheme,
|
||
port=default_port,
|
||
)
|
||
except Exception:
|
||
pass
|
||
return ""
|
||
|
||
|
||
def _rewrite_loopback_control_plane_url(raw_url: str) -> str:
|
||
normalized = str(raw_url or "").strip()
|
||
if not normalized:
|
||
return ""
|
||
parsed = urlsplit(normalized)
|
||
if not parsed.scheme or not parsed.netloc:
|
||
return normalized
|
||
if not _is_loopback_hostname(str(parsed.hostname or "").strip()):
|
||
return normalized
|
||
resolved = _resolve_public_control_plane_origin(normalized)
|
||
return resolved or normalized
|
||
|
||
|
||
def _build_absolute_release_package_url(base_url: str, raw_path: str) -> str:
|
||
normalized_base = _normalize_public_base_url(base_url)
|
||
normalized_path = str(raw_path or "").strip()
|
||
if not normalized_path:
|
||
return ""
|
||
if normalized_path.startswith("http://") or normalized_path.startswith("https://"):
|
||
return normalized_path
|
||
if not normalized_base:
|
||
return ""
|
||
if normalized_path.startswith("/"):
|
||
return f"{normalized_base}{normalized_path}"
|
||
return f"{normalized_base}/{normalized_path}"
|
||
|
||
|
||
def _read_latest_release_metadata_file() -> dict:
|
||
release_root = _release_package_root()
|
||
latest_json_path = release_root / "latest_release.json"
|
||
if not latest_json_path.exists():
|
||
return {}
|
||
try:
|
||
return json.loads(latest_json_path.read_text(encoding="utf-8"))
|
||
except Exception:
|
||
return {}
|
||
|
||
|
||
def _read_final_release_report_file() -> dict:
|
||
release_root = _release_package_root()
|
||
report_path = release_root / "final_release_report.json"
|
||
if not report_path.exists():
|
||
return {}
|
||
try:
|
||
return json.loads(report_path.read_text(encoding="utf-8"))
|
||
except Exception:
|
||
return {}
|
||
|
||
|
||
def _normalize_release_package_metadata(payload: dict) -> dict:
|
||
raw_payload = dict(payload or {})
|
||
package_name = str(raw_payload.get("package_name") or "").strip()
|
||
if not package_name:
|
||
return {}
|
||
|
||
archive_path_value = (
|
||
str(raw_payload.get("archive_path") or "").strip()
|
||
or str(raw_payload.get("zip_path") or "").strip()
|
||
)
|
||
sha256_path_value = str(raw_payload.get("sha256_path") or "").strip()
|
||
staging_path_value = str(raw_payload.get("staging_path") or "").strip()
|
||
package_type = str(raw_payload.get("package_type") or "").strip()
|
||
archive_filename = ""
|
||
if archive_path_value:
|
||
archive_filename = Path(archive_path_value).name
|
||
if not package_type:
|
||
if archive_filename.endswith(".tar.gz"):
|
||
package_type = "tar.gz"
|
||
elif archive_filename.endswith(".zip"):
|
||
package_type = "zip"
|
||
|
||
smoke_test_ok = bool(raw_payload.get("smoke_test_ok", False))
|
||
manifest_smoke = {}
|
||
commit_sha = str(raw_payload.get("commit_sha") or "").strip()
|
||
commit_ref = str(raw_payload.get("commit_ref") or "").strip()
|
||
manifest = {}
|
||
final_release_report = _read_final_release_report_file()
|
||
manifest_path = Path(staging_path_value) / "release_manifest.json" if staging_path_value else None
|
||
if manifest_path and manifest_path.exists():
|
||
try:
|
||
manifest = json.loads(manifest_path.read_text(encoding="utf-8"))
|
||
except Exception:
|
||
manifest = {}
|
||
manifest_smoke = manifest.get("smoke_test") or {}
|
||
if not commit_sha:
|
||
commit_sha = str(manifest.get("commit_sha") or "").strip()
|
||
if not commit_ref:
|
||
commit_ref = str(manifest.get("commit_ref") or "").strip()
|
||
smoke_test_report = str(
|
||
raw_payload.get("smoke_test_report")
|
||
or manifest_smoke.get("report")
|
||
or ""
|
||
).strip()
|
||
smoke_test_message = str(
|
||
raw_payload.get("smoke_test_message")
|
||
or manifest_smoke.get("message")
|
||
or ""
|
||
).strip()
|
||
if "\n" in smoke_test_message or len(smoke_test_message) > 240:
|
||
smoke_test_message = f"详见 {smoke_test_report}" if smoke_test_report else smoke_test_message.splitlines()[0][:240]
|
||
|
||
final_release_report_available = bool(final_release_report)
|
||
final_release_report_ok = bool(final_release_report.get("ok", False)) if final_release_report_available else False
|
||
final_release_latest = dict(final_release_report.get("latest_release") or {}) if final_release_report_available else {}
|
||
final_release_report_package_name = str(final_release_latest.get("package_name") or "").strip()
|
||
final_release_report_archive_path = str(
|
||
final_release_latest.get("archive_path") or final_release_latest.get("zip_path") or ""
|
||
).strip()
|
||
final_release_report_sha256 = str(final_release_latest.get("sha256") or "").strip().lower()
|
||
current_archive_path = archive_path_value.strip()
|
||
current_sha256 = str(raw_payload.get("sha256") or "").strip().lower()
|
||
final_release_report_matches_latest = False
|
||
final_release_report_stale_reason = ""
|
||
if final_release_report_available:
|
||
mismatches: list[str] = []
|
||
if final_release_report_package_name != package_name:
|
||
mismatches.append("package_name")
|
||
if final_release_report_archive_path != current_archive_path:
|
||
mismatches.append("archive_path")
|
||
if final_release_report_sha256 != current_sha256:
|
||
mismatches.append("sha256")
|
||
final_release_report_matches_latest = not mismatches
|
||
if mismatches:
|
||
final_release_report_stale_reason = "final_release_report does not match latest_release: " + ", ".join(mismatches)
|
||
|
||
final_release_gate = dict(final_release_report.get("release_gate") or {}) if final_release_report_available else {}
|
||
final_release_gate_decision = str(final_release_gate.get("decision") or "").strip()
|
||
final_release_gate_blocked_reasons = [
|
||
str(item).strip()
|
||
for item in list(final_release_gate.get("blocked_reasons") or [])
|
||
if str(item).strip()
|
||
]
|
||
final_release_gate_verify_ok = (
|
||
bool(final_release_gate.get("verify_ok")) if "verify_ok" in final_release_gate else None
|
||
)
|
||
final_release_gate_smoke_test_ok = (
|
||
bool(final_release_gate.get("smoke_test_ok")) if "smoke_test_ok" in final_release_gate else None
|
||
)
|
||
final_release_ok = bool(
|
||
final_release_report_available
|
||
and final_release_report_ok
|
||
and final_release_report_matches_latest
|
||
)
|
||
|
||
summary_parts = [f"导入自本机发布包 {archive_filename or package_name}"]
|
||
generated_at = str(raw_payload.get("generated_at") or "").strip()
|
||
if generated_at:
|
||
summary_parts.append(f"生成时间 {generated_at}")
|
||
if commit_sha:
|
||
summary_parts.append(f"提交 {commit_sha}")
|
||
summary_parts.append(f"自检{'通过' if smoke_test_ok else '未执行或未通过'}")
|
||
if final_release_ok:
|
||
summary_parts.append("最终签收通过")
|
||
elif final_release_report_available:
|
||
summary_parts.append("最终签收未通过")
|
||
|
||
return {
|
||
"available": True,
|
||
"package_name": package_name,
|
||
"package_type": package_type or "unknown",
|
||
"generated_at": generated_at,
|
||
"archive_path": archive_path_value,
|
||
"archive_filename": archive_filename,
|
||
"archive_present": bool(archive_path_value and Path(archive_path_value).exists()),
|
||
"sha256_path": sha256_path_value,
|
||
"sha256_present": bool(sha256_path_value and Path(sha256_path_value).exists()),
|
||
"sha256": str(raw_payload.get("sha256") or "").strip(),
|
||
"commit_sha": commit_sha,
|
||
"commit_ref": commit_ref,
|
||
"staging_path": staging_path_value,
|
||
"staging_present": bool(staging_path_value and Path(staging_path_value).exists()),
|
||
"smoke_test_ok": smoke_test_ok,
|
||
"smoke_test_message": smoke_test_message,
|
||
"smoke_test_report": smoke_test_report,
|
||
"final_release_ok": final_release_ok,
|
||
"final_release_report_available": final_release_report_available,
|
||
"final_release_report_ok": final_release_report_ok if final_release_report_available else None,
|
||
"final_release_report_matches_latest": final_release_report_matches_latest if final_release_report_available else None,
|
||
"final_release_report_package_name": final_release_report_package_name,
|
||
"final_release_report_stale_reason": final_release_report_stale_reason,
|
||
"final_release_gate_decision": final_release_gate_decision,
|
||
"final_release_gate_blocked_reasons": final_release_gate_blocked_reasons,
|
||
"final_release_gate_verify_ok": final_release_gate_verify_ok,
|
||
"final_release_gate_smoke_test_ok": final_release_gate_smoke_test_ok,
|
||
"download_path": f"/api/v1/ops/releases/packages/{package_name}/download",
|
||
"sha256_download_path": f"/api/v1/ops/releases/packages/{package_name}/sha256",
|
||
"release_version_suggestion": package_name,
|
||
"notes_suggestion": ";".join(summary_parts),
|
||
"metadata": {
|
||
"package_name": package_name,
|
||
"package_type": package_type or "unknown",
|
||
"generated_at": generated_at,
|
||
"archive_filename": archive_filename,
|
||
"archive_path": archive_path_value,
|
||
"sha256_path": sha256_path_value,
|
||
"commit_sha": commit_sha,
|
||
"commit_ref": commit_ref,
|
||
"staging_path": staging_path_value,
|
||
"smoke_test_ok": smoke_test_ok,
|
||
"smoke_test_message": smoke_test_message,
|
||
"smoke_test_report": smoke_test_report,
|
||
"final_release_report": final_release_report,
|
||
"manifest": manifest,
|
||
},
|
||
}
|
||
|
||
|
||
def get_latest_release_package_metadata() -> dict:
|
||
payload = _read_latest_release_metadata_file()
|
||
normalized = _normalize_release_package_metadata(payload)
|
||
if normalized:
|
||
return normalized
|
||
return {
|
||
"available": False,
|
||
"reason": "latest_release.json 不存在或不可解析",
|
||
}
|
||
|
||
|
||
def _validate_latest_release_package_final_gate(package_metadata: dict | None = None) -> tuple[bool, str, dict]:
|
||
normalized_package = dict(package_metadata or {})
|
||
if not bool(normalized_package.get("available")):
|
||
message = str(normalized_package.get("reason") or "当前没有可用的本机发布包")
|
||
return False, message, {
|
||
"package": normalized_package,
|
||
"final_release_gate": {
|
||
"ready": False,
|
||
"status": "package_unavailable",
|
||
"status_label": "发布包不可用",
|
||
"summary": message,
|
||
"blocked_reasons": ["package_unavailable"],
|
||
},
|
||
}
|
||
|
||
if bool(normalized_package.get("final_release_ok")):
|
||
return True, "", {
|
||
"package": normalized_package,
|
||
"final_release_gate": {
|
||
"ready": True,
|
||
"status": "ready",
|
||
"status_label": "最终签收通过",
|
||
"summary": "最新发布包已完成最终签收,可进入 Release 创建或智能 Rollout。",
|
||
"blocked_reasons": [],
|
||
"decision": str(normalized_package.get("final_release_gate_decision") or "").strip(),
|
||
},
|
||
}
|
||
|
||
blocked_reasons: list[str] = []
|
||
if not bool(normalized_package.get("smoke_test_ok")):
|
||
blocked_reasons.append("smoke_test_failed_or_skipped")
|
||
if not bool(normalized_package.get("final_release_report_available")):
|
||
blocked_reasons.append("final_release_report_missing")
|
||
if bool(normalized_package.get("final_release_report_available")) and not bool(
|
||
normalized_package.get("final_release_report_ok")
|
||
):
|
||
blocked_reasons.append("final_release_report_not_ok")
|
||
if normalized_package.get("final_release_report_matches_latest") is False:
|
||
blocked_reasons.append("final_release_report_not_matching_latest")
|
||
if normalized_package.get("final_release_gate_verify_ok") is False:
|
||
blocked_reasons.append("verify_failed")
|
||
if normalized_package.get("final_release_gate_smoke_test_ok") is False:
|
||
blocked_reasons.append("smoke_test_failed")
|
||
|
||
decision = str(normalized_package.get("final_release_gate_decision") or "").strip()
|
||
if decision and decision != "ready":
|
||
blocked_reasons.append(f"final_release_gate_{decision}")
|
||
|
||
for item in list(normalized_package.get("final_release_gate_blocked_reasons") or []):
|
||
normalized_item = str(item or "").strip()
|
||
if normalized_item:
|
||
blocked_reasons.append(normalized_item)
|
||
|
||
deduped_blocked_reasons: list[str] = []
|
||
for item in blocked_reasons:
|
||
if item and item not in deduped_blocked_reasons:
|
||
deduped_blocked_reasons.append(item)
|
||
|
||
hints: list[str] = []
|
||
stale_reason = str(normalized_package.get("final_release_report_stale_reason") or "").strip()
|
||
if stale_reason:
|
||
hints.append(stale_reason)
|
||
if not bool(normalized_package.get("final_release_report_available")):
|
||
hints.append("缺少 final_release_report.json")
|
||
elif not bool(normalized_package.get("final_release_report_ok")):
|
||
hints.append("final_release_report.json 尚未给出 ok=true")
|
||
if not bool(normalized_package.get("smoke_test_ok")):
|
||
hints.append("smoke test 未通过")
|
||
if normalized_package.get("final_release_gate_verify_ok") is False:
|
||
hints.append("verify 未通过")
|
||
if normalized_package.get("final_release_gate_smoke_test_ok") is False:
|
||
hints.append("最终签收门禁中的 smoke_test_ok=false")
|
||
|
||
summary = "最新发布包尚未完成最终签收,暂不允许创建 Release 或发起基于最新包的智能 Rollout。"
|
||
if hints:
|
||
summary = f"{summary} {hints[0]}"
|
||
|
||
return False, summary, {
|
||
"package": normalized_package,
|
||
"final_release_gate": {
|
||
"ready": False,
|
||
"status": "blocked",
|
||
"status_label": "最终签收未完成",
|
||
"summary": summary,
|
||
"blocked_reasons": deduped_blocked_reasons,
|
||
"decision": decision,
|
||
"report_available": bool(normalized_package.get("final_release_report_available")),
|
||
"report_ok": bool(normalized_package.get("final_release_report_ok")),
|
||
"report_matches_latest": normalized_package.get("final_release_report_matches_latest"),
|
||
"stale_reason": stale_reason,
|
||
"verify_ok": normalized_package.get("final_release_gate_verify_ok"),
|
||
"smoke_test_ok": normalized_package.get("final_release_gate_smoke_test_ok"),
|
||
},
|
||
}
|
||
|
||
|
||
def prepare_latest_release_package(*, requested_by: str = "api/release-prepare") -> tuple[bool, str, dict]:
|
||
repo_root = _release_repo_root()
|
||
script_path = repo_root / "prepare_final_release.sh"
|
||
if not script_path.exists():
|
||
return False, "prepare_final_release.sh 不存在", {
|
||
"requested_by": str(requested_by or "").strip() or "api/release-prepare",
|
||
"script_path": str(script_path),
|
||
}
|
||
|
||
completed = subprocess.run(
|
||
["bash", str(script_path)],
|
||
cwd=str(repo_root),
|
||
capture_output=True,
|
||
text=True,
|
||
timeout=1200,
|
||
)
|
||
package = get_latest_release_package_metadata()
|
||
gate_ok, gate_message, gate_data = _validate_latest_release_package_final_gate(package)
|
||
result = {
|
||
"requested_by": str(requested_by or "").strip() or "api/release-prepare",
|
||
"script_path": str(script_path),
|
||
"command": ["bash", str(script_path)],
|
||
"returncode": int(completed.returncode or 0),
|
||
"stdout_preview": _trim_release_command_output(completed.stdout, 12000),
|
||
"stderr_preview": _trim_release_command_output(completed.stderr, 4000),
|
||
"package": package,
|
||
"final_release_gate": dict(gate_data.get("final_release_gate") or {}),
|
||
}
|
||
if completed.returncode != 0:
|
||
message = _trim_release_command_output(completed.stderr or completed.stdout, 400) or "最终签收准备执行失败"
|
||
return False, message, result
|
||
if not gate_ok:
|
||
return False, gate_message, result
|
||
return True, "最新发布包最终签收已完成", result
|
||
|
||
|
||
def get_release_package_archive_path(package_name: str) -> Path | None:
|
||
normalized_name = _safe_release_package_name(package_name)
|
||
release_root = _release_package_root()
|
||
for suffix in (".tar.gz", ".zip"):
|
||
candidate = release_root / f"{normalized_name}{suffix}"
|
||
if candidate.exists() and candidate.is_file():
|
||
return candidate
|
||
metadata = _normalize_release_package_metadata(_read_latest_release_metadata_file())
|
||
if metadata and str(metadata.get("package_name") or "").strip() == normalized_name:
|
||
archive_path_value = str(metadata.get("archive_path") or "").strip()
|
||
if archive_path_value:
|
||
candidate = Path(archive_path_value)
|
||
if candidate.exists() and candidate.is_file():
|
||
return candidate
|
||
return None
|
||
|
||
|
||
def get_release_package_sha256_path(package_name: str) -> Path | None:
|
||
normalized_name = _safe_release_package_name(package_name)
|
||
release_root = _release_package_root()
|
||
candidate = release_root / f"{normalized_name}.sha256.txt"
|
||
if candidate.exists() and candidate.is_file():
|
||
return candidate
|
||
metadata = _normalize_release_package_metadata(_read_latest_release_metadata_file())
|
||
if metadata and str(metadata.get("package_name") or "").strip() == normalized_name:
|
||
sha256_path_value = str(metadata.get("sha256_path") or "").strip()
|
||
if sha256_path_value:
|
||
resolved = Path(sha256_path_value)
|
||
if resolved.exists() and resolved.is_file():
|
||
return resolved
|
||
return None
|
||
|
||
|
||
def _get_release_by_version(release_version: str) -> dict:
|
||
ensure_ops_release_schema()
|
||
normalized_release_version = str(release_version or "").strip()
|
||
if not normalized_release_version:
|
||
return {}
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
SELECT
|
||
id, release_version, channel, commit_sha, status, artifact_url, checksum, notes, metadata_json, created_by, created_at, activated_at, updated_at
|
||
FROM ops_releases
|
||
WHERE release_version = %s
|
||
LIMIT 1
|
||
""",
|
||
(normalized_release_version,),
|
||
)
|
||
row = cur.fetchone()
|
||
return _serialize_release_row(row) if row else {}
|
||
|
||
|
||
def _inspection_job_bucket_for_rollout(action: str, payload: dict) -> str:
|
||
normalized_action = str(action or "").strip()
|
||
if normalized_action in {"health.snapshot", "diagnostics.collect"}:
|
||
return normalized_action
|
||
if normalized_action != "logs.collect":
|
||
return ""
|
||
service_name = str((payload or {}).get("service_name") or "").strip().lower()
|
||
if "worker" not in service_name:
|
||
return ""
|
||
return "logs.collect"
|
||
|
||
|
||
def _rollout_inspection_status_label(status: str) -> str:
|
||
normalized_status = str(status or "").strip()
|
||
mapping = {
|
||
"healthy": "巡检通过",
|
||
"running": "巡检进行中",
|
||
"attention": "巡检待关注",
|
||
"missing": "缺少巡检",
|
||
}
|
||
return mapping.get(normalized_status, normalized_status or "未知")
|
||
|
||
|
||
def _rollout_agent_state_from_context(
|
||
*,
|
||
last_seen_at: str,
|
||
latest_token: dict,
|
||
cluster_status: str,
|
||
current_load: int,
|
||
) -> tuple[str, str, str]:
|
||
now = datetime.now()
|
||
last_seen_dt = _parse_time(last_seen_at)
|
||
token_status = str(latest_token.get("token_status") or "absent").strip() or "absent"
|
||
|
||
if last_seen_dt:
|
||
current_time = datetime.now(last_seen_dt.tzinfo) if last_seen_dt.tzinfo else now
|
||
age_seconds = max(0, int((current_time - last_seen_dt).total_seconds()))
|
||
if age_seconds <= 180:
|
||
if cluster_status == "busy" or int(current_load or 0) > 0:
|
||
return "online_busy", "已接管", f"Agent 最近心跳 {age_seconds} 秒前,节点当前负载 {int(current_load or 0)}。"
|
||
return "online", "已接管", f"Agent 最近心跳 {age_seconds} 秒前。"
|
||
return "stale", "心跳过期", f"Agent 最近心跳 {age_seconds} 秒前,已超过在线阈值。"
|
||
|
||
if token_status == "active":
|
||
return "pending_bootstrap", "待接入", "已签发有效 Agent Token,但节点尚未完成 register/heartbeat。"
|
||
if token_status == "expired":
|
||
return "token_expired", "Token 过期", "已签发 Agent Token,但该 Token 已过期。"
|
||
if token_status == "disabled":
|
||
return "token_disabled", "Token 已禁用", "最近一枚 Agent Token 已被禁用,需重新签发。"
|
||
if cluster_status in {"online", "busy"}:
|
||
return "runtime_only", "仅运行态在线", "控制面能看到 runtime 心跳,但 Node Agent 尚未接管。"
|
||
return "unmanaged", "未接管", "当前节点尚未生成有效 Agent 接入状态。"
|
||
|
||
|
||
def _rollout_execution_mode_label(execution_mode: str) -> str:
|
||
normalized_execution_mode = str(execution_mode or "remote-agent").strip() or "remote-agent"
|
||
return execution_mode_label(normalized_execution_mode)
|
||
|
||
|
||
def _rollout_delivery_queue_context(metadata: dict | None = None) -> tuple[str, str, str, dict]:
|
||
snapshot = dict((metadata or {}).get("delivery_queue") or {})
|
||
pending_count = int(snapshot.get("pending_count", 0) or 0)
|
||
dead_letter_count = int(snapshot.get("dead_letter_count", 0) or 0)
|
||
last_flush_at = str(snapshot.get("last_flush_at") or "").strip()
|
||
oldest_pending_at = str(snapshot.get("oldest_pending_at") or "").strip()
|
||
oldest_dead_letter_at = str(snapshot.get("oldest_dead_letter_at") or "").strip()
|
||
|
||
if not snapshot:
|
||
return "unknown", "未上报", "Node Agent 尚未上报回执队列状态。", {}
|
||
if dead_letter_count > 0:
|
||
reason = f"当前存在 {dead_letter_count} 条死信记录,说明该节点的 Agent 回执链路尚未完全收口。"
|
||
if oldest_dead_letter_at:
|
||
reason += f" 最早死信时间:{oldest_dead_letter_at}。"
|
||
return "dead_letter", f"死信 {dead_letter_count}", reason, snapshot
|
||
if pending_count > 0:
|
||
reason = f"当前存在 {pending_count} 条待重试回执,Node Agent 会继续自动回放。"
|
||
if oldest_pending_at:
|
||
reason += f" 最早积压时间:{oldest_pending_at}。"
|
||
return "retrying", f"待重试 {pending_count}", reason, snapshot
|
||
|
||
reason = "当前没有待重试回执,也没有死信记录。"
|
||
if last_flush_at:
|
||
reason += f" 最近一次队列冲刷:{last_flush_at}。"
|
||
return "healthy", "正常", reason, snapshot
|
||
|
||
|
||
def build_rollout_target_operational_readiness(
|
||
target_nodes: list[dict],
|
||
*,
|
||
execution_mode: str = "remote-agent",
|
||
desired_release: dict | None = None,
|
||
) -> dict:
|
||
from app.services.cluster_runtime_service import get_cluster_snapshot
|
||
from app.services.ops_agent_service import (
|
||
ensure_ops_agent_schema,
|
||
get_managed_node_onboarding,
|
||
list_managed_nodes_with_agent_state,
|
||
)
|
||
from app.services.ops_job_service import list_managed_nodes
|
||
|
||
ensure_ops_agent_schema()
|
||
normalized_execution_mode = str(execution_mode or "remote-agent").strip() or "remote-agent"
|
||
normalized_desired_release = dict(desired_release or {})
|
||
desired_release_version = str(normalized_desired_release.get("release_version") or "").strip()
|
||
normalized_target_nodes = [
|
||
dict(item or {})
|
||
for item in list(target_nodes or [])
|
||
if str((item or {}).get("node_code") or "").strip()
|
||
]
|
||
target_node_codes = [str(item.get("node_code") or "").strip() for item in normalized_target_nodes]
|
||
if not target_node_codes:
|
||
return {
|
||
"execution_mode": normalized_execution_mode,
|
||
"summary": {
|
||
"nodes_total": 0,
|
||
"managed_nodes": 0,
|
||
"managed_enabled_nodes": 0,
|
||
"agent_online_nodes": 0,
|
||
"ssh_ready_nodes": 0,
|
||
"remote_agent_ready_nodes": 0,
|
||
"execution_ready_nodes": 0,
|
||
"inspection_healthy_nodes": 0,
|
||
"inspection_running_nodes": 0,
|
||
"inspection_attention_nodes": 0,
|
||
"inspection_missing_nodes": 0,
|
||
},
|
||
"rows": [],
|
||
"blocking_reasons": [],
|
||
"warning_reasons": [],
|
||
"recommendations": [],
|
||
}
|
||
|
||
cluster_snapshot = get_cluster_snapshot()
|
||
cluster_map = {
|
||
str(item.get("node_code") or "").strip(): item
|
||
for item in list(cluster_snapshot.get("nodes") or [])
|
||
if str(item.get("node_code") or "").strip()
|
||
}
|
||
managed_nodes_payload = list_managed_nodes_with_agent_state()
|
||
managed_nodes = list_managed_nodes()
|
||
managed_map = {
|
||
str(item.get("node_code") or "").strip(): item
|
||
for item in managed_nodes
|
||
if str(item.get("node_code") or "").strip()
|
||
}
|
||
|
||
latest_tokens: dict[str, dict] = {}
|
||
latest_inspection_jobs: dict[tuple[str, str], dict] = {}
|
||
placeholders = ", ".join(["%s"] * len(target_node_codes))
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
f"""
|
||
SELECT DISTINCT ON (node_code)
|
||
node_code, id, is_enabled, expires_at, last_used_at, metadata_json, created_at, updated_at
|
||
FROM ops_node_tokens
|
||
WHERE purpose = 'agent' AND node_code IN ({placeholders})
|
||
ORDER BY node_code ASC, created_at DESC, id DESC
|
||
""",
|
||
tuple(target_node_codes),
|
||
)
|
||
for row in cur.fetchall():
|
||
node_code = str(row[0] or "").strip()
|
||
if not node_code:
|
||
continue
|
||
expires_at = _ts(row[3])
|
||
token_status = "active"
|
||
expires_dt = _parse_time(row[3])
|
||
if not bool(row[2]):
|
||
token_status = "disabled"
|
||
elif expires_dt:
|
||
current_time = datetime.now(expires_dt.tzinfo) if expires_dt.tzinfo else datetime.now()
|
||
if expires_dt < current_time:
|
||
token_status = "expired"
|
||
latest_tokens[node_code] = {
|
||
"record_id": int(row[1]),
|
||
"is_enabled": bool(row[2]),
|
||
"expires_at": expires_at,
|
||
"last_used_at": _ts(row[4]),
|
||
"metadata": _decode_json(row[5]),
|
||
"created_at": _ts(row[6]),
|
||
"updated_at": _ts(row[7]),
|
||
"token_status": token_status,
|
||
}
|
||
|
||
cur.execute(
|
||
f"""
|
||
SELECT
|
||
id, target_node_code, action, status, payload_json, created_at, updated_at
|
||
FROM ops_jobs
|
||
WHERE target_node_code IN ({placeholders})
|
||
AND action IN ('health.snapshot', 'logs.collect', 'diagnostics.collect')
|
||
ORDER BY created_at DESC, id DESC
|
||
""",
|
||
tuple(target_node_codes),
|
||
)
|
||
for row in cur.fetchall():
|
||
node_code = str(row[1] or "").strip()
|
||
if not node_code:
|
||
continue
|
||
payload = _decode_payload_json(row[4])
|
||
bucket = _inspection_job_bucket_for_rollout(str(row[2] or ""), payload)
|
||
if not bucket:
|
||
continue
|
||
key = (node_code, bucket)
|
||
if key in latest_inspection_jobs:
|
||
continue
|
||
latest_inspection_jobs[key] = {
|
||
"id": int(row[0]),
|
||
"status": str(row[3] or ""),
|
||
"created_at": _ts(row[5]),
|
||
"updated_at": _ts(row[6]),
|
||
}
|
||
|
||
rows: list[dict] = []
|
||
blocking_reasons: list[str] = []
|
||
warning_reasons: list[str] = []
|
||
recommendations: list[str] = []
|
||
|
||
unmanaged_nodes: list[str] = []
|
||
disabled_nodes: list[str] = []
|
||
agent_offline_nodes: list[str] = []
|
||
ssh_unready_nodes: list[str] = []
|
||
inspection_attention_nodes: list[str] = []
|
||
inspection_missing_nodes: list[str] = []
|
||
onboarding_bootstrap_pending_nodes: list[str] = []
|
||
onboarding_acceptance_ready_nodes: list[str] = []
|
||
onboarding_noop_nodes: list[str] = []
|
||
queue_dead_letter_nodes: list[str] = []
|
||
queue_retrying_nodes: list[str] = []
|
||
|
||
for target in normalized_target_nodes:
|
||
node_code = str(target.get("node_code") or "").strip()
|
||
cluster_node = cluster_map.get(node_code, {})
|
||
managed = managed_map.get(node_code, {})
|
||
latest_token = latest_tokens.get(node_code, {})
|
||
metadata = dict(managed.get("metadata") or {})
|
||
(
|
||
delivery_queue_state,
|
||
delivery_queue_label,
|
||
delivery_queue_reason,
|
||
delivery_queue_snapshot,
|
||
) = _rollout_delivery_queue_context(metadata)
|
||
last_seen_at = str(managed.get("last_seen_at") or "").strip() or str(metadata.get("last_seen_at") or "").strip()
|
||
cluster_status = str(cluster_node.get("status") or target.get("status") or "").strip()
|
||
current_load = int(cluster_node.get("current_load", target.get("current_load", 0)) or 0)
|
||
onboarding = (
|
||
get_managed_node_onboarding(node_code, nodes_payload=managed_nodes_payload)
|
||
if node_code
|
||
else {}
|
||
)
|
||
onboarding_stage = dict(onboarding.get("onboarding_stage") or {})
|
||
recovery_decision = dict(onboarding.get("recovery_decision") or {})
|
||
onboarding_stage_code = str(onboarding_stage.get("code") or "").strip()
|
||
onboarding_stage_label = str(onboarding_stage.get("label") or "").strip()
|
||
onboarding_summary = str(onboarding.get("summary") or onboarding_stage.get("summary") or "").strip()
|
||
recovery_action = str(recovery_decision.get("action") or "").strip()
|
||
recovery_label = str(recovery_decision.get("label") or "").strip()
|
||
recovery_summary = str(recovery_decision.get("summary") or "").strip()
|
||
recovery_command_hint = str(recovery_decision.get("command_hint") or "").strip()
|
||
recovery_window = str(recovery_decision.get("window") or "").strip()
|
||
agent_state, agent_state_label, agent_state_reason = _rollout_agent_state_from_context(
|
||
last_seen_at=last_seen_at,
|
||
latest_token=latest_token,
|
||
cluster_status=cluster_status,
|
||
current_load=current_load,
|
||
)
|
||
is_managed = bool(managed)
|
||
is_enabled = bool(managed.get("is_enabled", target.get("is_enabled", False)))
|
||
is_agent_online = agent_state in {"online", "online_busy"}
|
||
remote_agent_ready = is_managed and is_enabled and is_agent_online
|
||
has_ssh_access = bool(str(managed.get("ssh_host") or "").strip() and str(managed.get("ssh_user") or "").strip())
|
||
ssh_ready = is_managed and is_enabled and has_ssh_access
|
||
execution_ready = remote_agent_ready if normalized_execution_mode == "remote-agent" else ssh_ready
|
||
|
||
inspection_status_map = {
|
||
bucket: str((latest_inspection_jobs.get((node_code, bucket)) or {}).get("status") or "").strip()
|
||
for bucket in _ROLLOUT_INSPECTION_ACTION_KEYS
|
||
}
|
||
present_statuses = [status for status in inspection_status_map.values() if status]
|
||
if not present_statuses:
|
||
inspection_status = "missing"
|
||
inspection_reason = "当前缺少健康快照、Worker 日志和诊断包的巡检记录。"
|
||
elif any(status in {"failed", "blocked", "cancelled"} for status in present_statuses):
|
||
inspection_status = "attention"
|
||
inspection_reason = "最近巡检存在失败、阻断或取消记录。"
|
||
elif any(status in {"queued", "dispatching", "running", "awaiting_approval"} for status in present_statuses):
|
||
inspection_status = "running"
|
||
inspection_reason = "最近巡检仍在执行、排队或待审批中。"
|
||
elif all(inspection_status_map.get(bucket) == "success" for bucket in _ROLLOUT_INSPECTION_ACTION_KEYS):
|
||
inspection_status = "healthy"
|
||
inspection_reason = "最近一轮健康快照、Worker 日志和诊断包均已成功。"
|
||
else:
|
||
inspection_status = "attention"
|
||
missing_actions = [
|
||
bucket for bucket in _ROLLOUT_INSPECTION_ACTION_KEYS if not inspection_status_map.get(bucket)
|
||
]
|
||
inspection_reason = (
|
||
f"最近巡检仍缺少 {'、'.join(missing_actions)} 的成功记录。"
|
||
if missing_actions
|
||
else "最近巡检尚未形成完整健康结论。"
|
||
)
|
||
|
||
inspection_success_count = sum(1 for bucket in _ROLLOUT_INSPECTION_ACTION_KEYS if inspection_status_map.get(bucket) == "success")
|
||
inspection_ready = inspection_status == "healthy"
|
||
row_blocking_reasons: list[str] = []
|
||
row_warning_reasons: list[str] = []
|
||
|
||
if not is_managed:
|
||
row_blocking_reasons.append("当前节点尚未纳管,不能直接进入正式发布执行面。")
|
||
elif not is_enabled:
|
||
row_blocking_reasons.append("当前节点已纳管但未启用,需先恢复启用状态。")
|
||
elif normalized_execution_mode == "remote-agent" and not is_agent_online:
|
||
row_blocking_reasons.append("当前节点 Node Agent 未在线,暂不适合 remote-agent 发布。")
|
||
elif normalized_execution_mode == "ssh" and not has_ssh_access:
|
||
row_blocking_reasons.append("当前节点尚未配置完整 SSH 入口,暂不适合 SSH 发布。")
|
||
|
||
if delivery_queue_state == "dead_letter":
|
||
row_blocking_reasons.append("当前节点 Node Agent 回执队列存在死信,说明执行回执链路未收口,暂不适合正式发布。")
|
||
elif delivery_queue_state == "retrying":
|
||
row_warning_reasons.append("当前节点 Node Agent 回执队列仍在自动重试,建议等回执链路收口后再发布。")
|
||
|
||
if inspection_status == "attention":
|
||
row_warning_reasons.append("最近巡检存在失败、阻断或取消记录,建议先看巡检收口。")
|
||
elif inspection_status == "missing":
|
||
row_warning_reasons.append("当前缺少健康快照、Worker 日志和诊断包的完整巡检记录。")
|
||
elif inspection_status == "running":
|
||
row_warning_reasons.append("最近巡检仍在执行中,建议等待收口后再推进发布。")
|
||
|
||
if recovery_action == "bootstrap_run":
|
||
onboarding_bootstrap_pending_nodes.append(node_code)
|
||
if not execution_ready:
|
||
row_warning_reasons.append("当前节点还在接入阶段,建议先执行 onboarding.bootstrap。")
|
||
elif recovery_action == "run_acceptance":
|
||
onboarding_acceptance_ready_nodes.append(node_code)
|
||
if not execution_ready:
|
||
row_warning_reasons.append("当前节点已进入接管验收窗口,建议先执行 onboarding.acceptance。")
|
||
else:
|
||
onboarding_noop_nodes.append(node_code)
|
||
|
||
current_release_version = _extract_node_release_version(target, cluster_node, managed, metadata)
|
||
risk_level = _rollout_readiness_risk_level(
|
||
blocking_reasons=row_blocking_reasons,
|
||
inspection_status=inspection_status,
|
||
)
|
||
if row_blocking_reasons:
|
||
if recovery_action in {"bootstrap_run", "run_acceptance"}:
|
||
recommended_action_code = recovery_action
|
||
else:
|
||
recommended_action_code = "handover_first_gap" if not is_managed else "fix_managed_nodes"
|
||
elif inspection_status in {"attention", "missing", "running"}:
|
||
recommended_action_code = "run_standard_inspection"
|
||
else:
|
||
recommended_action_code = ""
|
||
|
||
if not is_managed:
|
||
unmanaged_nodes.append(node_code)
|
||
elif not is_enabled:
|
||
disabled_nodes.append(node_code)
|
||
elif not is_agent_online:
|
||
agent_offline_nodes.append(node_code)
|
||
if is_managed and is_enabled and not has_ssh_access:
|
||
ssh_unready_nodes.append(node_code)
|
||
|
||
if inspection_status == "attention":
|
||
inspection_attention_nodes.append(node_code)
|
||
elif inspection_status == "missing":
|
||
inspection_missing_nodes.append(node_code)
|
||
if delivery_queue_state == "dead_letter":
|
||
queue_dead_letter_nodes.append(node_code)
|
||
elif delivery_queue_state == "retrying":
|
||
queue_retrying_nodes.append(node_code)
|
||
|
||
rows.append(
|
||
{
|
||
"node_code": node_code,
|
||
"region": str(target.get("region") or cluster_node.get("region") or ""),
|
||
"role": str(target.get("role") or cluster_node.get("role") or ""),
|
||
"status": cluster_status,
|
||
"current_load": current_load,
|
||
"agent_managed": is_managed,
|
||
"is_managed": is_managed,
|
||
"is_enabled": is_enabled,
|
||
"agent_state": agent_state,
|
||
"agent_state_label": agent_state_label,
|
||
"agent_reason": agent_state_reason,
|
||
"is_agent_online": is_agent_online,
|
||
"has_ssh_access": has_ssh_access,
|
||
"ssh_ready": ssh_ready,
|
||
"remote_agent_ready": remote_agent_ready,
|
||
"execution_ready": execution_ready,
|
||
"execution_mode": normalized_execution_mode,
|
||
"execution_mode_label": _rollout_execution_mode_label(normalized_execution_mode),
|
||
"inspection_status": inspection_status,
|
||
"inspection_status_label": _rollout_inspection_status_label(inspection_status),
|
||
"inspection_reason": inspection_reason,
|
||
"inspection_success_count": inspection_success_count,
|
||
"inspection_ready": inspection_ready,
|
||
"inspection_status_map": inspection_status_map,
|
||
"delivery_queue_state": delivery_queue_state,
|
||
"delivery_queue_label": delivery_queue_label,
|
||
"delivery_queue_reason": delivery_queue_reason,
|
||
"delivery_queue_pending_count": int(delivery_queue_snapshot.get("pending_count", 0) or 0),
|
||
"delivery_queue_dead_letter_count": int(delivery_queue_snapshot.get("dead_letter_count", 0) or 0),
|
||
"current_release_version": current_release_version,
|
||
"desired_release_version": desired_release_version,
|
||
"onboarding_stage_code": onboarding_stage_code,
|
||
"onboarding_stage_label": onboarding_stage_label,
|
||
"onboarding_summary": onboarding_summary,
|
||
"recovery_action": recovery_action,
|
||
"recovery_label": recovery_label,
|
||
"recovery_summary": recovery_summary,
|
||
"recovery_command_hint": recovery_command_hint,
|
||
"recovery_window": recovery_window,
|
||
"risk_level": risk_level,
|
||
"blocking_reasons": row_blocking_reasons,
|
||
"warning_reasons": row_warning_reasons,
|
||
"recommended_action_code": recommended_action_code,
|
||
"focus_ref": _build_rollout_readiness_focus_ref(normalized_desired_release, node_code=node_code),
|
||
}
|
||
)
|
||
|
||
if normalized_execution_mode == "remote-agent":
|
||
if unmanaged_nodes:
|
||
blocking_reasons.append(f"有 {len(unmanaged_nodes)} 台目标节点尚未纳管,remote-agent rollout 无法下发。")
|
||
if disabled_nodes:
|
||
blocking_reasons.append(f"有 {len(disabled_nodes)} 台目标节点已纳管但未启用,需先恢复启用。")
|
||
if agent_offline_nodes:
|
||
blocking_reasons.append(f"有 {len(agent_offline_nodes)} 台目标节点 Agent 未在线,当前不适合直接 remote-agent 发布。")
|
||
elif normalized_execution_mode == "ssh":
|
||
if unmanaged_nodes:
|
||
blocking_reasons.append(f"有 {len(unmanaged_nodes)} 台目标节点尚未纳管,SSH 发布无法直接下发。")
|
||
if disabled_nodes:
|
||
blocking_reasons.append(f"有 {len(disabled_nodes)} 台目标节点已纳管但未启用,需先恢复启用。")
|
||
if ssh_unready_nodes:
|
||
blocking_reasons.append(f"有 {len(ssh_unready_nodes)} 台目标节点尚未配置完整 SSH 入口,当前不适合直接 SSH 发布。")
|
||
|
||
if inspection_attention_nodes:
|
||
warning_reasons.append(f"有 {len(inspection_attention_nodes)} 台目标节点最近巡检待关注,建议先看巡检结果再发布。")
|
||
if inspection_missing_nodes:
|
||
warning_reasons.append(f"有 {len(inspection_missing_nodes)} 台目标节点缺少完整巡检记录,建议先补齐标准巡检。")
|
||
if queue_dead_letter_nodes:
|
||
blocking_reasons.append(f"有 {len(queue_dead_letter_nodes)} 台目标节点存在 Node Agent 死信回执,当前不能进入正式发布。")
|
||
if queue_retrying_nodes:
|
||
warning_reasons.append(f"有 {len(queue_retrying_nodes)} 台目标节点仍在自动回放回执,建议等待回执链路收口。")
|
||
|
||
if normalized_execution_mode == "remote-agent" and (unmanaged_nodes or disabled_nodes or agent_offline_nodes):
|
||
recommendations.append("先把目标节点全部收敛到 已纳管 + 已启用 + Agent 在线,再推进 remote-agent rollout。")
|
||
if normalized_execution_mode == "ssh" and (unmanaged_nodes or disabled_nodes or ssh_unready_nodes):
|
||
recommendations.append("先把目标节点全部收敛到 已纳管 + 已启用 + SSH 已备好,再推进 SSH 发布。")
|
||
if inspection_attention_nodes or inspection_missing_nodes:
|
||
recommendations.append("先对待发布节点执行 inspection.standard,至少拿到 健康快照 + Worker 日志 + 诊断包 的最新回执。")
|
||
if queue_dead_letter_nodes:
|
||
recommendations.append("先处理 Node Agent 死信队列并确认回执链路恢复正常,再推进正式发布。")
|
||
elif queue_retrying_nodes:
|
||
recommendations.append("先等待 Node Agent 自动回放收口,确认待重试回执清空后再推进发布。")
|
||
if onboarding_bootstrap_pending_nodes:
|
||
recommendations.append(
|
||
f"有 {len(onboarding_bootstrap_pending_nodes)} 台目标节点仍停留在接入阶段,建议先执行 onboarding.bootstrap。"
|
||
)
|
||
if onboarding_acceptance_ready_nodes:
|
||
recommendations.append(
|
||
f"有 {len(onboarding_acceptance_ready_nodes)} 台目标节点已进入验收窗口,建议先执行 onboarding.acceptance。"
|
||
)
|
||
|
||
summary = {
|
||
"nodes_total": len(rows),
|
||
"managed_nodes": sum(1 for row in rows if bool(row.get("is_managed", False))),
|
||
"managed_enabled_nodes": sum(1 for row in rows if bool(row.get("is_managed", False)) and bool(row.get("is_enabled", False))),
|
||
"agent_online_nodes": sum(1 for row in rows if bool(row.get("is_agent_online", False))),
|
||
"ssh_ready_nodes": sum(1 for row in rows if bool(row.get("ssh_ready", False))),
|
||
"remote_agent_ready_nodes": sum(1 for row in rows if bool(row.get("remote_agent_ready", False))),
|
||
"execution_ready_nodes": sum(1 for row in rows if bool(row.get("execution_ready", False))),
|
||
"inspection_healthy_nodes": sum(1 for row in rows if str(row.get("inspection_status") or "") == "healthy"),
|
||
"inspection_running_nodes": sum(1 for row in rows if str(row.get("inspection_status") or "") == "running"),
|
||
"inspection_attention_nodes": sum(1 for row in rows if str(row.get("inspection_status") or "") == "attention"),
|
||
"inspection_missing_nodes": sum(1 for row in rows if str(row.get("inspection_status") or "") == "missing"),
|
||
"queue_retrying_nodes": len(queue_retrying_nodes),
|
||
"queue_dead_letter_nodes": len(queue_dead_letter_nodes),
|
||
"onboarding_bootstrap_pending_nodes": len(onboarding_bootstrap_pending_nodes),
|
||
"onboarding_acceptance_ready_nodes": len(onboarding_acceptance_ready_nodes),
|
||
"onboarding_noop_nodes": len(onboarding_noop_nodes),
|
||
}
|
||
|
||
rows.sort(
|
||
key=lambda item: (
|
||
0 if bool(item.get("execution_ready", False)) else 1,
|
||
0 if str(item.get("inspection_status") or "") == "healthy" else 1,
|
||
str(item.get("role") or ""),
|
||
str(item.get("node_code") or ""),
|
||
)
|
||
)
|
||
|
||
return {
|
||
"execution_mode": normalized_execution_mode,
|
||
"summary": summary,
|
||
"rows": rows,
|
||
"blocking_reasons": blocking_reasons,
|
||
"warning_reasons": warning_reasons,
|
||
"recommendations": recommendations,
|
||
}
|
||
|
||
|
||
def build_release_rollout_gate(release: dict, target_nodes: list[dict], *, execution_mode: str = "remote-agent") -> dict:
|
||
normalized_release = dict(release or {})
|
||
execution_mode_label = _rollout_execution_mode_label(execution_mode)
|
||
target_rows = [
|
||
dict(item or {})
|
||
for item in list(target_nodes or [])
|
||
if str((item or {}).get("node_code") or "").strip()
|
||
]
|
||
readiness = build_rollout_target_operational_readiness(
|
||
target_rows,
|
||
execution_mode=execution_mode,
|
||
desired_release=normalized_release,
|
||
)
|
||
readiness_summary = dict(readiness.get("summary") or {})
|
||
blocking_reasons = [str(item).strip() for item in list(readiness.get("blocking_reasons") or []) if str(item).strip()]
|
||
warning_reasons = [str(item).strip() for item in list(readiness.get("warning_reasons") or []) if str(item).strip()]
|
||
recommendations = [str(item).strip() for item in list(readiness.get("recommendations") or []) if str(item).strip()]
|
||
|
||
release_id = int(normalized_release.get("id") or 0)
|
||
release_status = str(normalized_release.get("status") or "").strip()
|
||
release_version = str(normalized_release.get("release_version") or "-").strip() or "-"
|
||
release_artifact_url = str(normalized_release.get("artifact_url") or "").strip()
|
||
is_preview_release = release_id <= 0 and bool(release_version not in {"", "-"} or release_status or release_artifact_url)
|
||
|
||
if release_id <= 0 and not is_preview_release:
|
||
gate_status = "missing_release"
|
||
gate_label = "缺少 Release"
|
||
gate_type = "info"
|
||
gate_summary = "当前还没有默认 Release,正式 Rollout 入口尚未建立。"
|
||
elif release_status not in {"ready", "active"}:
|
||
gate_status = "release_preview_not_ready" if is_preview_release else "release_not_ready"
|
||
gate_label = "Release 预览未就绪" if is_preview_release else "Release 未就绪"
|
||
gate_type = "warning"
|
||
gate_summary = (
|
||
f"预览 Release {release_version} 当前状态为 {release_status or 'draft'},建议先收口到 ready 或 active。"
|
||
if is_preview_release
|
||
else f"默认 Release {release_version} 当前状态为 {release_status or 'draft'},建议先收口到 ready 或 active。"
|
||
)
|
||
elif not release_artifact_url:
|
||
gate_status = "artifact_missing"
|
||
gate_label = "预览缺少制品" if is_preview_release else "缺少制品"
|
||
gate_type = "warning"
|
||
gate_summary = (
|
||
"当前预览 Release 尚未填写 artifact_url,Rollout 下发链路仍不完整。"
|
||
if is_preview_release
|
||
else "默认 Release 尚未填写 artifact_url,Rollout 下发链路仍不完整。"
|
||
)
|
||
elif not target_rows:
|
||
gate_status = "no_targets"
|
||
gate_label = "无默认目标"
|
||
gate_type = "warning"
|
||
gate_summary = (
|
||
"当前没有在线且可视为有效执行面的预览 Rollout 目标节点。"
|
||
if is_preview_release
|
||
else "当前没有在线且可视为有效执行面的默认 Rollout 目标节点。"
|
||
)
|
||
elif blocking_reasons:
|
||
gate_status = "blocked"
|
||
gate_label = "存在阻断"
|
||
gate_type = "danger"
|
||
gate_summary = (
|
||
f"{'预览 Release 已推导完成' if is_preview_release else '默认 Release 已选定'},"
|
||
f"但目标节点里仅 {int(readiness_summary.get('execution_ready_nodes', 0) or 0)}/"
|
||
f"{int(readiness_summary.get('nodes_total', 0) or 0)} 台 {execution_mode_label} 就绪,暂不适合直接发 Rollout。"
|
||
)
|
||
elif warning_reasons:
|
||
gate_status = "attention"
|
||
gate_label = "待补巡检"
|
||
gate_type = "warning"
|
||
gate_summary = (
|
||
f"{'预览 Release 已可用于 Rollout' if is_preview_release else '默认 Release 已可用于 Rollout'},"
|
||
f"但当前仅 {int(readiness_summary.get('inspection_healthy_nodes', 0) or 0)}/"
|
||
f"{int(readiness_summary.get('nodes_total', 0) or 0)} 台目标节点完成健康巡检。"
|
||
)
|
||
else:
|
||
gate_status = "preview_ready" if is_preview_release else "ready"
|
||
gate_label = "预览可发 Rollout" if is_preview_release else "可发 Rollout"
|
||
gate_type = "success"
|
||
gate_summary = (
|
||
"当前基于最新发布包推导出的 Release 预览与目标节点均已通过门禁,已具备进入定向或灰度 Rollout 的条件。"
|
||
if is_preview_release
|
||
else "默认 Release 与默认 Rollout 目标均已通过门禁,当前可直接进入定向或灰度 Rollout。"
|
||
)
|
||
|
||
gate = {
|
||
"status": gate_status,
|
||
"status_label": gate_label,
|
||
"status_type": gate_type,
|
||
"execution_mode": str(execution_mode or "remote-agent").strip() or "remote-agent",
|
||
"execution_mode_label": execution_mode_label,
|
||
"summary": gate_summary,
|
||
"release": {
|
||
"id": release_id,
|
||
"release_version": str(normalized_release.get("release_version") or ""),
|
||
"channel": str(normalized_release.get("channel") or ""),
|
||
"status": release_status,
|
||
"status_label": _release_status_label(release_status),
|
||
"artifact_url": release_artifact_url,
|
||
"is_preview": is_preview_release,
|
||
},
|
||
"is_preview_release": is_preview_release,
|
||
"target_nodes": target_rows,
|
||
"target_node_codes": [str(item.get("node_code") or "").strip() for item in target_rows if str(item.get("node_code") or "").strip()],
|
||
"blocking_reasons": blocking_reasons,
|
||
"warning_reasons": warning_reasons,
|
||
"recommendations": recommendations,
|
||
"operational_readiness": {
|
||
"summary": readiness_summary,
|
||
"rows": list(readiness.get("rows") or []),
|
||
},
|
||
}
|
||
gate["summary_text"] = str(gate.get("summary") or "").strip()
|
||
gate["focus_ref"] = _build_release_focus_ref(
|
||
normalized_release,
|
||
section="default_rollout_gate",
|
||
)
|
||
return gate
|
||
|
||
|
||
def _release_gate_status_rank(status: str) -> int:
|
||
normalized_status = str(status or "").strip()
|
||
if normalized_status in {"ready", "preview_ready"}:
|
||
return 0
|
||
if normalized_status == "attention":
|
||
return 1
|
||
if normalized_status == "blocked":
|
||
return 2
|
||
if normalized_status == "no_targets":
|
||
return 3
|
||
if normalized_status == "artifact_missing":
|
||
return 4
|
||
if normalized_status in {"release_not_ready", "release_preview_not_ready"}:
|
||
return 5
|
||
if normalized_status == "missing_release":
|
||
return 6
|
||
return 7
|
||
|
||
|
||
def _build_release_execution_mode_reason(
|
||
*,
|
||
recommended_mode: str,
|
||
recommended_gate: dict,
|
||
alternate_gate: dict,
|
||
) -> str:
|
||
recommended_label = _rollout_execution_mode_label(recommended_mode)
|
||
alternate_mode = str(alternate_gate.get("execution_mode") or "").strip()
|
||
alternate_label = _rollout_execution_mode_label(alternate_mode) if alternate_mode else ""
|
||
recommended_status = str(recommended_gate.get("status") or "").strip()
|
||
alternate_status = str(alternate_gate.get("status") or "").strip()
|
||
recommended_summary = dict((recommended_gate.get("operational_readiness") or {}).get("summary") or {})
|
||
alternate_summary = dict((alternate_gate.get("operational_readiness") or {}).get("summary") or {})
|
||
recommended_ready_nodes = int(recommended_summary.get("execution_ready_nodes", 0) or 0)
|
||
alternate_ready_nodes = int(alternate_summary.get("execution_ready_nodes", 0) or 0)
|
||
|
||
if recommended_status in {"missing_release", "release_not_ready", "release_preview_not_ready", "artifact_missing"}:
|
||
return "当前 Release 门禁本身尚未收口,执行模式推荐仅用于预览后续下发路径。"
|
||
if recommended_status == "no_targets":
|
||
return "当前默认目标节点为空,暂时还无法形成稳定的发布执行路径。"
|
||
|
||
if recommended_mode == "ssh":
|
||
if _release_gate_status_rank(recommended_status) < _release_gate_status_rank(alternate_status):
|
||
return "当前目标节点通过 SSH 路径的门禁更完整,建议先走 SSH 发布。"
|
||
if recommended_ready_nodes > alternate_ready_nodes:
|
||
return "当前目标节点的 SSH 就绪度高于 remote-agent,就此版本建议先走 SSH 发布。"
|
||
return "当前标准 Agent 路径尚未完全收口,SSH 更适合作为当前发布入口。"
|
||
|
||
if alternate_mode == "ssh" and recommended_ready_nodes >= alternate_ready_nodes:
|
||
return "当前目标节点已满足标准 remote-agent 发布链路,默认优先沿用 remote-agent。"
|
||
return f"当前 {recommended_label} 路径不劣于 {alternate_label or '备用路径'},默认继续使用 {recommended_label}。"
|
||
|
||
|
||
def build_release_execution_mode_recommendation(release: dict, target_nodes: list[dict]) -> dict:
|
||
mode_candidates = ("remote-agent", "ssh")
|
||
gates = {
|
||
mode: build_release_rollout_gate(release, target_nodes, execution_mode=mode)
|
||
for mode in mode_candidates
|
||
}
|
||
ranked_modes = sorted(
|
||
mode_candidates,
|
||
key=lambda mode: (
|
||
_release_gate_status_rank(gates[mode].get("status")),
|
||
-int(((gates[mode].get("operational_readiness") or {}).get("summary") or {}).get("execution_ready_nodes", 0) or 0),
|
||
len(list(gates[mode].get("blocking_reasons") or [])),
|
||
len(list(gates[mode].get("warning_reasons") or [])),
|
||
0 if mode == "remote-agent" else 1,
|
||
),
|
||
)
|
||
recommended_mode = ranked_modes[0] if ranked_modes else "remote-agent"
|
||
fallback_mode = ranked_modes[1] if len(ranked_modes) > 1 else ""
|
||
recommended_gate = dict(gates.get(recommended_mode) or {})
|
||
fallback_gate = dict(gates.get(fallback_mode) or {})
|
||
|
||
options = []
|
||
for mode in mode_candidates:
|
||
gate = dict(gates.get(mode) or {})
|
||
readiness_summary = dict((gate.get("operational_readiness") or {}).get("summary") or {})
|
||
options.append(
|
||
{
|
||
"mode": mode,
|
||
"mode_label": _rollout_execution_mode_label(mode),
|
||
"status": str(gate.get("status") or ""),
|
||
"status_label": str(gate.get("status_label") or ""),
|
||
"summary": str(gate.get("summary") or ""),
|
||
"execution_ready_nodes": int(readiness_summary.get("execution_ready_nodes", 0) or 0),
|
||
"nodes_total": int(readiness_summary.get("nodes_total", 0) or 0),
|
||
}
|
||
)
|
||
|
||
return {
|
||
"recommended_mode": recommended_mode,
|
||
"recommended_mode_label": _rollout_execution_mode_label(recommended_mode),
|
||
"fallback_mode": fallback_mode,
|
||
"fallback_mode_label": _rollout_execution_mode_label(fallback_mode) if fallback_mode else "",
|
||
"reason": _build_release_execution_mode_reason(
|
||
recommended_mode=recommended_mode,
|
||
recommended_gate=recommended_gate,
|
||
alternate_gate=fallback_gate,
|
||
),
|
||
"options": options,
|
||
"recommended_gate": recommended_gate,
|
||
"gates": gates,
|
||
}
|
||
|
||
|
||
def _resolve_rollout_targets(selector: dict) -> list[dict]:
|
||
from app.services.cluster_runtime_service import get_cluster_snapshot
|
||
from app.services.ops_job_service import list_managed_nodes
|
||
|
||
selector = selector or {}
|
||
requested_node_codes_list = [
|
||
str(item).strip()
|
||
for item in list(selector.get("node_codes") or [])
|
||
if str(item).strip()
|
||
]
|
||
requested_node_codes = set(requested_node_codes_list)
|
||
requested_node_code_order = {
|
||
node_code: index for index, node_code in enumerate(requested_node_codes_list)
|
||
}
|
||
requested_region = str(selector.get("region") or "").strip()
|
||
requested_role = str(selector.get("role") or "").strip()
|
||
only_enabled = bool(selector.get("only_enabled", True))
|
||
only_effective_workers = bool(selector.get("only_effective_workers", False))
|
||
only_online = bool(selector.get("only_online", False))
|
||
|
||
managed_nodes = list_managed_nodes()
|
||
managed_map = {
|
||
str(item.get("node_code") or ""): item
|
||
for item in managed_nodes
|
||
if str(item.get("node_code") or "")
|
||
}
|
||
cluster_nodes = list((get_cluster_snapshot().get("nodes") or []))
|
||
|
||
merged: dict[str, dict] = {}
|
||
for cluster_node in cluster_nodes:
|
||
node_code = str(cluster_node.get("node_code") or "").strip()
|
||
if not node_code:
|
||
continue
|
||
managed = managed_map.get(node_code, {})
|
||
merged[node_code] = {
|
||
"node_code": node_code,
|
||
"region": str(cluster_node.get("region") or managed.get("region") or ""),
|
||
"role": str(cluster_node.get("role") or managed.get("role") or ""),
|
||
"status": str(cluster_node.get("status") or ""),
|
||
"current_load": int(cluster_node.get("current_load", 0) or 0),
|
||
"is_effective_worker": bool(cluster_node.get("is_effective_worker", False)),
|
||
"detect_participating": bool(cluster_node.get("detect_participating", False)),
|
||
"is_enabled": bool(managed.get("is_enabled", True)),
|
||
"ssh_host": str(managed.get("ssh_host") or ""),
|
||
"ssh_user": str(managed.get("ssh_user") or ""),
|
||
"ssh_port": int(managed.get("ssh_port") or 22),
|
||
"auth_mode": str(managed.get("auth_mode") or "key"),
|
||
"deploy_channel": str(managed.get("deploy_channel") or "stable"),
|
||
"title": str(managed.get("title") or node_code),
|
||
}
|
||
for managed in managed_nodes:
|
||
node_code = str(managed.get("node_code") or "").strip()
|
||
if node_code and node_code not in merged:
|
||
merged[node_code] = {
|
||
"node_code": node_code,
|
||
"region": str(managed.get("region") or ""),
|
||
"role": str(managed.get("role") or ""),
|
||
"status": "",
|
||
"current_load": 0,
|
||
"is_effective_worker": False,
|
||
"detect_participating": False,
|
||
"is_enabled": bool(managed.get("is_enabled", True)),
|
||
"ssh_host": str(managed.get("ssh_host") or ""),
|
||
"ssh_user": str(managed.get("ssh_user") or ""),
|
||
"ssh_port": int(managed.get("ssh_port") or 22),
|
||
"auth_mode": str(managed.get("auth_mode") or "key"),
|
||
"deploy_channel": str(managed.get("deploy_channel") or "stable"),
|
||
"title": str(managed.get("title") or node_code),
|
||
}
|
||
|
||
targets: list[dict] = []
|
||
for item in merged.values():
|
||
if requested_node_codes and str(item.get("node_code") or "") not in requested_node_codes:
|
||
continue
|
||
if requested_region and str(item.get("region") or "") != requested_region:
|
||
continue
|
||
if requested_role and str(item.get("role") or "") != requested_role:
|
||
continue
|
||
if only_enabled and not bool(item.get("is_enabled", False)):
|
||
continue
|
||
if only_effective_workers and not bool(item.get("is_effective_worker", False)):
|
||
continue
|
||
if only_online and str(item.get("status") or "") not in {"online", "busy"}:
|
||
continue
|
||
targets.append(item)
|
||
|
||
role_order = {
|
||
"worker": 0,
|
||
"control": 1,
|
||
}
|
||
|
||
return sorted(
|
||
targets,
|
||
key=lambda item: (
|
||
requested_node_code_order.get(str(item.get("node_code") or ""), 10_000),
|
||
str(item.get("region") or ""),
|
||
role_order.get(str(item.get("role") or ""), 9),
|
||
str(item.get("node_code") or ""),
|
||
),
|
||
)
|
||
|
||
|
||
def create_release(payload: dict) -> tuple[bool, str, dict]:
|
||
ensure_ops_release_schema()
|
||
release_version = str(payload.get("release_version") or "").strip()
|
||
if not release_version:
|
||
return False, "release_version 不能为空", {}
|
||
|
||
channel = str(payload.get("channel") or "stable").strip() or "stable"
|
||
commit_sha = str(payload.get("commit_sha") or "").strip()
|
||
status = str(payload.get("status") or "draft").strip() or "draft"
|
||
artifact_url = str(payload.get("artifact_url") or "").strip()
|
||
checksum = str(payload.get("checksum") or "").strip()
|
||
notes = str(payload.get("notes") or "")
|
||
metadata = payload.get("metadata") or {}
|
||
created_by = str(payload.get("created_by") or "api").strip() or "api"
|
||
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
INSERT INTO ops_releases (
|
||
release_version, channel, commit_sha, status, artifact_url, checksum, notes, metadata_json, created_by, created_at, updated_at
|
||
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||
RETURNING
|
||
id, release_version, channel, commit_sha, status, artifact_url, checksum, notes, metadata_json, created_by, created_at, activated_at, updated_at
|
||
""",
|
||
(
|
||
release_version,
|
||
channel,
|
||
commit_sha,
|
||
status,
|
||
artifact_url,
|
||
checksum,
|
||
notes,
|
||
json.dumps(metadata, ensure_ascii=False),
|
||
created_by,
|
||
),
|
||
)
|
||
row = cur.fetchone()
|
||
conn.commit()
|
||
return True, "Release 已创建", {"release": _serialize_release_row(row)}
|
||
|
||
|
||
def _materialize_release_preview_payload(create_payload: dict) -> dict:
|
||
payload = dict(create_payload or {})
|
||
return {
|
||
"id": 0,
|
||
"release_version": str(payload.get("release_version") or "").strip(),
|
||
"channel": str(payload.get("channel") or "stable").strip() or "stable",
|
||
"commit_sha": str(payload.get("commit_sha") or "").strip(),
|
||
"status": str(payload.get("status") or "ready").strip() or "ready",
|
||
"artifact_url": str(payload.get("artifact_url") or "").strip(),
|
||
"checksum": str(payload.get("checksum") or "").strip(),
|
||
"notes": str(payload.get("notes") or ""),
|
||
"metadata": dict(payload.get("metadata") or {}),
|
||
"created_by": str(payload.get("created_by") or "api/package-import").strip() or "api/package-import",
|
||
"created_at": "",
|
||
"activated_at": "",
|
||
"updated_at": "",
|
||
"is_preview": True,
|
||
}
|
||
|
||
|
||
def _resolve_release_from_latest_package(
|
||
payload: dict | None = None,
|
||
*,
|
||
control_plane_base_url: str = "",
|
||
) -> tuple[bool, str, dict]:
|
||
normalized_payload = dict(payload or {})
|
||
package_metadata = get_latest_release_package_metadata()
|
||
if not bool(package_metadata.get("available")):
|
||
return False, str(package_metadata.get("reason") or "当前没有可用的本机发布包"), {"package": package_metadata}
|
||
|
||
release_version = (
|
||
str(normalized_payload.get("release_version") or "").strip()
|
||
or str(package_metadata.get("release_version_suggestion") or "").strip()
|
||
or str(package_metadata.get("package_name") or "").strip()
|
||
)
|
||
if not release_version:
|
||
return False, "最新发布包缺少 release_version 建议值", {"package": package_metadata}
|
||
|
||
existing_release = _get_release_by_version(release_version)
|
||
|
||
normalized_base_url = _normalize_public_base_url(
|
||
str(normalized_payload.get("control_plane_base_url") or "").strip() or control_plane_base_url
|
||
)
|
||
artifact_url = _build_absolute_release_package_url(normalized_base_url, str(package_metadata.get("download_path") or ""))
|
||
sha256_download_url = _build_absolute_release_package_url(
|
||
normalized_base_url, str(package_metadata.get("sha256_download_path") or "")
|
||
)
|
||
if existing_release and not artifact_url:
|
||
release_preview = dict(existing_release or {})
|
||
return True, "已解析最新发布包对应的 Release 信息", {
|
||
"release": dict(existing_release or {}),
|
||
"release_preview": release_preview,
|
||
"preview_source": "existing_release",
|
||
"deduplicated": True,
|
||
"package": package_metadata,
|
||
"resolved_artifact_url": "",
|
||
"resolved_sha256_download_url": sha256_download_url,
|
||
"release_payload": {},
|
||
}
|
||
if not artifact_url:
|
||
return False, "当前无法为最新发布包生成可下载的 artifact_url", {"package": package_metadata}
|
||
|
||
notes = str(normalized_payload.get("notes") or "").strip() or str(package_metadata.get("notes_suggestion") or "").strip()
|
||
notes_suffix = str(normalized_payload.get("notes_suffix") or "").strip()
|
||
if notes_suffix:
|
||
notes = f"{notes}\n{notes_suffix}".strip() if notes else notes_suffix
|
||
|
||
package_source_metadata = {
|
||
**dict(package_metadata.get("metadata") or {}),
|
||
"download_path": str(package_metadata.get("download_path") or "").strip(),
|
||
"sha256_download_path": str(package_metadata.get("sha256_download_path") or "").strip(),
|
||
"download_url": artifact_url,
|
||
"sha256_download_url": sha256_download_url,
|
||
}
|
||
release_metadata = {
|
||
**dict(normalized_payload.get("metadata") or {}),
|
||
"created_from": "latest_release_package",
|
||
"package_source": package_source_metadata,
|
||
}
|
||
|
||
create_payload = {
|
||
"release_version": release_version,
|
||
"channel": str(normalized_payload.get("channel") or "stable").strip() or "stable",
|
||
"status": str(normalized_payload.get("status") or "ready").strip() or "ready",
|
||
"commit_sha": (
|
||
str(normalized_payload.get("commit_sha") or "").strip()
|
||
or str(package_metadata.get("commit_sha") or "").strip()
|
||
or str(package_source_metadata.get("commit_sha") or "").strip()
|
||
),
|
||
"artifact_url": artifact_url,
|
||
"checksum": str(normalized_payload.get("checksum") or "").strip() or str(package_metadata.get("sha256") or "").strip(),
|
||
"notes": notes,
|
||
"metadata": release_metadata,
|
||
"created_by": str(normalized_payload.get("created_by") or "api/package-import").strip() or "api/package-import",
|
||
}
|
||
release_preview = existing_release or _materialize_release_preview_payload(create_payload)
|
||
return True, "已解析最新发布包对应的 Release 信息", {
|
||
"release": dict(existing_release or {}),
|
||
"release_preview": release_preview,
|
||
"preview_source": "existing_release" if existing_release else "latest_package_draft",
|
||
"deduplicated": bool(existing_release),
|
||
"package": package_metadata,
|
||
"resolved_artifact_url": artifact_url,
|
||
"resolved_sha256_download_url": sha256_download_url,
|
||
"release_payload": create_payload,
|
||
}
|
||
|
||
|
||
def create_release_from_latest_package(
|
||
payload: dict | None = None,
|
||
*,
|
||
control_plane_base_url: str = "",
|
||
) -> tuple[bool, str, dict]:
|
||
ok, message, resolved = _resolve_release_from_latest_package(
|
||
payload,
|
||
control_plane_base_url=control_plane_base_url,
|
||
)
|
||
if not ok:
|
||
return False, message, resolved
|
||
|
||
gate_ok, gate_message, gate_data = _validate_latest_release_package_final_gate(resolved.get("package") or {})
|
||
if not gate_ok:
|
||
return False, gate_message, {
|
||
**dict(resolved or {}),
|
||
**dict(gate_data or {}),
|
||
}
|
||
|
||
existing_release = dict(resolved.get("release") or {})
|
||
if existing_release:
|
||
return True, "Release 已存在,按幂等返回", {
|
||
"release": existing_release,
|
||
"deduplicated": True,
|
||
"package": dict(resolved.get("package") or {}),
|
||
"final_release_gate": dict(gate_data.get("final_release_gate") or {}),
|
||
"resolved_artifact_url": str(resolved.get("resolved_artifact_url") or ""),
|
||
"resolved_sha256_download_url": str(resolved.get("resolved_sha256_download_url") or ""),
|
||
"preview_source": "existing_release",
|
||
}
|
||
|
||
create_payload = dict(resolved.get("release_payload") or {})
|
||
ok, message, data = create_release(create_payload)
|
||
if ok:
|
||
data["deduplicated"] = False
|
||
data["package"] = dict(resolved.get("package") or {})
|
||
data["final_release_gate"] = dict(gate_data.get("final_release_gate") or {})
|
||
data["resolved_artifact_url"] = str(resolved.get("resolved_artifact_url") or "")
|
||
data["resolved_sha256_download_url"] = str(resolved.get("resolved_sha256_download_url") or "")
|
||
data["preview_source"] = "latest_package_draft"
|
||
return ok, message, data
|
||
|
||
|
||
def preview_release_from_latest_package(
|
||
payload: dict | None = None,
|
||
*,
|
||
control_plane_base_url: str = "",
|
||
) -> tuple[bool, str, dict]:
|
||
ok, message, resolved = _resolve_release_from_latest_package(
|
||
payload,
|
||
control_plane_base_url=control_plane_base_url,
|
||
)
|
||
if not ok:
|
||
return False, message, resolved
|
||
|
||
preview_release = dict(resolved.get("release_preview") or resolved.get("release") or {})
|
||
return True, "最新发布包预览已生成", {
|
||
**dict(resolved or {}),
|
||
"release": preview_release,
|
||
"release_created": False,
|
||
"release_deduplicated": bool(resolved.get("deduplicated", False)),
|
||
}
|
||
|
||
|
||
def _normalize_smart_rollout_mode(raw_value: object) -> str:
|
||
normalized = str(raw_value or "").strip().lower()
|
||
if normalized in {"control", "controller"}:
|
||
return "control"
|
||
if normalized in {"auto", ""}:
|
||
return "auto"
|
||
return "worker"
|
||
|
||
|
||
def _smart_rollout_mode_label(mode: str) -> str:
|
||
return "Control 发布" if str(mode or "").strip() == "control" else "Worker 灰度"
|
||
|
||
|
||
def _smart_rollout_node_role(node: dict) -> str:
|
||
return str(node.get("role") or node.get("cluster_role") or "").strip().lower()
|
||
|
||
|
||
def _smart_rollout_node_region(node: dict) -> str:
|
||
return str(node.get("region") or node.get("cluster_region") or "").strip().lower()
|
||
|
||
|
||
def _smart_rollout_node_status(node: dict) -> str:
|
||
return str(node.get("cluster_status") or node.get("status") or "").strip().lower()
|
||
|
||
|
||
def _is_smart_rollout_online_node(node: dict) -> bool:
|
||
return _smart_rollout_node_status(node) in {"online", "busy"}
|
||
|
||
|
||
def _smart_rollout_delivery_queue_state(node: dict) -> str:
|
||
return str(
|
||
node.get("delivery_queue_state")
|
||
or ((node.get("metadata") or {}).get("delivery_queue") or {}).get("state")
|
||
or ""
|
||
).strip().lower()
|
||
|
||
|
||
def _is_smart_rollout_remote_ready(node: dict) -> bool:
|
||
queue_state = _smart_rollout_delivery_queue_state(node)
|
||
return bool(
|
||
node.get("is_managed", False)
|
||
and node.get("is_enabled", False)
|
||
and node.get("is_agent_online", False)
|
||
and queue_state != "dead_letter"
|
||
)
|
||
|
||
|
||
def _is_smart_rollout_ssh_ready(node: dict) -> bool:
|
||
queue_state = _smart_rollout_delivery_queue_state(node)
|
||
return bool(
|
||
node.get("is_managed", False)
|
||
and node.get("is_enabled", False)
|
||
and node.get("has_ssh_access", False)
|
||
and queue_state != "dead_letter"
|
||
)
|
||
|
||
|
||
def _is_smart_rollout_execution_ready(node: dict, execution_mode: str) -> bool:
|
||
if str(execution_mode or "").strip() == "ssh":
|
||
return _is_smart_rollout_ssh_ready(node)
|
||
return _is_smart_rollout_remote_ready(node)
|
||
|
||
|
||
def _pick_smart_rollout_execution_mode(candidates: list[dict]) -> tuple[str, list[dict]]:
|
||
remote_ready_candidates = [item for item in list(candidates or []) if _is_smart_rollout_remote_ready(item)]
|
||
ssh_ready_candidates = [item for item in list(candidates or []) if _is_smart_rollout_ssh_ready(item)]
|
||
if ssh_ready_candidates and len(ssh_ready_candidates) > len(remote_ready_candidates):
|
||
return "ssh", ssh_ready_candidates
|
||
if remote_ready_candidates:
|
||
return "remote-agent", remote_ready_candidates
|
||
if ssh_ready_candidates:
|
||
return "ssh", ssh_ready_candidates
|
||
return "", []
|
||
|
||
|
||
def _is_smart_rollout_effective_worker(node: dict) -> bool:
|
||
if "cluster_is_effective_worker" in node:
|
||
return bool(node.get("cluster_is_effective_worker", False))
|
||
if "is_effective_worker" in node:
|
||
return bool(node.get("is_effective_worker", False))
|
||
return _smart_rollout_node_role(node) == "worker"
|
||
|
||
|
||
def _pick_smart_rollout_dominant_region(nodes: list[dict]) -> str:
|
||
counts: dict[str, int] = {}
|
||
for item in nodes:
|
||
region = _smart_rollout_node_region(item)
|
||
if not region:
|
||
continue
|
||
counts[region] = int(counts.get(region, 0) or 0) + 1
|
||
if not counts:
|
||
return ""
|
||
return sorted(counts.items(), key=lambda item: (-int(item[1]), str(item[0])))[0][0]
|
||
|
||
|
||
def _sort_smart_rollout_candidates(nodes: list[dict]) -> list[dict]:
|
||
def _candidate_key(node: dict) -> tuple:
|
||
status = _smart_rollout_node_status(node)
|
||
status_rank = 0 if status == "online" else (1 if status == "busy" else 9)
|
||
current_load = int(node.get("cluster_current_load", node.get("current_load", 0)) or 0)
|
||
participating = bool(node.get("cluster_detect_participating", node.get("detect_participating", False)))
|
||
access_rank = 0 if _is_smart_rollout_remote_ready(node) else (1 if _is_smart_rollout_ssh_ready(node) else 2)
|
||
return (
|
||
access_rank,
|
||
0 if not participating else 1,
|
||
status_rank,
|
||
current_load,
|
||
str(node.get("node_code") or ""),
|
||
)
|
||
|
||
return sorted(list(nodes or []), key=_candidate_key)
|
||
|
||
|
||
def _build_smart_rollout_role_policy(mode: str, *, execution_mode: str = "remote-agent") -> dict:
|
||
if str(mode or "").strip() == "control":
|
||
deploy_payload = {
|
||
"restart_services": ["domaincheck-api", "domaincheck-worker", "domaincheck-sync-agent"],
|
||
"health_check_urls": ["http://127.0.0.1:8100/health"],
|
||
"health_check_services": ["domaincheck-api", "domaincheck-worker", "domaincheck-sync-agent"],
|
||
# Control 节点启动期间会先经历较长的 import / startup hook,
|
||
# systemd 已经 active 但 /health 仍可能在 15-20 秒内拒绝连接。
|
||
# 这里把健康检查窗口放宽到约 40 秒,避免被误回滚。
|
||
"health_check_timeout_seconds": 20,
|
||
"health_check_retries": 9,
|
||
"health_check_interval_seconds": 4,
|
||
"rollback_on_failure": True,
|
||
"switch_current": True,
|
||
}
|
||
else:
|
||
deploy_payload = {
|
||
"restart_services": ["domaincheck-worker"],
|
||
"health_check_urls": [],
|
||
"health_check_services": ["domaincheck-worker"],
|
||
"health_check_timeout_seconds": 10,
|
||
"health_check_retries": 2,
|
||
"health_check_interval_seconds": 2,
|
||
"rollback_on_failure": True,
|
||
"switch_current": True,
|
||
}
|
||
return {
|
||
"auto_start": True,
|
||
"auto_dispatch": False,
|
||
"auto_advance": True,
|
||
"auto_approve": False,
|
||
"continue_on_failure": False,
|
||
"continue_on_partial": False,
|
||
"execution_mode": str(execution_mode or "remote-agent").strip() or "remote-agent",
|
||
"first_batch_size": 1,
|
||
"batch_size": 1,
|
||
"deploy_payload": deploy_payload,
|
||
}
|
||
|
||
|
||
def build_smart_release_rollout_plan(mode: str = "auto") -> dict:
|
||
from app.services.ops_agent_service import list_managed_nodes_with_agent_state
|
||
|
||
normalized_mode = _normalize_smart_rollout_mode(mode)
|
||
managed_nodes_payload = list_managed_nodes_with_agent_state()
|
||
source_nodes = list(managed_nodes_payload.get("nodes") or [])
|
||
online_enabled_nodes = [
|
||
dict(item or {})
|
||
for item in source_nodes
|
||
if bool((item or {}).get("is_enabled", False)) and _is_smart_rollout_online_node(dict(item or {}))
|
||
]
|
||
if not online_enabled_nodes:
|
||
return {
|
||
"available": False,
|
||
"mode": normalized_mode,
|
||
"mode_label": _smart_rollout_mode_label(normalized_mode),
|
||
"reason": "当前没有在线且启用的节点可用于智能 Rollout。",
|
||
"target_selector": {},
|
||
"policy": {},
|
||
"target_nodes": [],
|
||
}
|
||
|
||
effective_worker_nodes = [item for item in online_enabled_nodes if _is_smart_rollout_effective_worker(item)]
|
||
remote_ready_worker_nodes = [item for item in effective_worker_nodes if _is_smart_rollout_remote_ready(item)]
|
||
ssh_ready_worker_nodes = [item for item in effective_worker_nodes if _is_smart_rollout_ssh_ready(item)]
|
||
preferred_ready_worker_nodes = (
|
||
ssh_ready_worker_nodes if len(ssh_ready_worker_nodes) > len(remote_ready_worker_nodes) else remote_ready_worker_nodes
|
||
)
|
||
dominant_region = _pick_smart_rollout_dominant_region(preferred_ready_worker_nodes or effective_worker_nodes or online_enabled_nodes)
|
||
resolved_mode = normalized_mode
|
||
if resolved_mode == "auto":
|
||
resolved_mode = "worker" if effective_worker_nodes else "control"
|
||
|
||
def _matches_mode(node: dict) -> bool:
|
||
role = _smart_rollout_node_role(node)
|
||
if resolved_mode == "control":
|
||
return role == "control"
|
||
return role == "worker" and _is_smart_rollout_effective_worker(node)
|
||
|
||
regional_candidates = [
|
||
item
|
||
for item in online_enabled_nodes
|
||
if _matches_mode(item) and (not dominant_region or _smart_rollout_node_region(item) == dominant_region)
|
||
]
|
||
fallback_candidates = [item for item in online_enabled_nodes if _matches_mode(item)]
|
||
ranked_candidates = _sort_smart_rollout_candidates(regional_candidates or fallback_candidates)
|
||
execution_mode, ready_candidates = _pick_smart_rollout_execution_mode(ranked_candidates)
|
||
candidates = ready_candidates or ranked_candidates
|
||
if not candidates:
|
||
return {
|
||
"available": False,
|
||
"mode": resolved_mode,
|
||
"mode_label": _smart_rollout_mode_label(resolved_mode),
|
||
"reason": "当前没有满足角色条件的在线节点可用于智能 Rollout。",
|
||
"target_selector": {},
|
||
"policy": {},
|
||
"target_nodes": [],
|
||
}
|
||
if not ready_candidates:
|
||
dead_letter_nodes = [
|
||
str(item.get("node_code") or "").strip()
|
||
for item in ranked_candidates
|
||
if _smart_rollout_delivery_queue_state(item) == "dead_letter" and str(item.get("node_code") or "").strip()
|
||
]
|
||
reason = "当前命中的在线节点尚未进入可执行发布状态,remote-agent 与 SSH 两条路径都还不够稳定,暂不建议一键创建 Rollout。"
|
||
if dead_letter_nodes:
|
||
reason = (
|
||
f"当前命中的候选节点里有 {len(dead_letter_nodes)} 台存在 Node Agent 死信回执,"
|
||
"需先收口 delivery queue,再继续智能 Rollout。"
|
||
)
|
||
return {
|
||
"available": False,
|
||
"mode": resolved_mode,
|
||
"mode_label": _smart_rollout_mode_label(resolved_mode),
|
||
"reason": reason,
|
||
"target_selector": {},
|
||
"policy": {},
|
||
"target_nodes": [
|
||
{
|
||
"node_code": str(item.get("node_code") or ""),
|
||
"region": str(item.get("region") or ""),
|
||
"role": str(item.get("role") or ""),
|
||
"status": _smart_rollout_node_status(item),
|
||
"current_load": int(item.get("cluster_current_load", item.get("current_load", 0)) or 0),
|
||
"remote_agent_ready": _is_smart_rollout_remote_ready(item),
|
||
"ssh_ready": _is_smart_rollout_ssh_ready(item),
|
||
"execution_ready": False,
|
||
"delivery_queue_state": _smart_rollout_delivery_queue_state(item),
|
||
}
|
||
for item in ranked_candidates
|
||
],
|
||
}
|
||
|
||
node_codes = [str(item.get("node_code") or "").strip() for item in candidates if str(item.get("node_code") or "").strip()]
|
||
batch_size = 1 if len(node_codes) <= 2 else (2 if len(node_codes) <= 5 else 3)
|
||
policy = _build_smart_rollout_role_policy(resolved_mode, execution_mode=execution_mode)
|
||
policy["batch_size"] = batch_size
|
||
selector = {
|
||
"node_codes": node_codes,
|
||
"region": "",
|
||
"role": "",
|
||
"only_enabled": True,
|
||
"only_online": True,
|
||
"only_effective_workers": False,
|
||
"smart_generated": True,
|
||
"smart_rollout_mode": resolved_mode,
|
||
}
|
||
return {
|
||
"available": True,
|
||
"mode": resolved_mode,
|
||
"mode_label": _smart_rollout_mode_label(resolved_mode),
|
||
"reason": "",
|
||
"summary": (
|
||
f"已按 {_smart_rollout_mode_label(resolved_mode)} 预选 {len(node_codes)} 台节点,"
|
||
f"执行路径采用 {_rollout_execution_mode_label(execution_mode)};"
|
||
f"优先锁定 {dominant_region or _smart_rollout_node_region(candidates[0]) or '当前区域'};"
|
||
f"首批 1 台,后续每批 {batch_size} 台。"
|
||
),
|
||
"dominant_region": dominant_region or _smart_rollout_node_region(candidates[0]),
|
||
"execution_mode": execution_mode,
|
||
"execution_mode_label": _rollout_execution_mode_label(execution_mode),
|
||
"target_selector": selector,
|
||
"policy": policy,
|
||
"target_nodes": [
|
||
{
|
||
"node_code": str(item.get("node_code") or ""),
|
||
"region": str(item.get("region") or ""),
|
||
"role": str(item.get("role") or ""),
|
||
"status": _smart_rollout_node_status(item),
|
||
"current_load": int(item.get("cluster_current_load", item.get("current_load", 0)) or 0),
|
||
"remote_agent_ready": _is_smart_rollout_remote_ready(item),
|
||
"ssh_ready": _is_smart_rollout_ssh_ready(item),
|
||
"execution_ready": _is_smart_rollout_execution_ready(item, execution_mode),
|
||
"detect_participating": bool(item.get("cluster_detect_participating", False)),
|
||
"delivery_queue_state": _smart_rollout_delivery_queue_state(item),
|
||
}
|
||
for item in candidates
|
||
],
|
||
}
|
||
|
||
|
||
def _preview_smart_release_rollout_for_release(
|
||
release: dict,
|
||
normalized_payload: dict,
|
||
rollout_mode: str,
|
||
) -> tuple[bool, str, dict]:
|
||
from app.services.ops_policy_service import preview_ops_job_policy
|
||
|
||
safe_release = dict(release or {})
|
||
smart_rollout = build_smart_release_rollout_plan(rollout_mode)
|
||
data = {
|
||
"release": safe_release,
|
||
"release_created": False,
|
||
"rollout_created": False,
|
||
"requires_confirmation": False,
|
||
"blocked": False,
|
||
"smart_rollout": smart_rollout,
|
||
"preview": {},
|
||
"rollout_payload": {},
|
||
}
|
||
if not bool(smart_rollout.get("available", False)):
|
||
return True, "当前无法自动生成智能 Rollout", data
|
||
|
||
rollout_payload = {
|
||
"created_by": str(normalized_payload.get("rollout_created_by") or normalized_payload.get("created_by") or "api/smart-rollout").strip()
|
||
or "api/smart-rollout",
|
||
"target_selector": dict(smart_rollout.get("target_selector") or {}),
|
||
"policy": dict(smart_rollout.get("policy") or {}),
|
||
}
|
||
preview = preview_ops_job_policy(
|
||
{
|
||
"action": "deploy.release",
|
||
"target_type": "rollout",
|
||
"release_id": int(safe_release.get("id") or 0),
|
||
"release": safe_release,
|
||
"target_selector": rollout_payload["target_selector"],
|
||
"policy": rollout_payload["policy"],
|
||
}
|
||
)
|
||
data["preview"] = preview
|
||
data["rollout_payload"] = rollout_payload
|
||
if bool(preview.get("blocked", False)):
|
||
data["blocked"] = True
|
||
return True, "智能 Rollout 预检未通过", data
|
||
|
||
confirm_risky = bool(normalized_payload.get("confirm_risky", False))
|
||
if (list(preview.get("warnings") or []) or bool(preview.get("approval_required", False))) and not confirm_risky:
|
||
data["requires_confirmation"] = True
|
||
return True, "智能 Rollout 需要确认", data
|
||
|
||
return True, "智能 Rollout 预检通过", data
|
||
|
||
|
||
def _create_smart_release_rollout_for_release(
|
||
release: dict,
|
||
normalized_payload: dict,
|
||
rollout_mode: str,
|
||
) -> tuple[bool, str, dict]:
|
||
release_id = int((release or {}).get("id") or 0)
|
||
if release_id <= 0:
|
||
return False, "智能 Rollout 需要有效的 release_id", {}
|
||
|
||
preview_ok, preview_message, data = _preview_smart_release_rollout_for_release(release, normalized_payload, rollout_mode)
|
||
if not preview_ok:
|
||
return False, preview_message, data
|
||
if preview_message in {"当前无法自动生成智能 Rollout", "智能 Rollout 预检未通过", "智能 Rollout 需要确认"}:
|
||
return True, preview_message, data
|
||
|
||
rollout_payload = dict(data.get("rollout_payload") or {})
|
||
rollout_ok, rollout_message, rollout_data = create_release_rollout(release_id, rollout_payload)
|
||
data.update(dict(rollout_data or {}))
|
||
data["rollout_created"] = bool(rollout_ok and dict(rollout_data or {}).get("rollout"))
|
||
data["requires_confirmation"] = False
|
||
data["blocked"] = False
|
||
if rollout_ok:
|
||
return True, rollout_message, data
|
||
return False, rollout_message, data
|
||
|
||
|
||
def preview_smart_release_rollout(
|
||
release_id: int,
|
||
payload: dict | None = None,
|
||
) -> tuple[bool, str, dict]:
|
||
normalized_payload = dict(payload or {})
|
||
safe_release_id = int(release_id or normalized_payload.get("release_id") or 0)
|
||
if safe_release_id <= 0:
|
||
return False, "智能 Rollout 需要有效的 release_id", {}
|
||
|
||
release = get_release(safe_release_id)
|
||
if not release:
|
||
return False, "Release 不存在", {}
|
||
|
||
rollout_mode = _normalize_smart_rollout_mode(normalized_payload.get("mode"))
|
||
ok, message, data = _preview_smart_release_rollout_for_release(release, normalized_payload, rollout_mode)
|
||
merged_data = {
|
||
**dict(data or {}),
|
||
"release": release,
|
||
"release_created": False,
|
||
"release_deduplicated": False,
|
||
}
|
||
return ok, message, merged_data
|
||
|
||
|
||
def create_smart_release_rollout(
|
||
release_id: int,
|
||
payload: dict | None = None,
|
||
) -> tuple[bool, str, dict]:
|
||
normalized_payload = dict(payload or {})
|
||
safe_release_id = int(release_id or normalized_payload.get("release_id") or 0)
|
||
if safe_release_id <= 0:
|
||
return False, "智能 Rollout 需要有效的 release_id", {}
|
||
|
||
release = get_release(safe_release_id)
|
||
if not release:
|
||
return False, "Release 不存在", {}
|
||
|
||
rollout_mode = _normalize_smart_rollout_mode(normalized_payload.get("mode"))
|
||
ok, message, data = _create_smart_release_rollout_for_release(release, normalized_payload, rollout_mode)
|
||
merged_data = {
|
||
**dict(data or {}),
|
||
"release": release,
|
||
"release_created": False,
|
||
}
|
||
return ok, message, merged_data
|
||
|
||
|
||
def preview_smart_release_rollout_from_latest_package(
|
||
payload: dict | None = None,
|
||
*,
|
||
control_plane_base_url: str = "",
|
||
) -> tuple[bool, str, dict]:
|
||
normalized_payload = dict(payload or {})
|
||
preview_ok, preview_message, preview_data = preview_release_from_latest_package(
|
||
normalized_payload,
|
||
control_plane_base_url=control_plane_base_url,
|
||
)
|
||
if not preview_ok:
|
||
return False, preview_message, preview_data
|
||
|
||
release = dict(preview_data.get("release") or {})
|
||
rollout_mode = _normalize_smart_rollout_mode(normalized_payload.get("mode"))
|
||
rollout_ok, rollout_message, rollout_data = _preview_smart_release_rollout_for_release(
|
||
release,
|
||
normalized_payload,
|
||
rollout_mode,
|
||
)
|
||
data = {
|
||
**dict(preview_data or {}),
|
||
**dict(rollout_data or {}),
|
||
"release": release,
|
||
"release_created": False,
|
||
"release_deduplicated": bool(dict(preview_data or {}).get("release_deduplicated", False)),
|
||
"created_release_message": preview_message,
|
||
}
|
||
return rollout_ok, rollout_message, data
|
||
|
||
|
||
def create_release_and_smart_rollout_from_latest_package(
|
||
payload: dict | None = None,
|
||
*,
|
||
control_plane_base_url: str = "",
|
||
) -> tuple[bool, str, dict]:
|
||
normalized_payload = dict(payload or {})
|
||
rollout_mode = _normalize_smart_rollout_mode(normalized_payload.get("mode"))
|
||
|
||
release_ok, release_message, release_data = create_release_from_latest_package(
|
||
normalized_payload,
|
||
control_plane_base_url=control_plane_base_url,
|
||
)
|
||
if not release_ok:
|
||
return False, release_message, release_data
|
||
|
||
release = dict(release_data.get("release") or {})
|
||
if int(release.get("id") or 0) <= 0:
|
||
return False, "Release 创建结果缺少有效 release_id", release_data
|
||
|
||
rollout_ok, rollout_message, rollout_data = _create_smart_release_rollout_for_release(release, normalized_payload, rollout_mode)
|
||
data = {
|
||
**dict(release_data or {}),
|
||
**dict(rollout_data or {}),
|
||
"release": release,
|
||
"release_created": not bool(dict(release_data or {}).get("deduplicated", False)),
|
||
"release_deduplicated": bool(dict(release_data or {}).get("deduplicated", False)),
|
||
"created_release_message": release_message,
|
||
}
|
||
if not bool(data.get("rollout_created", False)) and rollout_message == "当前无法自动生成智能 Rollout":
|
||
return True, "Release 已创建,但当前无法自动生成智能 Rollout", data
|
||
if not bool(data.get("rollout_created", False)) and rollout_message == "智能 Rollout 预检未通过":
|
||
return True, "Release 已创建,但智能 Rollout 预检未通过", data
|
||
if bool(data.get("requires_confirmation", False)) and rollout_message == "智能 Rollout 需要确认":
|
||
return True, "Release 已创建,智能 Rollout 需要确认", data
|
||
return rollout_ok, rollout_message, data
|
||
|
||
|
||
def list_releases(limit: int = 20, *, channel: str | None = None) -> list[dict]:
|
||
ensure_ops_release_schema()
|
||
safe_limit = min(max(int(limit or 20), 1), 100)
|
||
conditions: list[str] = []
|
||
params: list[object] = []
|
||
if str(channel or "").strip():
|
||
conditions.append("channel = %s")
|
||
params.append(str(channel).strip())
|
||
where_clause = f"WHERE {' AND '.join(conditions)}" if conditions else ""
|
||
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
f"""
|
||
SELECT
|
||
id, release_version, channel, commit_sha, status, artifact_url, checksum, notes, metadata_json, created_by, created_at, activated_at, updated_at
|
||
FROM ops_releases
|
||
{where_clause}
|
||
ORDER BY created_at DESC, id DESC
|
||
LIMIT %s
|
||
""",
|
||
(*params, safe_limit),
|
||
)
|
||
rows = cur.fetchall()
|
||
return [_serialize_release_row(row) for row in rows]
|
||
|
||
|
||
def get_release(release_id: int) -> dict:
|
||
ensure_ops_release_schema()
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
SELECT
|
||
id, release_version, channel, commit_sha, status, artifact_url, checksum, notes, metadata_json, created_by, created_at, activated_at, updated_at
|
||
FROM ops_releases
|
||
WHERE id = %s
|
||
""",
|
||
(int(release_id),),
|
||
)
|
||
row = cur.fetchone()
|
||
return _serialize_release_row(row) if row else {}
|
||
|
||
|
||
def get_latest_release(*, channel: str = "stable") -> dict:
|
||
ensure_ops_release_schema()
|
||
normalized_channel = str(channel or "stable").strip() or "stable"
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
SELECT
|
||
id, release_version, channel, commit_sha, status, artifact_url, checksum, notes, metadata_json, created_by, created_at, activated_at, updated_at
|
||
FROM ops_releases
|
||
WHERE channel = %s
|
||
ORDER BY
|
||
CASE WHEN status = 'active' THEN 0 WHEN status = 'ready' THEN 1 ELSE 2 END,
|
||
created_at DESC,
|
||
id DESC
|
||
LIMIT 1
|
||
""",
|
||
(normalized_channel,),
|
||
)
|
||
row = cur.fetchone()
|
||
return _serialize_release_row(row) if row else {}
|
||
|
||
|
||
def _compact_package_metadata(payload: dict) -> dict:
|
||
package = dict(payload or {})
|
||
return {
|
||
"available": bool(package.get("available", False)),
|
||
"package_name": str(package.get("package_name") or ""),
|
||
"package_type": str(package.get("package_type") or ""),
|
||
"generated_at": str(package.get("generated_at") or ""),
|
||
"commit_sha": str(package.get("commit_sha") or ""),
|
||
"commit_ref": str(package.get("commit_ref") or ""),
|
||
"smoke_test_ok": bool(package.get("smoke_test_ok", False)),
|
||
"smoke_test_message": str(package.get("smoke_test_message") or ""),
|
||
"final_release_ok": bool(package.get("final_release_ok", False)),
|
||
"final_release_report_available": bool(package.get("final_release_report_available", False)),
|
||
"final_release_report_ok": package.get("final_release_report_ok"),
|
||
"final_release_report_matches_latest": package.get("final_release_report_matches_latest"),
|
||
"final_release_report_package_name": str(package.get("final_release_report_package_name") or ""),
|
||
"final_release_report_stale_reason": str(package.get("final_release_report_stale_reason") or ""),
|
||
"final_release_gate_decision": str(package.get("final_release_gate_decision") or ""),
|
||
"final_release_gate_blocked_reasons": list(package.get("final_release_gate_blocked_reasons") or []),
|
||
"final_release_gate_verify_ok": package.get("final_release_gate_verify_ok"),
|
||
"final_release_gate_smoke_test_ok": package.get("final_release_gate_smoke_test_ok"),
|
||
"download_path": str(package.get("download_path") or ""),
|
||
"sha256_download_path": str(package.get("sha256_download_path") or ""),
|
||
"reason": str(package.get("reason") or ""),
|
||
"manifest": _extract_release_manifest(package),
|
||
}
|
||
|
||
|
||
def _compact_launchpad_runtime_build_info(payload: dict) -> dict:
|
||
build_info = dict(payload or {})
|
||
route_surface = dict(build_info.get("route_surface") or {})
|
||
repository_capabilities = dict(build_info.get("repository_capabilities") or {})
|
||
route_surface_missing_keys = [
|
||
str(item).strip()
|
||
for item in list(route_surface.get("missing_keys") or [])
|
||
if str(item).strip()
|
||
]
|
||
route_surface_expected_paths = {
|
||
str(key).strip(): str(value).strip()
|
||
for key, value in dict(route_surface.get("expected_paths") or {}).items()
|
||
if str(key).strip()
|
||
}
|
||
repo_supports_install_command_block = bool(repository_capabilities.get("supports_install_command_block", False))
|
||
route_surface_declares_bootstrap_plan = "ops_node_handover_bootstrap_plan" in route_surface_expected_paths
|
||
runtime_schema_stale = bool(
|
||
repo_supports_install_command_block
|
||
and (
|
||
"ops_node_handover_bootstrap_plan" in route_surface_missing_keys
|
||
or not route_surface_declares_bootstrap_plan
|
||
)
|
||
)
|
||
return {
|
||
"source": str(build_info.get("source") or "").strip(),
|
||
"package_name": str(build_info.get("package_name") or "").strip(),
|
||
"commit_sha": str(build_info.get("commit_sha") or "").strip(),
|
||
"commit_ref": str(build_info.get("commit_ref") or "").strip(),
|
||
"generated_at": str(build_info.get("generated_at") or "").strip(),
|
||
"route_surface_complete": bool(route_surface.get("surface_complete", False)),
|
||
"route_surface_missing_keys": route_surface_missing_keys,
|
||
"route_surface_declares_bootstrap_plan": route_surface_declares_bootstrap_plan,
|
||
"repository_supports_install_command_block": repo_supports_install_command_block,
|
||
"runtime_schema_stale": runtime_schema_stale,
|
||
}
|
||
|
||
|
||
def _compact_release_payload(payload: dict) -> dict:
|
||
release = dict(payload or {})
|
||
compact = {
|
||
"id": int(release.get("id") or 0),
|
||
"release_version": str(release.get("release_version") or ""),
|
||
"channel": str(release.get("channel") or ""),
|
||
"commit_sha": str(release.get("commit_sha") or ""),
|
||
"status": str(release.get("status") or ""),
|
||
"artifact_url": str(release.get("artifact_url") or ""),
|
||
"artifact_url_present": bool(str(release.get("artifact_url") or "").strip()),
|
||
"checksum": str(release.get("checksum") or ""),
|
||
"created_by": str(release.get("created_by") or ""),
|
||
"created_at": str(release.get("created_at") or ""),
|
||
"updated_at": str(release.get("updated_at") or ""),
|
||
"is_preview": bool(release.get("is_preview", False)),
|
||
"manifest": _extract_release_manifest(release),
|
||
}
|
||
compact["status_label"] = str(release.get("status_label") or "").strip() or _release_status_label(compact.get("status"))
|
||
compact["summary"] = str(release.get("summary") or "").strip() or _build_release_summary_text(compact)
|
||
compact["summary_text"] = str(release.get("summary_text") or compact.get("summary") or "").strip()
|
||
compact["focus_ref"] = dict(release.get("focus_ref") or {}) or _build_release_focus_ref(compact, section="release_detail")
|
||
return compact
|
||
|
||
|
||
def _compact_policy_preview(preview: dict) -> dict:
|
||
data = dict(preview or {})
|
||
return {
|
||
"risk_level": str(data.get("risk_level") or ""),
|
||
"approval_required": bool(data.get("approval_required", False)),
|
||
"blocked": bool(data.get("blocked", False)),
|
||
"blocking_reasons": list(data.get("blocking_reasons") or []),
|
||
"approval_reasons": list(data.get("approval_reasons") or []),
|
||
"warnings": list(data.get("warnings") or []),
|
||
"recommendations": list(data.get("recommendations") or []),
|
||
"target_summary": dict(data.get("target_summary") or {}),
|
||
"batch_plan": dict(data.get("batch_plan") or {}),
|
||
"cluster_guardrails": dict(data.get("cluster_guardrails") or {}),
|
||
"release_gate": dict(data.get("release_gate") or {}),
|
||
}
|
||
|
||
|
||
def _compact_smart_rollout_plan(payload: dict) -> dict:
|
||
smart_rollout = dict(payload or {})
|
||
selector = dict(smart_rollout.get("target_selector") or {})
|
||
return {
|
||
"available": bool(smart_rollout.get("available", False)),
|
||
"mode": str(smart_rollout.get("mode") or ""),
|
||
"mode_label": str(smart_rollout.get("mode_label") or ""),
|
||
"execution_mode": str(smart_rollout.get("execution_mode") or ""),
|
||
"execution_mode_label": str(smart_rollout.get("execution_mode_label") or ""),
|
||
"reason": str(smart_rollout.get("reason") or ""),
|
||
"summary": str(smart_rollout.get("summary") or smart_rollout.get("reason") or ""),
|
||
"dominant_region": str(smart_rollout.get("dominant_region") or ""),
|
||
"target_node_codes": list(selector.get("node_codes") or []),
|
||
"target_nodes": [
|
||
{
|
||
"node_code": str(item.get("node_code") or ""),
|
||
"region": str(item.get("region") or ""),
|
||
"role": str(item.get("role") or ""),
|
||
"status": str(item.get("status") or ""),
|
||
"current_load": int(item.get("current_load", 0) or 0),
|
||
"remote_agent_ready": bool(item.get("remote_agent_ready", False)),
|
||
"ssh_ready": bool(item.get("ssh_ready", False)),
|
||
"execution_ready": bool(item.get("execution_ready", False)),
|
||
"detect_participating": bool(item.get("detect_participating", False)),
|
||
"delivery_queue_state": str(item.get("delivery_queue_state") or ""),
|
||
}
|
||
for item in list(smart_rollout.get("target_nodes") or [])
|
||
if str(item.get("node_code") or "").strip()
|
||
],
|
||
"policy": dict(smart_rollout.get("policy") or {}),
|
||
}
|
||
|
||
|
||
def _compact_latest_package_preview(message: str, data: dict) -> dict:
|
||
payload = dict(data or {})
|
||
return {
|
||
"message": str(message or ""),
|
||
"preview_source": str(payload.get("preview_source") or ""),
|
||
"deduplicated": bool(payload.get("release_deduplicated", payload.get("deduplicated", False))),
|
||
"release": _compact_release_payload(payload.get("release") or {}),
|
||
"package": _compact_package_metadata(payload.get("package") or {}),
|
||
"resolved_artifact_url": str(payload.get("resolved_artifact_url") or ""),
|
||
"resolved_sha256_download_url": str(payload.get("resolved_sha256_download_url") or ""),
|
||
}
|
||
|
||
|
||
def _compact_latest_smart_rollout_preview(message: str, data: dict) -> dict:
|
||
payload = dict(data or {})
|
||
preview = _compact_policy_preview(payload.get("preview") or {})
|
||
smart_rollout = _compact_smart_rollout_plan(payload.get("smart_rollout") or {})
|
||
release = _compact_release_payload(payload.get("release") or {})
|
||
requires_confirmation = bool(payload.get("requires_confirmation", False))
|
||
blocked = bool(payload.get("blocked", False) or preview.get("blocked", False))
|
||
preview_batch_plan = dict(preview.get("batch_plan") or {})
|
||
preview_target_summary = dict(preview.get("target_summary") or {})
|
||
rollout_preview = {
|
||
"release": release,
|
||
"execution_mode": str(smart_rollout.get("execution_mode") or ""),
|
||
"execution_mode_label": str(smart_rollout.get("execution_mode_label") or ""),
|
||
"target_nodes_total": int(preview_target_summary.get("nodes_total", len(list(smart_rollout.get("target_nodes") or []))) or 0),
|
||
"target_node_codes": list(smart_rollout.get("target_node_codes") or []),
|
||
"expected_batches_total": int(preview_batch_plan.get("batches_total", 0) or 0),
|
||
"expected_jobs_total": int(preview_target_summary.get("nodes_total", len(list(smart_rollout.get("target_nodes") or []))) or 0),
|
||
"gate": dict(preview.get("release_gate") or {}),
|
||
"batch_plan": dict(preview_batch_plan),
|
||
"summary": str(
|
||
(dict(preview.get("release_gate") or {}).get("summary") or "")
|
||
or smart_rollout.get("summary")
|
||
or message
|
||
or ""
|
||
).strip(),
|
||
"focus_ref": _build_release_focus_ref(
|
||
release,
|
||
section="rollout_preview",
|
||
),
|
||
}
|
||
rollout_preview["summary_text"] = str(rollout_preview.get("summary") or "")
|
||
return {
|
||
"message": str(message or ""),
|
||
"preview_source": str(payload.get("preview_source") or ""),
|
||
"release": release,
|
||
"package": _compact_package_metadata(payload.get("package") or {}),
|
||
"smart_rollout": smart_rollout,
|
||
"policy_preview": preview,
|
||
"rollout_preview": rollout_preview,
|
||
"requires_confirmation": requires_confirmation,
|
||
"blocked": blocked,
|
||
"rollout_created": bool(payload.get("rollout_created", False)),
|
||
"release_deduplicated": bool(payload.get("release_deduplicated", payload.get("deduplicated", False))),
|
||
}
|
||
|
||
|
||
def _launchpad_preview_execution_context(preview: dict) -> tuple[str, str]:
|
||
smart_rollout = dict((preview or {}).get("smart_rollout") or {})
|
||
execution_mode = str(smart_rollout.get("execution_mode") or "").strip()
|
||
execution_mode_label = str(smart_rollout.get("execution_mode_label") or "").strip()
|
||
if execution_mode or execution_mode_label:
|
||
return execution_mode, execution_mode_label or _rollout_execution_mode_label(execution_mode)
|
||
release_gate = dict((dict(preview or {}).get("policy_preview") or {}).get("release_gate") or {})
|
||
execution_mode = str(release_gate.get("execution_mode") or "").strip()
|
||
execution_mode_label = str(release_gate.get("execution_mode_label") or "").strip()
|
||
if execution_mode or execution_mode_label:
|
||
return execution_mode, execution_mode_label or _rollout_execution_mode_label(execution_mode)
|
||
return "", ""
|
||
|
||
|
||
def _decorate_release_launchpad_status(
|
||
status: dict,
|
||
*,
|
||
gap_row: dict,
|
||
onboarding_bootstrap_pending_nodes: int,
|
||
onboarding_acceptance_ready_nodes: int,
|
||
) -> dict:
|
||
normalized_status = dict(status or {})
|
||
normalized_gap_row = dict(gap_row or {})
|
||
normalized_status["onboarding_bootstrap_pending_nodes"] = int(onboarding_bootstrap_pending_nodes or 0)
|
||
normalized_status["onboarding_acceptance_ready_nodes"] = int(onboarding_acceptance_ready_nodes or 0)
|
||
normalized_status["onboarding_gap_nodes_total"] = int(
|
||
int(onboarding_bootstrap_pending_nodes or 0) + int(onboarding_acceptance_ready_nodes or 0)
|
||
)
|
||
normalized_status["recommended_target_node_code"] = str(
|
||
normalized_status.get("recommended_target_node_code") or normalized_gap_row.get("node_code") or ""
|
||
).strip()
|
||
normalized_status["recommended_recovery_label"] = str(
|
||
normalized_status.get("recommended_recovery_label") or normalized_gap_row.get("recovery_label") or ""
|
||
).strip()
|
||
normalized_status["recommended_recovery_summary"] = str(
|
||
normalized_status.get("recommended_recovery_summary")
|
||
or normalized_gap_row.get("recovery_summary")
|
||
or normalized_gap_row.get("onboarding_summary")
|
||
or ""
|
||
).strip()
|
||
normalized_status["recommended_onboarding_stage_code"] = str(
|
||
normalized_gap_row.get("onboarding_stage_code") or ""
|
||
).strip()
|
||
normalized_status["recommended_onboarding_stage_label"] = str(
|
||
normalized_gap_row.get("onboarding_stage_label") or ""
|
||
).strip()
|
||
return normalized_status
|
||
|
||
|
||
def _build_release_launchpad_status(
|
||
*,
|
||
package: dict,
|
||
latest_package_preview: dict,
|
||
worker_preview: dict,
|
||
control_preview: dict,
|
||
runtime_build_info: dict | None = None,
|
||
) -> dict:
|
||
if not bool(package.get("available", False)):
|
||
status = {
|
||
"status": "blocked",
|
||
"status_label": "缺少发布包",
|
||
"summary": "当前没有可用的最新发布包,请先在控制面本机完成打包。",
|
||
"recommended_action_code": "release_package",
|
||
"recommended_command": build_bash_command("drive_ops_center.sh", "release-package"),
|
||
}
|
||
status["summary_text"] = str(status.get("summary") or "")
|
||
status["focus_ref"] = _build_release_focus_ref({}, section="release_launchpad")
|
||
return status
|
||
|
||
package_final_gate_ok = package.get("final_release_ok")
|
||
if package_final_gate_ok is False:
|
||
stale_reason = str(package.get("final_release_report_stale_reason") or "").strip()
|
||
blocked_reasons = [
|
||
str(item).strip()
|
||
for item in list(package.get("final_release_gate_blocked_reasons") or [])
|
||
if str(item).strip()
|
||
]
|
||
gate_decision = str(package.get("final_release_gate_decision") or "").strip()
|
||
summary = "最新发布包尚未完成最终签收,当前应先完成最终验包与签收,再进入 Release / Rollout。"
|
||
if stale_reason:
|
||
summary = f"{summary} {stale_reason}"
|
||
elif blocked_reasons:
|
||
summary = f"{summary} 当前阻断:{' / '.join(blocked_reasons)}"
|
||
elif gate_decision:
|
||
summary = f"{summary} 当前门禁状态:{gate_decision}"
|
||
status = {
|
||
"status": "blocked",
|
||
"status_label": "待完成最终签收",
|
||
"summary": summary,
|
||
"recommended_action_code": "release_prepare",
|
||
"recommended_command": build_bash_command("drive_ops_center.sh", "release-prepare"),
|
||
"recommended_execution_mode": "",
|
||
"recommended_execution_mode_label": "",
|
||
"final_release_gate_ready": False,
|
||
"final_release_gate_decision": gate_decision,
|
||
"final_release_gate_blocked_reasons": blocked_reasons,
|
||
}
|
||
status["summary_text"] = str(status.get("summary") or "")
|
||
status["focus_ref"] = _build_release_focus_ref({}, section="release_launchpad")
|
||
return status
|
||
|
||
runtime_refresh_state = dict(runtime_build_info or {})
|
||
runtime_schema_stale = bool(runtime_refresh_state.get("runtime_schema_stale", False))
|
||
route_surface_complete = bool(runtime_refresh_state.get("route_surface_complete", True))
|
||
route_surface_missing_keys = [
|
||
str(item).strip()
|
||
for item in list(runtime_refresh_state.get("route_surface_missing_keys") or [])
|
||
if str(item).strip()
|
||
]
|
||
if runtime_schema_stale or not route_surface_complete:
|
||
summary = "运行中的控制面 API 还没有刷新到当前仓库的最新运维能力,当前 Release Launchpad 判断可能偏旧,建议先完成运行时刷新与复检。"
|
||
if runtime_schema_stale:
|
||
summary = "仓库已经具备新的节点接管与发布能力,但运行中的控制面 API 路由面仍旧,建议先执行运行时刷新,再回到 Launchpad 继续判断。"
|
||
elif route_surface_missing_keys:
|
||
summary = f"{summary} 当前缺少路由键:{', '.join(route_surface_missing_keys)}。"
|
||
status = {
|
||
"status": "attention",
|
||
"status_label": "待刷新运行时",
|
||
"summary": summary,
|
||
"recommended_action_code": "api-restart",
|
||
"recommended_command": build_bash_command("drive_ops_center.sh", "runtime-refresh-recover"),
|
||
"recommended_execution_mode": "",
|
||
"recommended_execution_mode_label": "",
|
||
"runtime_refresh_required": True,
|
||
"runtime_schema_stale": runtime_schema_stale,
|
||
"runtime_route_surface_complete": route_surface_complete,
|
||
"runtime_route_surface_missing_keys": route_surface_missing_keys,
|
||
}
|
||
status["summary_text"] = str(status.get("summary") or "")
|
||
status["focus_ref"] = _build_release_focus_ref({}, section="release_launchpad")
|
||
return status
|
||
|
||
worker_plan = dict(worker_preview.get("smart_rollout") or {})
|
||
worker_policy = dict(worker_preview.get("policy_preview") or {})
|
||
control_plan = dict(control_preview.get("smart_rollout") or {})
|
||
control_policy = dict(control_preview.get("policy_preview") or {})
|
||
worker_execution_mode, worker_execution_mode_label = _launchpad_preview_execution_context(worker_preview)
|
||
control_execution_mode, control_execution_mode_label = _launchpad_preview_execution_context(control_preview)
|
||
worker_release_gate = dict(worker_policy.get("release_gate") or {})
|
||
control_release_gate = dict(control_policy.get("release_gate") or {})
|
||
candidate_rows = [
|
||
dict(item or {})
|
||
for item in (
|
||
list(((worker_release_gate.get("operational_readiness") or {}).get("rows") or []))
|
||
+ list(((control_release_gate.get("operational_readiness") or {}).get("rows") or []))
|
||
)
|
||
if dict(item or {})
|
||
]
|
||
gap_row = next(
|
||
(
|
||
row
|
||
for row in candidate_rows
|
||
if str(row.get("recovery_action") or "").strip() in {"bootstrap_run", "run_acceptance"}
|
||
),
|
||
{},
|
||
)
|
||
gap_recovery_action = str(gap_row.get("recovery_action") or "").strip()
|
||
gap_recovery_label = str(gap_row.get("recovery_label") or "").strip()
|
||
gap_recovery_summary = str(
|
||
gap_row.get("recovery_summary")
|
||
or gap_row.get("onboarding_summary")
|
||
or gap_row.get("agent_reason")
|
||
or ""
|
||
).strip()
|
||
gap_node_code = str(gap_row.get("node_code") or "").strip()
|
||
onboarding_bootstrap_pending_nodes = sum(
|
||
1 for row in candidate_rows if str(row.get("recovery_action") or "").strip() == "bootstrap_run"
|
||
)
|
||
onboarding_acceptance_ready_nodes = sum(
|
||
1 for row in candidate_rows if str(row.get("recovery_action") or "").strip() == "run_acceptance"
|
||
)
|
||
|
||
if bool(worker_policy.get("blocked", False)) or bool(control_policy.get("blocked", False)):
|
||
recommended_action_code = "fix_rollout_blockers"
|
||
recommended_command = build_bash_command("drive_ops_center.sh", "release-preview-smart", "worker")
|
||
summary = "最新发布包已就绪,但当前 Rollout 预检存在阻断条件,需先修复节点接管、在线状态或巡检缺口。"
|
||
if gap_recovery_action in {"bootstrap_run", "run_acceptance"} and gap_node_code:
|
||
recommended_action_code = gap_recovery_action
|
||
if gap_recovery_action == "bootstrap_run":
|
||
recommended_command = build_bash_command("drive_ops_center.sh", "node-bootstrap-run", "http://127.0.0.1:8100", gap_node_code, "cli")
|
||
else:
|
||
recommended_command = build_bash_command("drive_ops_center.sh", "node-acceptance-run", "http://127.0.0.1:8100", gap_node_code, "cli")
|
||
summary = gap_recovery_summary or summary
|
||
status = {
|
||
"status": "blocked",
|
||
"status_label": "发布门禁阻断",
|
||
"summary": summary,
|
||
"recommended_action_code": recommended_action_code,
|
||
"recommended_command": recommended_command,
|
||
"recommended_execution_mode": worker_execution_mode or control_execution_mode,
|
||
"recommended_execution_mode_label": worker_execution_mode_label or control_execution_mode_label,
|
||
"recommended_target_node_code": gap_node_code,
|
||
"recommended_recovery_label": gap_recovery_label,
|
||
}
|
||
status["summary_text"] = str(status.get("summary") or "")
|
||
status["focus_ref"] = _build_release_focus_ref({}, section="release_launchpad")
|
||
return _decorate_release_launchpad_status(
|
||
status,
|
||
gap_row=gap_row,
|
||
onboarding_bootstrap_pending_nodes=onboarding_bootstrap_pending_nodes,
|
||
onboarding_acceptance_ready_nodes=onboarding_acceptance_ready_nodes,
|
||
)
|
||
|
||
if bool(worker_preview.get("requires_confirmation", False)) or bool(control_preview.get("requires_confirmation", False)):
|
||
status = {
|
||
"status": "attention",
|
||
"status_label": "待人工确认",
|
||
"summary": "最新发布包已经满足预发条件,但智能 Rollout 仍有告警或审批要求,建议先人工确认后再推进。",
|
||
"recommended_action_code": "review_smart_rollout_preview",
|
||
"recommended_command": build_bash_command("drive_ops_center.sh", "release-preview-smart", "worker"),
|
||
"recommended_execution_mode": worker_execution_mode or control_execution_mode,
|
||
"recommended_execution_mode_label": worker_execution_mode_label or control_execution_mode_label,
|
||
}
|
||
status["summary_text"] = str(status.get("summary") or "")
|
||
status["focus_ref"] = _build_release_focus_ref({}, section="release_launchpad")
|
||
return _decorate_release_launchpad_status(
|
||
status,
|
||
gap_row=gap_row,
|
||
onboarding_bootstrap_pending_nodes=onboarding_bootstrap_pending_nodes,
|
||
onboarding_acceptance_ready_nodes=onboarding_acceptance_ready_nodes,
|
||
)
|
||
|
||
if bool(worker_plan.get("available", False)):
|
||
status = {
|
||
"status": "ready",
|
||
"status_label": "可发 Worker 灰度",
|
||
"summary": "最新发布包、Release 记录与 Worker 智能 Rollout 预检均已收口,当前可以直接进入 Worker 灰度发布。",
|
||
"recommended_action_code": "publish_latest_worker",
|
||
"recommended_command": build_bash_command("drive_ops_center.sh", "publish-latest", "worker"),
|
||
"recommended_execution_mode": worker_execution_mode,
|
||
"recommended_execution_mode_label": worker_execution_mode_label,
|
||
}
|
||
status["summary_text"] = str(status.get("summary") or "")
|
||
status["focus_ref"] = _build_release_focus_ref({}, section="release_launchpad")
|
||
return _decorate_release_launchpad_status(
|
||
status,
|
||
gap_row=gap_row,
|
||
onboarding_bootstrap_pending_nodes=onboarding_bootstrap_pending_nodes,
|
||
onboarding_acceptance_ready_nodes=onboarding_acceptance_ready_nodes,
|
||
)
|
||
|
||
if bool(control_plan.get("available", False)):
|
||
status = {
|
||
"status": "attention",
|
||
"status_label": "仅 Control 可发",
|
||
"summary": "当前 Control 侧满足预发条件,但 Worker 侧还未形成可用灰度面,建议先收敛 Worker 节点状态。",
|
||
"recommended_action_code": "review_control_rollout",
|
||
"recommended_command": build_bash_command("drive_ops_center.sh", "release-preview-smart", "control"),
|
||
"recommended_execution_mode": control_execution_mode,
|
||
"recommended_execution_mode_label": control_execution_mode_label,
|
||
}
|
||
status["summary_text"] = str(status.get("summary") or "")
|
||
status["focus_ref"] = _build_release_focus_ref({}, section="release_launchpad")
|
||
return _decorate_release_launchpad_status(
|
||
status,
|
||
gap_row=gap_row,
|
||
onboarding_bootstrap_pending_nodes=onboarding_bootstrap_pending_nodes,
|
||
onboarding_acceptance_ready_nodes=onboarding_acceptance_ready_nodes,
|
||
)
|
||
|
||
status = {
|
||
"status": "attention",
|
||
"status_label": "待补执行面",
|
||
"summary": "最新发布包已经可用,但当前没有可用于智能 Rollout 的在线启用执行节点,需要先补齐托管或接入状态。",
|
||
"recommended_action_code": "fix_managed_nodes",
|
||
"recommended_command": build_bash_command("drive_ops_center.sh", "doctor"),
|
||
"recommended_execution_mode": worker_execution_mode or control_execution_mode,
|
||
"recommended_execution_mode_label": worker_execution_mode_label or control_execution_mode_label,
|
||
}
|
||
if gap_recovery_action in {"bootstrap_run", "run_acceptance"} and gap_node_code:
|
||
status["recommended_action_code"] = gap_recovery_action
|
||
status["recommended_command"] = (
|
||
build_bash_command("drive_ops_center.sh", "node-bootstrap-run", "http://127.0.0.1:8100", gap_node_code, "cli")
|
||
if gap_recovery_action == "bootstrap_run"
|
||
else build_bash_command("drive_ops_center.sh", "node-acceptance-run", "http://127.0.0.1:8100", gap_node_code, "cli")
|
||
)
|
||
status["summary"] = gap_recovery_summary or status["summary"]
|
||
status["recommended_target_node_code"] = gap_node_code
|
||
status["recommended_recovery_label"] = gap_recovery_label
|
||
status["summary_text"] = str(status.get("summary") or "")
|
||
status["focus_ref"] = _build_release_focus_ref({}, section="release_launchpad")
|
||
return _decorate_release_launchpad_status(
|
||
status,
|
||
gap_row=gap_row,
|
||
onboarding_bootstrap_pending_nodes=onboarding_bootstrap_pending_nodes,
|
||
onboarding_acceptance_ready_nodes=onboarding_acceptance_ready_nodes,
|
||
)
|
||
|
||
|
||
def get_release_launchpad(*, control_plane_base_url: str = "", channel: str = "stable") -> dict:
|
||
package = get_latest_release_package_metadata()
|
||
release_summary = get_release_summary()
|
||
latest_release = get_latest_release(channel=channel)
|
||
runtime_build_info = _compact_launchpad_runtime_build_info(get_runtime_build_info())
|
||
|
||
preview_ok, preview_message, preview_data = preview_release_from_latest_package(
|
||
{"channel": channel, "created_by": "api/release-launchpad"},
|
||
control_plane_base_url=control_plane_base_url,
|
||
)
|
||
worker_ok, worker_message, worker_data = preview_smart_release_rollout_from_latest_package(
|
||
{
|
||
"channel": channel,
|
||
"mode": "worker",
|
||
"created_by": "api/release-launchpad",
|
||
"rollout_created_by": "api/release-launchpad/worker",
|
||
},
|
||
control_plane_base_url=control_plane_base_url,
|
||
)
|
||
control_ok, control_message, control_data = preview_smart_release_rollout_from_latest_package(
|
||
{
|
||
"channel": channel,
|
||
"mode": "control",
|
||
"created_by": "api/release-launchpad",
|
||
"rollout_created_by": "api/release-launchpad/control",
|
||
},
|
||
control_plane_base_url=control_plane_base_url,
|
||
)
|
||
|
||
latest_package_preview = _compact_latest_package_preview(preview_message, preview_data if preview_ok else {"package": package})
|
||
worker_preview = _compact_latest_smart_rollout_preview(worker_message, worker_data if worker_ok else {"package": package})
|
||
control_preview = _compact_latest_smart_rollout_preview(control_message, control_data if control_ok else {"package": package})
|
||
compact_latest_release = _compact_release_payload(latest_release)
|
||
launchpad_status = _build_release_launchpad_status(
|
||
package=_compact_package_metadata(package),
|
||
latest_package_preview=latest_package_preview,
|
||
worker_preview=worker_preview,
|
||
control_preview=control_preview,
|
||
runtime_build_info=runtime_build_info,
|
||
)
|
||
launchpad_status["focus_ref"] = _build_release_focus_ref(
|
||
compact_latest_release,
|
||
section="release_launchpad",
|
||
)
|
||
|
||
return {
|
||
"channel": channel,
|
||
"control_plane_base_url": _normalize_public_base_url(control_plane_base_url),
|
||
"latest_package": _compact_package_metadata(package),
|
||
"latest_release": compact_latest_release,
|
||
"release_summary": release_summary,
|
||
"runtime_build_info": runtime_build_info,
|
||
"latest_package_preview": latest_package_preview,
|
||
"worker_rollout_preview": worker_preview,
|
||
"control_rollout_preview": control_preview,
|
||
"launchpad_status": launchpad_status,
|
||
}
|
||
|
||
|
||
def activate_release(release_id: int) -> tuple[bool, str, dict]:
|
||
ensure_ops_release_schema()
|
||
release = get_release(release_id)
|
||
if not release:
|
||
return False, "Release 不存在", {}
|
||
channel = str(release.get("channel") or "stable")
|
||
|
||
with get_db() as conn:
|
||
conn.autocommit = False
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
UPDATE ops_releases
|
||
SET
|
||
status = CASE WHEN id = %s THEN 'active' ELSE status END,
|
||
activated_at = CASE WHEN id = %s THEN CURRENT_TIMESTAMP ELSE activated_at END,
|
||
updated_at = CURRENT_TIMESTAMP
|
||
WHERE channel = %s AND id = %s
|
||
""",
|
||
(int(release_id), int(release_id), channel, int(release_id)),
|
||
)
|
||
cur.execute(
|
||
"""
|
||
UPDATE ops_releases
|
||
SET
|
||
status = CASE WHEN status = 'active' THEN 'ready' ELSE status END,
|
||
updated_at = CURRENT_TIMESTAMP
|
||
WHERE channel = %s AND id <> %s
|
||
""",
|
||
(channel, int(release_id)),
|
||
)
|
||
conn.commit()
|
||
|
||
return True, "Release 已激活", {"release": get_release(int(release_id))}
|
||
|
||
|
||
def get_release_summary() -> dict:
|
||
ensure_ops_release_schema()
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
SELECT status, count(*)
|
||
FROM ops_releases
|
||
GROUP BY status
|
||
"""
|
||
)
|
||
release_rows = cur.fetchall()
|
||
cur.execute(
|
||
"""
|
||
SELECT channel, release_version
|
||
FROM ops_releases
|
||
WHERE status = 'active'
|
||
ORDER BY channel ASC, activated_at DESC NULLS LAST, id DESC
|
||
"""
|
||
)
|
||
active_rows = cur.fetchall()
|
||
cur.execute(
|
||
"""
|
||
SELECT status, count(*)
|
||
FROM ops_release_rollouts
|
||
GROUP BY status
|
||
"""
|
||
)
|
||
rollout_rows = cur.fetchall()
|
||
status_counts = {str(row[0] or ""): int(row[1] or 0) for row in release_rows}
|
||
active_by_channel = {str(row[0] or ""): str(row[1] or "") for row in active_rows}
|
||
rollout_status_counts = {str(row[0] or ""): int(row[1] or 0) for row in rollout_rows}
|
||
return {
|
||
"status_counts": status_counts,
|
||
"total": sum(status_counts.values()),
|
||
"active_by_channel": active_by_channel,
|
||
"rollout_status_counts": rollout_status_counts,
|
||
"rollouts_total": sum(rollout_status_counts.values()),
|
||
}
|
||
|
||
|
||
def _get_rollout_row(rollout_id: int) -> tuple | None:
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
SELECT
|
||
id, release_id, rollout_code, target_selector_json, target_nodes_json, policy_json, status,
|
||
batch_cursor, batches_total, jobs_total, jobs_created, result_summary_json, created_by, created_at, updated_at
|
||
FROM ops_release_rollouts
|
||
WHERE id = %s
|
||
""",
|
||
(int(rollout_id),),
|
||
)
|
||
return cur.fetchone()
|
||
|
||
|
||
def _list_rollout_jobs(rollout_id: int, *, total_targets: int) -> list[dict]:
|
||
from app.services.ops_job_service import list_ops_jobs
|
||
|
||
safe_limit = max(50, min(max(total_targets * 2, total_targets + 10), 10000))
|
||
return list_ops_jobs(limit=safe_limit, rollout_id=int(rollout_id))
|
||
|
||
|
||
def _compute_rollout_summary(rollout: dict, jobs: list[dict]) -> dict:
|
||
target_nodes = list(rollout.get("target_nodes") or [])
|
||
target_node_codes = [str(item.get("node_code") or "").strip() for item in target_nodes if str(item.get("node_code") or "").strip()]
|
||
launched_node_codes = []
|
||
status_counts: dict[str, int] = {}
|
||
waiting_approval_targets = 0
|
||
|
||
for job in jobs:
|
||
status = str(job.get("status") or "")
|
||
target_node_code = str(job.get("target_node_code") or "").strip()
|
||
status_counts[status] = int(status_counts.get(status, 0) or 0) + 1
|
||
if target_node_code and target_node_code not in launched_node_codes:
|
||
launched_node_codes.append(target_node_code)
|
||
if status == "awaiting_approval":
|
||
waiting_approval_targets += 1
|
||
|
||
launched_target_codes = set(launched_node_codes)
|
||
remaining_targets = [item for item in target_nodes if str(item.get("node_code") or "").strip() not in launched_target_codes]
|
||
remaining_target_codes = [str(item.get("node_code") or "").strip() for item in remaining_targets]
|
||
|
||
active_jobs = sum(int(status_counts.get(status, 0) or 0) for status in _ACTIVE_JOB_STATUSES)
|
||
issue_jobs = sum(int(status_counts.get(status, 0) or 0) for status in _ISSUE_JOB_STATUSES)
|
||
success_jobs = int(status_counts.get("success", 0) or 0)
|
||
failed_jobs = int(status_counts.get("failed", 0) or 0)
|
||
cancelled_jobs = int(status_counts.get("cancelled", 0) or 0)
|
||
blocked_jobs = int(status_counts.get("blocked", 0) or 0)
|
||
partial_jobs = int(status_counts.get("partially_succeeded", 0) or 0)
|
||
|
||
policy = dict(rollout.get("policy") or {})
|
||
continue_on_failure = bool(policy.get("continue_on_failure", False))
|
||
continue_on_partial = bool(policy.get("continue_on_partial", False))
|
||
|
||
halted_by_failures = (failed_jobs > 0 or cancelled_jobs > 0 or blocked_jobs > 0) and not continue_on_failure
|
||
halted_by_partial = partial_jobs > 0 and not continue_on_partial
|
||
halted = halted_by_failures or halted_by_partial
|
||
next_batch_ready = active_jobs == 0 and bool(remaining_target_codes) and not halted
|
||
|
||
if active_jobs > 0:
|
||
status = "awaiting_approval" if active_jobs == waiting_approval_targets else "running"
|
||
elif not jobs and remaining_target_codes:
|
||
status = "planned"
|
||
elif next_batch_ready:
|
||
status = "ready_for_next_batch"
|
||
elif remaining_target_codes and halted:
|
||
status = "halted"
|
||
elif not remaining_target_codes:
|
||
if issue_jobs <= 0 and success_jobs > 0:
|
||
status = "completed"
|
||
elif success_jobs <= 0 and issue_jobs > 0:
|
||
status = "failed"
|
||
else:
|
||
status = "completed_with_issues"
|
||
else:
|
||
status = str(rollout.get("status") or "planned")
|
||
|
||
launched_count = len(launched_target_codes)
|
||
total_targets = len(target_node_codes)
|
||
batches_total = max(int(rollout.get("batches_total") or 0), _calculate_total_batches(total_targets, policy))
|
||
batch_size = max(1, int(policy.get("batch_size") or 1)) if total_targets > 0 else 0
|
||
current_batch_index = int(ceil(launched_count / batch_size)) if launched_count > 0 and batch_size > 0 else 0
|
||
|
||
return {
|
||
"status": status,
|
||
"status_counts": status_counts,
|
||
"jobs_created": int(len(jobs)),
|
||
"active_jobs": active_jobs,
|
||
"issue_jobs": issue_jobs,
|
||
"success_jobs": success_jobs,
|
||
"failed_jobs": failed_jobs,
|
||
"cancelled_jobs": cancelled_jobs,
|
||
"blocked_jobs": blocked_jobs,
|
||
"partial_jobs": partial_jobs,
|
||
"launched_targets": launched_count,
|
||
"remaining_targets": len(remaining_target_codes),
|
||
"remaining_target_codes": remaining_target_codes,
|
||
"launched_target_codes": list(launched_target_codes),
|
||
"next_batch_ready": next_batch_ready,
|
||
"halted": halted,
|
||
"current_batch_index": current_batch_index,
|
||
"batches_total": batches_total,
|
||
}
|
||
|
||
|
||
def refresh_release_rollout(rollout_id: int, *, advance_if_possible: bool = True) -> dict:
|
||
ensure_ops_release_schema()
|
||
row = _get_rollout_row(int(rollout_id))
|
||
if not row:
|
||
return {}
|
||
rollout = _serialize_rollout_row(row)
|
||
jobs = _list_rollout_jobs(int(rollout_id), total_targets=len(list(rollout.get("target_nodes") or [])))
|
||
summary = _compute_rollout_summary(rollout, jobs)
|
||
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
UPDATE ops_release_rollouts
|
||
SET
|
||
status = %s,
|
||
batch_cursor = %s,
|
||
batches_total = %s,
|
||
jobs_total = %s,
|
||
jobs_created = %s,
|
||
result_summary_json = %s,
|
||
updated_at = CURRENT_TIMESTAMP
|
||
WHERE id = %s
|
||
""",
|
||
(
|
||
str(summary.get("status") or "planned"),
|
||
int(summary.get("launched_targets") or 0),
|
||
int(summary.get("batches_total") or 0),
|
||
len(list(rollout.get("target_nodes") or [])),
|
||
int(summary.get("jobs_created") or 0),
|
||
json.dumps(summary, ensure_ascii=False),
|
||
int(rollout_id),
|
||
),
|
||
)
|
||
conn.commit()
|
||
|
||
refreshed_row = _get_rollout_row(int(rollout_id))
|
||
refreshed = _serialize_rollout_row(refreshed_row) if refreshed_row else {}
|
||
refreshed["jobs"] = jobs
|
||
if advance_if_possible and bool((refreshed.get("policy") or {}).get("auto_advance", False)) and bool((refreshed.get("result_summary") or {}).get("next_batch_ready", False)):
|
||
ok, _message, data = _enqueue_rollout_batch(int(rollout_id), created_by=str(refreshed.get("created_by") or "api"), reason="auto_advance")
|
||
if ok and isinstance(data, dict) and data.get("rollout"):
|
||
return data.get("rollout") or refreshed
|
||
return refreshed
|
||
|
||
|
||
def refresh_release_rollout_for_job(job_id: int) -> dict:
|
||
from app.services.ops_job_service import get_ops_job
|
||
|
||
job = get_ops_job(int(job_id))
|
||
rollout_id = int(job.get("rollout_id") or 0)
|
||
if rollout_id <= 0:
|
||
return {}
|
||
return refresh_release_rollout(rollout_id)
|
||
|
||
|
||
def _default_release_deploy_payload_for_target(target: dict) -> dict:
|
||
role = str((target or {}).get("role") or "").strip().lower()
|
||
if role == "control":
|
||
return {
|
||
"restart_services": ["domaincheck-api", "domaincheck-worker", "domaincheck-sync-agent"],
|
||
"health_check_urls": ["http://127.0.0.1:8100/health"],
|
||
"health_check_services": ["domaincheck-api", "domaincheck-worker", "domaincheck-sync-agent"],
|
||
"health_check_timeout_seconds": 20,
|
||
"health_check_retries": 9,
|
||
"health_check_interval_seconds": 4,
|
||
}
|
||
return {
|
||
"restart_services": ["domaincheck-worker"],
|
||
"health_check_urls": [],
|
||
"health_check_services": ["domaincheck-worker"],
|
||
"health_check_timeout_seconds": 10,
|
||
"health_check_retries": 2,
|
||
"health_check_interval_seconds": 2,
|
||
}
|
||
|
||
|
||
def _build_release_job_payload(release: dict, rollout: dict, *, target: dict | None = None) -> dict:
|
||
policy = dict(rollout.get("policy") or {})
|
||
deploy_payload = dict(policy.get("deploy_payload") or {})
|
||
target_defaults = _default_release_deploy_payload_for_target(target or {})
|
||
|
||
if not [str(item).strip() for item in list(deploy_payload.get("restart_services") or []) if str(item).strip()]:
|
||
deploy_payload["restart_services"] = list(target_defaults.get("restart_services") or [])
|
||
if not [str(item).strip() for item in list(deploy_payload.get("health_check_services") or []) if str(item).strip()]:
|
||
deploy_payload["health_check_services"] = list(target_defaults.get("health_check_services") or [])
|
||
if not [str(item).strip() for item in list(deploy_payload.get("health_check_urls") or []) if str(item).strip()]:
|
||
deploy_payload["health_check_urls"] = list(target_defaults.get("health_check_urls") or [])
|
||
if deploy_payload.get("health_check_timeout_seconds") in (None, "", 0, "0"):
|
||
deploy_payload["health_check_timeout_seconds"] = int(target_defaults.get("health_check_timeout_seconds") or 10)
|
||
if deploy_payload.get("health_check_retries") in (None, "", 0, "0"):
|
||
deploy_payload["health_check_retries"] = int(target_defaults.get("health_check_retries") or 2)
|
||
if deploy_payload.get("health_check_interval_seconds") in (None, "", 0, "0"):
|
||
deploy_payload["health_check_interval_seconds"] = int(target_defaults.get("health_check_interval_seconds") or 2)
|
||
|
||
artifact_url = _rewrite_loopback_control_plane_url(str(release.get("artifact_url") or "").strip())
|
||
return {
|
||
"release_id": int(release.get("id") or 0),
|
||
"rollout_id": int(rollout.get("id") or 0),
|
||
"release_version": str(release.get("release_version") or ""),
|
||
"artifact_url": artifact_url,
|
||
"checksum": str(release.get("checksum") or ""),
|
||
"channel": str(release.get("channel") or ""),
|
||
"commit_sha": str(release.get("commit_sha") or ""),
|
||
"notes": str(release.get("notes") or ""),
|
||
"target_node_role": str((target or {}).get("role") or "").strip(),
|
||
**deploy_payload,
|
||
}
|
||
|
||
|
||
def _enqueue_rollout_batch(rollout_id: int, *, created_by: str, reason: str = "manual_advance") -> tuple[bool, str, dict]:
|
||
ensure_ops_release_schema()
|
||
rollout = refresh_release_rollout(int(rollout_id), advance_if_possible=False)
|
||
if not rollout:
|
||
return False, "Rollout 不存在", {}
|
||
|
||
release = get_release(int(rollout.get("release_id") or 0))
|
||
if not release:
|
||
return False, "Release 不存在", {"rollout": rollout}
|
||
|
||
summary = dict(rollout.get("result_summary") or {})
|
||
if int(summary.get("active_jobs") or 0) > 0:
|
||
return False, "当前已有批次正在执行,暂不可推进下一批", {"rollout": rollout}
|
||
if bool(summary.get("halted", False)):
|
||
return False, "当前 rollout 已暂停,请先处理失败/阻断节点", {"rollout": rollout}
|
||
if int(summary.get("remaining_targets") or 0) <= 0:
|
||
return False, "已无剩余目标节点", {"rollout": rollout}
|
||
|
||
target_nodes = list(rollout.get("target_nodes") or [])
|
||
launched_node_codes = set(summary.get("launched_target_codes") or [])
|
||
remaining_targets = [
|
||
item
|
||
for item in target_nodes
|
||
if str(item.get("node_code") or "").strip() and str(item.get("node_code") or "").strip() not in launched_node_codes
|
||
]
|
||
if not remaining_targets:
|
||
return False, "已无剩余目标节点", {"rollout": rollout}
|
||
|
||
policy = dict(rollout.get("policy") or {})
|
||
batch_size = int(policy.get("first_batch_size") or 1) if int(summary.get("jobs_created") or 0) == 0 else int(policy.get("batch_size") or 1)
|
||
batch_targets = remaining_targets[: max(1, batch_size)]
|
||
|
||
from app.services.ops_job_service import create_ops_job, dispatch_ops_job
|
||
|
||
jobs: list[dict] = []
|
||
auto_dispatch = bool(policy.get("auto_dispatch", False))
|
||
auto_approve = bool(policy.get("auto_approve", False))
|
||
execution_mode = str(policy.get("execution_mode") or "remote-agent").strip() or "remote-agent"
|
||
for target in batch_targets:
|
||
target_node_code = str(target.get("node_code") or "").strip()
|
||
if not target_node_code:
|
||
continue
|
||
job_payload = _build_release_job_payload(release, rollout, target=target)
|
||
job_ok, _job_message, job_data = create_ops_job(
|
||
{
|
||
"action": "deploy.release",
|
||
"target_type": "node",
|
||
"target_node_code": target_node_code,
|
||
"requested_by": created_by,
|
||
"execution_mode": execution_mode,
|
||
"auto_approve": auto_approve,
|
||
"rollout_id": int(rollout.get("id") or 0),
|
||
"payload": {
|
||
**job_payload,
|
||
"target_node_code": target_node_code,
|
||
},
|
||
}
|
||
)
|
||
job = (job_data.get("job") or {}) if isinstance(job_data, dict) else {}
|
||
if job:
|
||
jobs.append(job)
|
||
if auto_dispatch and str(job.get("status") or "") == "queued":
|
||
dispatch_ops_job(int(job.get("id") or 0))
|
||
elif not job_ok:
|
||
continue
|
||
|
||
refreshed = refresh_release_rollout(int(rollout_id), advance_if_possible=False)
|
||
refreshed.setdefault("result_summary", {})
|
||
refreshed["result_summary"]["last_enqueue_reason"] = reason
|
||
return True, "已生成下一批发布任务", {"rollout": refreshed, "jobs": jobs}
|
||
|
||
|
||
def create_release_rollout(release_id: int, payload: dict) -> tuple[bool, str, dict]:
|
||
ensure_ops_release_schema()
|
||
release = get_release(release_id)
|
||
if not release:
|
||
return False, "Release 不存在", {}
|
||
|
||
selector = payload.get("target_selector") or {}
|
||
created_by = str(payload.get("created_by") or "api").strip() or "api"
|
||
target_nodes = _resolve_rollout_targets(selector)
|
||
if not target_nodes:
|
||
return False, "未命中任何目标节点", {"release": release}
|
||
|
||
normalized_policy = _normalize_rollout_policy(payload.get("policy") or {}, total_targets=len(target_nodes))
|
||
rollout_code = _build_rollout_code()
|
||
batches_total = _calculate_total_batches(len(target_nodes), normalized_policy)
|
||
|
||
with get_db() as conn:
|
||
conn.autocommit = False
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
"""
|
||
INSERT INTO ops_release_rollouts (
|
||
release_id, rollout_code, target_selector_json, target_nodes_json, policy_json, status,
|
||
batch_cursor, batches_total, jobs_total, jobs_created, result_summary_json, created_by, created_at, updated_at
|
||
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||
RETURNING
|
||
id, release_id, rollout_code, target_selector_json, target_nodes_json, policy_json, status,
|
||
batch_cursor, batches_total, jobs_total, jobs_created, result_summary_json, created_by, created_at, updated_at
|
||
""",
|
||
(
|
||
int(release_id),
|
||
rollout_code,
|
||
json.dumps(selector, ensure_ascii=False),
|
||
json.dumps({"items": target_nodes}, ensure_ascii=False),
|
||
json.dumps(normalized_policy, ensure_ascii=False),
|
||
"planned",
|
||
0,
|
||
batches_total,
|
||
len(target_nodes),
|
||
0,
|
||
json.dumps({}, ensure_ascii=False),
|
||
created_by,
|
||
),
|
||
)
|
||
row = cur.fetchone()
|
||
conn.commit()
|
||
|
||
rollout = _serialize_rollout_row(row)
|
||
if bool(normalized_policy.get("auto_start", True)):
|
||
return _enqueue_rollout_batch(int(rollout.get("id") or 0), created_by=created_by, reason="initial_batch")
|
||
|
||
refreshed = refresh_release_rollout(int(rollout.get("id") or 0), advance_if_possible=False)
|
||
return True, "Release rollout 已创建", {"rollout": refreshed, "jobs": []}
|
||
|
||
|
||
def advance_release_rollout(rollout_id: int, payload: dict | None = None) -> tuple[bool, str, dict]:
|
||
payload = payload or {}
|
||
created_by = str(payload.get("created_by") or "api").strip() or "api"
|
||
return _enqueue_rollout_batch(int(rollout_id), created_by=created_by, reason="manual_advance")
|
||
|
||
|
||
def list_release_rollouts(*, release_id: int | None = None, limit: int = 20) -> list[dict]:
|
||
ensure_ops_release_schema()
|
||
safe_limit = min(max(int(limit or 20), 1), 100)
|
||
conditions: list[str] = []
|
||
params: list[object] = []
|
||
if release_id is not None and int(release_id or 0) > 0:
|
||
conditions.append("release_id = %s")
|
||
params.append(int(release_id))
|
||
where_clause = f"WHERE {' AND '.join(conditions)}" if conditions else ""
|
||
with get_db() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(
|
||
f"""
|
||
SELECT
|
||
id, release_id, rollout_code, target_selector_json, target_nodes_json, policy_json, status,
|
||
batch_cursor, batches_total, jobs_total, jobs_created, result_summary_json, created_by, created_at, updated_at
|
||
FROM ops_release_rollouts
|
||
{where_clause}
|
||
ORDER BY created_at DESC, id DESC
|
||
LIMIT %s
|
||
""",
|
||
(*params, safe_limit),
|
||
)
|
||
rows = cur.fetchall()
|
||
return [_serialize_rollout_row(row) for row in rows]
|
||
|
||
|
||
def get_release_rollout(rollout_id: int, *, refresh: bool = True) -> dict:
|
||
ensure_ops_release_schema()
|
||
if refresh:
|
||
refreshed = refresh_release_rollout(int(rollout_id), advance_if_possible=False)
|
||
if refreshed:
|
||
return refreshed
|
||
row = _get_rollout_row(int(rollout_id))
|
||
return _serialize_rollout_row(row) if row else {}
|
||
|
||
|
||
def list_release_rollout_jobs(rollout_id: int, *, limit: int = 200) -> list[dict]:
|
||
from app.services.ops_job_service import list_ops_jobs
|
||
|
||
safe_limit = min(max(int(limit or 200), 1), 2000)
|
||
return list_ops_jobs(limit=safe_limit, rollout_id=int(rollout_id))
|