diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 825dd40..b63a932 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -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: diff --git a/README.md b/README.md index 8febb24..788f780 100644 --- a/README.md +++ b/README.md @@ -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. diff --git a/README_zh.md b/README_zh.md index 14d3325..f97f333 100644 --- a/README_zh.md +++ b/README_zh.md @@ -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)。 + ### 📖 基本概念 一个工作流代表了一个工作任务流程,包含了输入、输出以及工作流的触发方式。你可以任意定义输入是什么,输出是什么,以及输入是如何处理并到达输出结果的。 diff --git a/backend/api/workflow_api.py b/backend/api/workflow_api.py index a428083..f285c51 100644 --- a/backend/api/workflow_api.py +++ b/backend/api/workflow_api.py @@ -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)) diff --git a/backend/cli.py b/backend/cli.py new file mode 100644 index 0000000..8fd0837 --- /dev/null +++ b/backend/cli.py @@ -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()) diff --git a/backend/main.spec b/backend/main.spec index 6ce62cf..f1b84a3 100644 --- a/backend/main.spec +++ b/backend/main.spec @@ -20,7 +20,7 @@ block_cipher = None a = Analysis( - ["main.py"], + ["main.py", "cli.py"], pathex=[], binaries=[], datas=[ @@ -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", @@ -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, diff --git a/backend/models/workflow_models.py b/backend/models/workflow_models.py index f6a2bcc..a57aa1e 100644 --- a/backend/models/workflow_models.py +++ b/backend/models/workflow_models.py @@ -81,6 +81,7 @@ class WorkflowRunRecord(BaseModel): """用户工作流运行记录""" class RunFromTypes: + CLI = "CLI" WEB = "WEB" SCHEDULE = "SCHEDULE" API = "API" diff --git a/backend/pyproject.toml b/backend/pyproject.toml index a444539..e2b42b5 100644 --- a/backend/pyproject.toml +++ b/backend/pyproject.toml @@ -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" diff --git a/backend/tests/test_cli.py b/backend/tests/test_cli.py new file mode 100644 index 0000000..9ffd714 --- /dev/null +++ b/backend/tests/test_cli.py @@ -0,0 +1,59 @@ +from __future__ import annotations + +import json +import subprocess +import sys +from pathlib import Path + + +CLI = Path(__file__).resolve().parents[1] / "cli.py" + + +def command(workspace, *arguments, expected=0): + result = subprocess.run([sys.executable, "-X", "utf8", str(CLI), "--workspace", str(workspace), *map(str, arguments)], + capture_output=True, text=True, encoding="utf-8", timeout=120) + assert result.returncode == expected, result.stderr + result.stdout + return json.loads(result.stdout) + + +def test_cli_create_patch_run_status_and_failure(tmp_path): + workspace = tmp_path / "工作区" + file = tmp_path / "workflow.json" + fields = {"language": {"value": "python"}, "value": {"value": "hello", "type": "str"}, + "code": {"value": "def main(value): return value.upper()"}, "output": {"value": "", "is_output": True}} + definition = {"nodes": [{"id": "node", "type": "ProgrammingFunction", "category": "tools", "data": { + "task_name": "tools.programming_function", "template": fields}}], "edges": []} + file.write_text(json.dumps(definition), encoding="utf-8") + created = command(workspace, "workflow", "create", "--file", file, "--title", "CLI test") + wid = created["wid"] + patch = tmp_path / "inputs.json" + patch.write_text(json.dumps([{"node_id": "node", "field_name": "value", "value": "updated"}]), encoding="utf-8") + updated = command(workspace, "workflow", "patch", wid, "--inputs", patch, "--expected-version", "1") + assert updated["version"] == 2 + stale = command(workspace, "workflow", "patch", wid, "--inputs", patch, "--expected-version", "1", expected=2) + assert "version changed" in stale["message"] + result = command(workspace, "workflow", "run", wid) + assert result["status"] == "FINISHED" and result["outputs"]["node"]["output"] == "UPDATED" + status = command(workspace, "workflow", "status", result["rid"]) + assert status["workflow_version"] == 2 + patch.write_text(json.dumps([{"node_id": "node", "field_name": "code", "value": "def main(value): raise ValueError('sample')"}]), encoding="utf-8") + failure = command(workspace, "workflow", "run", wid, "--inputs", patch, expected=1) + assert failure["status"] == "FAILED" + saved = command(workspace, "workflow", "get", wid) + assert saved["data"]["nodes"][0]["data"]["template"]["code"]["value"] == "def main(value): return value.upper()" + + +def test_help_does_not_initialize_workspace(tmp_path): + workspace = tmp_path / "unused" + result = subprocess.run([sys.executable, str(CLI), "--workspace", str(workspace), "--help"], capture_output=True, timeout=10) + assert result.returncode == 0 and not workspace.exists() + + +def test_external_settings_errors_do_not_print_or_store_credentials(tmp_path): + settings = tmp_path / "settings.local.json" + secret = "test-secret-that-must-not-appear" + settings.write_text(json.dumps({"VERSION": "2", "endpoints": [{"api_key": secret}], "backends": {}}), encoding="utf-8") + workspace = tmp_path / "private-workspace" + result = command(workspace, "--settings", settings, "init", expected=2) + assert secret not in json.dumps(result) + assert secret.encode() not in (workspace / "data/my_database.db").read_bytes() diff --git a/backend/utilities/config/settings.py b/backend/utilities/config/settings.py index aa6bb6a..3d3c31d 100644 --- a/backend/utilities/config/settings.py +++ b/backend/utilities/config/settings.py @@ -1,5 +1,8 @@ # @Author: Bi Ying # @Date: 2024-04-29 16:50:17 +import json +import os +from pathlib import Path from copy import deepcopy from collections.abc import Mapping from typing import Any @@ -511,6 +514,14 @@ def __init__(self): self.load_setting() except Exception: self.data = dict() + external = os.environ.get("VECTORVEIN_LLM_SETTINGS_FILE") + if external: + raw = json.loads(Path(external).read_text(encoding="utf-8-sig")) + VvLlmSettings.model_validate(raw) + self.data["llm_settings"] = raw + output_dir = os.environ.get("VECTORVEIN_OUTPUT_DIR") + if output_dir: + self.data["output_folder"] = str(Path(output_dir).resolve()) def load_setting(self): from models import model_serializer diff --git a/backend/utilities/general/ratelimit.py b/backend/utilities/general/ratelimit.py index b94a140..6516e59 100644 --- a/backend/utilities/general/ratelimit.py +++ b/backend/utilities/general/ratelimit.py @@ -27,7 +27,9 @@ def add_request_record(product: str, cycle: int = 60) -> bool: request_records[product].popleft() if len(request_records[product]) > 0: - threading.Timer(interval=cycle, function=clear_expired_records, args=[product, cycle]).start() + timer = threading.Timer(interval=cycle, function=clear_expired_records, args=[product, cycle]) + timer.daemon = True + timer.start() return True @@ -64,7 +66,9 @@ def is_request_allowed(product: str, cycle: int, max_count: int, add_record: boo if add_record: request_records[product].append(current_time) - threading.Timer(cycle, clear_expired_records, [product, cycle]).start() + timer = threading.Timer(cycle, clear_expired_records, [product, cycle]) + timer.daemon = True + timer.start() return True diff --git a/backend/worker/tasks/llms/base_llm.py b/backend/worker/tasks/llms/base_llm.py index 4f32e2c..501bca9 100644 --- a/backend/worker/tasks/llms/base_llm.py +++ b/backend/worker/tasks/llms/base_llm.py @@ -141,6 +141,12 @@ def add_model_request_record(model: ModelSetting, endpoint: EndpointSetting) -> return add_request_record(product, cycle) +class StructuredOutputError(ValueError): + def __init__(self, message: str, result: ModelOutput): + super().__init__(message) + self.result = result + + class BaseLLMTask: MODEL_TYPE: BackendType NAME: str = "BaseLLMTask" @@ -176,6 +182,11 @@ def __init__(self, workflow_data: dict, node_id: str): self.thinking: ThinkingConfigParam | None | NotGiven = NOT_GIVEN self.reasoning_effort: ReasoningEffort | None | NotGiven = NOT_GIVEN self.extra_body: dict = {} + if self.workflow.get_node_field_value(node_id, "thinking_enabled", None) is False: + if self.MODEL_TYPE == BackendType.DeepSeek: + self.extra_body["thinking"] = {"type": "disabled"} + elif self.MODEL_TYPE == BackendType.Qwen: + self.extra_body["enable_thinking"] = False if self.MODEL_TYPE == BackendType.OpenAI and self.model.startswith(("gpt-5", "gpt-6")): self.temperature = NOT_GIVEN self.top_p = NOT_GIVEN @@ -193,6 +204,14 @@ def __init__(self, workflow_data: dict, node_id: str): ) self.model_settings = self.chat_client.backend_settings.models[self.model] + budget = self.workflow.get_node_field_value(node_id, "max_output_tokens", 0) + if budget: + budget = int(budget) + if not 1 <= budget <= self.model_settings.context_length: + raise ValueError("Output token limit must fit the model context") + self.model_settings = self.model_settings.model_copy(update={"max_output_tokens": budget}) + if self.workflow.get_node_field_value(node_id, "endpoint_policy", "all") == "first": + self.model_settings = self.model_settings.model_copy(update={"endpoints": self.model_settings.endpoints[:1]}) if isinstance(self.input_prompt, str): self.prompts = [self.input_prompt] @@ -367,7 +386,7 @@ def process_prompt( max_tokens, ) request_success = True - self.add_endpoint_request_record(endpoint) + # endpoint_available already reserves this request in the rate window. break except APIStatusError as e: if e.status_code == 429: @@ -501,7 +520,7 @@ def process_validated_prompt(self, prompt: str, index: int) -> ModelOutput: break except ValueError as exc: if validator is None or attempt == self.max_repairs: - raise + raise StructuredOutputError(str(exc), result) from exc request_prompt = prompt + "\nReturn only corrected JSON matching this schema:\n" + json.dumps(schema, ensure_ascii=False) request_prompt += "\nValidation error: " + str(exc) + "\nPrevious output:\n" + (result.content_output or "") result = result.model_copy(update={"prompt_tokens": prompt_tokens, "completion_tokens": completion_tokens}) @@ -528,7 +547,10 @@ def run(self): self.total_prompt_tokens += result.prompt_tokens self.total_completion_tokens += result.completion_tokens except Exception as exc: - failures.append({"index": index, "error": type(exc).__name__}) + failure = {"index": index, "error": type(exc).__name__} + if isinstance(exc, StructuredOutputError): + failure.update(message=str(exc), invalid_output=exc.result.content_output or "", completion_tokens=exc.result.completion_tokens, reasoning_characters=len(exc.result.reasoning_content or "")) + failures.append(failure) content_output = self.content_outputs[0] if isinstance(self.input_prompt, str) else self.content_outputs self.workflow.update_node_field_value(self.node_id, "output", content_output) @@ -554,4 +576,10 @@ def run(self): return self.workflow.data def get_max_concurrent_requests(self): - return max(vv_llm_settings.get_endpoint(get_endpoint_id(endpoint)).concurrent_requests for endpoint in self.model_settings.endpoints) + configured = max(vv_llm_settings.get_endpoint(get_endpoint_id(endpoint)).concurrent_requests for endpoint in self.model_settings.endpoints) + requested = self.workflow.get_node_field_value(self.node_id, "max_concurrent_requests", 0) + if requested: + if int(requested) < 1: + raise ValueError("Concurrency must be positive") + return min(configured, int(requested)) + return configured diff --git a/docs/cli.md b/docs/cli.md new file mode 100644 index 0000000..4452ec9 --- /dev/null +++ b/docs/cli.md @@ -0,0 +1,39 @@ +# Workflow CLI + +Run `pdm run cli` in backend, or use the console executable `VectorVeinCLI` included beside the desktop executable in CLI-enabled builds. Commands write UTF-8 JSON to stdout; progress goes to stderr. No desktop window or separate Celery worker is required. + +Global options precede the command: + +- `--workspace DIR`: directory containing config.json and data/. Defaults to the source backend or executable directory. Use the desktop workspace or an isolated experiment directory. +- `--settings FILE`: external vv-llm V2 configuration for this process; credentials are not imported into the settings database. +- `--output-dir DIR`: artifact directory for this process. +- `--result-file FILE`: also save the command result as JSON. + +```sh +python cli.py --workspace ./workspace init +python cli.py --workspace ./workspace workflow validate --file workflow.json +python cli.py --workspace ./workspace workflow create --file workflow.json --title "My workflow" +python cli.py --workspace ./workspace workflow list +python cli.py --workspace ./workspace workflow get WORKFLOW_ID +python cli.py --workspace ./workspace workflow update WORKFLOW_ID --file revised.json --expected-version 1 +python cli.py --workspace ./workspace workflow patch WORKFLOW_ID --inputs values.json --expected-version 2 +python cli.py --workspace ./workspace --settings settings.local.json workflow run WORKFLOW_ID --inputs values.json +python cli.py --workspace ./workspace workflow status RUN_ID +python cli.py --workspace ./workspace workflow records WORKFLOW_ID +``` + +A definition is a nodes/edges object or an exported workflow with its definition under data. Creation and updates validate unique IDs, registered tasks, edge handles, one incoming edge per field and an acyclic graph. Updates preserve metadata and do not call a model to generate Agent tool metadata. + +Input overrides and patches: + +```json +[{"node_id":"input-node","field_name":"text","value":"New input"}] +``` + +patch saves a new version. run uses a run-only copy, records the saved version and returns its run ID/status/output fields. --full on run/status includes the complete node snapshot. Unknown node/field references fail. + +Exit codes: 0 success, 1 workflow failure, 2 invalid arguments/configuration/input. Run start emits its ID to stderr, permitting status queries from another process. Python nodes execute the code in the supplied workflow. + +Model node templates also accept max_output_tokens and max_concurrent_requests (positive integers; 0 or absent inherits configuration), endpoint_policy="first", and thinking_enabled=false for DeepSeek/Qwen. Schema failures retain validation diagnostics and invalid output in run_stats; invalid outputs are never cached as successes. + +Reproduce CLI tests with `python -m pytest tests/test_cli.py`. Help/version do not initialize a workspace.