Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,16 @@ jobs:
working-directory: backend
run: pdm run python build.py -t p -v "${{ steps.meta.outputs.version }}"

- name: Smoke-test packaged CLI
working-directory: backend
run: |
if [ -f cli.py ]; then
CLI="./dist/vector-vein/VectorVeinCLI"
if [ "${RUNNER_OS}" = "Windows" ]; then CLI="${CLI}.exe"; fi
"$CLI" --help
"$CLI" --workspace "${RUNNER_TEMP}/vectorvein-cli-smoke" init
fi

- name: Upload packaged artifact
uses: actions/upload-artifact@v7
with:
Expand Down
4 changes: 4 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,10 @@ For detailed API documentation, visit `http://localhost:8787/docs` after startin

Version 0.4.16 adds source-preserving multi-sheet workbook input, typed XLSX output, structured DOCX reports and automatic PNG charts. LLM nodes support local JSON Schema validation, bounded repairs and per-item result reuse. Failed batches stop downstream execution. See the [workflow guide and reproducible example](docs/reliable-workflows.md).

### Workflow CLI

Use `pdm run cli` from source, or the `VectorVeinCLI` console executable in CLI-enabled desktop builds, to create, validate, update and run workflows or inspect results. See the [CLI guide](docs/cli.md).

### 📖 Basic Concepts

A workflow represents a work task process, including input, output, and how input is processed to reach the output result.
Expand Down
4 changes: 4 additions & 0 deletions README_zh.md
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,10 @@ print(result['data']) # 工作流输出

从 v0.4.16 起,文件读取支持保留来源行号的多工作表 JSON,文档节点支持多表 XLSX、结构化 DOCX 和自动 PNG 图表。LLM 节点可配置 JSON Schema 校验、有限修复和逐项结果缓存;批次失败会阻止后续步骤。详见 [使用与复现说明](docs/reliable-workflows.md)。

### 命令行工作流

源码环境可用 `pdm run cli`,支持创建、验证、修改、运行工作流和查询运行结果;含 CLI 的桌面包另提供 `VectorVeinCLI` 控制台程序。支持独立工作区和外部模型配置,使用方式见 [CLI 文档](docs/cli.md)。

### 📖 基本概念

一个工作流代表了一个工作任务流程,包含了输入、输出以及工作流的触发方式。你可以任意定义输入是什么,输出是什么,以及输入是如何处理并到达输出结果的。
Expand Down
3 changes: 2 additions & 1 deletion backend/api/workflow_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -153,7 +153,8 @@ def update(self, payload):
workflow.version = int(workflow.version or 1) + 1
workflow.save()

update_workflow_tool_call_data.delay(workflow_wid=workflow.wid.hex, force=title_changed)
if payload.get("refresh_tool_data", True):
update_workflow_tool_call_data.delay(workflow_wid=workflow.wid.hex, force=title_changed)

return JResponse(data=model_serializer(workflow, manytomany=True))

Expand Down
262 changes: 262 additions & 0 deletions backend/cli.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,262 @@
"""Headless workflow CLI sharing the desktop database and execution engine."""
from __future__ import annotations

import argparse
import io
import contextlib
import json
import os
import re
import sys
from copy import deepcopy
from datetime import datetime
from pathlib import Path

from bootstrap import APP_ROOT


def read_json(file: Path):
return json.loads(file.read_text(encoding="utf-8-sig"))


def version():
file = APP_ROOT / "version.txt"
if file.exists():
return file.read_text(encoding="utf-8").strip()
match = re.search(r'^version\s*=\s*"([^"]+)"', (APP_ROOT / "pyproject.toml").read_text(encoding="utf-8"), re.M)
return match.group(1) if match else "unknown"


def parser():
result = argparse.ArgumentParser(description=__doc__)
result.add_argument("--version", action="version", version=version())
result.add_argument("--workspace", type=Path, default=APP_ROOT, help="Desktop installation or isolated workspace directory")
result.add_argument("--settings", type=Path, help="External vv-llm settings for this process; never copied into the database")
result.add_argument("--output-dir", type=Path, help="Artifact directory for this process")
result.add_argument("--result-file", type=Path, help="Also save the command JSON result")
groups = result.add_subparsers(dest="group", required=True)
groups.add_parser("init", help="Initialize the workspace without opening a GUI")
commands = groups.add_parser("workflow").add_subparsers(dest="command", required=True)
listing = commands.add_parser("list")
listing.add_argument("--search", default="")
listing.add_argument("--limit", type=int, default=20)
get = commands.add_parser("get")
get.add_argument("wid")
create = commands.add_parser("create")
create.add_argument("--file", type=Path, required=True)
create.add_argument("--title")
update = commands.add_parser("update")
update.add_argument("wid")
update.add_argument("--file", type=Path, required=True)
update.add_argument("--title")
update.add_argument("--expected-version", type=int, help="Refuse to overwrite a different saved version")
patch = commands.add_parser("patch")
patch.add_argument("wid")
patch.add_argument("--inputs", type=Path, required=True, help="JSON array of node_id/field_name/value records")
patch.add_argument("--expected-version", type=int)
validate = commands.add_parser("validate")
validate.add_argument("--file", type=Path, required=True)
run = commands.add_parser("run")
run.add_argument("wid")
run.add_argument("--inputs", type=Path)
run.add_argument("--full", action="store_true", help="Include all intermediate node values")
status = commands.add_parser("status")
status.add_argument("rid")
status.add_argument("--full", action="store_true")
records = commands.add_parser("records")
records.add_argument("wid")
records.add_argument("--limit", type=int, default=20)
return result


