Metadata-Version: 2.5
Name: aigc-pipeline
Version: 0.1.6
Summary: AIGC 接口批量驱动 + 飞书审核工作流（图片生成 → 飞书推送 → 审核补生成 → 视频）
Project-URL: Homepage, https://github.com/your-username/aigc-pipeline
Project-URL: Repository, https://github.com/your-username/aigc-pipeline
Project-URL: Issues, https://github.com/your-username/aigc-pipeline/issues
Author-email: xuchaohui <hfyhui@126.com>
License-Expression: MIT
License-File: LICENSE
Keywords: aigc,feishu,image-generation,video-generation,workflow
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Topic :: Multimedia :: Graphics
Classifier: Topic :: Utilities
Requires-Python: >=3.10
Requires-Dist: httpx>=0.27
Provides-Extra: dev
Requires-Dist: build>=1.0; extra == 'dev'
Requires-Dist: pytest>=7; extra == 'dev'
Description-Content-Type: text/markdown

# aigc-pipeline

> AIGC 接口批量驱动 + 飞书审核工作流（图片生成 → 飞书推送 → 审核补生成 → 视频）

按用户的 8 步流程串联：

1. 读 `generation_tasks.json`
2. 按 `image_prompt_b` + `all_reference_images_d` 生成图片
3. 按命名规则保存到本地（`{片段名}-{日期}-{署名}/{片段名}-{日期}-（N）.png`）
4. 推送到飞书群 `@` 审核员（interactive 卡片 + 文件名列表）
5. 审核员把不合格图片移到 `<out_dir>/rejected/` 子目录
6. 程序检测 `rejected/` 触发补生成（最多 N 轮），重推全量
7. 通过图片 + `video_prompt_f` 生成视频（一图一视频）
8. 视频按命名规则存储 + 推送飞书

---

## 安装

### 方式一：从源码开发安装

```bash
git clone <repo>
cd aigc-pipeline
pip install -e .
```

### 方式二：从 whl 安装

```bash
pip install aigc_pipeline-0.1.0-py3-none-any.whl
```

---

## 快速开始

### 1. 安装

```bash
# 方式一：whl 安装（推荐集成方）
pip install aigc_pipeline-0.1.4-py3-none-any.whl

# 方式二：源码开发安装（推荐开发方）
git clone <repo> && cd aigc-pipeline
pip install -e ".[dev]"

# 方式三：PyPI（发布后）
pip install aigc-pipeline
```

### 2. 准备 `config.toml`

`config.toml` 必须与 `users.yaml` 放在同目录（默认是工作目录）：

```toml
[server]
base_url = "http://xx.xx.xx.xx"     # AIGC 服务地址

[auth]
username = "D-39JLlay"
password = "39JLlay"

[paths]
output_dir = "./outputs"                # 产物目录
db_path = "./aigc_pipeline.db"          # SQLite 状态库
token_file = "./token.txt"

[task]
mode = "image"                          # image | video
model_image = "auto-image"
model_video = "auto-video"
aspect_ratio = "1:1"

[feishu]
webhook_url = "https://open.feishu.cn/open-apis/bot/v2/hook/xxx"
reviewer_user_ids = ["ou_xxx", "ou_yyy"]
dry_run = true                           # true=只 log 不真发

[review]
review_timeout_hours = 72                # 超时→abandoned
max_retry = 3

[baidu]
access_token = ""                       # 百度网盘 OAuth（接入时填）
share_expires_days = 7
```

```bash
# users.yaml（账号→中文署名映射，可选）
"D-39JLlay": "lay"
"D-39haitong": "haitong"
```

### 3. 三种使用方式

#### 3.1 CLI 跑批

```bash
# 跑完整流程（图片+飞书审核+视频+推送），默认顺序模式
aigc-pipeline --tasks-json ./tasks.json

# 并发模式（方案 B：多进程 × 线程）
aigc-pipeline --parallel --max-processes 8 --tasks-json ./tasks.json

# 指定输出目录 / dry-run
aigc-pipeline --tasks-json ./tasks.json --output-dir ./my_out
aigc-pipeline --tasks-json ./tasks.json --skip-feishu    # 干跑，不真发飞书

# 跑单张图（ad-hoc）
aigc-pipeline --file 街头小店.jpg
```

#### 3.2 Python facade（嵌入业务代码，推荐）

