diff --git a/backend/app/api/simulation.py b/backend/app/api/simulation.py index 0df3f43c..4adf627e 100644 --- a/backend/app/api/simulation.py +++ b/backend/app/api/simulation.py @@ -337,8 +337,9 @@ def _check_simulation_prepared(simulation_id: str) -> tuple: # - running: 正在运行,说明准备早就完成了 # - completed: 运行完成,说明准备早就完成了 # - stopped: 已停止,说明准备早就完成了 + # - paused: 手动停止后会写入 paused,配置仍然可复用 # - failed: 运行失败(但准备是完成的) - prepared_statuses = ["ready", "preparing", "running", "completed", "stopped", "failed"] + prepared_statuses = ["ready", "preparing", "running", "completed", "stopped", "paused", "failed"] if status in prepared_statuses and config_generated: # 获取文件统计信息 profiles_file = os.path.join(simulation_dir, "reddit_profiles.json") @@ -1691,6 +1692,40 @@ def start_simulation(): "error": t('api.graphIdRequiredForMemory') }), 400 + existing_run_state = SimulationRunner.get_run_state(simulation_id) + restartable_statuses = { + RunnerStatus.IDLE, + RunnerStatus.STOPPED, + RunnerStatus.COMPLETED, + RunnerStatus.FAILED, + } + if ( + existing_run_state + and existing_run_state.runner_status in restartable_statuses + ): + if ZepGraphMemoryManager.get_updater(simulation_id) is not None: + return jsonify({ + "success": False, + "error": ( + "The previous simulation still has pending graph " + "memory updates; finalize or reset it before restarting" + ), + }), 409 + logger.info( + f"清理已结束的旧运行记录后重新启动: " + f"simulation_id={simulation_id}, runner_status={existing_run_state.runner_status.value}" + ) + cleanup_result = SimulationRunner.cleanup_simulation_logs(simulation_id) + if not cleanup_result.get("success"): + return jsonify({ + "success": False, + "error": ( + "Failed to clean previous simulation logs: " + f"{cleanup_result.get('errors')}" + ), + }), 500 + force_restarted = True + graph_guard = ( graph_lifecycle_lock(graph_id) if enable_graph_memory_update diff --git a/backend/scripts/run_parallel_simulation.py b/backend/scripts/run_parallel_simulation.py index 9ec72f72..a6411931 100644 --- a/backend/scripts/run_parallel_simulation.py +++ b/backend/scripts/run_parallel_simulation.py @@ -608,7 +608,7 @@ def load_config(config_path: str) -> Dict[str, Any]: # 需要过滤掉的非核心动作类型(这些动作对分析价值较低) -FILTERED_ACTIONS = {'refresh', 'sign_up'} +FILTERED_ACTIONS = {'refresh', 'sign_up', 'do_nothing'} # 动作类型映射表(数据库中的名称 -> 标准名称) ACTION_TYPE_MAP = { diff --git a/backend/tests/test_zep_simulation_barrier.py b/backend/tests/test_zep_simulation_barrier.py index 63289fce..cb6878e6 100644 --- a/backend/tests/test_zep_simulation_barrier.py +++ b/backend/tests/test_zep_simulation_barrier.py @@ -270,6 +270,119 @@ def test_force_restart_does_not_continue_while_old_ingestion_is_pending(monkeypa assert cleanup_called == [] +def test_terminal_restart_cleans_old_logs_before_start(monkeypatch): + simulation = SimpleNamespace( + simulation_id="sim-terminal", + project_id="proj-1", + graph_id=None, + status=SimulationStatus.READY, + ) + events = [] + monkeypatch.setattr( + simulation_api, + "SimulationManager", + lambda: SimpleNamespace(get_simulation=lambda _simulation_id: simulation), + ) + monkeypatch.setattr( + simulation_api.SimulationRunner, + "get_run_state", + classmethod( + lambda _cls, _simulation_id: SimulationRunState( + simulation_id="sim-terminal", + runner_status=RunnerStatus.COMPLETED, + ) + ), + ) + monkeypatch.setattr( + simulation_api.SimulationRunner, + "cleanup_simulation_logs", + classmethod( + lambda _cls, _simulation_id: ( + events.append("cleanup") or {"success": True, "errors": []} + ) + ), + ) + monkeypatch.setattr( + simulation_api.SimulationRunner, + "start_simulation", + classmethod( + lambda _cls, **_kwargs: ( + events.append("start") + or SimulationRunState( + simulation_id="sim-terminal", + runner_status=RunnerStatus.STARTING, + ) + ) + ), + ) + monkeypatch.setattr( + simulation_api.ZepGraphMemoryManager, + "get_updater", + classmethod(lambda _cls, _simulation_id: None), + ) + + app = Flask(__name__) + with app.test_request_context( + "/api/simulation/start", + method="POST", + json={"simulation_id": "sim-terminal"}, + ): + response = simulation_api.start_simulation() + + assert response.status_code == 200 + assert response.get_json()["data"]["force_restarted"] is True + assert events == ["cleanup", "start"] + + +def test_terminal_restart_rejects_pending_graph_updates(monkeypatch): + simulation = SimpleNamespace( + simulation_id="sim-pending-terminal", + project_id="proj-1", + graph_id=None, + status=SimulationStatus.READY, + ) + cleanup_called = [] + monkeypatch.setattr( + simulation_api, + "SimulationManager", + lambda: SimpleNamespace(get_simulation=lambda _simulation_id: simulation), + ) + monkeypatch.setattr( + simulation_api.SimulationRunner, + "get_run_state", + classmethod( + lambda _cls, _simulation_id: SimulationRunState( + simulation_id="sim-pending-terminal", + runner_status=RunnerStatus.FAILED, + ) + ), + ) + monkeypatch.setattr( + simulation_api.SimulationRunner, + "cleanup_simulation_logs", + classmethod( + lambda _cls, _simulation_id: cleanup_called.append(True) + ), + ) + monkeypatch.setattr( + simulation_api.ZepGraphMemoryManager, + "get_updater", + classmethod(lambda _cls, _simulation_id: object()), + ) + + app = Flask(__name__) + with app.test_request_context( + "/api/simulation/start", + method="POST", + json={"simulation_id": "sim-pending-terminal"}, + ): + response, status = simulation_api.start_simulation() + + assert status == 409 + assert "pending graph memory updates" in response.get_json()["error"] + assert cleanup_called == [] + + def test_monitor_start_failure_terminates_the_spawned_process(monkeypatch, tmp_path): simulation_id = "sim-start-failure" sim_dir = tmp_path / "runs" / simulation_id @@ -482,3 +595,45 @@ def test_shutdown_drain_failure_remains_failed_and_retryable(monkeypatch): SimulationRunner._cleanup_done = False SimulationRunner._graph_memory_enabled.pop(simulation_id, None) SimulationRunner._manual_stop_requests.discard(simulation_id) + + +def test_shutdown_preserves_completed_state_without_pending_resources(monkeypatch): + simulation_id = "sim-shutdown-completed" + state = SimulationRunState( + simulation_id=simulation_id, + runner_status=RunnerStatus.COMPLETED, + completed_at="2026-07-23T00:00:00", + ) + + class FinishedProcess: + def poll(self): + return 0 + + monkeypatch.setattr( + SimulationRunner, + "get_run_state", + classmethod(lambda _cls, _simulation_id: state), + ) + monkeypatch.setattr( + runner_module.ZepGraphMemoryManager, + "get_simulation_ids", + classmethod(lambda _cls: []), + ) + monkeypatch.setattr( + runner_module.ZepGraphMemoryManager, + "get_updater", + classmethod(lambda _cls, _simulation_id: None), + ) + + SimulationRunner._cleanup_done = False + SimulationRunner._processes[simulation_id] = FinishedProcess() + SimulationRunner._graph_memory_enabled.pop(simulation_id, None) + try: + SimulationRunner.cleanup_all_simulations() + assert state.runner_status == RunnerStatus.COMPLETED + assert state.completed_at == "2026-07-23T00:00:00" + assert state.error is None + finally: + SimulationRunner._cleanup_done = False + SimulationRunner._processes.pop(simulation_id, None) + SimulationRunner._manual_stop_requests.discard(simulation_id) diff --git a/frontend/src/components/Step3Simulation.vue b/frontend/src/components/Step3Simulation.vue index e9dd3dc3..834fd678 100644 --- a/frontend/src/components/Step3Simulation.vue +++ b/frontend/src/components/Step3Simulation.vue @@ -106,9 +106,9 @@
-
+
- TOTAL EVENTS: {{ allActions.length }} + TOTAL EVENTS: {{ chronologicalActions.length }}
- - -
+
{{ action.action_args.content }}
@@ -262,7 +279,7 @@
-
+
Waiting for agent actions...
@@ -326,19 +343,21 @@ const allActions = ref([]) // 所有动作(增量累积) const actionIds = ref(new Set()) // 用于去重的动作ID集合 const scrollContainer = ref(null) +const isDisplayableAction = (action) => action?.action_type !== 'DO_NOTHING' + // Computed // 按时间顺序显示动作(最新的在最后面,即底部) const chronologicalActions = computed(() => { - return allActions.value + return allActions.value.filter(isDisplayableAction) }) // 各平台动作计数 const twitterActionsCount = computed(() => { - return allActions.value.filter(a => a.platform === 'twitter').length + return chronologicalActions.value.filter(a => a.platform === 'twitter').length }) const redditActionsCount = computed(() => { - return allActions.value.filter(a => a.platform === 'reddit').length + return chronologicalActions.value.filter(a => a.platform === 'reddit').length }) // 格式化模拟流逝时间(根据轮次和每轮分钟数计算) @@ -379,15 +398,75 @@ const resetAllState = () => { stopPolling() // 停止之前可能存在的轮询 } +const hasExistingRunData = (data) => { + if (!data) return false + const runnerStatus = data.runner_status + const actionsCount = Number(data.total_actions_count || 0) + + Number(data.twitter_actions_count || 0) + + Number(data.reddit_actions_count || 0) + return Boolean(data.started_at) + || actionsCount > 0 + || ['starting', 'running', 'completed', 'stopped', 'paused', 'failed'].includes(runnerStatus) +} + +const isActiveRun = (data) => { + return data?.runner_status === 'starting' + || data?.runner_status === 'running' + || data?.twitter_running + || data?.reddit_running +} + +const restoreExistingRun = async (data) => { + runStatus.value = data + prevTwitterRound.value = data.twitter_current_round || 0 + prevRedditRound.value = data.reddit_current_round || 0 + await fetchRunStatusDetail() + + if (isActiveRun(data)) { + phase.value = 1 + emit('update-status', 'processing') + startStatusPolling() + startDetailPolling() + return + } + + phase.value = 2 + emit('update-status', data.runner_status === 'failed' ? 'error' : 'completed') + if (data.error) { + addLog(t('log.startFailed', { error: data.error })) + } +} + +const initializeSimulationView = async () => { + if (!props.simulationId) { + addLog(t('log.errorMissingSimId')) + return + } + + resetAllState() + + try { + const res = await getRunStatus(props.simulationId) + if (res.success && hasExistingRunData(res.data)) { + await restoreExistingRun(res.data) + return + } + } catch (err) { + console.warn('恢复模拟运行状态失败:', err) + } + + await doStartSimulation({ reset: false }) +} + // 启动模拟 -const doStartSimulation = async () => { +const doStartSimulation = async ({ reset = true } = {}) => { if (!props.simulationId) { addLog(t('log.errorMissingSimId')) return } // 先重置所有状态,确保不会受到上一次模拟的影响 - resetAllState() + if (reset) resetAllState() isStarting.value = true startError.value = null @@ -398,7 +477,7 @@ const doStartSimulation = async () => { const params = { simulation_id: props.simulationId, platform: 'parallel', - force: true, // 强制重新开始 + force: false, enable_graph_memory_update: true // 开启动态图谱更新 } @@ -598,8 +677,10 @@ const getActionTypeLabel = (type) => { 'CREATE_POST': 'POST', 'REPOST': 'REPOST', 'LIKE_POST': 'LIKE', + 'DISLIKE_POST': 'DISLIKE', 'CREATE_COMMENT': 'COMMENT', 'LIKE_COMMENT': 'LIKE', + 'DISLIKE_COMMENT': 'DISLIKE', 'DO_NOTHING': 'IDLE', 'FOLLOW': 'FOLLOW', 'SEARCH_POSTS': 'SEARCH', @@ -615,8 +696,10 @@ const getActionTypeClass = (type) => { 'CREATE_POST': 'badge-post', 'REPOST': 'badge-action', 'LIKE_POST': 'badge-action', + 'DISLIKE_POST': 'badge-action', 'CREATE_COMMENT': 'badge-comment', 'LIKE_COMMENT': 'badge-action', + 'DISLIKE_COMMENT': 'badge-action', 'QUOTE_POST': 'badge-post', 'FOLLOW': 'badge-meta', 'SEARCH_POSTS': 'badge-meta', @@ -691,7 +774,7 @@ watch(() => props.systemLogs?.length, () => { onMounted(() => { addLog(t('log.step3Init')) if (props.simulationId) { - doStartSimulation() + initializeSimulationView() } }) @@ -1125,7 +1208,7 @@ onUnmounted(() => { } /* Info Blocks (Quote, Repost, etc) */ -.quoted-block, .repost-content { +.quoted-block, .repost-content, .liked-content, .voted-content { background: #F9F9F9; border: 1px solid #EEE; padding: 10px 12px;