def validate_data(data: dict):
from utilities.workflow import Workflow
from background_task.workflow_tasks import task_functions
if not isinstance(data, dict) or not isinstance(data.get("nodes"), list) or not isinstance(data.get("edges"), list):
raise ValueError("Workflow must contain nodes and edges arrays")
nodes = {node["id"]: node for node in data["nodes"]}
if len(nodes) != len(data["nodes"]):
raise ValueError("Node IDs must be unique")
for node in nodes.values():
module, function = node["data"]["task_name"].split(".")
if function not in task_functions.get(module, {}):
raise ValueError(f"Unknown task: {module}.{function}")
targets = set()
for edge in data["edges"]:
if edge.get("ignored"):
continue
for side in ("source", "target"):
if edge[side] not in nodes or edge[side + "Handle"] not in nodes[edge[side]]["data"]["template"]:
raise ValueError("Edge references an unknown node or field")
target = (edge["target"], edge["targetHandle"])
if target in targets:
raise ValueError("A field may have only one incoming edge")
targets.add(target)
Workflow(deepcopy(data)).dag.topological_sort()


def apply_inputs(data, inputs):
if not isinstance(inputs, list):
raise ValueError("Inputs must be an array")
nodes = {node["id"]: node for node in data["nodes"]}
for item in inputs:
node = nodes.get(item["node_id"])
if node is None or item["field_name"] not in node["data"]["template"]:
raise ValueError("Input references an unknown node or field")
node["data"]["template"][item["field_name"]]["value"] = item["value"]


def initialize():
from models import create_tables, database
from utilities.config import config, Settings
from peewee_migrate import Router
Path(config.data_path).mkdir(parents=True, exist_ok=True)
router = Router(database, migrate_dir=str(APP_ROOT / "migrations"))
if not database.table_exists("workflow"):
create_tables()
router.run(fake=True)
else:
router.run()
create_tables()
Settings()


def checked(response):
if response.get("status") != 200:
raise ValueError(response.get("msg", "Command failed"))
return response["data"]


def record_result(record, full=False):
result = {"rid": record.rid.hex, "wid": record.workflow.wid.hex, "status": record.status,
"workflow_version": record.workflow_version, "outputs": {}}
for node in record.data.get("nodes", []):
fields = {name: item.get("value") for name, item in node["data"]["template"].items() if item.get("is_output")}
if fields:
result["outputs"][node["id"]] = fields
if full:
result["data"] = record.data
if record.status == "FAILED":
result["error_task"] = record.data.get("error_task", "")
return result


def execute(args):
from api.workflow_api import WorkflowAPI
from models import Workflow as WorkflowModel, WorkflowRunRecord
from utilities.workflow import WorkflowData
from celery_worker import app
from background_task.workflow_tasks import run_workflow
api = WorkflowAPI()
if args.group == "init":
return {"workspace": str(Path.cwd()), "version": version()}, 0
if args.command == "list":
data = checked(api.list({"search_text": args.search, "page_size": args.limit}))
return {"total": data["total"], "workflows": [{k: row[k] for k in ("wid", "title", "version", "status")} for row in data["workflows"]]}, 0
if args.command == "validate":
spec = read_json(args.file)
validate_data(spec.get("data", spec))
return {"valid": True}, 0
if args.command == "create":
spec = read_json(args.file)
data = spec.get("data", spec)
validate_data(data)
payload = {"title": args.title or spec.get("title") or args.file.stem, "brief": spec.get("brief", ""), "data": data}
created = checked(api.create(payload))
return {k: created[k] for k in ("wid", "title", "version")}, 0
if args.command == "status":
record = WorkflowRunRecord.get_or_none(WorkflowRunRecord.rid == args.rid)
if record is None:
raise ValueError("Run record not found")
return record_result(record, args.full), 0
if args.command == "records":
query = WorkflowRunRecord.select().join(WorkflowModel).where(WorkflowModel.wid == args.wid).order_by(WorkflowRunRecord.start_time.desc()).limit(args.limit)
return {"runs": [{"rid": record.rid.hex, "status": record.status, "workflow_version": record.workflow_version, "start_time": record.start_time.isoformat()} for record in query]}, 0
saved = checked(api.get({"wid": args.wid}))
if args.command == "get":
return saved, 0
if args.command in ("update", "patch"):
if args.expected_version is not None and saved["version"] != args.expected_version:
raise ValueError("Saved workflow version changed")
payload = {key: saved[key] for key in ("title", "brief", "images", "language")}
payload.update(wid=args.wid, tags=[tag["tid"] for tag in saved.get("tags", [])], refresh_tool_data=False)
if args.command == "update":
spec = read_json(args.file)
data = spec.get("data", spec)
payload["title"] = args.title or spec.get("title") or saved["title"]
payload["brief"] = spec.get("brief", saved["brief"])
else:
data = deepcopy(saved["data"])
apply_inputs(data, read_json(args.inputs))
validate_data(data)
payload["data"] = data
updated = checked(api.update(payload))
return {key: updated[key] for key in ("wid", "title", "version")}, 0
data = deepcopy(saved["data"])
if args.inputs:
apply_inputs(data, read_json(args.inputs))
validate_data(data)
data["related_workflows"] = WorkflowData(data).related_workflows
model = WorkflowModel.get_by_id(args.wid)
record = WorkflowRunRecord.create(workflow=model, workflow_version=model.version, data=data, status="RUNNING", run_from=WorkflowRunRecord.RunFromTypes.CLI)
data.update(wid=model.wid.hex, rid=record.rid.hex)
print(json.dumps({"event": "run_started", "rid": record.rid.hex}), file=sys.stderr, flush=True)
app.conf.update(task_always_eager=True, task_eager_propagates=False)
try:
# The same Celery chain runs synchronously; a separate broker worker is unnecessary.
run_workflow.run(data).get(propagate=True)
except (Exception, KeyboardInterrupt) as exc:
record = WorkflowRunRecord.get_by_id(record.rid)
if record.status != "FAILED":
record.status = "FAILED"
record.end_time = datetime.now()
record.data = {**record.data, "error_type": type(exc).__name__}
record.save()
record = WorkflowRunRecord.get_by_id(record.rid)
return record_result(record, args.full), 0 if record.status == "FINISHED" else 1