```python
from aigc_pipeline import AIGCPipeline

pipeline = AIGCPipeline.from_config("./config.toml")

# 跑批（指定 tasks_json + output_dir）
summary = pipeline.run_from_json(
    json_path="./tasks.json",
    output_dir="./my_outputs",
    parallel=True,           # 进程并发
    max_processes=8,
)

# ad-hoc：单张图 / 单个视频
paths = pipeline.generate_image(
    prompt="让门店表现廉价感",
    reference_paths=["./ref1.png", "./ref2.png"],
    quantity=3,
    segment_name="玉湖公园-1",
)
video = pipeline.generate_video(
    image_path="./output/玉湖公园-1-8.23-（1）.png",
    prompt="门店表现出城市烟火气",
)
```

#### 3.3 并发模式（方案 B）

```bash
# 默认进程数 = min(len(tasks), cpu_count)，硬上限 16
aigc-pipeline --parallel --tasks-json ./tasks.json
aigc-pipeline --parallel --max-processes 8 --tasks-json ./tasks.json
```

| task 量 | 推荐模式 |
|---|---|
| < 3,000 / 天 | 顺序模式（默认） |
| 3,000 ~ 30,000 / 天 | **并发模式**（方案 B） |
| 30,000 ~ 100,000 / 天 | 并发模式 + 多机 |
| > 100,000 / 天 | 升级到 [Dramatiq + RabbitMQ](docs/ARCHITECTURE.md) |

### 4. 审批工作流（重点）

整套工作流核心是 **state 驱动 + 外部程序审阅**：

```
 run_from_json()              ← 你跑
   ↓
 image pending_review         ← 状态机推到这里
   ↓
 get_pending_reviews()        ← 外部程序调
   ↓
 submit_review_by_token()    ← 客户审过
   ↓
 process_approved_images()   ← 视频生成
   ↓
 video pending_review
   ↓
 submit_review_by_token()    ← 客户审视频
   ↓
 process_approved_videos()   ← 百度网盘 + 分享链接
   ↓
 image / video / task → completed
```

**完整示例**（外部业务系统）：

```python
from aigc_pipeline import AIGCPipeline

pipeline = AIGCPipeline.from_config("./config.toml")

# Step 1: 跑批（生成 image + 写 state）
pipeline.run_from_json("./tasks.json", parallel=True, max_processes=8)

# Step 2: 定时轮询拉取待审批
pending = pipeline.get_pending_reviews()
# pending = [
#   {"target_type": "image", "target_id": "img_001",
#    "record_id": "rec_xxx", "review_token": "abcdef...",
#    "file_path": ".../玉湖公园-1-8.23-（1）.png",
#    "retry_count": 0, "image_id": None, "source_image_path": None},
#   {"target_type": "video", ...},
# ]

# Step 3: 把待审批列表推给客户（邮件/IM/自定义 UI）
# (这里你自己接 UI，比如返回前端页面 / 调 IM API)
for p in pending:
    notify_client(p)

# Step 4: 客户审完，外部程序把决策回传
result = pipeline.submit_review_by_token(
    review_token=pending[0]["review_token"],
    decision="approved",       # 或 "rejected"
    comment="构图可以",
)
# result: {image_id, decision, new_status, retry_count, next_action}

# Step 5: 触发视频生成（approved image → video pending_review）
pipeline.process_approved_images()

# Step 6: 视频审批通过后调百度网盘
result = pipeline.process_approved_videos()
# 空壳（NotImplementedError）→ skipped；接入后 → shared + share_url

# Step 7: 定时巡查超时（推荐每 10 分钟跑一次）
abandoned = pipeline.run_dispatcher()
# {abandoned_images, abandoned_videos, count}

# Step 8: 查询全状态
status = pipeline.get_task_status("rec_xxx")
# {task, images[{...image, videos}]}
```

**客户端只要看到 token 就能审批**（无需知道 image_id / video_id）：

```python
# 极简：客户点 "通过" 后调一行
pipeline.submit_review_by_token(
    review_token="abcdef1234...",  # 客户系统从 URL/邮件中拿到
    decision="approved",
)
```

**决策路由规则**（已固定）：

| 决策 | 后续动作 |
|---|---|
| image approved | 触发 video 生成（`next_action="video_generation"`） |
| image rejected (retry < 3) | 回到 generating 重生（`next_action="regenerate"`） |
| image rejected (retry ≥ 3) | 标 abandoned（`next_action="abandoned"`） |
| video approved | 触发百度网盘（`next_action="baidu_upload"`） |
| video rejected (retry < 3) | 回到 generating，**用同 image 重生** |
| video rejected (retry ≥ 3) | 标 abandoned |
| pending 超时（72h 默认） | dispatcher 巡查 → abandoned |

### 5. 配置文件 + 文档

