diff --git a/backend/app/api/v1/datasources.py b/backend/app/api/v1/datasources.py index 22337662..1428e073 100644 --- a/backend/app/api/v1/datasources.py +++ b/backend/app/api/v1/datasources.py @@ -3,7 +3,7 @@ from datetime import datetime, timedelta, timezone from typing import Optional from fastapi import APIRouter, Depends, HTTPException, Query -from sqlalchemy import func, select +from sqlalchemy import func, select, text from sqlalchemy.ext.asyncio import AsyncSession from app.core.time import to_iso8601_utc @@ -11,10 +11,16 @@ from app.core.security import get_current_user from app.core.data_sources import get_data_sources_config from app.db.session import get_db from app.models.collected_data import CollectedData +from app.models.data_snapshot import DataSnapshot from app.models.datasource import DataSource from app.models.task import CollectionTask from app.models.user import User -from app.services.scheduler import get_latest_task_id_for_datasource, run_collector_now, sync_datasource_job +from app.services.scheduler import ( + cancel_running_collector_now, + get_latest_task_id_for_datasource, + run_collector_now, + sync_datasource_job, +) router = APIRouter() STALE_RUNNING_TASK_TIMEOUT_MINUTES = 90 @@ -100,6 +106,70 @@ async def get_running_task(db: AsyncSession, datasource_id: int) -> Optional[Col return None +async def rollback_orphaned_running_task( + db: AsyncSession, + datasource: DataSource, + running_task: CollectionTask, +) -> None: + snapshot_result = await db.execute( + select(DataSnapshot) + .where( + DataSnapshot.datasource_id == datasource.id, + DataSnapshot.task_id == running_task.id, + ) + .order_by(DataSnapshot.id.desc()) + .limit(1) + ) + snapshot = snapshot_result.scalar_one_or_none() + + await db.execute(CollectedData.__table__.delete().where(CollectedData.task_id == running_task.id)) + + await db.execute( + text( + """ + UPDATE collected_data + SET is_current = FALSE + WHERE source = :source + """ + ), + {"source": datasource.source}, + ) + + if snapshot is not None: + snapshot.status = "cancelled" + snapshot.is_current = False + snapshot.completed_at = datetime.now(timezone.utc) + summary = dict(snapshot.summary or {}) + summary["rollback"] = True + summary["rollback_reason"] = "orphaned_running_task_after_backend_restart" + snapshot.summary = summary + + if snapshot.parent_snapshot_id is not None: + parent_snapshot = await db.get(DataSnapshot, snapshot.parent_snapshot_id) + if parent_snapshot: + parent_snapshot.is_current = True + await db.execute( + text( + """ + UPDATE collected_data + SET is_current = TRUE + WHERE snapshot_id = :snapshot_id + """ + ), + {"snapshot_id": snapshot.parent_snapshot_id}, + ) + + running_task.status = "cancelled" + running_task.phase = "cancelled" + running_task.completed_at = datetime.now(timezone.utc) + existing_error = (running_task.error_message or "").strip() + cancel_reason = "Cancelled after backend restart because the running task handle was lost; incomplete writes rolled back" + running_task.error_message = f"{existing_error}\n{cancel_reason}".strip() if existing_error else cancel_reason + datasource.last_status = "cancelled" + datasource.last_run_at = datetime.now(timezone.utc) + await db.commit() + + @router.get("") async def list_datasources( module: Optional[str] = None, @@ -346,6 +416,7 @@ async def get_datasource_stats( @router.post("/{source_id}/trigger") async def trigger_datasource( source_id: str, + force: bool = Query(False), current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db), ): @@ -356,6 +427,26 @@ async def trigger_datasource( if not datasource.is_active: raise HTTPException(status_code=400, detail="Data source is disabled") + running_task = await get_running_task(db, datasource.id) + if running_task is not None and not force: + raise HTTPException( + status_code=409, + detail={ + "reason": "running_task_in_progress", + "message": "当前采集任务尚未完成,重新触发会丢失本次未完成进度。是否强制重新采集?", + "task_id": running_task.id, + "phase": running_task.phase, + "progress": running_task.progress, + "records_processed": running_task.records_processed, + "total_records": running_task.total_records, + }, + ) + + if running_task is not None and force: + cancelled = await cancel_running_collector_now(datasource.source) + if not cancelled: + await rollback_orphaned_running_task(db, datasource, running_task) + previous_task_id = await get_latest_task_id_for_datasource(datasource.id) success = run_collector_now(datasource.source) if not success: @@ -375,6 +466,7 @@ async def trigger_datasource( "source_id": datasource.id, "task_id": task_id, "collector_name": datasource.source, + "force": force, "message": f"Collector '{datasource.source}' has been triggered", } diff --git a/backend/app/services/collectors/base.py b/backend/app/services/collectors/base.py index 64453974..fa2faaac 100644 --- a/backend/app/services/collectors/base.py +++ b/backend/app/services/collectors/base.py @@ -1,5 +1,6 @@ """Base collector class for all data sources""" +import asyncio from abc import ABC, abstractmethod from typing import Dict, List, Any, Optional from datetime import UTC, datetime @@ -166,6 +167,59 @@ class BaseCollector(ABC): await db.commit() return snapshot.id + async def _rollback_incomplete_run( + self, + db: AsyncSession, + *, + task_id: int, + snapshot_id: Optional[int], + reason: str, + ) -> None: + from app.models.collected_data import CollectedData + from app.models.data_snapshot import DataSnapshot + + await db.execute(CollectedData.__table__.delete().where(CollectedData.task_id == task_id)) + + parent_snapshot_id: Optional[int] = None + if snapshot_id is not None: + snapshot = await db.get(DataSnapshot, snapshot_id) + if snapshot: + parent_snapshot_id = snapshot.parent_snapshot_id + snapshot.status = "cancelled" + snapshot.is_current = False + snapshot.completed_at = datetime.now(UTC) + summary = dict(snapshot.summary or {}) + summary["rollback"] = True + summary["rollback_reason"] = reason + snapshot.summary = summary + + await db.execute( + text( + """ + UPDATE collected_data + SET is_current = FALSE + WHERE source = :source + """ + ), + {"source": self.name}, + ) + + if parent_snapshot_id is not None: + parent_snapshot = await db.get(DataSnapshot, parent_snapshot_id) + if parent_snapshot: + parent_snapshot.is_current = True + + await db.execute( + text( + """ + UPDATE collected_data + SET is_current = TRUE + WHERE snapshot_id = :snapshot_id + """ + ), + {"snapshot_id": parent_snapshot_id}, + ) + async def run(self, db: AsyncSession) -> Dict[str, Any]: """Full pipeline: fetch -> transform -> save""" from app.services.collectors.registry import collector_registry @@ -227,6 +281,21 @@ class BaseCollector(ABC): "records_processed": records_count, "execution_time_seconds": (datetime.now(UTC) - start_time).total_seconds(), } + except asyncio.CancelledError: + task.status = "cancelled" + task.phase = "cancelled" + task.error_message = "Collection cancelled by operator and rolled back" + task.completed_at = datetime.now(UTC) + if snapshot_id is not None: + await self._rollback_incomplete_run( + db, + task_id=task_id, + snapshot_id=snapshot_id, + reason="cancelled_by_operator", + ) + await db.commit() + await self._publish_task_update(force=True) + raise except Exception as e: task.status = "failed" task.phase = "failed" @@ -276,20 +345,34 @@ class BaseCollector(ABC): updated_count = 0 unchanged_count = 0 seen_entity_keys: set[str] = set() - previous_current_keys: set[str] = set() + progress_commit_interval = 1000 previous_current_result = await db.execute( - select(CollectedData.entity_key).where( + select(CollectedData) + .where( CollectedData.source == self.name, CollectedData.is_current == True, ) + .order_by(CollectedData.entity_key.asc(), CollectedData.collected_at.desc().nullslast(), CollectedData.id.desc()) ) - previous_current_keys = {row[0] for row in previous_current_result.fetchall() if row[0]} + previous_current_records = previous_current_result.scalars().all() + previous_current_keys = {record.entity_key for record in previous_current_records if record.entity_key} + previous_current_map: dict[str, CollectedData] = {} + stale_previous_records: list[CollectedData] = [] + + for existing_record in previous_current_records: + entity_key = existing_record.entity_key + if not entity_key: + continue + if entity_key not in previous_current_map: + previous_current_map[entity_key] = existing_record + continue + stale_previous_records.append(existing_record) + + for stale_record in stale_previous_records: + stale_record.is_current = False for i, item in enumerate(data): - print( - f"DEBUG: Saving item {i}: name={item.get('name')}, metadata={item.get('metadata', 'NOT FOUND')}" - ) raw_metadata = item.get("metadata", {}) extra_data = build_dynamic_metadata( raw_metadata, @@ -318,20 +401,9 @@ class BaseCollector(ABC): previous_record = None if entity_key and entity_key not in seen_entity_keys: - result = await db.execute( - select(CollectedData) - .where( - CollectedData.source == self.name, - CollectedData.entity_key == entity_key, - CollectedData.is_current == True, - ) - .order_by(CollectedData.collected_at.desc().nullslast(), CollectedData.id.desc()) - ) - previous_records = result.scalars().all() - if previous_records: - previous_record = previous_records[0] - for old_record in previous_records: - old_record.is_current = False + previous_record = previous_current_map.get(entity_key) + if previous_record is not None: + previous_record.is_current = False record = CollectedData( snapshot_id=snapshot_id, @@ -375,7 +447,7 @@ class BaseCollector(ABC): seen_entity_keys.add(entity_key) records_added += 1 - if i % 100 == 0: + if (i + 1) % progress_commit_interval == 0: await self.update_progress(i + 1, commit=True) if snapshot_id is not None: diff --git a/backend/app/services/scheduler.py b/backend/app/services/scheduler.py index ed0e45dd..0baf24be 100644 --- a/backend/app/services/scheduler.py +++ b/backend/app/services/scheduler.py @@ -19,6 +19,30 @@ logger = logging.getLogger(__name__) scheduler = AsyncIOScheduler() RUNNING_TASK_GUARD_TIMEOUT_MINUTES = 90 +RUNNING_COLLECTOR_TASKS: dict[str, asyncio.Task[Any]] = {} + + +def _collector_task_name(collector_name: str) -> str: + return f"collector:{collector_name}" + + +def get_running_collector_task(collector_name: str) -> asyncio.Task[Any] | None: + task = RUNNING_COLLECTOR_TASKS.get(collector_name) + if task is not None and not task.done(): + return task + + if task is not None and task.done(): + RUNNING_COLLECTOR_TASKS.pop(collector_name, None) + + target_name = _collector_task_name(collector_name) + for candidate in asyncio.all_tasks(): + if candidate.done(): + continue + if candidate.get_name() == target_name: + RUNNING_COLLECTOR_TASKS[collector_name] = candidate + return candidate + + return None async def _update_next_run_at(datasource: DataSource, session) -> None: @@ -133,6 +157,12 @@ async def run_collector_task(collector_name: str): datasource.last_status = task_result.get("status") await _update_next_run_at(datasource, db) logger.info("Collector %s completed: %s", collector_name, task_result) + except asyncio.CancelledError: + datasource.last_run_at = datetime.now(UTC) + datasource.last_status = "cancelled" + await db.commit() + logger.warning("Collector %s cancelled by operator", collector_name) + raise except Exception as exc: datasource.last_run_at = datetime.now(UTC) datasource.last_status = "failed" @@ -245,9 +275,31 @@ def run_collector_now(collector_name: str) -> bool: return False try: - asyncio.create_task(run_collector_task(collector_name)) + task = asyncio.create_task(run_collector_task(collector_name), name=_collector_task_name(collector_name)) + RUNNING_COLLECTOR_TASKS[collector_name] = task + + def _cleanup_task(done_task: asyncio.Task[Any]) -> None: + current = RUNNING_COLLECTOR_TASKS.get(collector_name) + if current is done_task: + RUNNING_COLLECTOR_TASKS.pop(collector_name, None) + + task.add_done_callback(_cleanup_task) logger.info("Triggered collector: %s", collector_name) return True except Exception as exc: logger.error("Failed to trigger collector %s: %s", collector_name, exc) - return False + return False + + +async def cancel_running_collector_now(collector_name: str) -> bool: + task = get_running_collector_task(collector_name) + if task is None or task.done(): + RUNNING_COLLECTOR_TASKS.pop(collector_name, None) + return False + + task.cancel() + try: + await task + except asyncio.CancelledError: + return True + return task.cancelled() diff --git a/frontend/src/pages/DataSources/DataSources.tsx b/frontend/src/pages/DataSources/DataSources.tsx index 70b65042..b560b44d 100644 --- a/frontend/src/pages/DataSources/DataSources.tsx +++ b/frontend/src/pages/DataSources/DataSources.tsx @@ -1,6 +1,6 @@ import { useCallback, useEffect, useRef, useState } from 'react' import { - Table, Tag, Space, Button, Form, Input, Select, Progress, Checkbox, message, + Table, Tag, Space, Button, Form, Input, Select, Progress, Checkbox, message, Modal, Drawer, Tabs, Empty, Tooltip, Popconfirm, Collapse, InputNumber, Row, Col, Card } from 'antd' import { @@ -382,30 +382,92 @@ function DataSources() { return () => clearInterval(interval) }, [builtInSources, taskProgress, taskSocketConnected, fetchData]) + const triggerDatasource = async (id: number, options?: { force?: boolean }) => { + const force = options?.force ?? false + const res = await axios.post(`/api/v1/datasources/${id}/trigger`, null, { + params: { force }, + }) + + if (res.data.task_id) { + setTaskProgress(prev => ({ + ...prev, + [id]: { + task_id: res.data.task_id, + progress: 0, + is_running: true, + phase: 'queued', + status: 'running', + }, + })) + } else { + window.setTimeout(() => { + fetchData() + }, 800) + } + + fetchData() + return res + } + + const confirmForceTrigger = (id: number, taskInfo?: { + progress?: number | null + phase?: string | null + records_processed?: number | null + total_records?: number | null + message?: string + }) => { + const progressText = typeof taskInfo?.progress === 'number' ? `${Math.round(taskInfo.progress)}%` : '未知' + const phaseText = taskInfo?.phase || 'running' + const processedText = typeof taskInfo?.records_processed === 'number' + ? `${taskInfo.records_processed}${typeof taskInfo?.total_records === 'number' && taskInfo.total_records > 0 ? ` / ${taskInfo.total_records}` : ''}` + : '未知' + + Modal.confirm({ + title: '当前任务未完成', + content: ( +
+

{taskInfo?.message || '当前采集任务仍在运行,重新触发会丢失本次未完成进度。'}

+

当前阶段: {phaseText}

+

当前进度: {progressText}

+

已处理记录: {processedText}

+

确认后会强制取消当前采集,并回滚未完成写入,然后重新开始采集。

+
+ ), + okText: '强制重新采集', + cancelText: '取消', + okButtonProps: { danger: true }, + onOk: async () => { + try { + await triggerDatasource(id, { force: true }) + messageApi.success('已强制重新触发,未完成采集将回滚') + } catch (error: unknown) { + const err = error as { response?: { data?: { detail?: string | { message?: string } } } } + const detail = err.response?.data?.detail + messageApi.error(typeof detail === 'string' ? detail : detail?.message || '强制重新采集失败') + } + }, + }) + } + const handleTrigger = async (id: number) => { try { - const res = await axios.post(`/api/v1/datasources/${id}/trigger`) + await triggerDatasource(id) messageApi.success('任务已触发') - if (res.data.task_id) { - setTaskProgress(prev => ({ - ...prev, - [id]: { - task_id: res.data.task_id, - progress: 0, - is_running: true, - phase: 'queued', - status: 'running', - }, - })) - } else { - window.setTimeout(() => { - fetchData() - }, 800) - } - fetchData() } catch (error: unknown) { - const err = error as { response?: { data?: { detail?: string } } } - messageApi.error(err.response?.data?.detail || '触发失败') + const err = error as { response?: { status?: number; data?: { detail?: string | { + reason?: string + message?: string + progress?: number | null + phase?: string | null + records_processed?: number | null + total_records?: number | null + } } } } + const detail = err.response?.data?.detail + if (err.response?.status === 409 && typeof detail === 'object' && detail?.reason === 'running_task_in_progress') { + confirmForceTrigger(id, detail) + return + } + messageApi.error(typeof detail === 'string' ? detail : detail?.message || '触发失败') } } @@ -527,12 +589,24 @@ function DataSources() { const handleUpdateSource = async () => { if (!viewingSource) return try { - await axios.post(`/api/v1/datasources/${viewingSource.id}/trigger`) + await triggerDatasource(viewingSource.id) messageApi.success('已触发更新') setViewDrawerVisible(false) } catch (error: unknown) { - const err = error as { response?: { data?: { detail?: string } } } - messageApi.error(err.response?.data?.detail || '更新失败') + const err = error as { response?: { status?: number; data?: { detail?: string | { + reason?: string + message?: string + progress?: number | null + phase?: string | null + records_processed?: number | null + total_records?: number | null + } } } } + const detail = err.response?.data?.detail + if (err.response?.status === 409 && typeof detail === 'object' && detail?.reason === 'running_task_in_progress') { + confirmForceTrigger(viewingSource.id, detail) + return + } + messageApi.error(typeof detail === 'string' ? detail : detail?.message || '更新失败') } } diff --git a/planet.sh b/planet.sh index 2e2ae5bc..4f109889 100755 --- a/planet.sh +++ b/planet.sh @@ -5,11 +5,30 @@ set -e SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" cd "$SCRIPT_DIR" -RED='\033[0;31m' -GREEN='\033[0;32m' -YELLOW='\033[1;33m' -BLUE='\033[0;34m' +RED='\033[38;5;203m' +GREEN='\033[38;5;114m' +YELLOW='\033[38;5;221m' +BLUE='\033[38;5;111m' +CYAN='\033[38;5;117m' +WHITE='\033[1;37m' +DIM='\033[2m' NC='\033[0m' +WAIT_SPINNER_FRAMES=( + '⠉⠁' + '⠈⠉' + '⠀⠙' + '⠀⠸' + '⠀⢰' + '⠀⣠' + '⢀⣀' + '⣀⡀' + '⣄⠀' + '⡆⠀' + '⠇⠀' + '⠋⠀' +) +WAIT_SPINNER_STEP=0 +WAIT_SPINNER_TICKS_PER_SECOND=8 BACKEND_MAX_RETRIES="${BACKEND_MAX_RETRIES:-3}" BACKEND_HEALTH_CHECK_ATTEMPTS="${BACKEND_HEALTH_CHECK_ATTEMPTS:-10}" @@ -37,6 +56,100 @@ compose_up() { fi } +render_wait_spinner() { + local message="$1" + local step="${2:-1}" + local frame_index=$(( (step - 1) % ${#WAIT_SPINNER_FRAMES[@]} )) + local frame="${WAIT_SPINNER_FRAMES[$frame_index]}" + printf "\r${DIM}·${NC} ${CYAN}%s${NC} ${WHITE}%s${NC}\033[K" "$frame" "$message" +} + +animate_wait_spinner() { + local message="$1" + local duration="${2:-1}" + local ticks=$(( duration * WAIT_SPINNER_TICKS_PER_SECOND )) + local tick=0 + + if [ "$ticks" -lt 1 ]; then + ticks=1 + fi + + while [ "$tick" -lt "$ticks" ]; do + WAIT_SPINNER_STEP=$((WAIT_SPINNER_STEP + 1)) + render_wait_spinner "$message" "$WAIT_SPINNER_STEP" + sleep 0.125 + tick=$((tick + 1)) + done +} + +clear_wait_spinner() { + printf "\r\033[K" +} + +finish_wait_spinner() { + log_line "done" "$GREEN" "$1" +} + +log_line() { + local label="$1" + local color="$2" + local message="$3" + clear_wait_spinner + printf "${DIM}·${NC} ${color}%-4s${NC} ${WHITE}%s${NC}\n" "$label" "$message" +} + +log_step() { + log_line "run" "$BLUE" "$1" +} + +log_warn() { + log_line "warn" "$YELLOW" "$1" +} + +log_error() { + log_line "fail" "$RED" "$1" +} + +log_note() { + printf "${DIM} %s${NC}\n" "$1" +} + +log_success() { + log_line "done" "$GREEN" "$1" +} + +print_splash() { + clear_wait_spinner + printf "%b" "$CYAN" + cat <<'EOF' + ___ _ _ _ _ _ + |_ _|_ __ | |_ ___| | (_) __ _ ___ _ __ | |_ + | || '_ \| __/ _ \ | | |/ _` |/ _ \ '_ \| __| + | || | | | || __/ | | | (_| | __/ | | | |_ + |___|_| |_|\__\___|_|_|_|\__, |\___|_| |_|\__| + |___/ +EOF + printf "%b\n" "$NC" + printf "%b" "$BLUE" + cat <<'EOF' + ____ _ _ + | _ \| | __ _ _ __ ___ | |_ + | |_) | |/ _` | '_ \ / _ \| __| + | __/| | (_| | | | | __/| |_ + |_| |_|\__,_|_| |_|\___| \__| +EOF + printf "%b\n" "$NC" + printf "%b" "$WHITE" + cat <<'EOF' + ____ _ + | _ \| | __ _ _ __ + | |_) | |/ _` | '_ \ + | __/| | (_| | | | | + |_| |_|\__,_|_| |_| +EOF + printf "%b\n\n" "$NC" +} + run_with_retry() { local max_retries="$1" local interval="$2" @@ -52,12 +165,12 @@ run_with_retry() { fi if [ "$attempt" -eq "$max_retries" ]; then - echo -e "${RED}❌ ${failure_message}${NC}" + clear_wait_spinner + log_error "${failure_message}" return 1 fi - echo -e "${YELLOW}⚠️ ${retry_label} 第 ${attempt}/${max_retries} 次失败,${interval} 秒后重试...${NC}" - sleep "$interval" + animate_wait_spinner "${retry_label} 第 ${attempt}/${max_retries} 次失败,${interval} 秒后重试" "$interval" attempt=$((attempt + 1)) done @@ -65,17 +178,17 @@ run_with_retry() { } ensure_uv_backend_deps() { - echo -e "${BLUE}📦 检查后端 uv 环境...${NC}" + log_step "检查后端 uv 环境" if ! command -v uv >/dev/null 2>&1; then - echo -e "${RED}❌ 未找到 uv,请先安装 uv 并加入 PATH${NC}" + log_error "未找到 uv,请先安装 uv 并加入 PATH" exit 1 fi cd "$SCRIPT_DIR" if [ ! -x "$SCRIPT_DIR/.venv/bin/python" ]; then - echo -e "${YELLOW}⚠️ 未检测到 .venv,正在执行 uv sync...${NC}" + log_warn "未检测到 .venv,正在执行 uv sync" if ! run_with_retry \ "$DEPENDENCY_INSTALL_MAX_RETRIES" \ "$DEPENDENCY_INSTALL_RETRY_INTERVAL" \ @@ -87,23 +200,23 @@ ensure_uv_backend_deps() { fi if [ ! -x "$SCRIPT_DIR/.venv/bin/python" ]; then - echo -e "${RED}❌ uv 环境初始化失败,未找到 .venv/bin/python${NC}" + log_error "uv 环境初始化失败,未找到 .venv/bin/python" exit 1 fi } ensure_frontend_deps() { - echo -e "${BLUE}📦 检查前端依赖...${NC}" + log_step "检查前端依赖" if ! command -v bun >/dev/null 2>&1; then - echo -e "${RED}❌ 未找到 bun,请先安装或加载 bun 到 PATH${NC}" + log_error "未找到 bun,请先安装或加载 bun 到 PATH" exit 1 fi cd "$SCRIPT_DIR/frontend" if [ ! -x "$SCRIPT_DIR/frontend/node_modules/.bin/vite" ]; then - echo -e "${YELLOW}⚠️ 前端依赖缺失,正在执行 bun install...${NC}" + log_warn "前端依赖缺失,正在执行 bun install" if ! run_with_retry \ "$DEPENDENCY_INSTALL_MAX_RETRIES" \ "$DEPENDENCY_INSTALL_RETRY_INTERVAL" \ @@ -115,7 +228,7 @@ ensure_frontend_deps() { fi if [ ! -x "$SCRIPT_DIR/frontend/node_modules/.bin/vite" ]; then - echo -e "${RED}❌ 前端依赖安装失败,未找到 vite${NC}" + log_error "前端依赖安装失败,未找到 vite" exit 1 fi } @@ -129,14 +242,15 @@ wait_for_http() { while [ "$attempt" -le "$attempts" ]; do if curl -s "$url" > /dev/null 2>&1; then + finish_wait_spinner "${service_name}已就绪" return 0 fi - echo -e "${YELLOW}⏳ 等待${service_name}就绪...${NC}" - sleep "$interval" + animate_wait_spinner "等待${service_name}就绪" "$interval" attempt=$((attempt + 1)) done + clear_wait_spinner return 1 } @@ -152,14 +266,15 @@ wait_for_container_health() { status="$(docker inspect --format '{{if .State.Health}}{{.State.Health.Status}}{{else}}{{.State.Status}}{{end}}' "$container_name" 2>/dev/null || true)" if [ "$status" = "healthy" ] || [ "$status" = "running" ]; then + finish_wait_spinner "${service_name}容器健康检查已通过" return 0 fi - echo -e "${YELLOW}⏳ 等待${service_name}容器健康检查通过...${NC}" - sleep "$interval" + animate_wait_spinner "等待${service_name}容器健康检查通过" "$interval" attempt=$((attempt + 1)) done + clear_wait_spinner return 1 } @@ -173,15 +288,15 @@ wait_for_postgres_health() { } start_database_services() { - docker start planet_postgres planet_redis 2>/dev/null || compose_up up -d postgres redis + docker start planet_postgres planet_redis >/dev/null 2>&1 || compose_up up -d postgres redis >/dev/null 2>&1 } restart_database_services() { - docker restart planet_postgres planet_redis 2>/dev/null || compose_up up -d postgres redis + docker restart planet_postgres planet_redis >/dev/null 2>&1 || compose_up up -d postgres redis >/dev/null 2>&1 } start_postgres_service() { - docker start planet_postgres 2>/dev/null || compose_up up -d postgres + docker start planet_postgres >/dev/null 2>&1 || compose_up up -d postgres >/dev/null 2>&1 } start_backend_with_retry() { @@ -198,11 +313,12 @@ start_backend_with_retry() { return 0 fi - echo -e "${YELLOW}⚠️ 后端第 ${retry}/${BACKEND_MAX_RETRIES} 次启动未就绪,准备重试...${NC}" kill "$BACKEND_PID" 2>/dev/null || true + animate_wait_spinner "后端第 ${retry}/${BACKEND_MAX_RETRIES} 次启动未就绪,准备重试" "$BACKEND_HEALTH_CHECK_INTERVAL" retry=$((retry + 1)) done + clear_wait_spinner return 1 } @@ -210,10 +326,10 @@ start_ai_provider_service() { local ai_provider_port="${1:-$DEFAULT_AI_PROVIDER_PORT}" local retry=1 - echo -e "${BLUE}🧠 启动 AI Provider...${NC}" + log_step "启动 AI Provider" while [ "$retry" -le "$AI_PROVIDER_START_MAX_RETRIES" ]; do - if docker start planet_aiprovider 2>/dev/null || compose_up up -d aiprovider; then + if docker start planet_aiprovider >/dev/null 2>&1 || compose_up up -d aiprovider >/dev/null 2>&1; then if wait_for_container_health "planet_aiprovider" "$AI_PROVIDER_HEALTH_CHECK_ATTEMPTS" "$AI_PROVIDER_HEALTH_CHECK_INTERVAL" "AI Provider" && wait_for_http "http://localhost:${ai_provider_port}/health" "$AI_PROVIDER_HEALTH_CHECK_ATTEMPTS" "$AI_PROVIDER_HEALTH_CHECK_INTERVAL" "AI Provider"; then return 0 @@ -221,18 +337,19 @@ start_ai_provider_service() { fi if [ "$retry" -eq "$AI_PROVIDER_START_MAX_RETRIES" ]; then - echo -e "${RED}❌ AI Provider 启动失败,已重试 ${AI_PROVIDER_START_MAX_RETRIES} 次${NC}" + clear_wait_spinner + log_error "AI Provider 启动失败,已重试 ${AI_PROVIDER_START_MAX_RETRIES} 次" docker logs --tail 20 planet_aiprovider 2>/dev/null || true exit 1 fi - echo -e "${YELLOW}⚠️ AI Provider 第 ${retry}/${AI_PROVIDER_START_MAX_RETRIES} 次启动后仍未健康,正在重启容器...${NC}" - docker restart planet_aiprovider 2>/dev/null || true - sleep "$AI_PROVIDER_RETRY_INTERVAL" + docker restart planet_aiprovider >/dev/null 2>&1 || true + animate_wait_spinner "AI Provider 第 ${retry}/${AI_PROVIDER_START_MAX_RETRIES} 次启动后仍未健康,正在重启容器" "$AI_PROVIDER_RETRY_INTERVAL" retry=$((retry + 1)) done - echo -e "${RED}❌ AI Provider 启动失败${NC}" + clear_wait_spinner + log_error "AI Provider 启动失败" docker logs --tail 20 planet_aiprovider 2>/dev/null || true exit 1 } @@ -252,11 +369,12 @@ start_frontend_with_retry() { return 0 fi - echo -e "${YELLOW}⚠️ 前端第 ${retry}/${FRONTEND_MAX_RETRIES} 次启动未就绪,准备重试...${NC}" kill "$FRONTEND_PID" 2>/dev/null || true + animate_wait_spinner "前端第 ${retry}/${FRONTEND_MAX_RETRIES} 次启动未就绪,准备重试" "$FRONTEND_HEALTH_CHECK_INTERVAL" retry=$((retry + 1)) done + clear_wait_spinner return 1 } @@ -269,15 +387,15 @@ ensure_database_services_healthy() { fi if [ "$retry" -eq "$DATABASE_START_MAX_RETRIES" ]; then - echo -e "${RED}❌ 数据库启动失败,已重试 ${DATABASE_START_MAX_RETRIES} 次${NC}" + clear_wait_spinner + log_error "数据库启动失败,已重试 ${DATABASE_START_MAX_RETRIES} 次" docker logs --tail 20 planet_postgres 2>/dev/null || true docker logs --tail 20 planet_redis 2>/dev/null || true exit 1 fi - echo -e "${YELLOW}⚠️ 数据库第 ${retry}/${DATABASE_START_MAX_RETRIES} 次启动后仍未健康,正在重启容器...${NC}" - docker restart planet_postgres planet_redis 2>/dev/null || true - sleep "$DATABASE_RETRY_INTERVAL" + docker restart planet_postgres planet_redis >/dev/null 2>&1 || true + animate_wait_spinner "数据库第 ${retry}/${DATABASE_START_MAX_RETRIES} 次启动后仍未健康,正在重启容器" "$DATABASE_RETRY_INTERVAL" retry=$((retry + 1)) done } @@ -291,14 +409,14 @@ ensure_postgres_service_healthy() { fi if [ "$retry" -eq "$DATABASE_START_MAX_RETRIES" ]; then - echo -e "${RED}❌ PostgreSQL 启动失败,已重试 ${DATABASE_START_MAX_RETRIES} 次${NC}" + clear_wait_spinner + log_error "PostgreSQL 启动失败,已重试 ${DATABASE_START_MAX_RETRIES} 次" docker logs --tail 20 planet_postgres 2>/dev/null || true exit 1 fi - echo -e "${YELLOW}⚠️ PostgreSQL 第 ${retry}/${DATABASE_START_MAX_RETRIES} 次启动后仍未健康,正在重启容器...${NC}" - docker restart planet_postgres 2>/dev/null || true - sleep "$DATABASE_RETRY_INTERVAL" + docker restart planet_postgres >/dev/null 2>&1 || true + animate_wait_spinner "PostgreSQL 第 ${retry}/${DATABASE_START_MAX_RETRIES} 次启动后仍未健康,正在重启容器" "$DATABASE_RETRY_INTERVAL" retry=$((retry + 1)) done } @@ -306,7 +424,7 @@ ensure_postgres_service_healthy() { restart_database_service() { local retry=1 - echo -e "${BLUE}🗄️ 重启数据库...${NC}" + log_step "重启数据库" while [ "$retry" -le "$DATABASE_START_MAX_RETRIES" ]; do if restart_database_services && wait_for_database_health; then @@ -315,14 +433,14 @@ restart_database_service() { fi if [ "$retry" -eq "$DATABASE_START_MAX_RETRIES" ]; then - echo -e "${RED}❌ 数据库重启失败,已重试 ${DATABASE_START_MAX_RETRIES} 次${NC}" + clear_wait_spinner + log_error "数据库重启失败,已重试 ${DATABASE_START_MAX_RETRIES} 次" docker logs --tail 20 planet_postgres 2>/dev/null || true docker logs --tail 20 planet_redis 2>/dev/null || true exit 1 fi - echo -e "${YELLOW}⚠️ 数据库第 ${retry}/${DATABASE_START_MAX_RETRIES} 次重启后仍未健康,继续重试...${NC}" - sleep "$DATABASE_RETRY_INTERVAL" + animate_wait_spinner "数据库第 ${retry}/${DATABASE_START_MAX_RETRIES} 次重启后仍未健康,继续重试" "$DATABASE_RETRY_INTERVAL" retry=$((retry + 1)) done } @@ -332,7 +450,7 @@ start_backend_service() { local backend_port_requested="$2" local ai_provider_port="${3:-$DEFAULT_AI_PROVIDER_PORT}" - echo -e "${BLUE}🗄️ 启动数据库...${NC}" + log_step "启动数据库" ensure_database_services_healthy sleep 3 @@ -342,10 +460,10 @@ start_backend_service() { kill_port_if_requested "$backend_port" "后端" fi - echo -e "${BLUE}🔧 启动后端...${NC}" + log_step "启动后端" ensure_uv_backend_deps if ! start_backend_with_retry "$backend_port"; then - echo -e "${RED}❌ 后端启动失败,已重试 ${BACKEND_MAX_RETRIES} 次${NC}" + log_error "后端启动失败,已重试 ${BACKEND_MAX_RETRIES} 次" tail -10 /tmp/planet_backend.log exit 1 fi @@ -359,10 +477,10 @@ start_frontend_service() { kill_port_if_requested "$frontend_port" "前端" fi - echo -e "${BLUE}🌐 启动前端...${NC}" + log_step "启动前端" ensure_frontend_deps if ! start_frontend_with_retry "$frontend_port"; then - echo -e "${RED}❌ 前端启动失败,已重试 ${FRONTEND_MAX_RETRIES} 次${NC}" + log_error "前端启动失败,已重试 ${FRONTEND_MAX_RETRIES} 次" tail -10 /tmp/planet_frontend.log exit 1 fi @@ -372,7 +490,7 @@ validate_port() { local port="$1" if ! [[ "$port" =~ ^[0-9]+$ ]] || [ "$port" -lt 1 ] || [ "$port" -gt 65535 ]; then - echo -e "${RED}❌ 非法端口: ${port}${NC}" + log_error "非法端口: ${port}" exit 1 fi } @@ -381,10 +499,10 @@ kill_port_if_requested() { local port="$1" local service_name="$2" - echo -e "${YELLOW}🧹 检测 ${service_name} 端口 ${port} 占用...${NC}" + log_warn "检测 ${service_name} 端口 ${port} 占用" if command -v fuser >/dev/null 2>&1 && fuser "${port}/tcp" >/dev/null 2>&1; then - echo -e "${BLUE}🔌 发现端口 ${port} 占用,正在终止...${NC}" + log_step "发现端口 ${port} 占用,正在终止" fuser -k "${port}/tcp" >/dev/null 2>&1 || true sleep 1 return 0 @@ -394,14 +512,14 @@ kill_port_if_requested() { local pids pids="$(lsof -ti tcp:"${port}" 2>/dev/null || true)" if [ -n "$pids" ]; then - echo -e "${BLUE}🔌 发现端口 ${port} 占用,正在终止...${NC}" + log_step "发现端口 ${port} 占用,正在终止" kill $pids 2>/dev/null || true sleep 1 return 0 fi fi - echo -e "${GREEN}✅ 端口 ${port} 未被占用${NC}" + log_success "端口 ${port} 未被占用" } parse_service_args() { @@ -447,7 +565,7 @@ parse_service_args() { shift 1 ;; *) - echo -e "${RED}❌ 未知参数: $1${NC}" + log_error "未知参数: $1" exit 1 ;; esac @@ -462,9 +580,9 @@ cleanup_exit_containers() { local exit_containers exit_containers="$(docker ps -a --filter status=exited -q 2>/dev/null || true)" if [ -n "$exit_containers" ]; then - echo -e "${BLUE}🗑️ 清理残留 Exit 容器...${NC}" + log_step "清理残留 Exit 容器" echo "$exit_containers" | xargs -r docker rm -f >/dev/null 2>&1 || true - echo -e "${GREEN}✅ 残留容器已清理${NC}" + log_success "残留容器已清理" fi } @@ -492,15 +610,15 @@ create_user() { ensure_uv_backend_deps - echo -e "${BLUE}🗄️ 启动数据库...${NC}" + log_step "启动数据库" ensure_postgres_service_healthy sleep 2 - echo -e "${BLUE}👤 创建用户${NC}" + log_step "创建用户" read -r -p "用户名: " username if [ -z "$username" ]; then - echo -e "${RED}❌ 用户名不能为空${NC}" + log_error "用户名不能为空" exit 1 fi @@ -516,17 +634,17 @@ create_user() { echo "" if [ -z "$password" ]; then - echo -e "${RED}❌ 密码不能为空${NC}" + log_error "密码不能为空" exit 1 fi if [ "${#password}" -lt 8 ]; then - echo -e "${RED}❌ 密码长度不能少于 8 位${NC}" + log_error "密码长度不能少于 8 位" exit 1 fi if [ "$password" != "$password_confirm" ]; then - echo -e "${RED}❌ 两次输入的密码不一致${NC}" + log_error "两次输入的密码不一致" exit 1 fi @@ -580,35 +698,36 @@ async def main(): asyncio.run(main()) PY - echo -e "${GREEN}✅ 用户创建成功${NC}" - echo " 用户名: ${username}" - echo " 角色: ${role}" - echo " 邮箱: ${email}" + log_success "用户创建成功" + log_note "用户名: ${username}" + log_note "角色: ${role}" + log_note "邮箱: ${email}" } start() { parse_service_args "$@" cleanup_exit_containers - echo -e "${BLUE}🚀 启动智能星球计划...${NC}" + print_splash + log_step "启动智能星球计划" start_backend_service "$BACKEND_PORT" "$BACKEND_PORT_REQUESTED" "$AI_PROVIDER_PORT" start_frontend_service "$FRONTEND_PORT" "$FRONTEND_PORT_REQUESTED" echo "" - echo -e "${GREEN}✅ 启动完成!${NC}" - echo " 前端: http://localhost:${FRONTEND_PORT}" - echo " 后端: http://localhost:${BACKEND_PORT}" - echo " AI Provider: http://localhost:${AI_PROVIDER_PORT}" + log_success "启动完成" + log_note "前端: http://localhost:${FRONTEND_PORT}" + log_note "后端: http://localhost:${BACKEND_PORT}" + log_note "AI Provider: http://localhost:${AI_PROVIDER_PORT}" } stop() { - echo -e "${YELLOW}🛑 停止服务...${NC}" + log_warn "停止服务" stop_backend_service stop_ai_provider_service stop_frontend_service docker stop planet_postgres planet_redis 2>/dev/null || true - echo -e "${GREEN}✅ 已停止${NC}" + log_success "已停止" } restart() { @@ -622,7 +741,7 @@ restart() { return 0 fi - echo -e "${YELLOW}🔄 按需重启服务...${NC}" + log_warn "按需重启服务" if [ "$DATABASE_REQUESTED" -eq 1 ]; then restart_database_service @@ -647,43 +766,43 @@ restart() { fi echo "" - echo -e "${GREEN}✅ 重启完成!${NC}" + log_success "重启完成" if [ "$DATABASE_REQUESTED" -eq 1 ]; then - echo " 数据库: planet_postgres, planet_redis" + log_note "数据库: planet_postgres, planet_redis" fi if [ "$AI_PROVIDER_REQUESTED" -eq 1 ]; then - echo " AI Provider: http://localhost:${AI_PROVIDER_PORT}" + log_note "AI Provider: http://localhost:${AI_PROVIDER_PORT}" fi if [ "$BACKEND_PORT_REQUESTED" -eq 1 ]; then - echo " 后端: http://localhost:${BACKEND_PORT}" + log_note "后端: http://localhost:${BACKEND_PORT}" fi if [ "$FRONTEND_PORT_REQUESTED" -eq 1 ]; then - echo " 前端: http://localhost:${FRONTEND_PORT}" + log_note "前端: http://localhost:${FRONTEND_PORT}" fi } health() { - echo "📊 容器状态:" + echo -e "${DIM}·${NC} ${BLUE}view${NC} ${WHITE}容器状态${NC}" docker ps --filter "name=planet_" --format "table {{.Names}}\t{{.Status}}\t{{.Ports}}" echo "" - echo "🔍 服务状态:" + echo -e "${DIM}·${NC} ${BLUE}view${NC} ${WHITE}服务状态${NC}" if curl -s http://localhost:8000/health > /dev/null 2>&1; then - echo -e " 后端: ${GREEN}✅ 运行中${NC}" + echo -e "${DIM} 后端:${NC} ${GREEN}online${NC}" else - echo -e " 后端: ${RED}❌ 未运行${NC}" + echo -e "${DIM} 后端:${NC} ${RED}offline${NC}" fi if curl -s "http://localhost:${DEFAULT_AI_PROVIDER_PORT}/health" > /dev/null 2>&1; then - echo -e " AI Provider: ${GREEN}✅ 运行中${NC}" + echo -e "${DIM} AI Provider:${NC} ${GREEN}online${NC}" else - echo -e " AI Provider: ${RED}❌ 未运行${NC}" + echo -e "${DIM} AI Provider:${NC} ${RED}offline${NC}" fi if curl -s http://localhost:3000 > /dev/null 2>&1; then - echo -e " 前端: ${GREEN}✅ 运行中${NC}" + echo -e "${DIM} 前端:${NC} ${GREEN}online${NC}" else - echo -e " 前端: ${RED}❌ 未运行${NC}" + echo -e "${DIM} 前端:${NC} ${RED}offline${NC}" fi } diff --git a/uv.lock b/uv.lock index 2604d004..0ff80823 100644 --- a/uv.lock +++ b/uv.lock @@ -475,7 +475,7 @@ wheels = [ [[package]] name = "planet" -version = "0.23.0" +version = "0.23.1" source = { virtual = "." } dependencies = [ { name = "aiofiles" },