- 在线 Worker {{ runtime.detect?.capacity_plan?.online_worker_nodes || 0 }}
+ 有效执行节点 {{ runtime.detect?.capacity_plan?.online_worker_nodes || 0 }}
当前吞吐 {{ runtime.detect?.capacity_plan?.current_processed_per_hour || 0 }}/小时
预计剩余 {{ runtime.detect?.capacity_plan?.estimated_hours_remaining || 0 }} 小时
建议总 Worker {{ runtime.detect?.capacity_plan?.recommended_total_workers || 1 }}
@@ -377,7 +382,7 @@
在线控制面 {{ runtime.cluster?.summary?.online_control_nodes || 0 }}
- 在线 Worker {{ runtime.cluster?.summary?.online_worker_nodes || 0 }}
+ 有效执行节点 {{ runtime.cluster?.summary?.online_worker_nodes || 0 }}
独立 Worker {{ runtime.cluster?.summary?.dedicated_online_worker_nodes || 0 }}
忙碌 {{ runtime.cluster?.summary?.status_counts?.busy || 0 }}
失活 {{ runtime.cluster?.summary?.status_counts?.stale || 0 }}
@@ -803,17 +808,19 @@ const clusterIssueCount = computed(() => {
const clusterAlertText = computed(() => {
const summary = runtime.cluster?.summary || {};
const busy = Array.isArray(summary.busy_nodes) ? summary.busy_nodes.join("、") : "";
- const stale = Array.isArray(summary.stale_nodes) ? summary.stale_nodes.join("、") : "";
- const offline = Array.isArray(summary.offline_nodes) ? summary.offline_nodes.join("、") : "";
- const parts = [];
+ const staleNodes = Array.isArray(summary.stale_nodes) ? summary.stale_nodes : [];
+ const offlineNodes = Array.isArray(summary.offline_nodes) ? summary.offline_nodes : [];
+ const parts = [
+ `当前有效执行节点 ${summary.online_worker_nodes || 0},其中独立 Worker ${summary.dedicated_online_worker_nodes || 0}`
+ ];
if (busy) {
parts.push(`忙碌节点:${busy}`);
}
- if (stale) {
- parts.push(`失活节点:${stale}`);
+ if (staleNodes.length) {
+ parts.push(`失活节点 ${staleNodes.length} 个${staleNodes.length <= 2 ? `:${staleNodes.join("、")}` : ",详见下方节点表"}`);
}
- if (offline) {
- parts.push(`离线节点:${offline}`);
+ if (offlineNodes.length) {
+ parts.push(`离线节点 ${offlineNodes.length} 个${offlineNodes.length <= 2 ? `:${offlineNodes.join("、")}` : ",详见下方节点表"}`);
}
return parts.join(";");
});
@@ -824,7 +831,8 @@ const participatingNodesAlertText = computed(() => {
return "当前没有节点正在领任务、执行任务或产生近窗吞吐。";
}
const names = rows.map((item: Record) => String(item.node_code || "-")).join("、");
- return `当前参与检测节点 ${rows.length} 台:${names}`;
+ const effectiveCount = rows.filter((item: Record) => Boolean(item.is_effective_worker)).length;
+ return `当前参与检测节点 ${rows.length} 台,其中有效执行节点 ${effectiveCount} 台:${names}`;
});
const syncAlertText = computed(() => {
@@ -1034,18 +1042,36 @@ const invokeAction = async (action: "start_worker" | "stop_worker" | "restart_ap
stop_sync_agent: "停止 Sync Agent",
push_sync: "立即同步"
};
+ const responseMessage = String(response.message || "");
+ const uiLevel = String(response.data?.ui_level || "").trim().toLowerCase();
+ const actionType: "success" | "warning" | "error" =
+ uiLevel === "error" ? "error" : uiLevel === "warning" ? "warning" : "success";
lastAction.value = {
label: labelMap[action],
- message: response.message,
+ message: responseMessage,
at: new Date().toLocaleString("zh-CN", { hour12: false }),
- type: "success"
+ type: actionType
};
persistLastAction();
appendActionHistory(lastAction.value);
- ElMessage.success(response.message);
- window.setTimeout(() => {
- loadAll(false);
- }, Number(response.data?.poll_after_seconds || 2) * 1000);
+ if (actionType === "error") {
+ ElMessage.error(responseMessage);
+ } else if (actionType === "warning") {
+ ElMessage.warning(responseMessage);
+ } else {
+ ElMessage.success(responseMessage);
+ }
+ await loadAll(false);
+ const pollSchedule = Array.isArray(response.data?.poll_schedule_seconds)
+ ? response.data.poll_schedule_seconds
+ : [Number(response.data?.poll_after_seconds || 2)];
+ [...new Set(pollSchedule.map((item: unknown) => Number(item || 0)).filter((item: number) => item > 0))]
+ .sort((a, b) => a - b)
+ .forEach((seconds) => {
+ window.setTimeout(() => {
+ loadAll(false);
+ }, seconds * 1000);
+ });
} catch (error: any) {
const message = error?.message || "运行时动作执行失败";
const labelMap = {
diff --git a/domain-web/src/views/settings/SettingsView.vue b/domain-web/src/views/settings/SettingsView.vue
index e003190..1b0b25b 100644
--- a/domain-web/src/views/settings/SettingsView.vue
+++ b/domain-web/src/views/settings/SettingsView.vue
@@ -54,6 +54,25 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
@@ -232,7 +251,9 @@ const runtimeSettings = ref({
worker_mode: "windows-local",
worker_service_name: "domaincheck-worker",
api_service_name: "domaincheck-api",
- sync_agent_service_name: "domaincheck-sync-agent"
+ sync_agent_service_name: "domaincheck-sync-agent",
+ worker_log_sync_enabled: false,
+ worker_log_sync_mode: "key"
});
const importInput = ref(null);
const juziseoLoggingIn = ref(false);
@@ -316,6 +337,10 @@ const normalizeNodeThreadOverrides = (payload: Record)
.sort((a, b) => a.node_code.localeCompare(b.node_code));
};
+const normalizeWorkerLogSyncMode = (value: unknown) => {
+ return String(value || "").trim().toLowerCase() === "full" ? "full" : "key";
+};
+
const addNodeThreadOverride = () => {
nodeThreadOverrides.value.push({
node_code: "",
@@ -371,7 +396,9 @@ const loadSettings = async () => {
worker_mode: response.data.runtime_settings?.worker_mode || "windows-local",
worker_service_name: response.data.runtime_settings?.worker_service_name || "domaincheck-worker",
api_service_name: response.data.runtime_settings?.api_service_name || "domaincheck-api",
- sync_agent_service_name: response.data.runtime_settings?.sync_agent_service_name || "domaincheck-sync-agent"
+ sync_agent_service_name: response.data.runtime_settings?.sync_agent_service_name || "domaincheck-sync-agent",
+ worker_log_sync_enabled: Boolean(response.data.runtime_settings?.worker_log_sync_enabled),
+ worker_log_sync_mode: normalizeWorkerLogSyncMode(response.data.runtime_settings?.worker_log_sync_mode)
};
} catch {
ElMessage.error("读取系统设置失败");
diff --git a/domainCheck/detect_worker.py b/domainCheck/detect_worker.py
index 5e33531..2e1fe4f 100644
--- a/domainCheck/detect_worker.py
+++ b/domainCheck/detect_worker.py
@@ -41,6 +41,7 @@ CONFIG_UPDATE_CHANNEL = "domain_tool:config_update"
CONTROL_CHANNEL = "domain_tool:worker_control"
RUNTIME_STATE_KEY = "domain_tool:detect_runtime_state"
PENDING_CONTROL_KEY = "domain_tool:worker_pending_command"
+RUNTIME_SETTINGS_KEY = "domain_tool:runtime_settings"
class DetectThread(QThread):
"""
@@ -814,6 +815,11 @@ class DetectWorker:
self.current_cycle_token = ""
self.current_job_id = None
self.current_job_code = ""
+ self.runtime_settings = {
+ "worker_log_sync_enabled": False,
+ "worker_log_sync_mode": "key",
+ }
+ self._last_synced_worker_log = ""
# 初始化数据库连接
self.db = Database()
@@ -868,6 +874,7 @@ class DetectWorker:
self.detect_options = self.load_detect_options()
self.proxy_config = self.load_proxy_config()
self.thread_count = self.load_thread_count() # 从配置文件加载线程数
+ self.runtime_settings = self.load_runtime_settings()
# 初始化代理池
self.proxy_pool = []
@@ -1042,6 +1049,7 @@ class DetectWorker:
def _set_active_cycle_context(self, control_payload=None):
payload = control_payload or {}
+ self._last_synced_worker_log = ""
self.current_cycle_token = str(payload.get("cycle_token") or "").strip()
job_id = payload.get("job_id")
try:
@@ -1051,6 +1059,7 @@ class DetectWorker:
self.current_job_code = str(payload.get("job_code") or "").strip()
def _clear_active_cycle_context(self):
+ self._last_synced_worker_log = ""
self.current_cycle_token = ""
self.current_job_id = None
self.current_job_code = ""
@@ -1091,6 +1100,7 @@ class DetectWorker:
self.detecting = True
try:
logger.info(f"开始执行远程检测任务,来源: {source}")
+ self._sync_worker_log_event(f"开始执行检测任务,来源: {source}", payload={"source": source}, mode='key')
self._update_runtime_state(
"starting",
f"开始执行检测任务,来源: {source}",
@@ -1136,6 +1146,7 @@ class DetectWorker:
self._update_runtime_state("idle", "当前没有运行中的检测任务")
return False
logger.info(f"收到停止检测指令,来源: {source}")
+ self._sync_worker_log_event(f"收到停止检测指令,来源: {source}", payload={"source": source}, mode='key')
self.stop_requested = True
self._update_runtime_state(
"stopping",
@@ -1351,6 +1362,79 @@ class DetectWorker:
logger.info(f"使用默认检测线程数: {default_thread_count}")
return default_thread_count
+ def load_runtime_settings(self):
+ default_settings = {
+ "worker_log_sync_enabled": False,
+ "worker_log_sync_mode": "key",
+ }
+ try:
+ if self.use_redis:
+ runtime_settings_raw = self.redis_client.get(RUNTIME_SETTINGS_KEY)
+ if runtime_settings_raw:
+ runtime_settings = default_settings.copy()
+ runtime_settings.update(json.loads(runtime_settings_raw))
+ runtime_settings["worker_log_sync_enabled"] = bool(runtime_settings.get("worker_log_sync_enabled", False))
+ runtime_settings["worker_log_sync_mode"] = "full" if str(runtime_settings.get("worker_log_sync_mode", "key")).strip().lower() == "full" else "key"
+ logger.info(f"从Redis加载运行时设置成功: {runtime_settings}")
+ return runtime_settings
+
+ candidate_paths = [
+ 'runtime_settings.json',
+ os.path.join('runtime', 'runtime_settings.json'),
+ ]
+ for candidate_path in candidate_paths:
+ if os.path.exists(candidate_path):
+ with open(candidate_path, 'r', encoding='utf-8') as f:
+ runtime_settings = default_settings.copy()
+ runtime_settings.update(json.load(f))
+ runtime_settings["worker_log_sync_enabled"] = bool(runtime_settings.get("worker_log_sync_enabled", False))
+ runtime_settings["worker_log_sync_mode"] = "full" if str(runtime_settings.get("worker_log_sync_mode", "key")).strip().lower() == "full" else "key"
+ logger.info(f"从本地文件加载运行时设置成功: {runtime_settings}")
+ return runtime_settings
+ except Exception as e:
+ logger.error(f"加载运行时设置失败: {e}")
+ return default_settings
+
+ def _worker_log_sync_mode(self):
+ if not bool((self.runtime_settings or {}).get("worker_log_sync_enabled", False)):
+ return "off"
+ return "full" if str((self.runtime_settings or {}).get("worker_log_sync_mode", "key")).strip().lower() == "full" else "key"
+
+ def _sync_worker_log_event(self, message, level='info', payload=None, *, mode='key', job_item_id=None):
+ current_mode = self._worker_log_sync_mode()
+ if current_mode == "off":
+ return
+ if mode == 'full' and current_mode != "full":
+ return
+ normalized_message = str(message or '').strip()
+ if not normalized_message:
+ return
+ if normalized_message == self._last_synced_worker_log and mode != 'full':
+ return
+ job_id = self.current_job_id
+ if not job_id:
+ return
+ event_payload = dict(payload or {})
+ if self.current_cycle_token and not event_payload.get("cycle_token"):
+ event_payload["cycle_token"] = self.current_cycle_token
+ if self.current_job_code and not event_payload.get("job_code"):
+ event_payload["job_code"] = self.current_job_code
+ event_payload["log_mode"] = mode
+ event_payload["synced_by"] = "worker_log_callback"
+ try:
+ self.db.append_detect_run_event(
+ job_id,
+ job_item_id,
+ config.NODE_CODE,
+ event_type='worker_log',
+ message=normalized_message,
+ level=level,
+ payload=event_payload,
+ )
+ self._last_synced_worker_log = normalized_message
+ except Exception as e:
+ logger.debug(f"回传Worker日志事件失败: {e}")
+
def test_proxy(self, proxy_item, result_queue):
"""
测试单个代理的可用性
@@ -1495,6 +1579,7 @@ class DetectWorker:
"idle" if not self.detecting else "running",
"代理未启用,使用直接连接",
)
+ self._sync_worker_log_event("代理未启用,使用直接连接", mode='key')
self.proxy_refresh_lock.release()
return
@@ -1510,6 +1595,7 @@ class DetectWorker:
"refreshing_proxy" if self.detecting else "idle",
f"代理池刷新冷却中,{wait_seconds} 秒后再试",
)
+ self._sync_worker_log_event(f"代理池刷新冷却中,{wait_seconds} 秒后再试", mode='full')
return
proxy_api_urls = self.proxy_config.get('proxy_urls') or []
@@ -1622,6 +1708,16 @@ class DetectWorker:
"running" if self.detecting else "idle",
f"代理池刷新完成,共 {len(new_proxies)} 个可用代理,来源链接 {len(proxy_api_urls)} 个,原始 {self.proxy_last_refresh_total_items} 个,验证 {self.proxy_last_validated_count} 个",
)
+ self._sync_worker_log_event(
+ f"代理池刷新完成,共 {len(new_proxies)} 个可用代理,来源链接 {len(proxy_api_urls)} 个,原始 {self.proxy_last_refresh_total_items} 个,验证 {self.proxy_last_validated_count} 个",
+ payload={
+ "available_proxy_count": len(new_proxies),
+ "source_count": len(proxy_api_urls),
+ "raw_items": self.proxy_last_refresh_total_items,
+ "validated_count": self.proxy_last_validated_count,
+ },
+ mode='key',
+ )
else:
with self.proxy_pool_lock:
self.proxy_pool = []
@@ -1636,6 +1732,7 @@ class DetectWorker:
"refreshing_proxy" if self.detecting else "idle",
"所有代理池链接均未返回可用代理数据",
)
+ self._sync_worker_log_event("所有代理池链接均未返回可用代理数据", level='warning', mode='key')
else:
with self.proxy_pool_lock:
self.proxy_pool = []
@@ -1649,6 +1746,7 @@ class DetectWorker:
"idle" if not self.detecting else "running",
"未配置代理池链接",
)
+ self._sync_worker_log_event("未配置代理池链接", level='warning', mode='key')
except Exception as e:
with self.proxy_pool_lock:
self.proxy_pool = []
@@ -1663,6 +1761,7 @@ class DetectWorker:
"refreshing_proxy" if self.detecting else "failed",
f"刷新代理池失败: {e}",
)
+ self._sync_worker_log_event(f"刷新代理池失败: {e}", level='error', mode='key')
finally:
self.proxy_refresh_lock.release()
@@ -2383,6 +2482,15 @@ class DetectWorker:
"job_code": job_code,
},
)
+ self._sync_worker_log_event(
+ f"开始检测域名: {domain_name}",
+ payload={
+ "domain_id": domain_id,
+ "domain": domain_name,
+ },
+ mode='full',
+ job_item_id=job_item_id,
+ )
# 使用共享的敏感词列表
sensitive_words = self.sensitive_words
@@ -2433,11 +2541,32 @@ class DetectWorker:
self._complete_detection(domain_id, domain_name)
finalize_job_item('completed', f"域名检测完成: {domain_name}")
logger.info(f"域名检测完成: {domain_name}")
+ self._sync_worker_log_event(
+ f"域名检测完成: {domain_name}",
+ payload={
+ "domain_id": domain_id,
+ "domain": domain_name,
+ "status": "completed",
+ },
+ mode='full',
+ job_item_id=job_item_id,
+ )
except Exception as e:
logger.error(f"检测域名出错: {domain_name}, 错误: {e}")
self.db.update_domain_detect_status(domain_id, DETECT_STATUS_FAILED)
finalize_job_item('failed', str(e))
+ self._sync_worker_log_event(
+ f"检测域名出错: {domain_name}, 错误: {e}",
+ level='error',
+ payload={
+ "domain_id": domain_id,
+ "domain": domain_name,
+ "status": "failed",
+ },
+ mode='full',
+ job_item_id=job_item_id,
+ )
def start_detection(self):
"""
@@ -2447,6 +2576,7 @@ class DetectWorker:
self.detecting = True
self.stop_requested = False
self._mark_detection_phase("preparing", "开始执行域名检测任务,正在加载配置")
+ self._sync_worker_log_event("开始执行域名检测任务,正在加载配置", mode='key')
# 重新加载配置,确保获取最新的配置
try:
@@ -2468,6 +2598,7 @@ class DetectWorker:
import traceback
logger.error(traceback.format_exc())
self._mark_detection_phase("failed", f"加载配置失败: {e}")
+ self._sync_worker_log_event(f"加载配置失败: {e}", level='error', mode='key')
return
# 重新加载cookies
@@ -2486,6 +2617,7 @@ class DetectWorker:
if self.proxy_config.get('proxy_enable', False):
logger.info("开始检测,刷新代理池")
self._mark_detection_phase("refreshing_proxy", "开始检测,正在刷新代理池")
+ self._sync_worker_log_event("开始检测,正在刷新代理池", mode='key')
self.refresh_proxy_pool()
# 更新JC和Juziseo实例的代理设置
@@ -2516,6 +2648,7 @@ class DetectWorker:
if self.stop_requested:
logger.info("检测任务收到停止请求,停止继续领取任务")
self._mark_detection_phase("stopping", "检测任务收到停止请求,准备安全退出")
+ self._sync_worker_log_event("检测任务收到停止请求,准备安全退出", level='warning', mode='key')
break
recycled_job_items = self.db.recycle_expired_detect_job_items()
if recycled_job_items:
@@ -2538,10 +2671,19 @@ class DetectWorker:
batch_size=current_batch_size,
queue_source="detect_job_items" if using_job_queue else "domains",
)
+ self._sync_worker_log_event(
+ f"从{queue_label}获取到 {current_batch_size} 个需要检测的域名",
+ payload={
+ "batch_size": current_batch_size,
+ "queue_source": "detect_job_items" if using_job_queue else "domains",
+ },
+ mode='full',
+ )
if not domains:
logger.info("没有需要检测的域名")
self._mark_detection_phase("idle", "当前没有需要检测的域名,Worker 等待下一次启动")
+ self._sync_worker_log_event("当前没有需要检测的域名,Worker 等待下一次启动", mode='key')
break
# 旧链路直接扫 domains 表时,少于 batch_size 代表已接近尾批;
@@ -2560,6 +2702,11 @@ class DetectWorker:
batch_size=current_batch_size,
max_threads=max_threads,
)
+ self._sync_worker_log_event(
+ f"开始创建线程,当前批次域名数: {current_batch_size},最大线程数: {max_threads}",
+ payload={"batch_size": current_batch_size, "max_threads": max_threads},
+ mode='full',
+ )
for i, domain in enumerate(domains):
# 检查是否需要停止
if not self.running or self.stop_requested:
@@ -2627,6 +2774,15 @@ class DetectWorker:
max_threads=max_threads,
batch_size=current_batch_size,
)
+ self._sync_worker_log_event(
+ f"当前实际线程数量: {current_active}/{max_threads}",
+ payload={
+ "active_threads": current_active,
+ "max_threads": max_threads,
+ "batch_size": current_batch_size,
+ },
+ mode='full',
+ )
# 更新GUI显示
if self.detect_thread:
@@ -2705,9 +2861,19 @@ class DetectWorker:
f"当前批次域名数量 ({current_batch_size}) 少于 {batch_size},本轮检测即将完成",
processed=total_processed,
)
+ self._sync_worker_log_event(
+ f"当前批次域名数量 ({current_batch_size}) 少于 {batch_size},本轮检测即将完成",
+ payload={"processed": total_processed},
+ mode='key',
+ )
break
logger.info(f"当前批次检测完成,累计处理 {total_processed} 个域名")
self._mark_detection_phase("batch_completed", f"当前批次检测完成,累计处理 {total_processed} 个域名", processed=total_processed)
+ self._sync_worker_log_event(
+ f"当前批次检测完成,累计处理 {total_processed} 个域名",
+ payload={"processed": total_processed},
+ mode='key',
+ )
# 完成进度
if self.detect_thread:
@@ -2721,10 +2887,12 @@ class DetectWorker:
logger.info("域名检测任务完成")
self._mark_detection_phase("completed", "域名检测任务完成")
+ self._sync_worker_log_event("域名检测任务完成", payload={"processed": total_processed}, mode='key')
except Exception as e:
logger.error(f"执行检测任务出错: {e}")
self._mark_detection_phase("failed", f"执行检测任务出错: {e}")
+ self._sync_worker_log_event(f"执行检测任务出错: {e}", level='error', mode='key')
finally:
self.detecting = False
@@ -2922,6 +3090,7 @@ class DetectWorker:
self.detect_options = self.load_detect_options()
self.proxy_config = self.load_proxy_config()
self.thread_count = self.load_thread_count()
+ self.runtime_settings = self.load_runtime_settings()
self.load_cookies_from_remote()
self.update_config_labels()
logger.debug("配置已更新")