- **配置项大全**：看 `config.toml`（带详细注释）
- **架构设计**：看 [docs/ARCHITECTURE.md](docs/ARCHITECTURE.md)（10万/天的演进路径）
- **CLI 参数**：见下表
- **公开 API**：见 [aigc_pipeline/__init__.py](src/aigc_pipeline/__init__.py) 和各模块文档


## CLI 参数

| 参数 | 作用 |
|---|---|
| `--config PATH` | 配置文件路径（TOML），默认 `./config.toml` |
| `--tasks-json PATH` | 批量任务 JSON 路径，覆盖 config |
| `--skip-video` | 跳过视频生成阶段 |
| `--skip-feishu` | 不实际推送飞书（dry_run=true） |
| `--max-rounds N` | 最多几轮“审核-补生成”循环（默认 5） |
| `--wait-seconds N` | 推送后阻塞 N 秒让审核员操作（默认 0） |
| `--record-id ID` | 只跑指定 `record_id` 的 task |
| `--mode image\|video` | 覆盖 `cfg.task.mode` |
| `--quantity-mode loop\|single` | 覆盖 `cfg.batch.quantity_mode` |
| `--parallel` | **启用并发模式（方案 B）**：多进程跑 tasks + 内置 ThreadPool(8) 并发 jobs |
| `--max-processes N` | 并发进程数上限（默认 = min(len(tasks), cpu_count)，硬上限 16） |

---

## 对外集成：脱敏策略

本包采用 **脱敏 facade** 设计：集成方调 `AIGCPipeline` 看不到底层细节。

| 暴露面 | 状态 |
|---|---|
| `from aigc_pipeline import AIGCPipeline` | ✅ 默认导出 |
| `AIGCPipeline.run_from_json(...)` | ✅ 公开方法 |
| `AIGCPipeline.generate_image(...)` | ✅ 公开方法 |
| `AIGCPipeline.output_dir` 属性 | ✅ 可读可写 |
| `AIGCPipeline.tasks_json_path` 属性 | ✅ 可读可写 |
| `base_url` / `cfg` / `httpx.Client` | ❌ 不暴露 |
| `core` / `workflow` / `feishu` 等内部模块 | 需明确 import 才能用 |

**日志脱敏（默认开启）**：
```
原始 URL: POST http://14.103.124.30/api/auth/login
脱敏后:   POST /api/auth/login              ← host 完全消失
Authorization: Bearer abc123def
       → Bearer ********                  ← token 被 mask
```

脱敏**仅限日志输出**，实际发到服务端的 URL 仍是完整地址（保证功能）。

高级用户如需直接调底层（明确知道自己在做什么）：

```python
from aigc_pipeline.config import AppConfig, load_config
from aigc_pipeline.core import run, load_tasks, make_logged_client
from aigc_pipeline.feishu import FeishuBot
from aigc_pipeline.workflow import WorkflowContext, run_full_workflow
```

---

## 典型场景

### 1. 第一次接入：dry-run 测通流程

```bash
aigc-pipeline --tasks-json /path/to/generation_tasks.json --skip-feishu
```

不真发飞书，跑通图片+审核+视频，看日志和本地产物。

### 2. 给审核员时间标记不合格图

```bash
aigc-pipeline --tasks-json /path/to/generation_tasks.json --wait-seconds 60
```

推送后阻塞 60 秒，期间审核员把不合格图移到 `<output_dir>/<片段名>/rejected/`。程序自动检测并补生成。

### 3. 只生成图片，不生成视频

```bash
aigc-pipeline --skip-video
```

用于审图阶段快速迭代。

### 4. 切换视频模式

在 `config.toml [task]` 里设 `mode = "video"`，或 CLI 一次性覆盖：

```bash
aigc-pipeline --mode video
```

### 5. 调试单条 task

```bash
aigc-pipeline --record-id recvsxt4rtJ4Wd --skip-feishu
```

只跑指定 record_id，节省时间。

### 6. 嵌入你的业务脚本

```python
import json
from pathlib import Path
from aigc_pipeline import AppConfig, load_config, load_tasks, run
from aigc_pipeline.workflow import run_full_workflow, WorkflowContext
from aigc_pipeline.feishu import FeishuBot
from aigc_pipeline.naming import today_str

cfg = load_config(Path("./config.toml"))
bot = FeishuBot(
    webhook_url=cfg.feishu.webhook_url,
    reviewer_user_ids=list(cfg.feishu.reviewer_user_ids),
    dry_run=cfg.feishu.dry_run,
)
ctx = WorkflowContext(
    cfg=cfg,
    bot=bot,
    user_mapping=dict(cfg.users),
    output_root=Path(cfg.paths.output_dir),
    date_str=today_str(),
)
summary = run_full_workflow(ctx, Path("./your_tasks.json"))
with open("summary.json", "w", encoding="utf-8") as f:
    json.dump(summary, f, ensure_ascii=False, indent=2)
```

