fix runtime ingest effective worker projection
This commit is contained in:
@@ -79,6 +79,14 @@ def _refresh_remote_runtime_node(*, source_region: str, projection: dict, receiv
|
|||||||
"phase_detail": projection.get("phase_detail", ""),
|
"phase_detail": projection.get("phase_detail", ""),
|
||||||
"proxy_runtime_label": projection.get("proxy_runtime_label", ""),
|
"proxy_runtime_label": projection.get("proxy_runtime_label", ""),
|
||||||
"proxy_runtime_reason": projection.get("proxy_runtime_reason", ""),
|
"proxy_runtime_reason": projection.get("proxy_runtime_reason", ""),
|
||||||
|
"worker_online": bool(projection.get("worker_online", False)),
|
||||||
|
"detect_participating": bool(projection.get("detect_participating", False)),
|
||||||
|
"active_job_code": str((projection.get("active_job") or {}).get("job_code") or ""),
|
||||||
|
"active_job_status": str((projection.get("active_job") or {}).get("status") or ""),
|
||||||
|
"job_items_total": int(((projection.get("active_job") or {}).get("items_total", 0) or 0)),
|
||||||
|
"job_items_claimed": int(((projection.get("active_job") or {}).get("items_claimed", 0) or 0)),
|
||||||
|
"job_items_running": int(((projection.get("active_job") or {}).get("items_running", 0) or 0)),
|
||||||
|
"job_items_completed": int(((projection.get("active_job") or {}).get("items_completed", 0) or 0)),
|
||||||
"updated_at": _format_time(received_at or datetime.now()),
|
"updated_at": _format_time(received_at or datetime.now()),
|
||||||
}
|
}
|
||||||
register_node_heartbeat(
|
register_node_heartbeat(
|
||||||
|
|||||||
@@ -444,6 +444,12 @@ def append_runtime_projection_if_changed(
|
|||||||
normalized_source_region = _normalize_region(source_region, _normalize_region(settings.sync_source_region, settings.node_region))
|
normalized_source_region = _normalize_region(source_region, _normalize_region(settings.sync_source_region, settings.node_region))
|
||||||
normalized_target_region = _normalize_region(target_region, _normalize_region(settings.sync_target_region, "overseas"))
|
normalized_target_region = _normalize_region(target_region, _normalize_region(settings.sync_target_region, "overseas"))
|
||||||
active_job = detect.get("active_job") or {}
|
active_job = detect.get("active_job") or {}
|
||||||
|
local_participating = False
|
||||||
|
for node in list(cluster.get("nodes") or []):
|
||||||
|
if str(node.get("node_code") or "").strip() != settings.node_code:
|
||||||
|
continue
|
||||||
|
local_participating = bool(node.get("detect_participating", False) or node.get("current_load", 0))
|
||||||
|
break
|
||||||
projection = {
|
projection = {
|
||||||
"node": {
|
"node": {
|
||||||
"node_code": settings.node_code,
|
"node_code": settings.node_code,
|
||||||
@@ -453,6 +459,7 @@ def append_runtime_projection_if_changed(
|
|||||||
"ip": _resolve_local_ip(),
|
"ip": _resolve_local_ip(),
|
||||||
},
|
},
|
||||||
"worker_online": bool(detect.get("worker_online", False)),
|
"worker_online": bool(detect.get("worker_online", False)),
|
||||||
|
"detect_participating": local_participating,
|
||||||
"worker_mode": detect.get("worker_mode", ""),
|
"worker_mode": detect.get("worker_mode", ""),
|
||||||
"phase_label": detect.get("phase_label", ""),
|
"phase_label": detect.get("phase_label", ""),
|
||||||
"phase_detail": detect.get("phase_detail", ""),
|
"phase_detail": detect.get("phase_detail", ""),
|
||||||
|
|||||||
Reference in New Issue
Block a user