This commit is contained in:
Your Name
2026-04-17 12:30:31 +08:00
parent 954a6a2c65
commit dc34a4e294
4 changed files with 472 additions and 1 deletions

View File

@@ -161,6 +161,35 @@ def register_node_heartbeat(
conn.commit()
def cleanup_imported_runtime_nodes(*, region: str, role: str, keep_node_code: str) -> None:
normalized_region = str(region or "").strip() or "unknown"
normalized_role = str(role or "").strip() or "unknown"
preserved_node_code = str(keep_node_code or "").strip()
if not preserved_node_code:
return
with get_db() as conn:
with conn.cursor() as cur:
cur.execute(
"""
DELETE FROM detect_worker_nodes
WHERE region = %s
AND role = %s
AND node_code <> %s
AND (
node_code = %s
OR (metadata_json->>'service') = 'runtime-ingest'
)
""",
(
normalized_region,
normalized_role,
preserved_node_code,
f"{normalized_region}-{normalized_role}-imported",
),
)
conn.commit()
def register_local_control_heartbeat() -> None:
register_node_heartbeat(
node_code=settings.node_code,

View File

@@ -10,7 +10,7 @@ from uuid import uuid4
from app.core.config import settings
from app.core.db import get_db
from app.services.cluster_runtime_service import register_node_heartbeat
from app.services.cluster_runtime_service import cleanup_imported_runtime_nodes, register_node_heartbeat
from app.services.sync_record_service import _decode_json, _normalize_region
@@ -91,6 +91,43 @@ def _refresh_remote_runtime_node(*, source_region: str, projection: dict, receiv
hostname_override=hostname,
ip_override=ip,
)
cleanup_imported_runtime_nodes(region=region, role=role, keep_node_code=node_code)
active_job = projection.get("active_job") or {}
for node_stat in list(active_job.get("node_stats") or []):
worker_node_code = str(node_stat.get("node_code") or "").strip()
if not worker_node_code or worker_node_code == "unassigned":
continue
items_running = int(node_stat.get("items_running", 0) or 0)
items_claimed = int(node_stat.get("items_claimed", 0) or 0)
items_total = int(node_stat.get("items_total", 0) or 0)
worker_status = "busy" if (items_running > 0 or items_claimed > 0) else "online"
worker_load = max(items_running, items_claimed, 0)
worker_metadata = {
"service": "runtime-ingest",
"projection_source_region": source_region,
"worker_mode": projection.get("worker_mode", ""),
"phase_label": projection.get("phase_label", ""),
"phase_detail": projection.get("phase_detail", ""),
"proxy_runtime_label": projection.get("proxy_runtime_label", ""),
"proxy_runtime_reason": projection.get("proxy_runtime_reason", ""),
"updated_at": _format_time(received_at or datetime.now()),
"job_items_total": items_total,
"job_items_running": items_running,
"job_items_claimed": items_claimed,
"derived_from": node_code,
}
register_node_heartbeat(
node_code=worker_node_code,
region=region,
role="worker",
status=worker_status,
current_load=worker_load,
metadata=worker_metadata,
hostname_override=hostname,
ip_override=ip,
)
cleanup_imported_runtime_nodes(region=region, role="worker", keep_node_code=worker_node_code)
def _load_latest_projection(sync_type: str) -> dict | None:

View File

@@ -435,6 +435,7 @@ def append_runtime_projection_if_changed(
"items_pending": active_job.get("items_pending", 0),
"items_running": active_job.get("items_running", 0),
"items_failed": active_job.get("items_failed", 0),
"node_stats": list(active_job.get("node_stats") or []),
},
"cluster_summary": {
"nodes_total": int(cluster.get("nodes_total", 0) or 0),