### 7. 并发批量跑（高吞吐场景）

```bash
# 启动 8 worker 进程跑一批 task
aigc-pipeline --parallel --max-processes 8 --tasks-json /path/to/big_tasks.json

# 跳过飞书推送（防止误推）
aigc-pipeline --parallel --skip-feishu --tasks-json /path/to/tasks.json

# 指定输出目录：CLI 参数或环境变量（当前版本需改 config.toml [paths].output_dir）
```

Python 嵌入：

```python
from pathlib import Path
from aigc_pipeline import parallel, core

cfg = core.load_config(Path("./config.toml"))
tasks = core.load_tasks(Path("./tasks.json"))

# 修改输出目录（运行时覆盖）
cfg.paths.output_dir = "/data/aigc_outputs"

summaries = parallel.run_tasks_parallel(
    tasks=tasks,
    cfg=cfg,
    config_path=Path("./config.toml"),
    max_processes=8,
)

# 汇总
ok = sum(1 for s in summaries if s.get("run_status", "ok") == "ok")
err = sum(1 for s in summaries if s.get("run_status") == "error")
print(f"ok={ok}  err={err}")
```

### 8. 指定输出目录的三种方式

| 场景 | 做法 |
|---|---|
| **临时改单次** | CLI：当前需改 `config.toml`；Python：`cfg.paths.output_dir = "/path"` |
| **永久改** | 改 `config.toml [paths].output_dir` |
| **AIGCPipeline facade** | `pipeline.output_dir = "/path"` (setter 已实现) |

> 注意：当前 `--parallel` CLI 不接 `--output-dir` 参数（如需可后续加）。

---

## 配置文件（`config.toml`）

```toml
[server]
base_url = "http://14.103.124.30"

[auth]
username = "D-39JLlay"
password = "39JLlay"

[task]
mode = "image"            # image | video
prompt = "让门店表现廉价感"
model_image = "auto-image"
model_video = "auto-video"

[batch]
tasks_json_path = "/path/to/generation_tasks.json"
quantity_mode = "loop"    # loop | single

[feishu]
webhook_url = "https://open.feishu.cn/open-apis/bot/v2/hook/xxx"
reviewer_user_ids = ["ou_xxx", "ou_yyy"]
dry_run = false           # true = 只打印 payload 不真发

[review]
rejected_subdir = "rejected"
wait_seconds = 0          # 推送后阻塞秒数
```

任何字段缺失会回退到代码默认值。详见 `config.toml`。

---

## 用户账号映射（`config.toml` 的 `[users]` section）

```toml
[users]
"D-39JLlay": "张三"
"D-99xxx": "李四"
```

把 `generation_tasks.json` 的 `assignee` 字段映射到中文署名（用于产物命名）。

**fallback 规则**：如果账号未在 `[users]` 中，自动提取账号字符串中的中文字符（如 `D-39姜` → `姜`）。

---

## 命名规则

| 类型 | 文件夹 | 文件 |
|---|---|---|
| 通用素材 | `{片段名}-{日期}-{署名}` | `{片段名}-{日期}-（N）.{ext}` |
| 门店素材 | `{片段名}-{门店编号}-{日期}-{署名}` | `{片段名}-{门店编号}-{日期}-（N）.{ext}` |
| 视频目录 | `{图片目录}_video` | `{图片 basename}.mp4` |

- 日期格式 `8.23` / `12.5`（不带前导 0）
- 序号用全角括号 `（1）` `（2）`
- 通用 vs 门店由 JSON 的 `category` 字段判断（`门店素材` / `门店` / `store` 任一关键词）
- 门店素材的 segment_name 第一段是门店编号（如 `130014WL-玉湖公园`）

示例：
```
outputs/
├── 玉湖公园-1-8.23-姜/
│   ├── 玉湖公园-1-8.23-（1）.png
│   ├── 玉湖公园-1-8.23-（2）.png
│   └── rejected/                # 审核员把不合格图放这里
│       └── 玉湖公园-1-8.23-（3）.png
└── 玉湖公园-1-8.23-姜_video/
    └── 玉湖公园-1-8.23-（1）.mp4
```

---

## 审核反馈机制

- **方案 C（文件标记）**：审核员把不合格图移到 `<image_dir>/rejected/` 子目录
- 程序下次运行（或本次 `--wait-seconds` 到期后）扫描 `rejected/`
- 发现不合格 → 删除对应位置 → 补生成 N 张 → 重新推送全量
- 最多循环 `--max-rounds` 轮（默认 5）