def main(argv=None):
for stream in (sys.stdout, sys.stderr):
if isinstance(stream, io.TextIOWrapper):
stream.reconfigure(encoding="utf-8")
args = parser().parse_args(argv)
for name in ("workspace", "settings", "output_dir", "result_file", "file", "inputs"):
value = getattr(args, name, None)
if value is not None:
setattr(args, name, value.resolve())
args.workspace.mkdir(parents=True, exist_ok=True)
if args.settings:
os.environ["VECTORVEIN_LLM_SETTINGS_FILE"] = str(args.settings)
if args.output_dir:
os.environ["VECTORVEIN_OUTPUT_DIR"] = str(args.output_dir)
os.chdir(args.workspace)
code = 2
try:
with contextlib.redirect_stdout(sys.stderr):
initialize()
result, code = execute(args)
except Exception as exc:
# Provider errors may contain credentials; only configuration errors expose a message.
result = {"error": type(exc).__name__}
if type(exc).__name__ == "ValidationError":
result["message"] = "Configuration validation failed; check the selected settings and schema"
elif isinstance(exc, (ValueError, FileNotFoundError)):
result["message"] = str(exc)
finally:
if "utilities.config" in sys.modules:
from utilities.config import cache, config
from models import database
cache.close()
database.close()
config.close()
encoded = json.dumps(result, ensure_ascii=False, default=str, indent=2)
if args.result_file:
args.result_file.parent.mkdir(parents=True, exist_ok=True)
args.result_file.write_text(encoded + "\n", encoding="utf-8")
print(encoded)
return code


if __name__ == "__main__":
raise SystemExit(main())
16 changes: 14 additions & 2 deletions backend/main.spec
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ block_cipher = None


a = Analysis(
["main.py"],
["main.py", "cli.py"],
pathex=[],
binaries=[],
datas=[
Expand Down Expand Up @@ -61,7 +61,7 @@ pyz = PYZ(a.pure, a.zipped_data, cipher=block_cipher)

exe = EXE(
pyz,
a.scripts,
[script for script in a.scripts if script[0] != "cli"],
[],
exclude_binaries=True,
name="VectorVein",
Expand All @@ -77,8 +77,20 @@ exe = EXE(
entitlements_file=None,
icon="web/assets/favicon.ico",
)
cli_exe = EXE(
pyz,
[script for script in a.scripts if script[0] != "main"],
[],
exclude_binaries=True,
name="VectorVeinCLI",
console=True,
strip=False,
upx=True,
)

coll = COLLECT(
exe,
cli_exe,
a.binaries,
a.zipfiles,
a.datas,
Expand Down
1 change: 1 addition & 0 deletions backend/models/workflow_models.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ class WorkflowRunRecord(BaseModel):
"""用户工作流运行记录"""

class RunFromTypes:
CLI = "CLI"
WEB = "WEB"
SCHEDULE = "SCHEDULE"
API = "API"
Expand Down
1 change: 1 addition & 0 deletions backend/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ build.cmd = "python build.py -t p"
build.env_file = ".env"
dev.cmd = "python main.py"
dev.env_file = ".env"
cli = "python cli.py"
fullstack-dev.cmd = "python run_fullstack_dev.py"
fullstack-dev.env_file = ".env"
test = "pytest tests"
Expand Down
Loading