#!/usr/bin/env python3
"""
千问AI代码评审系统 - 异步分批评审版本
支持大MR分多批处理，避免超时
"""

import os
import sys
import time
import functools

# 检查必要的依赖
try:
    import requests
    import json
    from datetime import datetime
except ImportError as e:
    print(f"❌ 缺少必要的Python库: {e}", flush=True)
    print("请运行: pip3 install requests", flush=True)
    sys.exit(1)

# 强制无缓冲输出，确保CI中日志实时可见
print = functools.partial(print, flush=True)


class QwenCodeReviewerAsync:
    def __init__(self):
        # 诊断环境
        self.mr_iid = os.getenv('CI_MERGE_REQUEST_IID')
        self.project_id = os.getenv('CI_PROJECT_ID')
        self.gitlab_url = os.getenv('CI_SERVER_URL', 'http://192.168.22.227')
        self.gitlab_token = os.getenv('AI_GITLAB_TOKEN')
        self.pipeline_id = os.getenv('CI_PIPELINE_ID', 'unknown')

        # 千问API配置
        self.api_url = "http://192.168.22.230:8000/v1/chat/completions"
        self.model_name = "Qwen3.5-35B-A3B-FP8"

        # 分批配置
        self.batch_size = int(os.getenv('REVIEW_BATCH_SIZE', '3'))  # 每批文件数
        self.max_total_files = int(os.getenv('MAX_REVIEW_FILES', '50'))  # 最多评审文件数
        self.max_diff_per_file = int(os.getenv('MAX_DIFF_PER_FILE', '8000'))  # 单文件diff限制

        print(f"🔍 环境检测: MR#{self.mr_iid}, 项目#{self.project_id}")
        print(f"⚙️  分批配置: 每批{self.batch_size}个文件, 最多{self.max_total_files}个文件")

    def get_mr_changes(self):
        """获取MR的真实代码差异"""
        if not all([self.mr_iid, self.project_id, self.gitlab_token]):
            print("❌ 缺少必要的MR环境变量或Token")
            return None

        print(f"📥 获取MR #{self.mr_iid} 的代码变更...")
        url = f"{self.gitlab_url}/api/v4/projects/{self.project_id}/merge_requests/{self.mr_iid}/changes"
        headers = {"PRIVATE-TOKEN": self.gitlab_token}
        print(f"   请求地址: {url}")

        try:
            print(f"   正在连接GitLab API...", end=" ")
            response = requests.get(url, headers=headers, timeout=(10, 30))
            print("连接成功")
            response.raise_for_status()
            data = response.json()
            changes = data.get('changes', [])

            if not changes:
                print("⚠️  MR中没有代码变更")
                return []

            print(f"✅ 获取到 {len(changes)} 个文件变更")
            return changes
        except requests.exceptions.ConnectTimeout:
            print(f"❌ 连接GitLab超时(10秒)，请检查网络: {self.gitlab_url}")
            return None
        except requests.exceptions.ReadTimeout:
            print(f"❌ 读取GitLab响应超时(30秒)")
            return None
        except requests.exceptions.ConnectionError as e:
            print(f"❌ 无法连接GitLab: {e}")
            return None
        except Exception as e:
            print(f"❌ 获取MR变更失败: {e}")
            return None

    def prioritize_files(self, changes):
        """按优先级排序文件，核心代码优先，跳过不需要评审的文件"""
        priority_map = {
            '.c': 1, '.h': 1, '.cpp': 1,  # 核心代码最高优先级
            '.py': 2, '.sh': 2, '.js': 2,  # 脚本次之
            '.json': 3,  # 配置
        }

        # 跳过不需要评审的文件类型
        skip_exts = {'.yml', '.yaml', '.md', '.txt', '.rst', '.log'}

        def get_priority(change):
            path = change.get('new_path') or change.get('old_path', '')
            ext = os.path.splitext(path)[1].lower()
            if ext in skip_exts:
                return 99  # 跳过的文件给最低优先级
            return priority_map.get(ext, 2)  # 默认优先级2

        # 过滤并排序（排除优先级99的文件）
        sorted_changes = [c for c in sorted(changes, key=get_priority) if get_priority(c) != 99]

        # 限制总数
        limited_changes = sorted_changes[:self.max_total_files]

        # 打印文件列表
        print(f"\n📋 文件优先级排序（前{len(limited_changes)}个）:")
        for i, change in enumerate(limited_changes, 1):
            path = change.get('new_path') or change.get('old_path', '未知')
            priority = get_priority(change)
            priority_label = {1: '🔴核心', 2: '🟡脚本', 3: '🔵配置'}.get(priority, '?')
            print(f"   {i:2d}. {priority_label} {path}")

        skipped = len(changes) - len(limited_changes)
        if skipped > 0:
            print(f"   ... 还有 {skipped} 个文件未列入评审（超出限制）")

        return limited_changes

    def trim_diff_context(self, diff_text, context_lines=1):
        """精简diff的上下文行数，将默认3行压缩为指定行数
        
        git diff默认-U3，每处变更前后各3行上下文。
        压缩后减少token消耗，同时保留基本语境。
        """
        lines = diff_text.split('\n')
        result = []
        i = 0
        while i < len(lines):
            line = lines[i]
            # 保留非上下文行（hunk头、+行、-行、diff头等）
            if line.startswith('@@') or line.startswith('diff ') or \
               line.startswith('index ') or line.startswith('---') or \
               line.startswith('+++') or line.startswith('+') or \
               line.startswith('-') or line.startswith('\\'):
                result.append(line)
                i += 1
            else:
                # 上下文行（空格开头或空行）：收集连续的上下文块
                context_block_start = i
                while i < len(lines):
                    cl = lines[i]
                    if cl.startswith(('+', '-', '@@', 'diff ', 'index ', '---', '+++', '\\')):
                        break
                    i += 1
                context_block = lines[context_block_start:i]
                # 只保留首尾各context_lines行
                if len(context_block) <= context_lines * 2:
                    result.extend(context_block)
                else:
                    result.extend(context_block[:context_lines])
                    result.append(f'  ... ({len(context_block) - context_lines * 2}行上下文省略) ...')
                    result.extend(context_block[-context_lines:])
        return '\n'.join(result)

    def create_batch_prompt(self, batch_changes, batch_num, total_batches, project_context):
        """为一批文件创建评审提示词"""

        # 系统提示词：详细评审标准（放在system消息中，模型遵循度更高）
        system_prompt = """你是资深嵌入式代码架构师，按以下标准评审代码：

## 问题优先级
🔴 P0— 必须立即修复（崩溃/安全漏洞/数据损坏）
🟠 P1— 强烈建议修复（潜在bug/性能/并发问题）

## 评审检查清单（逐项检查，发现即报告）

### A. 内存安全（最高优先级）
- A1. 危险字符串函数: sprintf→snprintf, strcpy→strncpy, strcat→strncat, gets→fgets; 注意strncpy不保证null终止
- A2. 动态内存管理: malloc/calloc/realloc/strdup返回值必须检查NULL; 失败路径需释放已分配资源(goto统一清理); free后置NULL防double-free/use-after-free
- A3. 数组/指针越界: 下标边界校验; 外部输入(网络报文/配置)作索引必须有合法性校验; 检查off-by-one
- A4. 栈资源安全: 栈数组>1KB警告, >4KB必须改堆分配或static; 禁止VLA变长数组
- A5. 未初始化变量: 局部变量声明后必须赋值再使用; struct局部变量需memset清零; 注意条件分支中某路径可能未初始化

### B. 错误处理完整性
- B1. 返回值检查: fopen/fclose/socket/connect/bind/listen/accept/send/recv/pthread_create/pthread_mutex_lock/strdup/sqlite3_*/json_*/OpenSSL全部函数必须检查返回值
- B2. 资源泄漏: 所有退出路径(return/break/goto)都必须释放资源; 使用goto统一清理模式(Linux内核风格)
- B3. 错误日志: 每个error path至少一行LOG_ERROR/LOG_WARN; 日志必须含上下文(错误码/文件名/关键变量); 禁止空日志和静默吞错误

### C. 多线程安全（核心检查项）
- C1. 共享可变状态: 全局/static变量(非const)读写必须有同步保护; 函数内static局部变量在多线程下不安全; 回调中全局写入需加锁或消息队列
- C2. 锁使用正确性: early return必须unlock(用goto统一出口); 持锁期间禁止阻塞操作(sleep/网络IO); 缩小临界区; pthread_mutex_lock返回值也必须检查
- C3. 死锁检测: 检查ABBA死锁; signal/longjmp/pthread_cancel越过unlock需pthread_cleanup_push
- C4. 非线程安全函数替换: strtok→strtok_r, gmtime→gmtime_r, localtime→localtime_r, ctime→ctime_r, rand→rand_r, strerror→strerror_r, inet_ntoa→inet_ntop
- C5. 条件变量规范: pthread_cond_wait必须在while循环中防虚假唤醒; signal时持同一个mutex

### D. 安全漏洞
- D1. 命令注入: 禁止system/popen处理外部输入; 参数来自网络/配置/文件必须白名单校验
- D2. 整数溢出: malloc(a*b)乘法溢出检查; 有符号/无符号混用时负数变超大正数
- D3. 敏感信息硬编码: 扫描密码/密钥/token/IP硬编码, 应从配置或安全存储读取
- D4. 不安全类型转换: 64位系统指针截断用intptr_t; 有符号→无符号隐式转换需先检查符号

### E. 嵌入式特有问题
- E1. 信号处理: handler中只调用async-signal-safe函数; 用volatile sig_atomic_t通知主循环
- E2. 长时间阻塞: recv/read需设SO_RCVTIMEO超时; sleep影响响应
- E3. 资源限制: 单次分配>1MB需标注; 循环频繁malloc/free建议预分配

### F. 代码质量
- F1-F2. 函数>80行/圈复杂度>10/嵌套>4层→标注; 相似代码≥2次→提取公共函数
- F3. 魔法数字: 除0/1/-1外裸数字→#define/enum/const
- F4-F5. 命名规范 & 注释质量(公共API缺功能/参数/返回值注释)

### G. 项目特定检查
- G1. 设备操作安全: 设备ID作下标前校验范围; 状态机非法跳转检测
- G2. 协议解析鲁棒性: 报文字段长度≤包体; JSON字段缺失/类型错误处理; Modbus/DLT645/IEC104的CRC校验和长度合法性
- G3. 数据库操作安全: SQL参数化查询防注入; 失败回滚; 长事务持锁时长控制


## 输出格式

按以下结构输出评审报告（简洁精炼，避免啰嗦）：

📊 评审总结
（一两句话概括变更内容和评价）

🔴 P0 问题（必须修复）
问题N: 问题简述
位置: 文件#L行
影响: 一句话说明直接后果，客观描述不推测过远
原因: 一句话说明原因

```c
// 问题代码
原代码
```

```c
// 修复代码
修复代码
```

🟠 P1 问题（强烈建议）
（同上格式）

要求：
- 每个问题精简为：影响1句 + 原因1句 + 问题代码 + 修复代码
- 只报告P0/P1级别的问题
- 每个优先级最多5个问题
- 同一问题只报告一次，不要以不同优先级重复报告"""

        # 用户消息：项目上下文 + 待评审代码
        user_prompt = f"""---
项目上下文：
- 分支: {project_context.get('branch', 'Unknown')}
- 作者: {project_context.get('author', 'Unknown')}
- 批次: 第{batch_num}批/共{total_batches}批
- 本批文件数: {len(batch_changes)}
---

待评审代码变更：
"""

        total_diff_len = 0
        for i, change in enumerate(batch_changes, 1):
            file_path = change.get('new_path') or change.get('old_path', '未知文件')
            original_diff = change.get('diff', '')

            # 截断单文件diff
            if len(original_diff) > self.max_diff_per_file:
                diff = original_diff[:self.max_diff_per_file]
                diff += f"\n... (截断，原长度{len(original_diff)}) ..."
            else:
                diff = original_diff

            # 精简diff：将3行上下文压缩为1行，减少token消耗
            diff = self.trim_diff_context(diff, context_lines=1)

            total_diff_len += len(diff)

            user_prompt += f"""
## 文件{i}: `{file_path}`
```diff
{diff}
```
"""

        user_prompt += "\n请按上述标准进行专业评审："

        print(f"\n📝 第{batch_num}批提示词: system={len(system_prompt)}字符, user={len(user_prompt)}字符 (diff内容: {total_diff_len} 字符)")
        print(f"{'='*60}")
        print(f"[SYSTEM]\n{system_prompt}")
        print(f"{'='*60}")
        print(f"[USER]\n{user_prompt}")
        print(f"{'='*60}")
        return system_prompt, user_prompt

    def call_qwen_api(self, system_prompt, user_prompt, max_retries=2, base_timeout=120):
        """调用千问API（流式），带重试机制
        
        关闭思维链(chat_template_kwargs)后速度从~400s降到几秒。
        仍保留流式模式以确保稳定性。
        """
        headers = {"Content-Type": "application/json"}
        data = {
            "model": self.model_name,
            "messages": [
                {"role": "system", "content": system_prompt},
                {"role": "user", "content": user_prompt}
            ],
            "temperature": 0.2,
            "max_tokens": 4096,
            "stream": True,
            "chat_template_kwargs": {"enable_thinking": False}
        }

        for attempt in range(1, max_retries + 1):
            try:
                print(f"   🔄 API调用尝试 {attempt}/{max_retries} (流式, max_tokens=4096)...")
                t0 = time.time()

                response = requests.post(
                    self.api_url, headers=headers, json=data,
                    timeout=(10, base_timeout), stream=True
                )
                response.raise_for_status()

                # 流式收集
                reasoning_chunks = []
                content_chunks = []
                first_token_time = None
                last_progress = 0

                for line in response.iter_lines():
                    if not line:
                        continue
                    line = line.decode("utf-8")
                    if not line.startswith("data: "):
                        continue
                    data_str = line[6:]
                    if data_str == "[DONE]":
                        break

                    try:
                        chunk = json.loads(data_str)
                        delta = chunk["choices"][0].get("delta", {})

                        if first_token_time is None and (delta.get("content") or delta.get("reasoning") or delta.get("reasoning_content")):
                            first_token_time = time.time() - t0

                        if delta.get("reasoning"):
                            reasoning_chunks.append(delta["reasoning"])
                        if delta.get("reasoning_content"):
                            reasoning_chunks.append(delta["reasoning_content"])
                        if delta.get("content"):
                            content_chunks.append(delta["content"])
                            print(delta["content"], end="")  # 实时打印正文

                        # 每30秒输出一次进度
                        now = time.time()
                        if now - last_progress >= 30:
                            elapsed = now - t0
                            r_len = sum(len(c) for c in reasoning_chunks)
                            c_len = sum(len(c) for c in content_chunks)
                            print(f"\n   ⏳ 已接收 {elapsed:.0f}s: 思维链{r_len}字符, 正文{c_len}字符")
                            last_progress = now
                    except json.JSONDecodeError:
                        pass

                elapsed = time.time() - t0
                reasoning = "".join(reasoning_chunks)
                content = "".join(content_chunks)
                if content_chunks:
                    print()  # 流式打印结束换行

                # 优先用content，没有则用reasoning
                final_content = content or reasoning
                content_source = "content" if content else ("reasoning" if reasoning else "empty")

                print(f"   ⏱️  总耗时: {elapsed:.1f}s (首token: {first_token_time:.1f}s)" if first_token_time else f"   ⏱️  总耗时: {elapsed:.1f}s")
                print(f"   📊 思维链: {len(reasoning)}字符 | 正文: {len(content)}字符 → 来源: {content_source}")

                if not final_content:
                    print(f"   ⚠️  API返回空内容")
                    if attempt < max_retries:
                        wait_time = min(2 ** attempt, 10)
                        print(f"   ⏳ 等待 {wait_time} 秒后重试...")
                        time.sleep(wait_time)
                        continue
                    return "❌ API返回空内容"

                print(f"   ✅ API响应成功 ({len(final_content)} 字符, 来源: {content_source}, 耗时: {elapsed:.1f}s)")
                return final_content

            except requests.exceptions.ConnectTimeout:
                print(f"   ❌ 第 {attempt} 次连接超时(10s)，AI服务可能不可达")
                if attempt < max_retries:
                    wait_time = min(2 ** attempt, 10)
                    print(f"   ⏳ 等待 {wait_time} 秒后重试...")
                    time.sleep(wait_time)
                else:
                    return "❌ AI服务连接超时，请检查API服务是否正常运行"
            except requests.exceptions.ReadTimeout:
                print(f"   ⏱️  第 {attempt} 次读取超时")
                if attempt < max_retries:
                    wait_time = min(2 ** attempt, 10)
                    print(f"   ⏳ 等待 {wait_time} 秒后重试...")
                    time.sleep(wait_time)
                else:
                    return "❌ 本批评审超时，请稍后重试或联系管理员"
            except requests.exceptions.ConnectionError as e:
                print(f"   ❌ 连接失败: {e}")
                return f"❌ AI服务不可达: {e}"
            except Exception as e:
                return f"❌ API调用失败: {str(e)}"

    def post_pending_notice(self):
        """发布"评审中"通知"""
        if not all([self.mr_iid, self.project_id, self.gitlab_token]):
            return None

        url = f"{self.gitlab_url}/api/v4/projects/{self.project_id}/merge_requests/{self.mr_iid}/notes"
        headers = {
            "PRIVATE-TOKEN": self.gitlab_token,
            "Content-Type": "application/json"
        }

        notice = f"""## 🤖 千问AI代码评审进行中

⏳ **状态**: 正在分析代码变更...
📊 **配置**: 每批 {self.batch_size} 个文件，最多评审 {self.max_total_files} 个文件
⏱️ **预计**: 2-5 分钟完成（根据文件数量）

> 本评审由 AI 自动生成，请勿重复触发。
> 流水线: {self.pipeline_id}"""

        try:
            print("   正在发布评审中通知...", end=" ")
            response = requests.post(url, headers=headers, json={"body": notice}, timeout=(10, 30))
            response.raise_for_status()
            note_id = response.json().get('id')
            print(f"成功 (note_id: {note_id})")
            return note_id
        except Exception as e:
            print(f"失败: {e}")
            return None

    def update_review_result(self, note_id, review_content, batch_info):
        """更新评审结果（编辑原评论）"""
        if not all([self.mr_iid, self.project_id, self.gitlab_token, note_id]):
            # 如果没有note_id，创建新评论
            return self.post_new_review(review_content, batch_info)

        url = f"{self.gitlab_url}/api/v4/projects/{self.project_id}/merge_requests/{self.mr_iid}/notes/{note_id}"
        headers = {
            "PRIVATE-TOKEN": self.gitlab_token,
            "Content-Type": "application/json"
        }

        full_content = f"""## 🤖 千问AI代码评审报告 ({self.model_name})

{batch_info}

{review_content}

---
**评审信息**
- **MR**: !{self.mr_iid}
- **模型**: {self.model_name}
- **完成时间**: {datetime.now().strftime("%Y-%m-%d %H:%M:%S")}
- **触发**: GitLab CI/CD 流水线

> 💡 本评审由公司内部千问大模型自动生成，旨在辅助代码质量提升。请结合人工评审做出最终决策。"""

        try:
            response = requests.put(url, headers=headers, json={"body": full_content}, timeout=30)
            response.raise_for_status()
            print("✅ 已更新评审结果")
            return True
        except Exception as e:
            print(f"⚠️  更新评论失败: {e}，尝试创建新评论")
            return self.post_new_review(review_content, batch_info)

    def post_new_review(self, review_content, batch_info):
        """创建新的评审评论"""
        url = f"{self.gitlab_url}/api/v4/projects/{self.project_id}/merge_requests/{self.mr_iid}/notes"
        headers = {
            "PRIVATE-TOKEN": self.gitlab_token,
            "Content-Type": "application/json"
        }

        full_content = f"""## 🤖 千问AI代码评审报告 ({self.model_name})

{batch_info}

{review_content}

---
**评审信息**
- **MR**: !{self.mr_iid}
- **模型**: {self.model_name}
- **完成时间**: {datetime.now().strftime("%Y-%m-%d %H:%M:%S")}
- **触发**: GitLab CI/CD 流水线

> 💡 本评审由公司内部千问大模型自动生成，旨在辅助代码质量提升。"""

        try:
            response = requests.post(url, headers=headers, json={"body": full_content}, timeout=30)
            response.raise_for_status()
            print("✅ 已创建新评审评论")
            return True
        except Exception as e:
            print(f"❌ 发布评审失败: {e}")
            return False

    def check_api_available(self):
        """预检AI API是否可达，避免白等超时"""
        print(f"🔍 预检AI API连通性: {self.api_url}")
        try:
            response = requests.get(self.api_url.replace('/v1/chat/completions', '/v1/models'),
                                    timeout=(5, 10))
            if response.status_code == 200:
                models = response.json().get('data', [])
                model_ids = [m.get('id', '') for m in models]
                print(f"✅ API可达，可用模型: {model_ids[:5]}{'...' if len(model_ids) > 5 else ''}")
                if self.model_name not in model_ids:
                    print(f"⚠️  模型 {self.model_name} 不在可用列表中，可能调用失败")
            else:
                print(f"⚠️  API返回非200状态码: {response.status_code}")
        except requests.exceptions.ConnectTimeout:
            print(f"❌ AI API连接超时(5秒)，服务不可达: {self.api_url}")
            return False
        except requests.exceptions.ConnectionError as e:
            print(f"❌ AI API无法连接: {e}")
            return False
        except Exception as e:
            print(f"⚠️  模型列表查询异常: {e}")

        # 最小化推理测试 — 验证completions端点是否能响应
        print(f"🧪 最小化推理测试...", end=" ")
        try:
            test_data = {
                "model": self.model_name,
                "messages": [{"role": "user", "content": "Say OK"}],
                "max_tokens": 50,
                "temperature": 0,
                "chat_template_kwargs": {"enable_thinking": False}
            }
            t0 = time.time()
            test_resp = requests.post(self.api_url,
                                      headers={"Content-Type": "application/json"},
                                      json=test_data, timeout=(10, 60))
            elapsed = time.time() - t0
            if test_resp.status_code == 200:
                resp_json = test_resp.json()
                choices = resp_json.get('choices') or []
                if choices:
                    msg = choices[0].get('message', {})
                    content = msg.get('content') or ''
                    reasoning = msg.get('reasoning') or msg.get('reasoning_content') or ''
                    if content:
                        print(f"✅ 推理正常 (耗时{elapsed:.1f}s, 响应: {content[:30]})")
                    elif reasoning:
                        print(f"✅ 推理正常 (耗时{elapsed:.1f}s, 思维链模式)")
                    else:
                        # max_tokens可能不够思维链输出，但API确实响应了
                        print(f"✅ API响应正常 (耗时{elapsed:.1f}s, 思维链输出被截断)")
                    # 只要HTTP 200且有choices，就算API可用
                    return True
                else:
                    print(f"⚠️  API返回无choices: {json.dumps(resp_json, ensure_ascii=False)[:200]}")
                    return False
            else:
                print(f"❌ 推理失败 (HTTP {test_resp.status_code}): {test_resp.text[:200]}")
                return False
        except requests.exceptions.Timeout:
            print(f"❌ 推理超时(60s)，模型可能未加载或GPU资源不足")
            return False
        except Exception as e:
            print(f"⚠️  推理测试异常: {e}，仍将尝试调用")
            return True

    def run(self):
        """主执行流程 - 异步分批评审"""
        print("=" * 60)
        print("🚀 启动千问AI代码评审（异步分批模式）")
        print("=" * 60)

        # 0. 预检AI API可用性
        if not self.check_api_available():
            print("❌ AI API不可达，跳过评审（避免浪费CI时间）")
            return False

        # 1. 获取MR变更
        changes = self.get_mr_changes()
        if changes is None:
            print("❌ 无法获取MR变更，退出")
            return False
        if not changes:
            print("⚠️  MR中没有代码变更")
            return True

        # 2. 文件优先级排序
        prioritized_changes = self.prioritize_files(changes)

        # 3. 先发布"评审中"通知
        note_id = self.post_pending_notice()

        # 4. 分批处理
        total_files = len(prioritized_changes)
        batch_size = self.batch_size
        total_batches = (total_files + batch_size - 1) // batch_size

        all_reviews = []
        project_context = {
            'branch': os.getenv('CI_COMMIT_REF_NAME', 'Unknown'),
            'author': os.getenv('GITLAB_USER_NAME', 'Unknown'),
        }

        print(f"\n{'=' * 60}")
        print(f"📦 开始分批评审: 共 {total_files} 个文件，分 {total_batches} 批")
        print(f"{'=' * 60}")

        for batch_num in range(1, total_batches + 1):
            start_idx = (batch_num - 1) * batch_size
            end_idx = min(start_idx + batch_size, total_files)
            batch_changes = prioritized_changes[start_idx:end_idx]

            # 显示批次信息：共X批，当前第Y批，文件a到b
            print(f"\n📌 批次进度: 共{total_batches}批 | 当前第{batch_num}批 | 文件{start_idx+1}-{end_idx}/{total_files}")
            print("-" * 60)

            # 显示本批文件列表
            for i, change in enumerate(batch_changes, start_idx + 1):
                file_path = change.get('new_path') or change.get('old_path', '未知')
                print(f"   [{i:2d}] {os.path.basename(file_path)}")

            # 创建提示词
            system_prompt, user_prompt = self.create_batch_prompt(batch_changes, batch_num, total_batches, project_context)

            # 调用AI
            review_result = self.call_qwen_api(system_prompt, user_prompt)

            # 保存结果
            all_reviews.append({
                'batch_num': batch_num,
                'files': [c.get('new_path') or c.get('old_path') for c in batch_changes],
                'result': review_result
            })

            # 批次间短暂休息，避免压垮模型
            if batch_num < total_batches:
                print("   😴 批次间休息 2 秒...")
                time.sleep(2)

        # 5. 合并所有评审结果
        print(f"\n{'=' * 60}")
        print("📊 合并所有批次评审结果...")
        print(f"{'=' * 60}")

        merged_review = self.merge_reviews(all_reviews, total_files, len(changes))

        # 6. 发布最终结果
        batch_info = f"📦 评审范围: 共 {len(changes)} 个文件变更，详细评审了 {total_files} 个文件（分 {total_batches} 批）"
        success = self.update_review_result(note_id, merged_review, batch_info)

        if success:
            print("\n🎉 AI代码评审流程完成！")
        else:
            print("\n⚠️  评审完成但发布失败，请查看日志")

        return success

    def merge_reviews(self, all_reviews, reviewed_count, total_count):
        """合并多批评审结果"""
        merged = []

        # 批次详情
        merged.append("### 📋 分批评审详情\n")
        for review in all_reviews:
            file_list = ', '.join([os.path.basename(f) for f in review['files'][:3]])
            if len(review['files']) > 3:
                file_list += f" 等{len(review['files'])}个文件"
            merged.append(f"- **第{review['batch_num']}批**: {file_list}")
        merged.append("")

        # 合并各批次内容
        merged.append("---\n")
        for review in all_reviews:
            merged.append(f"\n### 第 {review['batch_num']} 批评审结果\n")
            merged.append(review['result'])
            merged.append("\n")

        # 统计信息
        skipped = total_count - reviewed_count
        if skipped > 0:
            merged.append(f"\n> ⚠️  因长度限制，{skipped} 个文件未详细评审")

        return '\n'.join(merged)


if __name__ == "__main__":
    reviewer = QwenCodeReviewerAsync()
    success = reviewer.run()
    sys.exit(0 if success else 1)