---

## 项目结构

```
src/aigc_pipeline/
├── __init__.py         # 公开 API
├── __main__.py         # python -m aigc_pipeline
├── cli.py              # CLI 入口（argparse）
├── config.py           # AppConfig + load_config (TOML)
├── core.py             # AIGC 底层：upload/submit/poll/download
├── naming.py           # 命名规则 + users.yaml 加载
├── feishu.py           # 飞书 webhook 推送
└── workflow.py         # 审核工作流编排
```

---

## 开发

```bash
# 安装 dev 依赖
pip install -e ".[dev]"

# 跑测试
pytest

# 构建 whl + sdist
python -m build

# 产出文件
ls dist/
# aigc_pipeline-0.1.0-py3-none-any.whl
# aigc_pipeline-0.1.0.tar.gz
```

---

## 完整生命周期（Phase 4）

下面是本包支持的端到端流程，以客户审批为主轴：

```
 +-------------------+   +--------------------+   +-------------------+
 | 1. 读 generation_ |   | 3. 客户审批         |   | 5. 视频生成       |
 |    tasks.json     |   |                    |   |                    |
 +-------------------+   +--------------------+   +-------------------+
        |                       |                      |
        v                       |                      v
 +-------------------+          |             +-------------------+
 | 2. 创建 image     |          |             | 6. 视频审批         |
 |   记录 (generating)|          |             |   (pending_review) |
 |    并发跑 core.run |          |             +-------------------+
 +-------------------+          |                      |
        |                      |                      v
        v                      |             +-------------------+
 +-------------------+          |             | 7. 百度网盘上传     |
 | image pending_    |          |             |   + create_share    |
 |   review          |          |             +-------------------+
 +-------------------+          |                      |
        |                      |                      v
        +----------------------+             +-------------------+
                                  |              |  8. video.shared   |
                                  +------------->+ + image.completed |
                                                 | + task.completed  |
                                                 +-------------------+
```

### 外部程序集成示例

```python
from aigc_pipeline import AIGCPipeline

pipeline = AIGCPipeline.from_config("./config.toml")

# 1. 读 JSON + 创建任务/图片记录 + 跑 AIGC
summary = pipeline.run_from_json(
    json_path="./tasks.json",
    output_dir="./outputs",
    parallel=True,                 # 进程级并发
    max_processes=8,
)

# 2. 外部程序轮询（定时调）拉取待审批
pending = pipeline.get_pending_reviews()
for p in pending:
    print(f"待审批: type={p['target_type']}, id={p['target_id']}")
    # 展示给客户（自己接 UI/邮件/IM），获得决策

# 3. 客户提交决策（外部程序调）
result = pipeline.submit_review_by_token(
    review_token=pending[0]["review_token"],
    decision="approved",       # or "rejected"
    comment="looks good",
)
# result: {image_id, new_status, next_action, retry_count, ...}

# 4. 触发视频生成（approved image → video pending_review）
pipeline.process_approved_images()

# 5. 视频审批（类似上面 step 2-3）

# 6. 上传百度网盘（approved video → shared）
# baidu 需先实现 `BaiduPanUploader.upload` 和 `create_share_link`（空壳在 baidu_uploader.py）
pipeline.process_approved_videos()

# 7. 巡查超时（定时调，如每 10 分钟一次）
abandoned = pipeline.run_dispatcher()
# abandoned: {abandoned_images, abandoned_videos, count}

# 8. 查询 task 全状态
status = pipeline.get_task_status("recvsxt4rtJ4Wd")
# status: {task, images[{...image, videos}]}
```

### 状态机一览

| 对象 | 状态转换 |
|---|---|
| **image** | generating → pending_review → (approved\|rejected) → (generating\[重试\]) \| abandoned |
| **video** | pending → generating → pending_review → (approved\|rejected) → (generating\[重试\]) \| abandoned → uploaded → shared |
| **task** | pending → generating → reviewing → completed \| abandoned |

### 决策路由

| 决策 | 路由 |
|---|---|
| image approved | → 触发 video 生成（调 process_approved_images） |
| image rejected | retry < 3 → 回到 generating 重生；>= 3 → abandoned |
| video approved | → 触发百度网盘上传（调 process_approved_videos） |
| video rejected | retry < 3 → 回到 generating（用**同 image**重生）；>= 3 → abandoned |
| pending 超时 | dispatcher 巡查 → abandoned |

---

## License

MIT

## 进阶文档

- [docs/ARCHITECTURE.md](docs/ARCHITECTURE.md) — 架构改造设计文档（从单进程到 10万 task/天的演进路线）