diff --git a/domain-api/app/services/sync_push_service.py b/domain-api/app/services/sync_push_service.py index 5880178..6dfe74a 100644 --- a/domain-api/app/services/sync_push_service.py +++ b/domain-api/app/services/sync_push_service.py @@ -79,6 +79,14 @@ def _refresh_remote_runtime_node(*, source_region: str, projection: dict, receiv "phase_detail": projection.get("phase_detail", ""), "proxy_runtime_label": projection.get("proxy_runtime_label", ""), "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()), } register_node_heartbeat( diff --git a/domain-api/app/services/sync_record_service.py b/domain-api/app/services/sync_record_service.py index 2cc5cee..8792dc4 100644 --- a/domain-api/app/services/sync_record_service.py +++ b/domain-api/app/services/sync_record_service.py @@ -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_target_region = _normalize_region(target_region, _normalize_region(settings.sync_target_region, "overseas")) 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 = { "node": { "node_code": settings.node_code, @@ -453,6 +459,7 @@ def append_runtime_projection_if_changed( "ip": _resolve_local_ip(), }, "worker_online": bool(detect.get("worker_online", False)), + "detect_participating": local_participating, "worker_mode": detect.get("worker_mode", ""), "phase_label": detect.get("phase_label", ""), "phase_detail": detect.get("phase_detail", ""),